1use std::path::{Path, PathBuf};
41#[cfg(unix)]
42use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
43use std::sync::Arc;
44#[cfg(unix)]
45use std::time::Duration;
46
47use async_trait::async_trait;
48use khive_db::{StorageBackend, WalCeilingPolicy};
49use khive_storage::event::{EventPageQuery, EventPageWindow, IdempotentEventBatchResult};
50use khive_storage::{
51 BatchWriteSummary, Event, EventFilter, EventStore, Page, PageRequest, StorageError,
52 StorageResult,
53};
54use serde::{Deserialize, Serialize};
55
56#[cfg(unix)]
57use tokio::net::{UnixListener, UnixStream};
58use uuid::Uuid;
59
60#[cfg(unix)]
61use crate::daemon::{read_frame, write_frame};
62
63pub const EVENTS_PROTOCOL_VERSION: u32 = 4;
79
80#[cfg(unix)]
84pub const DEFAULT_APPEND_QUEUE_BATCHES: usize = 4096;
85
86#[cfg(unix)]
90pub const DEFAULT_APPEND_QUEUE_BYTES: usize = 32 * 1024 * 1024;
91
92#[cfg(unix)]
97const REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
98
99#[cfg(unix)]
101const FORWARDER_BACKOFF: Duration = Duration::from_secs(2);
102
103#[cfg(unix)]
110const MAX_EVENTS_CONNECTIONS: usize = 128;
111
112#[cfg(unix)]
120const CONN_IO_TIMEOUT: Duration = Duration::from_secs(60);
121
122#[cfg(unix)]
132const MAX_INFLIGHT_REQUEST_BYTES: usize = 64 * 1024 * 1024;
133
134#[cfg(unix)]
139const MAX_CACHED_NAMESPACE_STORES: usize = 1024;
140
141const EVENTS_SYMLINK_HOP_BOUND: u32 = 40;
146
147pub const MAX_QUERY_EVENTS_PAGE_ROWS: u32 = 4096;
165
166pub fn events_db_path_beside(main_db: &Path) -> PathBuf {
181 let resolved = std::fs::canonicalize(main_db)
188 .unwrap_or_else(|_| resolve_dangling_final_component(main_db));
189 let mut name = resolved
190 .file_name()
191 .map(std::ffi::OsStr::to_os_string)
192 .unwrap_or_else(|| std::ffi::OsString::from("khive.db"));
193 name.push(".events.db");
194 let path = match resolved.parent().filter(|dir| !dir.as_os_str().is_empty()) {
195 Some(dir) => std::fs::canonicalize(dir)
196 .unwrap_or_else(|_| dir.to_path_buf())
197 .join(&name),
198 None => PathBuf::from(name),
199 };
200 absolutize(&path)
201}
202
203fn resolve_dangling_final_component(path: &Path) -> PathBuf {
210 let mut current = path.to_path_buf();
211 for _ in 0..EVENTS_SYMLINK_HOP_BOUND {
216 match std::fs::read_link(¤t) {
217 Ok(target) => {
218 current = if target.is_absolute() {
219 target
220 } else {
221 match current.parent().filter(|dir| !dir.as_os_str().is_empty()) {
222 Some(dir) => dir.join(target),
223 None => target,
224 }
225 };
226 }
227 Err(_) => break,
228 }
229 }
230 current
231}
232
233fn absolutize(path: &Path) -> PathBuf {
239 if path.is_absolute() {
240 return path.to_path_buf();
241 }
242 std::env::current_dir()
243 .map(|cwd| cwd.join(path))
244 .unwrap_or_else(|_| path.to_path_buf())
245}
246
247pub fn events_socket_path_beside(events_db: &Path) -> PathBuf {
257 events_db.with_extension("sock")
258}
259
260#[cfg(unix)]
261#[path = "events_split_socket_path.rs"]
262pub(crate) mod socket_path;
263#[cfg(unix)]
264pub use socket_path::validate_events_socket_path;
265
266#[derive(Debug, Clone)]
268pub struct EventsSplitConfig {
269 pub db_path: PathBuf,
272 pub socket_path: Option<PathBuf>,
276}
277
278#[cfg(unix)]
288type ClientMap = std::collections::HashMap<PathBuf, Arc<EventsSplitClient>>;
289type BackendMap = std::collections::HashMap<PathBuf, (bool, Arc<StorageBackend>)>;
293
294#[cfg(unix)]
295fn client_registry() -> &'static std::sync::Mutex<ClientMap> {
296 static REGISTRY: std::sync::OnceLock<std::sync::Mutex<ClientMap>> = std::sync::OnceLock::new();
297 REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
298}
299
300fn direct_backend_registry() -> &'static std::sync::Mutex<BackendMap> {
301 static REGISTRY: std::sync::OnceLock<std::sync::Mutex<BackendMap>> = std::sync::OnceLock::new();
302 REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
303}
304
305#[cfg(any(test, feature = "test-internals"))]
312#[doc(hidden)]
313pub struct TestRegistryGuard<'a> {
314 root: &'a Path,
315 canonical_root: PathBuf,
316}
317
318#[cfg(any(test, feature = "test-internals"))]
319impl<'a> TestRegistryGuard<'a> {
320 pub fn new(root: &'a Path) -> Self {
322 Self {
323 root,
324 canonical_root: std::fs::canonicalize(root).expect("existing unique fixture root"),
325 }
326 }
327
328 fn remove_entries<T>(
329 &self,
330 registry: &std::sync::Mutex<std::collections::HashMap<PathBuf, T>>,
331 ) -> Vec<T> {
332 let mut entries = registry
333 .lock()
334 .unwrap_or_else(std::sync::PoisonError::into_inner);
335 let keys: Vec<_> = entries
336 .keys()
337 .filter(|path| path.starts_with(self.root) || path.starts_with(&self.canonical_root))
338 .cloned()
339 .collect();
340 keys.into_iter()
341 .filter_map(|key| entries.remove(&key))
342 .collect()
343 }
344}
345
346#[cfg(any(test, feature = "test-internals"))]
347impl Drop for TestRegistryGuard<'_> {
348 fn drop(&mut self) {
349 let backends = self.remove_entries(direct_backend_registry());
352 #[cfg(unix)]
353 let clients = self.remove_entries(client_registry());
354 drop(backends);
355 #[cfg(unix)]
356 drop(clients);
357 }
358}
359
360#[cfg(unix)]
363pub fn client_for(socket_path: &Path) -> crate::error::RuntimeResult<Arc<EventsSplitClient>> {
364 let mut registry = client_registry()
365 .lock()
366 .unwrap_or_else(std::sync::PoisonError::into_inner);
367 if let Some(existing) = registry.get(socket_path) {
368 return Ok(Arc::clone(existing));
369 }
370 let client = EventsSplitClient::new(socket_path.to_path_buf())?;
371 registry.insert(socket_path.to_path_buf(), Arc::clone(&client));
372 Ok(client)
373}
374
375pub fn direct_backend_for(db_path: &Path) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
378 direct_backend_with_max_readers(db_path, false, None)
379}
380
381pub fn direct_backend_read_only_for(
386 db_path: &Path,
387) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
388 direct_backend_with_max_readers(db_path, true, None)
389}
390
391pub(crate) fn direct_backend_with_max_readers(
393 db_path: &Path,
394 read_only: bool,
395 max_readers: Option<usize>,
396) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
397 let mut config = crate::RuntimeConfig {
398 db_path: Some(db_path.to_path_buf()),
399 ..crate::RuntimeConfig::no_embeddings()
400 };
401 let wal_ceiling = config.resolve_wal_ceiling_policy(read_only)?;
402 let disk_guard = config.resolve_disk_guard_policy(read_only)?;
403 direct_backend_with_policies(
404 db_path,
405 read_only,
406 max_readers,
407 wal_ceiling,
408 disk_guard,
409 config.volume_lock_dir,
410 )
411}
412
413pub(crate) fn direct_backend_with_policies(
417 db_path: &Path,
418 read_only: bool,
419 max_readers: Option<usize>,
420 wal_ceiling: WalCeilingPolicy,
421 disk_guard: Option<khive_db::EffectiveDiskGuardConfig>,
422 volume_lock_dir: Option<PathBuf>,
423) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
424 wal_ceiling.validate_static(true, true, read_only)?;
425 if !read_only {
426 disk_guard
427 .ok_or_else(|| {
428 crate::error::RuntimeError::Internal("missing events disk policy".into())
429 })?
430 .validate()?;
431 khive_db::require_volume_lock_dir(volume_lock_dir.clone())?;
432 }
433 let mut registry = direct_backend_registry()
434 .lock()
435 .unwrap_or_else(std::sync::PoisonError::into_inner);
436 let absolute_path = absolutize(db_path);
437 let db_path = absolute_path.as_path();
438 refuse_events_db_symlinks(db_path)
439 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
440 #[cfg(unix)]
441 if !read_only {
442 if let Some(parent) = db_path.parent().filter(|p| !p.as_os_str().is_empty()) {
443 let _ = std::fs::create_dir_all(parent);
444 }
445 }
446 #[cfg(unix)]
452 ensure_events_db_parent_trusted(db_path)
453 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
454 let key = std::fs::canonicalize(db_path)
455 .or_else(|error| {
456 if error.kind() != std::io::ErrorKind::NotFound {
457 return Err(error);
458 }
459 let parent = db_path.parent().ok_or(error)?;
460 let name = db_path.file_name().ok_or_else(|| {
461 std::io::Error::new(
462 std::io::ErrorKind::InvalidInput,
463 "events path has no file name",
464 )
465 })?;
466 std::fs::canonicalize(parent).map(|parent| parent.join(name))
467 })
468 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
469 if let Some((existing_read_only, existing)) = registry.get(&key) {
470 if *existing_read_only != read_only {
471 let existing_mode = if *existing_read_only {
472 "read-only"
473 } else {
474 "writable"
475 };
476 let requested_mode = if read_only { "read-only" } else { "writable" };
477 return Err(crate::error::RuntimeError::InvalidInput(format!(
478 "events database {} is already open {existing_mode} in this process; cannot \
479 open it {requested_mode}; read-only events access requires a separate frozen snapshot",
480 key.display()
481 )));
482 }
483 let numbers = |p: Option<khive_db::EffectiveDiskGuardConfig>| {
484 p.map(|p| (p.reserve_bytes, p.guard_deadline_ms))
485 };
486 if numbers(existing.pool().effective_disk_guard_config()) != numbers(disk_guard) {
487 return Err(crate::error::RuntimeError::Internal(
488 "events database is already open with a different disk reserve/deadline policy"
489 .into(),
490 ));
491 }
492 if !read_only && existing.pool().config().volume_lock_dir != volume_lock_dir {
493 return Err(crate::error::RuntimeError::Internal(
494 "events database is already open with a different volume-lock directory".into(),
495 ));
496 }
497 let existing_bytes = existing.pool().config().wal_ceiling.bytes;
498 if existing_bytes != wal_ceiling.bytes {
499 return Err(khive_db::SqliteError::InvalidConfig(format!(
500 "events database {} is already open with wal_ceiling_bytes={existing_bytes}; \
501 requested {}; drain and restart before changing the WAL ceiling",
502 key.display(),
503 wal_ceiling.bytes
504 ))
505 .into());
506 }
507 return Ok(Arc::clone(existing));
508 }
509 let db_path = key.as_path();
510 #[cfg(unix)]
516 if !read_only {
517 use std::os::unix::fs::OpenOptionsExt;
518 let _ = std::fs::OpenOptions::new()
519 .write(true)
520 .create_new(true)
521 .mode(0o600)
522 .open(db_path);
523 let _before_open = harden_events_db_sidecars(db_path)
526 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
527 }
528 let backend = Arc::new(if read_only {
529 StorageBackend::sqlite_read_only_with_max_readers_and_wal_ceiling(
530 db_path,
531 max_readers,
532 wal_ceiling,
533 )?
534 } else {
535 StorageBackend::sqlite_with_max_readers_and_policies(
536 db_path,
537 max_readers,
538 wal_ceiling,
539 disk_guard.ok_or_else(|| {
540 crate::error::RuntimeError::Internal("missing events disk policy".into())
541 })?,
542 khive_db::require_volume_lock_dir(volume_lock_dir)?,
543 )?
544 });
545 registry.insert(key, (read_only, Arc::clone(&backend)));
546 Ok(backend)
547}
548
549#[cfg(unix)]
552pub fn forwarding_metrics(socket_path: &Path) -> Option<EventsForwardingMetrics> {
553 let registry = client_registry()
554 .lock()
555 .unwrap_or_else(std::sync::PoisonError::into_inner);
556 registry.get(socket_path).map(|client| client.metrics())
557}
558
559#[derive(Debug, Serialize, Deserialize)]
565#[serde(tag = "op", rename_all = "snake_case")]
566pub enum EventsRequest {
567 AppendEvents {
571 protocol_version: u32,
572 namespace: String,
573 events: Vec<Event>,
574 },
575 AppendEventsIdempotent {
577 protocol_version: u32,
578 namespace: String,
579 events: Vec<Event>,
580 },
581 GetEvent {
582 protocol_version: u32,
583 namespace: String,
584 id: Uuid,
585 },
586 QueryEvents {
587 protocol_version: u32,
588 namespace: String,
589 filter: EventFilter,
590 page: PageRequest,
591 },
592 QueryEventPage {
593 protocol_version: u32,
594 namespace: String,
595 query: EventPageQuery,
596 },
597 CountEvents {
598 protocol_version: u32,
599 namespace: String,
600 filter: EventFilter,
601 },
602}
603
604impl EventsRequest {
605 fn protocol_version(&self) -> u32 {
606 match self {
607 Self::AppendEvents {
608 protocol_version, ..
609 }
610 | Self::AppendEventsIdempotent {
611 protocol_version, ..
612 }
613 | Self::GetEvent {
614 protocol_version, ..
615 }
616 | Self::QueryEvents {
617 protocol_version, ..
618 }
619 | Self::QueryEventPage {
620 protocol_version, ..
621 }
622 | Self::CountEvents {
623 protocol_version, ..
624 } => *protocol_version,
625 }
626 }
627
628 fn namespace(&self) -> &str {
629 match self {
630 Self::AppendEvents { namespace, .. }
631 | Self::AppendEventsIdempotent { namespace, .. }
632 | Self::GetEvent { namespace, .. }
633 | Self::QueryEvents { namespace, .. }
634 | Self::QueryEventPage { namespace, .. }
635 | Self::CountEvents { namespace, .. } => namespace,
636 }
637 }
638}
639
640#[derive(Debug, Serialize, Deserialize)]
642#[serde(tag = "kind", rename_all = "snake_case")]
643pub enum EventsResponse {
644 Appended {
645 summary: BatchWriteSummary,
646 },
647 Idempotent {
648 result: IdempotentEventBatchResult,
649 },
650 Event {
651 event: Option<Event>,
652 },
653 Pageful {
654 page: Page<Event>,
655 },
656 EventPageWindow {
657 window: EventPageWindow,
658 },
659 Count {
660 count: u64,
661 },
662 Error {
672 message: String,
673 retryable: bool,
674 #[serde(default)]
675 writer_task_failure: Option<WireWriterTaskFailure>,
676 },
677}
678
679#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
682#[serde(tag = "kind", rename_all = "snake_case")]
683pub enum WireWriterTaskFailure {
684 RequestFailed {
685 request_state: WireWriterTaskState,
686 },
687 TaskTerminated {
688 request_state: WireWriterTaskState,
689 #[serde(default, skip_serializing_if = "Option::is_none")]
690 sqlite_full_codes: Option<(i32, i32)>,
691 },
692}
693
694#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
698#[serde(rename_all = "snake_case")]
699pub enum WireWriterTaskState {
700 NotStarted,
701 TransactionRolledBack,
702 SideEffectsUnknown,
703}
704
705impl From<khive_storage::WriterTaskRequestState> for WireWriterTaskState {
706 fn from(state: khive_storage::WriterTaskRequestState) -> Self {
707 use khive_storage::WriterTaskRequestState as S;
708 match state {
709 S::NotStarted => Self::NotStarted,
710 S::TransactionRolledBack => Self::TransactionRolledBack,
711 S::SideEffectsUnknown => Self::SideEffectsUnknown,
712 }
713 }
714}
715
716impl From<WireWriterTaskState> for khive_storage::WriterTaskRequestState {
717 fn from(state: WireWriterTaskState) -> Self {
718 use khive_storage::WriterTaskRequestState as S;
719 match state {
720 WireWriterTaskState::NotStarted => S::NotStarted,
721 WireWriterTaskState::TransactionRolledBack => S::TransactionRolledBack,
722 WireWriterTaskState::SideEffectsUnknown => S::SideEffectsUnknown,
723 }
724 }
725}
726
727#[cfg(unix)]
735pub struct EventsDaemonGuard {
736 _file: std::fs::File,
737}
738
739#[cfg(unix)]
743pub fn try_acquire_events_daemon_guard(socket_path: &Path) -> Option<EventsDaemonGuard> {
744 match acquire_events_daemon_guard_outcome(socket_path) {
745 EventsDaemonGuardAcquisition::Held(guard) => Some(guard),
746 _ => None,
747 }
748}
749
750#[cfg(unix)]
753enum EventsDaemonGuardAcquisition {
754 Held(EventsDaemonGuard),
755 Contended,
756 OpenFailed(std::io::Error),
757 HardeningRefused(std::io::Error),
758}
759
760#[cfg(unix)]
761impl std::fmt::Debug for EventsDaemonGuardAcquisition {
762 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
763 match self {
764 Self::Held(_) => f.write_str("Held"),
765 Self::Contended => f.write_str("Contended"),
766 Self::OpenFailed(error) => f.debug_tuple("OpenFailed").field(error).finish(),
767 Self::HardeningRefused(error) => {
768 f.debug_tuple("HardeningRefused").field(error).finish()
769 }
770 }
771 }
772}
773
774#[cfg(unix)]
780fn acquire_events_daemon_guard_outcome(socket_path: &Path) -> EventsDaemonGuardAcquisition {
781 use EventsDaemonGuardAcquisition::{Contended, HardeningRefused, Held, OpenFailed};
782
783 let lock_path = socket_path.with_extension("lock");
784 if let Some(parent) = lock_path.parent() {
785 if let Err(error) = std::fs::create_dir_all(parent) {
786 return OpenFailed(error);
787 }
788 }
789 use std::os::unix::fs::OpenOptionsExt;
790 let file = match std::fs::OpenOptions::new()
794 .create(true)
795 .truncate(false)
796 .write(true)
797 .mode(0o600)
798 .custom_flags(libc::O_NOFOLLOW)
799 .open(&lock_path)
800 {
801 Ok(file) => file,
802 Err(error) if error.raw_os_error() == Some(libc::ELOOP) => {
803 return HardeningRefused(error);
804 }
805 Err(error) => return OpenFailed(error),
806 };
807 if let Err(error) = file.set_permissions(
813 <std::fs::Permissions as std::os::unix::fs::PermissionsExt>::from_mode(0o600),
814 ) {
815 return HardeningRefused(error);
816 }
817 use std::os::fd::AsRawFd;
818 let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
821 if rc == 0 {
822 Held(EventsDaemonGuard { _file: file })
823 } else {
824 let error = std::io::Error::last_os_error();
825 if error.raw_os_error() == Some(libc::EWOULDBLOCK)
826 || error.raw_os_error() == Some(libc::EAGAIN)
827 {
828 Contended
829 } else {
830 HardeningRefused(error)
831 }
832 }
833}
834
835fn events_db_targets(db_path: &Path) -> [PathBuf; 3] {
836 ["", "-wal", "-shm"].map(|suffix| {
837 let mut name = db_path.as_os_str().to_os_string();
838 name.push(suffix);
839 PathBuf::from(name)
840 })
841}
842
843fn refuse_events_db_symlinks(db_path: &Path) -> anyhow::Result<()> {
853 for path in events_db_targets(db_path) {
854 match std::fs::symlink_metadata(&path) {
855 Ok(meta) if meta.file_type().is_symlink() => anyhow::bail!(
856 "refusing to serve events: {} is a symlink; the events database and its \
857 sidecars must be regular files",
858 path.display()
859 ),
860 Ok(_) => {}
861 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
862 Err(e) => anyhow::bail!(
863 "refusing to serve events: cannot inspect {}: {e}",
864 path.display()
865 ),
866 }
867 }
868 Ok(())
869}
870
871#[cfg(unix)]
872#[derive(Clone, Copy, Debug, PartialEq, Eq)]
873struct EventsFileIdentity {
874 device: u64,
875 inode: u64,
876}
877
878#[cfg(unix)]
879impl EventsFileIdentity {
880 fn from_metadata(metadata: &std::fs::Metadata) -> Self {
881 use std::os::unix::fs::MetadataExt;
882 Self {
883 device: metadata.dev(),
884 inode: metadata.ino(),
885 }
886 }
887}
888
889#[cfg(unix)]
892type EventsDbIdentities = [Option<EventsFileIdentity>; 3];
893
894#[cfg(unix)]
903fn ensure_events_db_parent_trusted(db_path: &Path) -> anyhow::Result<()> {
904 let parent = absolutize(db_path);
905 let parent = parent.parent().filter(|p| !p.as_os_str().is_empty());
906 match parent {
907 Some(dir) => crate::daemon::ensure_socket_dir_is_trusted(dir).map_err(|e| {
908 anyhow::anyhow!("refusing to serve events from an untrusted directory: {e}")
909 }),
910 None => Ok(()),
911 }
912}
913
914#[cfg(unix)]
920fn ensure_events_db_owner_only(db_path: &Path) -> anyhow::Result<()> {
921 use std::os::unix::fs::OpenOptionsExt;
922 refuse_events_db_symlinks(db_path)?;
923 if let Some(parent) = db_path.parent() {
924 std::fs::create_dir_all(parent)?;
925 }
926 ensure_events_db_parent_trusted(db_path)?;
929 match std::fs::OpenOptions::new()
930 .write(true)
931 .create_new(true)
932 .mode(0o600)
933 .open(db_path)
934 {
935 Ok(_created) => Ok(()),
937 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
938 Err(e) => Err(anyhow::anyhow!(
939 "refusing to serve events: cannot create {} owner-only: {e}",
940 db_path.display()
941 )),
942 }
943}
944
945#[cfg(unix)]
958fn harden_events_db_sidecars(db_path: &Path) -> anyhow::Result<EventsDbIdentities> {
959 use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
960 let mut identities = [None; 3];
961 for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
962 let file = match std::fs::OpenOptions::new()
969 .read(true)
970 .custom_flags(libc::O_NOFOLLOW)
971 .open(&path)
972 {
973 Ok(file) => file,
974 Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
975 Err(e) => {
976 anyhow::bail!(
977 "refusing to serve events: cannot open {} without following symlinks: \
978 {e}. The events database and its sidecars must be regular files.",
979 path.display()
980 )
981 }
982 };
983 let metadata = file.metadata()?;
986 if !metadata.file_type().is_file() {
987 anyhow::bail!(
988 "refusing to serve events: {} is not a regular file. The events \
989 database and its sidecars must be regular files.",
990 path.display()
991 );
992 }
993 file.set_permissions(std::fs::Permissions::from_mode(0o600))
994 .map_err(|e| {
995 anyhow::anyhow!(
996 "refusing to serve events: cannot chmod 0600 {}: {e}. The events \
997 database and its sidecars must be owner-only.",
998 path.display()
999 )
1000 })?;
1001 identities[index] = Some(EventsFileIdentity::from_metadata(&metadata));
1004 }
1005 Ok(identities)
1006}
1007
1008#[cfg(unix)]
1026fn verify_events_db_owner_only_unopened(
1027 db_path: &Path,
1028 before: &EventsDbIdentities,
1029) -> anyhow::Result<()> {
1030 use std::os::unix::fs::PermissionsExt;
1031 for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
1032 let metadata = match std::fs::symlink_metadata(&path) {
1033 Ok(metadata) => metadata,
1034 Err(e) if e.kind() == std::io::ErrorKind::NotFound && index != 0 => continue,
1035 Err(e) if e.kind() == std::io::ErrorKind::NotFound => anyhow::bail!(
1036 "refusing to serve events: main database {} disappeared after SQLite open; \
1037 refusing to serve an unlinked database",
1038 path.display()
1039 ),
1040 Err(e) => {
1041 anyhow::bail!(
1042 "refusing to serve events: cannot stat {}: {e}",
1043 path.display()
1044 )
1045 }
1046 };
1047 if !metadata.file_type().is_file() {
1048 anyhow::bail!(
1049 "refusing to serve events: {} is not a regular file. The events database \
1050 and its sidecars must be regular files.",
1051 path.display()
1052 );
1053 }
1054 let mode = metadata.permissions().mode() & 0o777;
1055 if mode & 0o077 != 0 {
1056 anyhow::bail!(
1057 "refusing to serve events: {} is mode {mode:03o}, not owner-only. The events \
1058 database and its sidecars must be owner-only.",
1059 path.display()
1060 );
1061 }
1062 if before[index]
1063 .is_some_and(|identity| identity != EventsFileIdentity::from_metadata(&metadata))
1064 {
1065 anyhow::bail!(
1066 "refusing to serve events: {} changed identity between pre-open hardening \
1067 and post-open verification",
1068 path.display()
1069 );
1070 }
1071 }
1072 Ok(())
1073}
1074
1075#[cfg(unix)]
1101pub async fn supervise_events_daemon(db_path: PathBuf, socket_path: PathBuf) {
1102 match standalone_daemon_policies(&db_path) {
1103 Ok((wal_ceiling, disk_guard, volume_lock_dir)) => {
1104 supervise_events_daemon_with_policies(
1105 db_path,
1106 socket_path,
1107 wal_ceiling,
1108 disk_guard,
1109 volume_lock_dir,
1110 )
1111 .await;
1112 }
1113 Err(error) => {
1114 tracing::warn!(%error, "invalid events daemon policy; events supervisor not started");
1115 }
1116 }
1117}
1118
1119#[cfg(unix)]
1121pub async fn supervise_events_daemon_with_policies(
1122 db_path: PathBuf,
1123 socket_path: PathBuf,
1124 wal_ceiling: WalCeilingPolicy,
1125 disk_guard: khive_db::EffectiveDiskGuardConfig,
1126 volume_lock_dir: PathBuf,
1127) {
1128 if let Err(error) = disk_guard.validate() {
1129 tracing::warn!(%error, "invalid disk policy; events supervisor not started");
1130 return;
1131 }
1132 if let Err(error) = wal_ceiling.validate_static(true, true, false) {
1133 tracing::warn!(error = %error, "invalid WAL ceiling; events supervisor not started");
1134 return;
1135 }
1136 const PROBE_INTERVAL: Duration = Duration::from_secs(15);
1137 let shutdown = crate::daemon::daemon_shutdown_token();
1138 let mut child: Option<std::process::Child> = None;
1139 let mut respawns: u64 = 0;
1140 loop {
1141 if let Some(c) = child.as_mut() {
1145 match c.try_wait() {
1146 Ok(Some(status)) => {
1147 tracing::info!(%status, "events daemon child exited");
1148 child = None;
1149 }
1150 Ok(None) => {}
1151 Err(error) => {
1152 tracing::warn!(error = %error, "cannot poll events daemon child; dropping handle");
1153 child = None;
1154 }
1155 }
1156 }
1157
1158 let reachable = match connect_verified(&socket_path).await {
1159 Ok(_stream) => true,
1160 Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
1161 tracing::warn!(
1162 socket = %socket_path.display(),
1163 error = %error,
1164 "events socket answered by a foreign uid; treating as unreachable"
1165 );
1166 false
1167 }
1168 Err(_) => false,
1169 };
1170 if !reachable && child.is_none() {
1171 match std::env::current_exe() {
1172 Ok(exe) => {
1173 let spawned = events_daemon_command(
1174 &exe,
1175 &db_path,
1176 &socket_path,
1177 wal_ceiling,
1178 disk_guard,
1179 &volume_lock_dir,
1180 )
1181 .spawn();
1182 match spawned {
1183 Ok(spawned_child) => {
1184 respawns += 1;
1185 tracing::info!(
1186 pid = spawned_child.id(),
1187 respawns,
1188 socket = %socket_path.display(),
1189 "spawned events daemon"
1190 );
1191 child = Some(spawned_child);
1192 }
1193 Err(error) => {
1194 tracing::warn!(error = %error, "failed to spawn events daemon");
1195 }
1196 }
1197 }
1198 Err(error) => {
1199 tracing::warn!(error = %error, "cannot resolve current executable for events daemon spawn");
1200 }
1201 }
1202 }
1203 tokio::select! {
1204 _ = shutdown.cancelled() => break,
1205 _ = tokio::time::sleep(PROBE_INTERVAL) => {}
1206 }
1207 }
1208 if let Some(mut c) = child.take() {
1209 let _ = c.kill();
1213 let _ = c.wait();
1214 tracing::info!("events daemon child stopped with supervisor shutdown");
1215 }
1216}
1217
1218#[cfg(unix)]
1219fn events_daemon_command(
1220 executable: &Path,
1221 db_path: &Path,
1222 socket_path: &Path,
1223 wal_ceiling: WalCeilingPolicy,
1224 disk_guard: khive_db::EffectiveDiskGuardConfig,
1225 volume_lock_dir: &Path,
1226) -> std::process::Command {
1227 let source = match wal_ceiling.source {
1228 khive_db::WalCeilingSource::BackendField => "backend_field",
1229 khive_db::WalCeilingSource::Environment => "environment",
1230 khive_db::WalCeilingSource::Default => "default",
1231 };
1232 let mut command = std::process::Command::new(executable);
1233 command
1234 .arg("events-daemon")
1235 .arg("--db")
1236 .arg(db_path)
1237 .arg("--socket")
1238 .arg(socket_path)
1239 .arg("--wal-ceiling-bytes")
1240 .arg(wal_ceiling.bytes.to_string())
1241 .arg("--wal-ceiling-source")
1242 .arg(source)
1243 .arg("--disk-reserve-bytes")
1244 .arg(disk_guard.reserve_bytes.to_string())
1245 .arg("--disk-guard-deadline-ms")
1246 .arg(disk_guard.guard_deadline_ms.to_string())
1247 .arg("--disk-reserve-source")
1248 .arg(disk_guard.reserve_source.as_str())
1249 .arg("--disk-deadline-source")
1250 .arg(disk_guard.deadline_source.as_str())
1251 .arg("--disk-legacy-environment-present")
1252 .arg(disk_guard.legacy_environment_present.to_string())
1253 .arg("--volume-lock-dir")
1254 .arg(volume_lock_dir)
1255 .stdin(std::process::Stdio::null())
1256 .stdout(std::process::Stdio::null())
1257 .stderr(std::process::Stdio::null());
1258 command
1259}
1260
1261#[cfg(unix)]
1280pub async fn run_events_daemon(db_path: &Path, socket_path: &Path) -> anyhow::Result<()> {
1281 let (wal_ceiling, disk_guard, volume_lock_dir) = standalone_daemon_policies(db_path)?;
1282 run_events_daemon_with_policies(
1283 db_path,
1284 socket_path,
1285 wal_ceiling,
1286 disk_guard,
1287 volume_lock_dir,
1288 )
1289 .await
1290}
1291
1292#[cfg(unix)]
1294pub async fn run_events_daemon_with_policies(
1295 db_path: &Path,
1296 socket_path: &Path,
1297 wal_ceiling: WalCeilingPolicy,
1298 disk_guard: khive_db::EffectiveDiskGuardConfig,
1299 volume_lock_dir: PathBuf,
1300) -> anyhow::Result<()> {
1301 disk_guard.validate()?;
1302 wal_ceiling
1303 .validate_static(true, true, false)
1304 .map_err(crate::error::RuntimeError::from)?;
1305 let db_path = &absolutize(db_path);
1308 let socket_path = &absolutize(socket_path);
1309 validate_events_socket_path(socket_path)?;
1310 if let Some(parent) = socket_path.parent() {
1314 std::fs::create_dir_all(parent)?;
1315 crate::daemon::ensure_socket_dir_is_trusted(parent)?;
1316 }
1317 let _guard = match acquire_events_daemon_guard_outcome(socket_path) {
1318 EventsDaemonGuardAcquisition::Held(guard) => guard,
1319 refusal => {
1320 tracing::info!(
1321 socket = %socket_path.display(),
1322 reason = ?refusal,
1323 "events daemon lock unavailable; exiting"
1324 );
1325 return Ok(());
1326 }
1327 };
1328 ensure_events_db_owner_only(db_path)?;
1329 let before_open = harden_events_db_sidecars(db_path)?;
1342 let backend = Arc::new(
1343 StorageBackend::sqlite_with_max_readers_and_policies(
1344 db_path,
1345 None,
1346 wal_ceiling,
1347 disk_guard,
1348 volume_lock_dir,
1349 )
1350 .map_err(crate::error::RuntimeError::from)?,
1351 );
1352 backend.events()?;
1354 verify_events_db_owner_only_unopened(db_path, &before_open)?;
1355
1356 if socket_path.exists() {
1357 std::fs::remove_file(socket_path)?;
1358 }
1359 let listener = UnixListener::bind(socket_path)?;
1360 {
1361 use std::os::unix::fs::PermissionsExt;
1362 if let Err(e) =
1366 std::fs::set_permissions(socket_path, std::fs::Permissions::from_mode(0o600))
1367 {
1368 drop(listener);
1369 let _ = std::fs::remove_file(socket_path);
1370 anyhow::bail!(
1371 "refusing to serve events: cannot chmod 0600 {}: {e}. The events socket must \
1372 be owner-only.",
1373 socket_path.display()
1374 );
1375 }
1376 }
1377 let daemon_euid = unsafe { libc::geteuid() };
1378 let stores: NamespaceStores = Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
1384 let connections = Arc::new(tokio::sync::Semaphore::new(MAX_EVENTS_CONNECTIONS));
1385 let frame_budget = Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES));
1386 tracing::info!(
1387 socket = %socket_path.display(),
1388 db = %db_path.display(),
1389 wal_ceiling_configured_bytes = wal_ceiling.bytes,
1390 wal_ceiling_effective_bytes = wal_ceiling.effective_bytes(backend.is_read_only()),
1391 wal_ceiling_source = ?wal_ceiling.source,
1392 "events daemon listening"
1393 );
1394
1395 loop {
1396 let (stream, _addr) = match listener.accept().await {
1397 Ok(pair) => pair,
1398 Err(error) => {
1399 tracing::warn!(error = %error, "events daemon accept failed");
1400 continue;
1401 }
1402 };
1403 match crate::daemon::peer_uid(&stream) {
1404 Ok(uid) if crate::daemon::uid_is_permitted(uid, daemon_euid) => {}
1405 Ok(uid) => {
1406 tracing::warn!(peer_uid = uid, "events daemon rejected foreign-uid peer");
1407 continue;
1408 }
1409 Err(error) => {
1410 tracing::warn!(error = %error, "events daemon could not read peer credentials");
1411 continue;
1412 }
1413 }
1414 let permit = match Arc::clone(&connections).try_acquire_owned() {
1415 Ok(permit) => permit,
1416 Err(_) => {
1417 tracing::warn!(
1420 cap = MAX_EVENTS_CONNECTIONS,
1421 "events daemon at connection cap; dropping new connection"
1422 );
1423 continue;
1424 }
1425 };
1426 let backend = Arc::clone(&backend);
1427 let stores = Arc::clone(&stores);
1428 let frame_budget = Arc::clone(&frame_budget);
1429 crate::daemon::spawn_named_tracked_task("events_connection", async move {
1430 serve_events_conn(stream, backend, stores, frame_budget).await;
1431 drop(permit);
1432 });
1433 }
1434}
1435
1436#[cfg(unix)]
1437#[path = "events_split_policy_entries.rs"]
1438mod policy_entries;
1439#[cfg(unix)]
1440use policy_entries::standalone_daemon_policies;
1441#[cfg(unix)]
1442pub use policy_entries::{
1443 run_events_daemon_with_wal_ceiling, supervise_events_daemon_with_wal_ceiling,
1444};
1445
1446#[cfg(all(test, unix))]
1447#[path = "events_wal_policy_tests.rs"]
1448mod wal_policy_tests;
1449
1450#[cfg(unix)]
1458async fn read_frame_budgeted(
1459 stream: &mut UnixStream,
1460 budget: &Arc<tokio::sync::Semaphore>,
1461) -> std::io::Result<(Vec<u8>, tokio::sync::OwnedSemaphorePermit)> {
1462 use tokio::io::AsyncReadExt;
1463 let mut len_buf = [0u8; 4];
1464 stream.read_exact(&mut len_buf).await?;
1465 let len = u32::from_be_bytes(len_buf) as usize;
1466 if len > crate::daemon::MAX_FRAME_BYTES {
1467 return Err(std::io::Error::new(
1468 std::io::ErrorKind::InvalidData,
1469 format!(
1470 "daemon frame of {len} bytes exceeds {} cap",
1471 crate::daemon::MAX_FRAME_BYTES
1472 ),
1473 ));
1474 }
1475 let permit = Arc::clone(budget)
1476 .acquire_many_owned(len as u32)
1477 .await
1478 .map_err(|_| std::io::Error::other("frame budget closed"))?;
1479 let mut buf = vec![0u8; len];
1480 stream.read_exact(&mut buf).await?;
1481 Ok((buf, permit))
1482}
1483
1484#[cfg(unix)]
1485type NamespaceStores =
1486 Arc<std::sync::Mutex<std::collections::HashMap<String, Arc<dyn EventStore>>>>;
1487
1488#[cfg(unix)]
1500fn namespace_store(
1501 backend: &StorageBackend,
1502 stores: &NamespaceStores,
1503 namespace: &str,
1504) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1505 namespace_store_with_cap(backend, stores, namespace, MAX_CACHED_NAMESPACE_STORES)
1506}
1507
1508#[cfg(unix)]
1509fn namespace_store_with_cap(
1510 backend: &StorageBackend,
1511 stores: &NamespaceStores,
1512 namespace: &str,
1513 cap: usize,
1514) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1515 let key = namespace.trim();
1516 if let Some(store) = stores
1517 .lock()
1518 .unwrap_or_else(std::sync::PoisonError::into_inner)
1519 .get(key)
1520 {
1521 return Ok(Arc::clone(store));
1522 }
1523 let store = backend.events_for_namespace(namespace)?;
1524 let mut map = stores
1525 .lock()
1526 .unwrap_or_else(std::sync::PoisonError::into_inner);
1527 if map.len() >= cap && !map.contains_key(key) {
1528 if let Some(victim) = map.keys().next().cloned() {
1536 map.remove(&victim);
1537 }
1538 }
1539 map.insert(key.to_string(), Arc::clone(&store));
1540 Ok(store)
1541}
1542
1543#[cfg(unix)]
1544async fn serve_events_conn(
1545 mut stream: UnixStream,
1546 backend: Arc<StorageBackend>,
1547 stores: NamespaceStores,
1548 frame_budget: Arc<tokio::sync::Semaphore>,
1549) {
1550 loop {
1551 let (payload, budget_permit) = match tokio::time::timeout(
1557 CONN_IO_TIMEOUT,
1558 read_frame_budgeted(&mut stream, &frame_budget),
1559 )
1560 .await
1561 {
1562 Ok(Ok(pair)) => pair,
1563 Ok(Err(_)) | Err(_) => return,
1565 };
1566 let response = match serde_json::from_slice::<EventsRequest>(&payload) {
1567 Ok(request) => dispatch_events_request(request, &backend, &stores).await,
1568 Err(error) => EventsResponse::Error {
1569 message: format!("events daemon could not parse request frame: {error}"),
1570 retryable: false,
1571 writer_task_failure: None,
1572 },
1573 };
1574 let bytes = match serde_json::to_vec(&response) {
1575 Ok(bytes) => bytes,
1576 Err(error) => {
1577 tracing::error!(error = %error, "events daemon response serialization failed");
1578 return;
1579 }
1580 };
1581 let bytes = if bytes.len() > crate::daemon::MAX_FRAME_BYTES {
1582 let refusal = EventsResponse::Error {
1586 message: format!(
1587 "response_frame_size_limit: {} response bytes exceed the {}-byte events IPC frame cap; request a narrower page",
1588 bytes.len(),
1589 crate::daemon::MAX_FRAME_BYTES
1590 ),
1591 retryable: false,
1592 writer_task_failure: None,
1593 };
1594 match serde_json::to_vec(&refusal) {
1595 Ok(bytes) => bytes,
1596 Err(error) => {
1597 tracing::error!(error = %error, "events daemon frame-size refusal serialization failed");
1598 return;
1599 }
1600 }
1601 } else {
1602 bytes
1603 };
1604 match tokio::time::timeout(CONN_IO_TIMEOUT, write_frame(&mut stream, &bytes)).await {
1607 Ok(Ok(())) => {}
1608 Ok(Err(_)) | Err(_) => return,
1609 }
1610 drop(budget_permit);
1614 }
1615}
1616
1617#[cfg(unix)]
1618async fn dispatch_events_request(
1619 request: EventsRequest,
1620 backend: &StorageBackend,
1621 stores: &NamespaceStores,
1622) -> EventsResponse {
1623 if request.protocol_version() != EVENTS_PROTOCOL_VERSION {
1624 return EventsResponse::Error {
1625 message: format!(
1626 "events protocol version mismatch: daemon speaks {}, client sent {}",
1627 EVENTS_PROTOCOL_VERSION,
1628 request.protocol_version()
1629 ),
1630 retryable: false,
1631 writer_task_failure: None,
1632 };
1633 }
1634 if let Err(error) = khive_types::Namespace::parse(request.namespace().trim()) {
1640 return EventsResponse::Error {
1641 message: format!("events request namespace rejected: {error}"),
1642 retryable: false,
1643 writer_task_failure: None,
1644 };
1645 }
1646 let store = match namespace_store(backend, stores, request.namespace()) {
1647 Ok(store) => store,
1648 Err(error) => {
1649 return EventsResponse::Error {
1650 message: format!("events store unavailable: {error}"),
1651 retryable: true,
1652 writer_task_failure: None,
1653 };
1654 }
1655 };
1656 match request {
1657 EventsRequest::AppendEvents { events, .. } => match store.append_events(events).await {
1658 Ok(summary) => EventsResponse::Appended { summary },
1659 Err(error) => storage_error_response(&error),
1660 },
1661 EventsRequest::AppendEventsIdempotent { events, .. } => {
1662 match store.append_events_idempotent(events).await {
1663 Ok(result) => EventsResponse::Idempotent { result },
1664 Err(error) => storage_error_response(&error),
1665 }
1666 }
1667 EventsRequest::GetEvent { id, .. } => match store.get_event(id).await {
1668 Ok(event) => EventsResponse::Event { event },
1669 Err(error) => storage_error_response(&error),
1670 },
1671 EventsRequest::QueryEvents { filter, page, .. } => {
1672 if page.limit > MAX_QUERY_EVENTS_PAGE_ROWS {
1673 return EventsResponse::Error {
1674 message: format!(
1675 "events query page limit {} exceeds the daemon cap of {} rows; \
1676 request narrower pages",
1677 page.limit, MAX_QUERY_EVENTS_PAGE_ROWS
1678 ),
1679 retryable: false,
1680 writer_task_failure: None,
1681 };
1682 }
1683 match store.query_events(filter, page).await {
1684 Ok(page) => EventsResponse::Pageful { page },
1685 Err(error) => storage_error_response(&error),
1686 }
1687 }
1688 EventsRequest::QueryEventPage { query, .. } => {
1689 if let Err(error) = crate::event_page::validate_page_query(&query) {
1690 return storage_error_response(&error);
1691 }
1692 match store.query_event_page(query).await {
1693 Ok(window) => EventsResponse::EventPageWindow { window },
1694 Err(error) => storage_error_response(&error),
1695 }
1696 }
1697 EventsRequest::CountEvents { filter, .. } => match store.count_events(filter).await {
1698 Ok(count) => EventsResponse::Count { count },
1699 Err(error) => storage_error_response(&error),
1700 },
1701 }
1702}
1703
1704#[cfg(unix)]
1705fn storage_error_response(error: &StorageError) -> EventsResponse {
1706 let (message, writer_task_failure) = match error {
1707 StorageError::WriterTaskRequestFailed {
1708 request_state,
1709 source,
1710 } => (
1711 source.to_string(),
1712 Some(WireWriterTaskFailure::RequestFailed {
1713 request_state: (*request_state).into(),
1714 }),
1715 ),
1716 StorageError::WriterTaskTerminated {
1717 request_state,
1718 sqlite_full_codes,
1719 } => (
1720 error.to_string(),
1721 Some(WireWriterTaskFailure::TaskTerminated {
1722 request_state: (*request_state).into(),
1723 sqlite_full_codes: *sqlite_full_codes,
1724 }),
1725 ),
1726 _ => (error.to_string(), None),
1727 };
1728 EventsResponse::Error {
1729 message,
1730 retryable: error.is_retryable(),
1735 writer_task_failure,
1736 }
1737}
1738
1739#[cfg(unix)]
1750async fn connect_verified(socket_path: &Path) -> std::io::Result<UnixStream> {
1751 let stream = UnixStream::connect(socket_path).await?;
1752 let own_euid = unsafe { libc::geteuid() } as u32;
1754 let peer = crate::daemon::peer_uid(&stream)?;
1755 if !crate::daemon::uid_is_permitted(peer, own_euid) {
1756 return Err(std::io::Error::new(
1757 std::io::ErrorKind::PermissionDenied,
1758 format!(
1759 "events socket {} answered by uid {peer}, not this process's uid {own_euid}; \
1760 refusing to exchange event frames with an unowned daemon",
1761 socket_path.display()
1762 ),
1763 ));
1764 }
1765 Ok(stream)
1766}
1767
1768#[cfg(unix)]
1772#[derive(Debug, Clone, Copy, Serialize)]
1773pub struct EventsForwardingMetrics {
1774 pub forwarded_batches: u64,
1775 pub forwarded_events: u64,
1776 pub dropped_batches: u64,
1777 pub dropped_events: u64,
1778 pub queued_bytes: usize,
1780}
1781
1782#[cfg(unix)]
1783#[derive(Debug, Default)]
1784struct ForwardingCounters {
1785 forwarded_batches: AtomicU64,
1786 forwarded_events: AtomicU64,
1787 dropped_batches: AtomicU64,
1788 dropped_events: AtomicU64,
1789 queued_bytes: AtomicUsize,
1790}
1791
1792#[cfg(unix)]
1795#[derive(Debug)]
1796struct AppendByteReservation {
1797 counters: Arc<ForwardingCounters>,
1798 bytes: usize,
1799}
1800
1801#[cfg(unix)]
1802impl Drop for AppendByteReservation {
1803 fn drop(&mut self) {
1804 self.counters
1805 .queued_bytes
1806 .fetch_sub(self.bytes, Ordering::AcqRel);
1807 }
1808}
1809
1810#[cfg(unix)]
1811#[derive(Debug)]
1812struct QueuedAppend {
1813 request: EventsRequest,
1814 count: u64,
1815 _reservation: AppendByteReservation,
1816}
1817
1818#[cfg(all(test, unix))]
1819impl QueuedAppend {
1820 fn unmetered(namespace: &str, events: Vec<Event>, counters: Arc<ForwardingCounters>) -> Self {
1821 let count = events.len() as u64;
1822 Self {
1823 request: EventsRequest::AppendEvents {
1824 protocol_version: EVENTS_PROTOCOL_VERSION,
1825 namespace: namespace.to_string(),
1826 events,
1827 },
1828 count,
1829 _reservation: AppendByteReservation { counters, bytes: 0 },
1830 }
1831 }
1832}
1833
1834#[cfg(unix)]
1837#[derive(Default)]
1838struct BoundedFrameCounter {
1839 bytes: usize,
1840 exceeded: bool,
1841}
1842
1843#[cfg(unix)]
1844impl std::io::Write for BoundedFrameCounter {
1845 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
1846 let next = self.bytes.saturating_add(buf.len());
1847 if next > crate::daemon::MAX_FRAME_BYTES {
1848 self.exceeded = true;
1849 return Err(std::io::Error::new(
1850 std::io::ErrorKind::InvalidData,
1851 "events append request exceeds IPC frame cap",
1852 ));
1853 }
1854 self.bytes = next;
1855 Ok(buf.len())
1856 }
1857
1858 fn flush(&mut self) -> std::io::Result<()> {
1859 Ok(())
1860 }
1861}
1862
1863#[cfg(unix)]
1864fn reserve_append_bytes(
1865 counters: &Arc<ForwardingCounters>,
1866 bytes: usize,
1867 budget: usize,
1868) -> Option<AppendByteReservation> {
1869 let mut used = counters.queued_bytes.load(Ordering::Acquire);
1870 loop {
1871 let next = used.checked_add(bytes)?;
1872 if next > budget {
1873 return None;
1874 }
1875 match counters.queued_bytes.compare_exchange_weak(
1876 used,
1877 next,
1878 Ordering::AcqRel,
1879 Ordering::Acquire,
1880 ) {
1881 Ok(_) => {
1882 return Some(AppendByteReservation {
1883 counters: Arc::clone(counters),
1884 bytes,
1885 });
1886 }
1887 Err(observed) => used = observed,
1888 }
1889 }
1890}
1891
1892#[cfg(unix)]
1900pub struct EventsSplitClient {
1901 socket_path: PathBuf,
1902 append_tx: tokio::sync::mpsc::Sender<QueuedAppend>,
1903 append_queue_byte_budget: usize,
1904 counters: Arc<ForwardingCounters>,
1905 outage_logged: Arc<AtomicBool>,
1908 preflight_store: Arc<dyn EventStore>,
1910}
1911
1912#[cfg(unix)]
1913impl std::fmt::Debug for EventsSplitClient {
1914 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1915 f.debug_struct("EventsSplitClient")
1916 .field("socket_path", &self.socket_path)
1917 .finish_non_exhaustive()
1918 }
1919}
1920
1921#[cfg(unix)]
1922impl EventsSplitClient {
1923 pub fn new(socket_path: PathBuf) -> crate::error::RuntimeResult<Arc<Self>> {
1925 Self::new_with_queue_depth(socket_path, DEFAULT_APPEND_QUEUE_BATCHES)
1926 }
1927
1928 pub fn new_with_queue_depth(
1931 socket_path: PathBuf,
1932 queue_depth: usize,
1933 ) -> crate::error::RuntimeResult<Arc<Self>> {
1934 Self::new_with_queue_depth_and_delivery_timeout(socket_path, queue_depth, REQUEST_TIMEOUT)
1935 }
1936
1937 fn new_with_queue_depth_and_delivery_timeout(
1941 socket_path: PathBuf,
1942 queue_depth: usize,
1943 delivery_timeout: Duration,
1944 ) -> crate::error::RuntimeResult<Arc<Self>> {
1945 Self::new_with_limits_and_delivery_timeout(
1946 socket_path,
1947 queue_depth,
1948 DEFAULT_APPEND_QUEUE_BYTES,
1949 delivery_timeout,
1950 )
1951 }
1952
1953 fn new_with_limits_and_delivery_timeout(
1954 socket_path: PathBuf,
1955 queue_depth: usize,
1956 byte_budget: usize,
1957 delivery_timeout: Duration,
1958 ) -> crate::error::RuntimeResult<Arc<Self>> {
1959 validate_events_socket_path(&socket_path)?;
1960 let preflight_backend = StorageBackend::memory()?;
1961 let preflight_store = preflight_backend.events()?;
1962 let (append_tx, append_rx) = tokio::sync::mpsc::channel::<QueuedAppend>(queue_depth.max(1));
1966 let counters = Arc::new(ForwardingCounters::default());
1967 let outage_logged = Arc::new(AtomicBool::new(false));
1968
1969 let client = Arc::new(Self {
1970 socket_path: socket_path.clone(),
1971 append_tx,
1972 append_queue_byte_budget: byte_budget,
1973 counters: Arc::clone(&counters),
1974 outage_logged: Arc::clone(&outage_logged),
1975 preflight_store,
1976 });
1977
1978 crate::daemon::spawn_named_tracked_task(
1979 "events_forwarder",
1980 run_forwarder(
1981 socket_path,
1982 append_rx,
1983 counters,
1984 outage_logged,
1985 delivery_timeout,
1986 crate::daemon::daemon_shutdown_token(),
1987 ),
1988 );
1989 Ok(client)
1990 }
1991
1992 pub fn metrics(&self) -> EventsForwardingMetrics {
1994 EventsForwardingMetrics {
1995 forwarded_batches: self.counters.forwarded_batches.load(Ordering::Relaxed),
1996 forwarded_events: self.counters.forwarded_events.load(Ordering::Relaxed),
1997 dropped_batches: self.counters.dropped_batches.load(Ordering::Relaxed),
1998 dropped_events: self.counters.dropped_events.load(Ordering::Relaxed),
1999 queued_bytes: self.counters.queued_bytes.load(Ordering::Acquire),
2000 }
2001 }
2002
2003 fn enqueue(&self, namespace: &str, events: Vec<Event>) {
2006 let count = events.len() as u64;
2007 let request = EventsRequest::AppendEvents {
2008 protocol_version: EVENTS_PROTOCOL_VERSION,
2009 namespace: namespace.to_string(),
2010 events,
2011 };
2012 let mut frame_size = BoundedFrameCounter::default();
2013 if let Err(error) = serde_json::to_writer(&mut frame_size, &request) {
2014 self.count_append_drop(
2015 count,
2016 if frame_size.exceeded {
2017 "frame cap"
2018 } else {
2019 "serialization"
2020 },
2021 );
2022 tracing::warn!(error = %error, "events append batch could not fit a wire frame");
2023 return;
2024 }
2025 let Some(reservation) = reserve_append_bytes(
2026 &self.counters,
2027 frame_size.bytes,
2028 self.append_queue_byte_budget,
2029 ) else {
2030 self.count_append_drop(count, "queue byte budget");
2031 return;
2032 };
2033 let batch = QueuedAppend {
2034 request,
2035 count,
2036 _reservation: reservation,
2037 };
2038 match self.append_tx.try_send(batch) {
2039 Ok(()) => {}
2040 Err(_) => {
2041 self.count_append_drop(count, "queue batch limit");
2043 }
2044 }
2045 }
2046
2047 fn count_append_drop(&self, count: u64, reason: &'static str) {
2048 self.counters
2049 .dropped_batches
2050 .fetch_add(1, Ordering::Relaxed);
2051 self.counters
2052 .dropped_events
2053 .fetch_add(count, Ordering::Relaxed);
2054 if !self.outage_logged.swap(true, Ordering::Relaxed) {
2055 tracing::warn!(
2056 dropped_events = count,
2057 reason,
2058 "events append queue rejected loss-tolerant batch"
2059 );
2060 }
2061 }
2062
2063 async fn round_trip(&self, request: &EventsRequest) -> StorageResult<EventsResponse> {
2068 let op = "events daemon round-trip";
2069 let payload = serde_json::to_vec(request).map_err(|error| StorageError::Serialization {
2070 capability: khive_storage::StorageCapability::Events,
2071 message: format!("events request serialization failed: {error}"),
2072 })?;
2073
2074 let attempt = async {
2075 let mut stream = connect_verified(&self.socket_path)
2076 .await
2077 .map_err(|error| ("connect", error))?;
2078 write_frame(&mut stream, &payload)
2079 .await
2080 .map_err(|error| ("write", error))?;
2081 read_frame(&mut stream)
2082 .await
2083 .map_err(|error| ("read", error))
2084 };
2085 let bytes = match tokio::time::timeout(REQUEST_TIMEOUT, attempt).await {
2086 Ok(Ok(bytes)) => bytes,
2087 Ok(Err(("read", error))) => {
2088 return Err(StorageError::Serialization {
2089 capability: khive_storage::StorageCapability::Events,
2090 message: format!(
2091 "events daemon closed or broke the response after connection at {}: {error}",
2092 self.socket_path.display()
2093 ),
2094 });
2095 }
2096 Ok(Err((stage, error))) => {
2097 return Err(StorageError::Pool {
2098 operation: op.into(),
2099 message: format!(
2100 "events daemon {stage} failed at {}: {error}",
2101 self.socket_path.display()
2102 ),
2103 });
2104 }
2105 Err(_elapsed) => {
2106 return Err(StorageError::Timeout {
2107 operation: op.into(),
2108 });
2109 }
2110 };
2111
2112 serde_json::from_slice::<EventsResponse>(&bytes).map_err(|error| {
2113 StorageError::Serialization {
2114 capability: khive_storage::StorageCapability::Events,
2115 message: format!("events response deserialization failed: {error}"),
2116 }
2117 })
2118 }
2119}
2120
2121#[cfg(unix)]
2128fn drain_dropped_queue(
2129 rx: &mut tokio::sync::mpsc::Receiver<QueuedAppend>,
2130 counters: &ForwardingCounters,
2131) {
2132 let mut dropped_batches = 0u64;
2133 let mut dropped_events = 0u64;
2134 while let Ok(batch) = rx.try_recv() {
2135 dropped_batches += 1;
2136 dropped_events += batch.count;
2137 }
2139 if dropped_batches > 0 {
2140 counters
2141 .dropped_batches
2142 .fetch_add(dropped_batches, Ordering::Relaxed);
2143 counters
2144 .dropped_events
2145 .fetch_add(dropped_events, Ordering::Relaxed);
2146 tracing::warn!(
2147 dropped_batches,
2148 dropped_events,
2149 "events forwarder shutting down; dropping queued loss-tolerant batches"
2150 );
2151 }
2152}
2153
2154#[cfg(unix)]
2159async fn run_forwarder(
2160 socket_path: PathBuf,
2161 mut rx: tokio::sync::mpsc::Receiver<QueuedAppend>,
2162 counters: Arc<ForwardingCounters>,
2163 outage_logged: Arc<AtomicBool>,
2164 delivery_timeout: Duration,
2165 shutdown: tokio_util::sync::CancellationToken,
2166) {
2167 let mut conn: Option<UnixStream> = None;
2175 loop {
2176 let batch = tokio::select! {
2177 _ = shutdown.cancelled() => {
2178 drain_dropped_queue(&mut rx, &counters);
2179 break;
2180 },
2181 received = rx.recv() => match received {
2182 Some(batch) => batch,
2183 None => break,
2184 },
2185 };
2186 let count = batch.count;
2187 let payload = match serde_json::to_vec(&batch.request) {
2188 Ok(payload) => payload,
2189 Err(error) => {
2190 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2191 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2192 tracing::error!(error = %error, "events forwarder serialization failed; batch dropped");
2193 continue;
2194 }
2195 };
2196
2197 let delivered = tokio::select! {
2205 _ = shutdown.cancelled() => {
2206 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2207 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2208 tracing::warn!(
2209 dropped_events = count,
2210 "events forwarder shutting down mid-delivery; in-flight batch dropped"
2211 );
2212 drain_dropped_queue(&mut rx, &counters);
2213 break;
2214 },
2215 outcome = tokio::time::timeout(
2216 delivery_timeout,
2217 deliver_batch(&socket_path, &mut conn, &payload),
2218 ) => match outcome {
2219 Ok(delivered) => delivered,
2220 Err(_elapsed) => {
2221 conn = None;
2222 false
2223 }
2224 },
2225 };
2226 if delivered {
2227 counters.forwarded_batches.fetch_add(1, Ordering::Relaxed);
2228 counters
2229 .forwarded_events
2230 .fetch_add(count, Ordering::Relaxed);
2231 if outage_logged.swap(false, Ordering::Relaxed) {
2232 tracing::info!("events daemon reachable again; forwarding resumed");
2233 }
2234 } else {
2235 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2236 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2237 if !outage_logged.swap(true, Ordering::Relaxed) {
2238 tracing::warn!(
2239 socket = %socket_path.display(),
2240 "events daemon unreachable; dropping loss-tolerant events until it returns"
2241 );
2242 }
2243 tokio::select! {
2244 _ = shutdown.cancelled() => {
2245 drain_dropped_queue(&mut rx, &counters);
2246 break;
2247 },
2248 _ = tokio::time::sleep(FORWARDER_BACKOFF) => {}
2249 }
2250 }
2251 }
2252}
2253
2254#[cfg(unix)]
2259async fn deliver_batch(socket_path: &Path, conn: &mut Option<UnixStream>, payload: &[u8]) -> bool {
2260 for _attempt in 0..2u8 {
2261 if conn.is_none() {
2262 match connect_verified(socket_path).await {
2263 Ok(stream) => *conn = Some(stream),
2264 Err(_) => return false,
2265 }
2266 }
2267 let stream = conn.as_mut().expect("connection populated above");
2268 let ok = async {
2269 write_frame(stream, payload).await?;
2270 let bytes = read_frame(stream).await?;
2271 std::io::Result::Ok(bytes)
2272 }
2273 .await;
2274 match ok {
2275 Ok(bytes) => {
2276 return !matches!(
2277 serde_json::from_slice::<EventsResponse>(&bytes),
2278 Ok(EventsResponse::Error { .. }) | Err(_)
2279 );
2280 }
2281 Err(_) => {
2282 *conn = None;
2285 }
2286 }
2287 }
2288 false
2289}
2290
2291#[cfg(unix)]
2300#[derive(Debug)]
2301pub struct ForwardingEventStore {
2302 namespace: String,
2303 client: Arc<EventsSplitClient>,
2304}
2305
2306#[cfg(unix)]
2307impl ForwardingEventStore {
2308 pub fn new(namespace: impl Into<String>, client: Arc<EventsSplitClient>) -> Self {
2309 Self {
2310 namespace: namespace.into(),
2311 client,
2312 }
2313 }
2314
2315 fn unexpected(&self, op: &'static str, response: EventsResponse) -> StorageError {
2316 match response {
2317 EventsResponse::Error {
2318 message,
2319 retryable,
2320 writer_task_failure,
2321 } => {
2322 if let Some(failure) = writer_task_failure {
2323 match failure {
2324 WireWriterTaskFailure::RequestFailed { request_state } => {
2325 let source = if retryable {
2326 StorageError::Pool {
2327 operation: op.into(),
2328 message,
2329 }
2330 } else {
2331 StorageError::InvalidInput {
2332 capability: khive_storage::StorageCapability::Events,
2333 operation: op.into(),
2334 message,
2335 }
2336 };
2337 StorageError::WriterTaskRequestFailed {
2338 request_state: request_state.into(),
2339 source: Box::new(source),
2340 }
2341 }
2342 WireWriterTaskFailure::TaskTerminated {
2343 request_state,
2344 sqlite_full_codes,
2345 } => StorageError::WriterTaskTerminated {
2346 request_state: request_state.into(),
2347 sqlite_full_codes,
2348 },
2349 }
2350 } else if retryable {
2351 StorageError::Pool {
2352 operation: op.into(),
2353 message,
2354 }
2355 } else {
2356 StorageError::InvalidInput {
2357 capability: khive_storage::StorageCapability::Events,
2358 operation: op.into(),
2359 message,
2360 }
2361 }
2362 }
2363 other => StorageError::Serialization {
2364 capability: khive_storage::StorageCapability::Events,
2365 message: format!("events daemon returned mismatched response for {op}: {other:?}"),
2366 },
2367 }
2368 }
2369}
2370
2371#[cfg(unix)]
2372#[async_trait]
2373impl EventStore for ForwardingEventStore {
2374 async fn append_event(&self, event: Event) -> StorageResult<()> {
2375 self.client.enqueue(&self.namespace, vec![event]);
2376 Ok(())
2377 }
2378
2379 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2380 let attempted = events.len() as u64;
2381 self.client.enqueue(&self.namespace, events);
2382 Ok(BatchWriteSummary {
2385 attempted,
2386 affected: attempted,
2387 ..BatchWriteSummary::default()
2388 })
2389 }
2390
2391 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2392 let request = EventsRequest::GetEvent {
2393 protocol_version: EVENTS_PROTOCOL_VERSION,
2394 namespace: self.namespace.clone(),
2395 id,
2396 };
2397 match self.client.round_trip(&request).await? {
2398 EventsResponse::Event { event } => Ok(event),
2399 other => Err(self.unexpected("get_event", other)),
2400 }
2401 }
2402
2403 async fn query_events(
2404 &self,
2405 filter: EventFilter,
2406 page: PageRequest,
2407 ) -> StorageResult<Page<Event>> {
2408 let request = EventsRequest::QueryEvents {
2409 protocol_version: EVENTS_PROTOCOL_VERSION,
2410 namespace: self.namespace.clone(),
2411 filter,
2412 page,
2413 };
2414 match self.client.round_trip(&request).await? {
2415 EventsResponse::Pageful { page } => Ok(page),
2416 other => Err(self.unexpected("query_events", other)),
2417 }
2418 }
2419
2420 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2421 let request = EventsRequest::CountEvents {
2422 protocol_version: EVENTS_PROTOCOL_VERSION,
2423 namespace: self.namespace.clone(),
2424 filter,
2425 };
2426 match self.client.round_trip(&request).await? {
2427 EventsResponse::Count { count } => Ok(count),
2428 other => Err(self.unexpected("count_events", other)),
2429 }
2430 }
2431
2432 async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
2433 crate::event_page::validate_page_query(&query)?;
2434 let request = EventsRequest::QueryEventPage {
2435 protocol_version: EVENTS_PROTOCOL_VERSION,
2436 namespace: self.namespace.clone(),
2437 query: query.clone(),
2438 };
2439 match self.client.round_trip(&request).await? {
2440 EventsResponse::EventPageWindow { window } => {
2441 crate::event_page::validate_window(&query, Some(&self.namespace), &window)?;
2442 Ok(window)
2443 }
2444 error @ EventsResponse::Error { .. } => Err(self.unexpected("query_event_page", error)),
2445 _ => Err(crate::event_page::page_error(
2446 "events page response kind mismatch",
2447 )),
2448 }
2449 }
2450
2451 fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2452 self.client.preflight_store.preflight_event(event)
2454 }
2455
2456 async fn append_events_idempotent(
2457 &self,
2458 events: Vec<Event>,
2459 ) -> StorageResult<IdempotentEventBatchResult> {
2460 let request = EventsRequest::AppendEventsIdempotent {
2461 protocol_version: EVENTS_PROTOCOL_VERSION,
2462 namespace: self.namespace.clone(),
2463 events,
2464 };
2465 match self.client.round_trip(&request).await? {
2466 EventsResponse::Idempotent { result } => Ok(result),
2467 other => Err(self.unexpected("append_events_idempotent", other)),
2468 }
2469 }
2470
2471 fn supports_idempotent_audit_batch(&self) -> bool {
2472 true
2473 }
2474}
2475
2476pub struct SplitEventStore {
2493 legacy: Arc<dyn EventStore>,
2494 lane: Arc<dyn EventStore>,
2495}
2496
2497impl std::fmt::Debug for SplitEventStore {
2498 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2499 f.debug_struct("SplitEventStore").finish_non_exhaustive()
2500 }
2501}
2502
2503impl SplitEventStore {
2504 pub const MAX_MERGED_WINDOW_ROWS: u64 = 100_000;
2512
2513 pub fn new(legacy: Arc<dyn EventStore>, lane: Arc<dyn EventStore>) -> Self {
2514 Self { legacy, lane }
2515 }
2516}
2517
2518#[async_trait]
2519impl EventStore for SplitEventStore {
2520 async fn append_event(&self, event: Event) -> StorageResult<()> {
2521 self.legacy.append_event(event).await
2522 }
2523
2524 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2525 self.legacy.append_events(events).await
2526 }
2527
2528 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2529 if let Some(event) = self.legacy.get_event(id).await? {
2532 return Ok(Some(event));
2533 }
2534 self.lane.get_event(id).await
2535 }
2536
2537 async fn query_events(
2538 &self,
2539 filter: EventFilter,
2540 page: PageRequest,
2541 ) -> StorageResult<Page<Event>> {
2542 let window = page.offset.saturating_add(u64::from(page.limit));
2556 if window > Self::MAX_MERGED_WINDOW_ROWS {
2557 return Err(StorageError::InvalidInput {
2558 capability: khive_storage::StorageCapability::Events,
2559 operation: "query_events".into(),
2560 message: format!(
2561 "offset+limit ({window}) exceeds the merged event plane's window bound \
2562 of {}; page deep windows with a `before` cursor at offset 0, or narrow \
2563 the filter",
2564 Self::MAX_MERGED_WINDOW_ROWS
2565 ),
2566 });
2567 }
2568 let prefix = PageRequest {
2569 offset: 0,
2570 limit: window.min(u64::from(u32::MAX)) as u32,
2571 };
2572 let legacy = self
2573 .legacy
2574 .query_events(filter.clone(), prefix.clone())
2575 .await?;
2576 let lane = self.lane.query_events(filter, prefix).await?;
2577 let total = match (legacy.total, lane.total) {
2578 (Some(a), Some(b)) => Some(a + b),
2579 _ => None,
2580 };
2581 let mut items = legacy.items;
2582 items.extend(lane.items);
2583 items.sort_by(|a, b| {
2584 b.created_at
2585 .cmp(&a.created_at)
2586 .then_with(|| b.id.cmp(&a.id))
2587 });
2588 let items = items
2589 .into_iter()
2590 .skip(page.offset as usize)
2591 .take(page.limit as usize)
2592 .collect();
2593 Ok(Page { items, total })
2594 }
2595
2596 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2597 let legacy = self.legacy.count_events(filter.clone()).await?;
2598 let lane = self.lane.count_events(filter).await?;
2599 Ok(legacy + lane)
2600 }
2601
2602 async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
2603 crate::event_page::split_page(self.legacy.as_ref(), self.lane.as_ref(), query).await
2604 }
2605
2606 fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2607 self.lane.preflight_event(event)
2610 }
2611
2612 async fn append_events_idempotent(
2613 &self,
2614 events: Vec<Event>,
2615 ) -> StorageResult<IdempotentEventBatchResult> {
2616 if events.is_empty() {
2626 return self.lane.append_events_idempotent(events).await;
2627 }
2628 let ids: Vec<Uuid> = events.iter().map(|event| event.id).collect();
2629 let probe_limit = u32::try_from(ids.len()).map_err(|_| StorageError::InvalidInput {
2630 capability: khive_storage::StorageCapability::Events,
2631 operation: "append_events_idempotent".into(),
2632 message: format!(
2633 "batch of {} rows exceeds the legacy probe window",
2634 ids.len()
2635 ),
2636 })?;
2637 let existing = self
2638 .legacy
2639 .query_events(
2640 EventFilter {
2641 ids,
2642 ..EventFilter::default()
2643 },
2644 PageRequest {
2645 offset: 0,
2646 limit: probe_limit,
2647 },
2648 )
2649 .await?;
2650 let legacy_resident: std::collections::HashSet<Uuid> =
2651 existing.items.iter().map(|event| event.id).collect();
2652 if legacy_resident.is_empty() {
2653 return self.lane.append_events_idempotent(events).await;
2654 }
2655 let mut legacy_rows = Vec::new();
2656 let mut lane_rows = Vec::new();
2657 let mut routed_to_legacy = Vec::with_capacity(events.len());
2658 for event in events {
2659 if legacy_resident.contains(&event.id) {
2660 routed_to_legacy.push(true);
2661 legacy_rows.push(event);
2662 } else {
2663 routed_to_legacy.push(false);
2664 lane_rows.push(event);
2665 }
2666 }
2667 let legacy_expected = legacy_rows.len();
2668 let lane_expected = lane_rows.len();
2669 let legacy_result = self.legacy.append_events_idempotent(legacy_rows).await?;
2670 let lane_result = if lane_expected == 0 {
2671 IdempotentEventBatchResult { rows: Vec::new() }
2672 } else {
2673 self.lane.append_events_idempotent(lane_rows).await?
2674 };
2675 if legacy_result.rows.len() != legacy_expected || lane_result.rows.len() != lane_expected {
2676 return Err(StorageError::Driver {
2677 capability: khive_storage::StorageCapability::Events,
2678 operation: "append_events_idempotent".into(),
2679 source: format!(
2680 "idempotent sub-batch result length mismatch: legacy {}/{}, lane {}/{}",
2681 legacy_result.rows.len(),
2682 legacy_expected,
2683 lane_result.rows.len(),
2684 lane_expected
2685 )
2686 .into(),
2687 });
2688 }
2689 let mut legacy_iter = legacy_result.rows.into_iter();
2690 let mut lane_iter = lane_result.rows.into_iter();
2691 let rows = routed_to_legacy
2692 .into_iter()
2693 .map(|to_legacy| {
2694 if to_legacy {
2695 legacy_iter.next().expect("length checked above")
2696 } else {
2697 lane_iter.next().expect("length checked above")
2698 }
2699 })
2700 .collect();
2701 Ok(IdempotentEventBatchResult { rows })
2702 }
2703
2704 fn supports_idempotent_audit_batch(&self) -> bool {
2705 self.lane.supports_idempotent_audit_batch() && self.legacy.supports_idempotent_audit_batch()
2709 }
2710}
2711
2712#[cfg(test)]
2713#[path = "events_split_sqlite_full_tests.rs"]
2714mod sqlite_full_tests;
2715
2716#[cfg(test)]
2717mod tests {
2718 use super::*;
2719
2720 #[test]
2721 fn fixture_registry_cleanup_preserves_other_roots() {
2722 let live_dir = tempfile::tempdir().unwrap();
2723 let _live_guard = TestRegistryGuard::new(live_dir.path());
2724 let live_path = live_dir.path().join("events.db");
2725 let live_backend = direct_backend_for(&live_path).unwrap();
2726
2727 let finished_dir = tempfile::tempdir().unwrap();
2728 let finished_backend = {
2729 let _finished_guard = TestRegistryGuard::new(finished_dir.path());
2730 let backend = direct_backend_for(&finished_dir.path().join("events.db")).unwrap();
2731 let weak = Arc::downgrade(&backend);
2732 drop(backend);
2733 assert!(
2734 weak.upgrade().is_some(),
2735 "registry retains the live fixture"
2736 );
2737 weak
2738 };
2739
2740 assert!(
2741 finished_backend.upgrade().is_none(),
2742 "FD_FIXTURE_REGISTRY_RELEASED"
2743 );
2744 let same_live_backend = direct_backend_for(&live_path).unwrap();
2745 assert!(
2746 Arc::ptr_eq(&live_backend, &same_live_backend),
2747 "FD_FIXTURE_OTHER_ROOT_RETAINED"
2748 );
2749 }
2750
2751 #[cfg(unix)]
2752 mod target_filter_tests {
2753 include!("events_split_target_tests.rs");
2754 }
2755
2756 include!("events_sidecar_path_tests.rs");
2761
2762 #[cfg(unix)]
2763 #[test]
2764 fn deep_symlink_chains_within_the_kernel_bound_derive_the_target_sidecar() {
2765 let dir = tempfile::tempdir().unwrap();
2770 let real = dir.path().join("real.db");
2771 let mut prev = real.clone();
2772 for i in 0..35 {
2773 let link = dir.path().join(format!("hop{i}.db"));
2774 std::os::unix::fs::symlink(&prev, &link).unwrap();
2775 prev = link;
2776 }
2777 assert_eq!(
2778 events_db_path_beside(&prev),
2779 events_db_path_beside(&real),
2780 "a 35-hop dangling chain must derive the target's sidecar"
2781 );
2782 assert_ne!(
2784 events_db_path_beside(&dir.path().join("unrelated.db")),
2785 events_db_path_beside(&real)
2786 );
2787 }
2788
2789 include!("events_sidecar_hardening_tests.rs");
2790
2791 #[cfg(unix)]
2792 fn write_owner_only_test_file(path: &Path) {
2793 use std::os::unix::fs::PermissionsExt;
2794 std::fs::write(path, b"").unwrap();
2795 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).unwrap();
2796 }
2797
2798 #[cfg(unix)]
2799 #[test]
2800 fn unopened_check_refuses_same_mode_regular_file_replacement() {
2801 for target_index in 0..3 {
2805 let dir = tempfile::tempdir().unwrap();
2806 let db = dir.path().join("events.db");
2807 let targets = events_db_targets(&db);
2808 for path in &targets {
2809 write_owner_only_test_file(path);
2810 }
2811 let before_open = harden_events_db_sidecars(&db).unwrap();
2812 let replacement = dir.path().join("replacement");
2813 write_owner_only_test_file(&replacement);
2814 let replacement_id = EventsFileIdentity::from_metadata(
2815 &std::fs::symlink_metadata(&replacement).unwrap(),
2816 );
2817 assert_ne!(before_open[target_index], Some(replacement_id));
2818 let _connection = rusqlite::Connection::open(&db).unwrap();
2821 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2822 std::fs::rename(&replacement, &targets[target_index]).unwrap();
2823 let error = verify_events_db_owner_only_unopened(&db, &before_open)
2824 .unwrap_err()
2825 .to_string();
2826 assert!(error.contains("changed identity"), "{error}");
2827 assert!(
2828 error.contains(&targets[target_index].display().to_string()),
2829 "{error}"
2830 );
2831 }
2832 }
2833
2834 #[cfg(unix)]
2835 #[test]
2836 fn unopened_check_accepts_sidecars_missing_at_check_time() {
2837 let dir = tempfile::tempdir().unwrap();
2838 let db = dir.path().join("events.db");
2839 let targets = events_db_targets(&db);
2840 for path in &targets {
2841 write_owner_only_test_file(path);
2842 }
2843 let before_open = harden_events_db_sidecars(&db).unwrap();
2844 let _connection = rusqlite::Connection::open(&db).unwrap();
2845 for path in targets.iter().skip(1) {
2846 std::fs::remove_file(path).unwrap();
2847 verify_events_db_owner_only_unopened(&db, &before_open)
2848 .expect("SQLite sidecars may disappear without losing the main database");
2849 }
2850 }
2851
2852 #[cfg(unix)]
2853 #[test]
2854 fn unopened_check_refuses_main_database_removed_after_sqlite_open() {
2855 let dir = tempfile::tempdir().unwrap();
2856 let db = dir.path().join("events.db");
2857 ensure_events_db_owner_only(&db).unwrap();
2858 let before_open = harden_events_db_sidecars(&db).unwrap();
2859 let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2860 backend.events().unwrap();
2861 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2862 std::fs::remove_file(&db).unwrap();
2863 let error = verify_events_db_owner_only_unopened(&db, &before_open)
2866 .unwrap_err()
2867 .to_string();
2868 assert!(
2869 error.contains("main database") && error.contains("disappeared"),
2870 "{error}"
2871 );
2872 }
2873
2874 #[cfg(unix)]
2875 #[test]
2876 fn unopened_check_accepts_fresh_sqlite_sidecars_and_unchanged_main() {
2877 let dir = tempfile::tempdir().unwrap();
2878 let db = dir.path().join("new.events.db");
2879 assert!(!db.exists());
2880 ensure_events_db_owner_only(&db).unwrap();
2881 let before_open = harden_events_db_sidecars(&db).unwrap();
2882 assert!(before_open[0].is_some());
2883 assert_eq!(&before_open[1..], &[None, None]);
2884 let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2885 backend.events().unwrap();
2886 for path in events_db_targets(&db).iter().skip(1) {
2887 assert!(path.exists(), "SQLite creates {}", path.display());
2888 }
2889 verify_events_db_owner_only_unopened(&db, &before_open)
2890 .expect("new SQLite sidecars have no pre-open identity to contradict");
2891 }
2892
2893 #[test]
2894 fn socket_derives_beside_the_sidecar() {
2895 let db = events_db_path_beside(Path::new("/data/khive.db"));
2896 let socket = events_socket_path_beside(&db);
2897 assert!(socket.ends_with("khive.db.events.sock"), "got {socket:?}");
2898 }
2899
2900 #[test]
2901 fn bare_relative_main_db_yields_an_absolute_sidecar() {
2902 let path = events_db_path_beside(Path::new("khive.db"));
2906 assert!(path.is_absolute(), "got relative {path:?}");
2907 assert!(path.ends_with("khive.db.events.db"), "got {path:?}");
2908 }
2909
2910 #[cfg(unix)]
2911 #[tokio::test]
2912 async fn over_cap_query_page_limit_is_refused_before_materialization() {
2913 let backend = Arc::new(StorageBackend::memory().unwrap());
2918 let stores: NamespaceStores =
2919 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
2920 let request = |limit: u32| EventsRequest::QueryEvents {
2921 protocol_version: EVENTS_PROTOCOL_VERSION,
2922 namespace: "local".to_string(),
2923 filter: EventFilter::default(),
2924 page: PageRequest { offset: 0, limit },
2925 };
2926 let refused =
2927 dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS + 1), &backend, &stores)
2928 .await;
2929 match refused {
2930 EventsResponse::Error {
2931 message,
2932 retryable,
2933 writer_task_failure,
2934 } => {
2935 assert!(
2936 message.contains(&MAX_QUERY_EVENTS_PAGE_ROWS.to_string()),
2937 "refusal must name the cap: {message}"
2938 );
2939 assert!(!retryable, "an over-cap page is not transient");
2940 assert!(writer_task_failure.is_none());
2941 }
2942 other => panic!("over-cap query must be refused, got {other:?}"),
2943 }
2944 let at_cap =
2947 dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS), &backend, &stores).await;
2948 assert!(
2949 matches!(at_cap, EventsResponse::Pageful { .. }),
2950 "at-cap query must reach the store, got {at_cap:?}"
2951 );
2952 }
2953
2954 #[cfg(unix)]
2955 #[test]
2956 fn side_effects_unknown_state_crosses_the_wire_as_terminal_failure() {
2957 let response = storage_error_response(&StorageError::writer_task_terminated(
2958 khive_storage::WriterTaskRequestState::SideEffectsUnknown,
2959 ));
2960 assert!(
2961 matches!(
2962 response,
2963 EventsResponse::Error {
2964 writer_task_failure: Some(WireWriterTaskFailure::TaskTerminated {
2965 request_state: WireWriterTaskState::SideEffectsUnknown,
2966 sqlite_full_codes: None,
2967 }),
2968 ..
2969 }
2970 ),
2971 "an unknown-commit-state termination must carry its state on the wire"
2972 );
2973 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
2975 assert!(matches!(
2976 busy,
2977 EventsResponse::Error {
2978 writer_task_failure: None,
2979 ..
2980 }
2981 ));
2982 }
2983
2984 #[cfg(unix)]
2985 #[test]
2986 fn proven_rollback_crosses_the_wire_without_claiming_task_termination() {
2987 let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
2988 request_state: khive_storage::WriterTaskRequestState::TransactionRolledBack,
2989 source: Box::new(StorageError::Pool {
2990 operation: "writer_task_commit".into(),
2991 message: "commit refused".into(),
2992 }),
2993 });
2994 assert!(matches!(
2995 response,
2996 EventsResponse::Error {
2997 writer_task_failure: Some(WireWriterTaskFailure::RequestFailed {
2998 request_state: WireWriterTaskState::TransactionRolledBack,
2999 }),
3000 ..
3001 }
3002 ));
3003 }
3004
3005 #[test]
3006 fn error_frames_without_writer_failure_still_parse() {
3007 let bytes = br#"{"kind":"error","message":"boom","retryable":true}"#;
3011 let parsed: EventsResponse = serde_json::from_slice(bytes).expect("stateless frame parses");
3012 assert!(matches!(
3013 parsed,
3014 EventsResponse::Error {
3015 retryable: true,
3016 writer_task_failure: None,
3017 ..
3018 }
3019 ));
3020 }
3021
3022 fn split_retry_event(verb: &str) -> Event {
3023 Event::new(
3024 "test",
3025 verb,
3026 khive_types::EventKind::RecallExecuted,
3027 khive_types::SubstrateKind::Note,
3028 "agent:test",
3029 )
3030 }
3031
3032 fn store_pair(dir: &Path) -> (Arc<dyn EventStore>, Arc<dyn EventStore>) {
3033 let legacy = direct_backend_for(&dir.join("legacy.db"))
3034 .expect("legacy backend")
3035 .events_for_namespace("test")
3036 .expect("legacy store");
3037 let lane = direct_backend_for(&dir.join("lane.db"))
3038 .expect("lane backend")
3039 .events_for_namespace("test")
3040 .expect("lane store");
3041 (legacy, lane)
3042 }
3043
3044 #[tokio::test]
3049 async fn idempotent_retry_of_legacy_resident_rows_does_not_duplicate() {
3050 use khive_storage::event::EventAppendDisposition;
3051
3052 let dir = tempfile::tempdir().unwrap();
3053 let _registry_guard = TestRegistryGuard::new(dir.path());
3054 let (legacy, lane) = store_pair(dir.path());
3055
3056 let resident = split_retry_event("recall");
3057 legacy
3058 .append_events_idempotent(vec![resident.clone()])
3059 .await
3060 .expect("pre-cutover landing");
3061
3062 let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3063 let fresh = split_retry_event("search");
3064 let result = split
3065 .append_events_idempotent(vec![resident.clone(), fresh.clone()])
3066 .await
3067 .expect("mixed retry batch");
3068 assert_eq!(
3069 result.rows,
3070 vec![
3071 EventAppendDisposition::AlreadyPresentIdentical,
3072 EventAppendDisposition::Inserted,
3073 ],
3074 "input order must be preserved across the two sub-batches"
3075 );
3076 assert_eq!(
3077 lane.count_events(EventFilter::default()).await.unwrap(),
3078 1,
3079 "the legacy-resident row must not reach the lane"
3080 );
3081 assert_eq!(
3082 split.count_events(EventFilter::default()).await.unwrap(),
3083 2,
3084 "merged count must not double-count the retried id"
3085 );
3086
3087 let mut mutated = resident.clone();
3090 mutated.verb = "other".to_string();
3091 let conflict = split
3092 .append_events_idempotent(vec![mutated])
3093 .await
3094 .expect("conflicting retry");
3095 assert_eq!(
3096 conflict.rows,
3097 vec![EventAppendDisposition::IdentityConflict]
3098 );
3099 assert_eq!(split.count_events(EventFilter::default()).await.unwrap(), 2);
3100 }
3101
3102 #[tokio::test]
3105 async fn merged_offset_window_is_bounded() {
3106 let dir = tempfile::tempdir().unwrap();
3107 let _registry_guard = TestRegistryGuard::new(dir.path());
3108 let (legacy, lane) = store_pair(dir.path());
3109 let split = SplitEventStore::new(legacy, lane);
3110
3111 let err = split
3112 .query_events(
3113 EventFilter::default(),
3114 PageRequest {
3115 offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS,
3116 limit: 1,
3117 },
3118 )
3119 .await
3120 .expect_err("a window past the bound must be refused");
3121 assert!(
3122 matches!(err, StorageError::InvalidInput { .. }),
3123 "got {err:?}"
3124 );
3125 assert!(
3126 err.to_string().contains("before"),
3127 "the refusal must name the cursor remedy: {err}"
3128 );
3129
3130 let at_bound = split
3131 .query_events(
3132 EventFilter::default(),
3133 PageRequest {
3134 offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS - 1,
3135 limit: 1,
3136 },
3137 )
3138 .await
3139 .expect("a window at the bound is admitted");
3140 assert!(at_bound.items.is_empty());
3141 }
3142
3143 #[cfg(unix)]
3144 #[tokio::test]
3145 async fn stateless_non_retryable_error_maps_terminal_not_writer_terminated() {
3146 let client = EventsSplitClient::new(std::path::PathBuf::from("/tmp/never-bound.sock"))
3153 .expect("client builds");
3154 let store = ForwardingEventStore::new("test", client);
3155 let bytes = br#"{"kind":"error","message":"refused","retryable":false}"#;
3156 let parsed: EventsResponse = serde_json::from_slice(bytes).expect("frame parses");
3157 assert!(matches!(
3158 store.unexpected("append", parsed),
3159 StorageError::InvalidInput { .. }
3160 ));
3161 let with_state = br#"{"kind":"error","message":"died","retryable":false,"writer_task_failure":{"kind":"task_terminated","request_state":"side_effects_unknown"}}"#;
3165 let parsed: EventsResponse =
3166 serde_json::from_slice(with_state).expect("stateful frame parses");
3167 assert!(matches!(
3168 store.unexpected("append", parsed),
3169 StorageError::WriterTaskTerminated { .. }
3170 ));
3171 }
3172
3173 #[cfg(unix)]
3174 #[tokio::test]
3175 async fn writer_task_states_cross_the_wire_verbatim() {
3176 use khive_storage::WriterTaskRequestState as S;
3182 let dir = tempfile::tempdir().unwrap();
3183 let client =
3184 EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3185 let store = ForwardingEventStore::new("test", client);
3186 for state in [
3187 S::NotStarted,
3188 S::TransactionRolledBack,
3189 S::SideEffectsUnknown,
3190 ] {
3191 let response = storage_error_response(&StorageError::writer_task_terminated(state));
3192 let bytes = serde_json::to_vec(&response).unwrap();
3194 let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3195 let err = store.unexpected("append_events_idempotent", parsed);
3196 assert!(
3197 matches!(
3198 err,
3199 StorageError::WriterTaskTerminated { request_state, .. } if request_state == state
3200 ),
3201 "state {state:?} did not survive the socket: got {err:?}"
3202 );
3203 }
3204 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3206 assert!(matches!(
3207 busy,
3208 EventsResponse::Error {
3209 writer_task_failure: None,
3210 ..
3211 }
3212 ));
3213 }
3214
3215 #[cfg(unix)]
3216 #[tokio::test]
3217 async fn proven_rollback_round_trips_as_non_terminal_request_failure() {
3218 use khive_storage::WriterTaskRequestState as S;
3219
3220 let dir = tempfile::tempdir().unwrap();
3221 let client =
3222 EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3223 let store = ForwardingEventStore::new("test", client);
3224 let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
3225 request_state: S::TransactionRolledBack,
3226 source: Box::new(StorageError::Pool {
3227 operation: "writer_task_commit".into(),
3228 message: "commit refused".into(),
3229 }),
3230 });
3231 let bytes = serde_json::to_vec(&response).unwrap();
3232 let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3233 let err = store.unexpected("append_events_idempotent", parsed);
3234 assert!(
3235 matches!(
3236 err,
3237 StorageError::WriterTaskRequestFailed {
3238 request_state: S::TransactionRolledBack,
3239 ..
3240 }
3241 ),
3242 "proven rollback must not reconstruct as a terminal writer: {err:?}"
3243 );
3244 }
3245
3246 #[cfg(unix)]
3247 #[test]
3248 fn direct_backend_hardens_preexisting_db_and_sidecars() {
3249 use std::os::unix::fs::PermissionsExt;
3250 let dir = tempfile::tempdir().unwrap();
3251 let _registry_guard = TestRegistryGuard::new(dir.path());
3252 let db = dir.path().join("pre-existing.events.db");
3253 let wal = dir.path().join("pre-existing.events.db-wal");
3254 std::fs::write(&db, b"").unwrap();
3255 std::fs::write(&wal, b"").unwrap();
3256 for path in [&db, &wal] {
3257 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o644)).unwrap();
3258 }
3259 assert_eq!(
3261 std::fs::metadata(&db).unwrap().permissions().mode() & 0o777,
3262 0o644
3263 );
3264 direct_backend_for(&db).expect("writable open succeeds");
3265 for path in [&db, &wal] {
3266 assert_eq!(
3267 std::fs::metadata(path).unwrap().permissions().mode() & 0o777,
3268 0o600,
3269 "pre-existing {} must be tightened to owner-only",
3270 path.display()
3271 );
3272 }
3273 }
3274
3275 #[cfg(unix)]
3276 #[test]
3277 fn namespace_store_cache_is_bounded_and_trim_normalized() {
3278 let backend = StorageBackend::memory().unwrap();
3279 let stores: NamespaceStores =
3280 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3281 namespace_store_with_cap(&backend, &stores, "alpha", 1).expect("first store");
3284 namespace_store_with_cap(&backend, &stores, "beta", 1).expect("second store evicts");
3285 {
3286 let map = stores.lock().unwrap();
3287 assert_eq!(map.len(), 1, "cache must not grow past its cap");
3288 assert!(
3289 map.contains_key("beta"),
3290 "the newest namespace must be admitted at the cap"
3291 );
3292 }
3293 namespace_store_with_cap(&backend, &stores, "beta", 1).expect("cached beta");
3296 namespace_store_with_cap(&backend, &stores, " beta ", 1).expect("trimmed spelling");
3297 {
3298 let map = stores.lock().unwrap();
3299 assert_eq!(
3300 map.len(),
3301 1,
3302 "spellings of one namespace must share one entry"
3303 );
3304 assert!(map.contains_key("beta"), "trimmed key is the cache key");
3305 }
3306 }
3307
3308 #[cfg(unix)]
3309 #[tokio::test]
3310 async fn frame_budget_admits_before_allocating_and_releases_after() {
3311 let (mut client, mut server) = UnixStream::pair().expect("socketpair");
3312 let budget = Arc::new(tokio::sync::Semaphore::new(1024));
3313 let payload = vec![7u8; 100];
3314 write_frame(&mut client, &payload).await.expect("write");
3315 let (bytes, permit) = read_frame_budgeted(&mut server, &budget)
3316 .await
3317 .expect("budgeted read");
3318 assert_eq!(bytes, payload);
3319 assert_eq!(budget.available_permits(), 1024 - 100);
3321 drop(permit);
3323 assert_eq!(budget.available_permits(), 1024);
3324 let mut oversized = Vec::from(((crate::daemon::MAX_FRAME_BYTES + 1) as u32).to_be_bytes());
3327 oversized.extend_from_slice(&[0u8; 8]);
3328 use tokio::io::AsyncWriteExt;
3329 client.write_all(&oversized).await.expect("raw prefix");
3330 let err = read_frame_budgeted(&mut server, &budget)
3331 .await
3332 .expect_err("oversized declaration must refuse");
3333 assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
3334 assert_eq!(
3335 budget.available_permits(),
3336 1024,
3337 "refusal must not consume budget"
3338 );
3339 }
3340
3341 #[cfg(unix)]
3342 #[tokio::test]
3343 async fn dispatch_rejects_invalid_wire_namespace() {
3344 let backend = Arc::new(StorageBackend::memory().unwrap());
3345 let stores: NamespaceStores =
3346 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3347 let huge = "n".repeat(64 * 1024);
3351 let response = dispatch_events_request(
3352 EventsRequest::CountEvents {
3353 protocol_version: EVENTS_PROTOCOL_VERSION,
3354 namespace: huge,
3355 filter: Default::default(),
3356 },
3357 &backend,
3358 &stores,
3359 )
3360 .await;
3361 assert!(
3362 matches!(
3363 &response,
3364 EventsResponse::Error {
3365 retryable: false,
3366 writer_task_failure: None,
3367 ..
3368 }
3369 ),
3370 "oversized namespace must be a typed refusal, got {response:?}"
3371 );
3372 assert_eq!(
3373 stores.lock().unwrap().len(),
3374 0,
3375 "a rejected namespace must never enter the cache"
3376 );
3377 let ok = dispatch_events_request(
3379 EventsRequest::CountEvents {
3380 protocol_version: EVENTS_PROTOCOL_VERSION,
3381 namespace: "local".to_string(),
3382 filter: Default::default(),
3383 },
3384 &backend,
3385 &stores,
3386 )
3387 .await;
3388 assert!(
3389 matches!(ok, EventsResponse::Count { .. }),
3390 "valid namespace must dispatch, got {ok:?}"
3391 );
3392 }
3393
3394 include!("events_split_shutdown_tests.rs");
3395
3396 #[cfg(unix)]
3397 #[tokio::test]
3398 async fn forwarder_abandons_hung_delivery() {
3399 let dir = tempfile::tempdir().unwrap();
3403 let socket = dir.path().join("hung.sock");
3404 let listener = tokio::net::UnixListener::bind(&socket).unwrap();
3405 let _server = tokio::spawn(async move {
3407 let mut held = Vec::new();
3408 loop {
3409 if let Ok((stream, _)) = listener.accept().await {
3410 held.push(stream);
3411 }
3412 }
3413 });
3414 let client = EventsSplitClient::new_with_queue_depth_and_delivery_timeout(
3415 socket.clone(),
3416 4,
3417 Duration::from_millis(100),
3418 )
3419 .expect("client builds");
3420 let event = Event::new(
3421 "test",
3422 "noop",
3423 khive_types::EventKind::Audit,
3424 khive_types::SubstrateKind::Event,
3425 "tester",
3426 );
3427 client.enqueue("test", vec![event]);
3428 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
3429 loop {
3430 if client.metrics().dropped_batches >= 1 {
3431 break;
3432 }
3433 assert!(
3434 tokio::time::Instant::now() < deadline,
3435 "forwarder never abandoned the hung delivery: {:?}",
3436 client.metrics()
3437 );
3438 tokio::time::sleep(Duration::from_millis(20)).await;
3439 }
3440 }
3441
3442 #[cfg(unix)]
3443 #[test]
3444 fn wire_retryability_follows_the_storage_classifier() {
3445 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3446 assert!(
3447 matches!(
3448 busy,
3449 EventsResponse::Error {
3450 retryable: true,
3451 ..
3452 }
3453 ),
3454 "transient writer contention must stay retryable across the socket"
3455 );
3456 let terminated = storage_error_response(&StorageError::writer_task_terminated(
3459 khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3460 ));
3461 assert!(
3462 matches!(
3463 terminated,
3464 EventsResponse::Error {
3465 retryable: false,
3466 ..
3467 }
3468 ),
3469 "a terminated writer must not be reported transient"
3470 );
3471 }
3472
3473 use khive_storage::event::EventAppendDisposition;
3474 use khive_types::{EventKind, SubstrateKind};
3475
3476 fn test_event(namespace: &str) -> Event {
3477 Event::new(
3478 namespace,
3479 "test.verb",
3480 EventKind::Audit,
3481 SubstrateKind::Event,
3482 "actor:test",
3483 )
3484 }
3485
3486 #[cfg(unix)]
3488 async fn boot_daemon(dir: &tempfile::TempDir) -> (PathBuf, PathBuf) {
3489 let db = dir.path().join("events.db");
3490 let socket = dir.path().join("events.sock");
3491 let (db_clone, socket_clone) = (db.clone(), socket.clone());
3492 tokio::spawn(async move {
3493 let _ = run_events_daemon(&db_clone, &socket_clone).await;
3494 });
3495 for _ in 0..100 {
3496 if UnixStream::connect(&socket).await.is_ok() {
3497 return (db, socket);
3498 }
3499 tokio::time::sleep(Duration::from_millis(20)).await;
3500 }
3501 panic!("events daemon did not come up on {}", socket.display());
3502 }
3503
3504 #[tokio::test]
3509 async fn split_store_routes_plain_to_legacy_idempotent_to_lane_and_merges_reads() {
3510 let dir = tempfile::tempdir().expect("tempdir");
3511 let _registry_guard = TestRegistryGuard::new(dir.path());
3512 let legacy_backend =
3513 direct_backend_for(&dir.path().join("legacy.db")).expect("legacy backend");
3514 let lane_backend = direct_backend_for(&dir.path().join("lane.db")).expect("lane backend");
3515 let legacy = legacy_backend
3516 .events_for_namespace("local")
3517 .expect("legacy store");
3518 let lane = lane_backend
3519 .events_for_namespace("local")
3520 .expect("lane store");
3521 let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3522
3523 let plain = test_event("local");
3524 let plain_id = plain.id;
3525 split.append_event(plain).await.expect("plain append");
3526
3527 let audit = test_event("local");
3528 let audit_id = audit.id;
3529 let result = split
3530 .append_events_idempotent(vec![audit])
3531 .await
3532 .expect("idempotent append");
3533 assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3534
3535 assert!(legacy
3537 .get_event(plain_id)
3538 .await
3539 .expect("legacy get")
3540 .is_some());
3541 assert!(lane.get_event(plain_id).await.expect("lane get").is_none());
3542 assert!(legacy
3543 .get_event(audit_id)
3544 .await
3545 .expect("legacy get")
3546 .is_none());
3547 assert!(lane.get_event(audit_id).await.expect("lane get").is_some());
3548
3549 assert!(split
3551 .get_event(plain_id)
3552 .await
3553 .expect("split get")
3554 .is_some());
3555 assert!(split
3556 .get_event(audit_id)
3557 .await
3558 .expect("split get")
3559 .is_some());
3560 assert_eq!(
3561 split
3562 .count_events(EventFilter::default())
3563 .await
3564 .expect("split count"),
3565 2
3566 );
3567 let page = split
3568 .query_events(
3569 EventFilter::default(),
3570 PageRequest {
3571 offset: 0,
3572 limit: 10,
3573 },
3574 )
3575 .await
3576 .expect("split query");
3577 let ids: Vec<Uuid> = page.items.iter().map(|e| e.id).collect();
3578 assert!(ids.contains(&plain_id) && ids.contains(&audit_id));
3579
3580 let mut seen = Vec::new();
3583 for offset in 0..2 {
3584 let page = split
3585 .query_events(EventFilter::default(), PageRequest { offset, limit: 1 })
3586 .await
3587 .expect("windowed query");
3588 assert_eq!(page.items.len(), 1);
3589 seen.push(page.items[0].id);
3590 }
3591 seen.sort();
3592 let mut expected = vec![plain_id, audit_id];
3593 expected.sort();
3594 assert_eq!(seen, expected);
3595 }
3596
3597 #[tokio::test]
3610 async fn raw_sql_consumers_of_the_legacy_events_table_still_see_plain_appends() {
3611 use khive_storage::types::{SqlStatement, SqlValue};
3612
3613 let dir = tempfile::tempdir().expect("tempdir");
3614 let _registry_guard = TestRegistryGuard::new(dir.path());
3615 let legacy_backend =
3616 direct_backend_for(&dir.path().join("legacy-guard.db")).expect("legacy backend");
3617 let lane_backend =
3618 direct_backend_for(&dir.path().join("lane-guard.db")).expect("lane backend");
3619 let split = SplitEventStore::new(
3620 legacy_backend
3621 .events_for_namespace("local")
3622 .expect("legacy store"),
3623 lane_backend
3624 .events_for_namespace("local")
3625 .expect("lane store"),
3626 );
3627
3628 let provenance = test_event("local");
3631 let provenance_id = provenance.id;
3632 split.append_event(provenance).await.expect("plain append");
3633 let audit = test_event("local");
3634 let audit_id = audit.id;
3635 split
3636 .append_events_idempotent(vec![audit])
3637 .await
3638 .expect("idempotent append");
3639
3640 let count_by_raw_sql = |backend: Arc<StorageBackend>, id: Uuid| async move {
3641 let mut reader = backend.sql().reader().await.expect("sql reader");
3642 let rows = reader
3643 .query_all(SqlStatement {
3644 sql: "SELECT actor FROM events WHERE id = ?1".to_string(),
3645 params: vec![SqlValue::Text(id.to_string())],
3646 label: None,
3647 })
3648 .await
3649 .expect("raw events query");
3650 rows.len()
3651 };
3652
3653 assert_eq!(
3656 count_by_raw_sql(Arc::clone(&legacy_backend), provenance_id).await,
3657 1,
3658 "plain appends must stay visible to raw-SQL consumers of the legacy events table"
3659 );
3660 assert_eq!(
3663 count_by_raw_sql(Arc::clone(&legacy_backend), audit_id).await,
3664 0,
3665 "audit-lane rows must not land in the legacy events table"
3666 );
3667 assert_eq!(
3671 count_by_raw_sql(Arc::clone(&lane_backend), audit_id).await,
3672 1,
3673 "the raw query must prove it can find rows where they actually live"
3674 );
3675 }
3676
3677 #[cfg(unix)]
3680 #[tokio::test]
3681 async fn idempotent_append_round_trips_through_the_daemon() {
3682 let dir = tempfile::tempdir().expect("tempdir");
3683 let (_db, socket) = boot_daemon(&dir).await;
3684 let client = EventsSplitClient::new(socket).expect("client");
3685 let store = ForwardingEventStore::new("local", client);
3686
3687 let event = test_event("local");
3688 let id = event.id;
3689 let result = store
3690 .append_events_idempotent(vec![event.clone()])
3691 .await
3692 .expect("idempotent append over socket");
3693 assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3694
3695 let retry = store
3698 .append_events_idempotent(vec![event])
3699 .await
3700 .expect("idempotent retry");
3701 assert_eq!(
3702 retry.rows,
3703 vec![EventAppendDisposition::AlreadyPresentIdentical]
3704 );
3705
3706 let fetched = store.get_event(id).await.expect("get over socket");
3707 assert_eq!(fetched.map(|e| e.id), Some(id));
3708 let count = store
3709 .count_events(EventFilter::default())
3710 .await
3711 .expect("count over socket");
3712 assert_eq!(count, 1);
3713 }
3714
3715 #[cfg(unix)]
3718 #[tokio::test]
3719 async fn fire_and_forget_append_lands_in_the_daemon_store() {
3720 let dir = tempfile::tempdir().expect("tempdir");
3721 let (_db, socket) = boot_daemon(&dir).await;
3722 let client = EventsSplitClient::new(socket).expect("client");
3723 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3724
3725 store
3726 .append_event(test_event("local"))
3727 .await
3728 .expect("append_event is fire-and-forget");
3729
3730 let (count, metrics) = tokio::time::timeout(Duration::from_secs(30), async {
3734 loop {
3735 let count = store
3736 .count_events(EventFilter::default())
3737 .await
3738 .expect("count over socket");
3739 let metrics = client.metrics();
3740 if (count == 1 && metrics.forwarded_events >= 1) || metrics.dropped_events > 0 {
3741 break (count, metrics);
3742 }
3743 tokio::time::sleep(Duration::from_millis(20)).await;
3744 }
3745 })
3746 .await
3747 .expect("forwarder must deliver or report a drop within 30 seconds");
3748 assert_eq!(
3749 metrics.dropped_events, 0,
3750 "forwarder dropped an event before it reached the daemon store"
3751 );
3752 assert_eq!(count, 1, "forwarded event must land in the daemon store");
3753 assert!(metrics.forwarded_events >= 1, "delivery must be counted");
3754 }
3755
3756 #[cfg(unix)]
3760 #[tokio::test]
3761 async fn queue_overflow_drops_and_counts_but_never_errors() {
3762 let dir = tempfile::tempdir().expect("tempdir");
3763 let dead_socket = dir.path().join("nobody-home.sock");
3764 let client = EventsSplitClient::new_with_queue_depth(dead_socket, 2).expect("client");
3765 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3766
3767 for _ in 0..20 {
3768 store
3769 .append_event(test_event("local"))
3770 .await
3771 .expect("append_event must not error under overflow");
3772 }
3773 let metrics = client.metrics();
3774 assert!(
3775 metrics.dropped_events > 0,
3776 "overflow must register in the drop counter, got {metrics:?}"
3777 );
3778 }
3779
3780 #[cfg(unix)]
3784 #[tokio::test]
3785 async fn dead_socket_reads_fail_typed_and_preflight_stays_local() {
3786 let dir = tempfile::tempdir().expect("tempdir");
3787 let dead_socket = dir.path().join("nobody-home.sock");
3788 let client = EventsSplitClient::new(dead_socket).expect("client");
3789 let store = ForwardingEventStore::new("local", client);
3790
3791 let error = store
3792 .get_event(Uuid::new_v4())
3793 .await
3794 .expect_err("read against a dead socket must fail");
3795 assert!(
3796 matches!(error, StorageError::Pool { .. }),
3797 "expected the typed unreachable error, got {error:?}"
3798 );
3799
3800 assert!(store.supports_idempotent_audit_batch());
3803 store
3804 .preflight_event(&test_event("local"))
3805 .expect("offline preflight validates a well-formed event");
3806 }
3807
3808 #[cfg(unix)]
3811 #[tokio::test]
3812 async fn connected_events_peer_closing_without_response_is_not_unreachable() {
3813 let dir = tempfile::tempdir().unwrap();
3814 let socket = dir.path().join("closes-after-request.sock");
3815 let listener = UnixListener::bind(&socket).unwrap();
3816 let server = tokio::spawn(async move {
3817 let (mut stream, _) = listener.accept().await.unwrap();
3818 read_frame(&mut stream)
3819 .await
3820 .expect("complete request arrives");
3821 });
3823 let client = EventsSplitClient::new(socket).unwrap();
3824 let store = ForwardingEventStore::new("local", client);
3825 let error = store.get_event(Uuid::new_v4()).await.unwrap_err();
3826 assert!(
3827 matches!(error, StorageError::Serialization { .. }),
3828 "a post-connect response failure must not be Pool/unreachable: {error:?}"
3829 );
3830 server.await.unwrap();
3831 }
3832
3833 #[cfg(unix)]
3837 #[tokio::test]
3838 async fn query_page_over_frame_cap_returns_typed_size_refusal() {
3839 let dir = tempfile::tempdir().unwrap();
3840 let socket = dir.path().join("large-query.sock");
3841 let backend = Arc::new(StorageBackend::memory().unwrap());
3842 let event_store = backend.events_for_namespace("local").unwrap();
3843 let events: Vec<_> = (0..MAX_QUERY_EVENTS_PAGE_ROWS)
3844 .map(|_| {
3845 test_event("local").with_payload(serde_json::json!({"data": "x".repeat(3 * 1024)}))
3846 })
3847 .collect();
3848 event_store
3849 .append_events(events)
3850 .await
3851 .expect("seed large page");
3852
3853 let listener = UnixListener::bind(&socket).unwrap();
3854 let stores: NamespaceStores =
3855 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3856 let server = tokio::spawn(async move {
3857 let (stream, _) = listener.accept().await.unwrap();
3858 serve_events_conn(
3859 stream,
3860 backend,
3861 stores,
3862 Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES)),
3863 )
3864 .await;
3865 });
3866 let client = EventsSplitClient::new(socket).unwrap();
3867 let store = ForwardingEventStore::new("local", client);
3868 let error = store
3869 .query_events(
3870 EventFilter::default(),
3871 PageRequest {
3872 offset: 0,
3873 limit: MAX_QUERY_EVENTS_PAGE_ROWS,
3874 },
3875 )
3876 .await
3877 .expect_err("the full page cannot fit one frame");
3878 assert!(
3879 matches!(error, StorageError::InvalidInput { .. }),
3880 "expected non-retryable frame-size refusal, got {error:?}"
3881 );
3882 assert!(
3883 error.to_string().contains("response_frame_size_limit"),
3884 "{error}"
3885 );
3886 assert!(
3887 error.to_string().contains("request a narrower page"),
3888 "{error}"
3889 );
3890 server.abort();
3891 }
3892
3893 #[cfg(unix)]
3897 #[tokio::test]
3898 async fn forwarding_queue_enforces_serialized_byte_and_frame_limits_with_drop_metrics() {
3899 let dir = tempfile::tempdir().unwrap();
3900 let socket = dir.path().join("stalled-forwarder.sock");
3901 let listener = UnixListener::bind(&socket).unwrap();
3902 let server = tokio::spawn(async move {
3903 let (stream, _) = listener.accept().await.unwrap();
3904 let _hold_open = stream;
3905 std::future::pending::<()>().await;
3906 });
3907 let event = test_event("local").with_payload(serde_json::json!({"data": "x".repeat(4096)}));
3908 let first = vec![event.clone(), event.clone()];
3909 let request_bytes = serde_json::to_vec(&EventsRequest::AppendEvents {
3910 protocol_version: EVENTS_PROTOCOL_VERSION,
3911 namespace: "local".into(),
3912 events: first.clone(),
3913 })
3914 .unwrap()
3915 .len();
3916 let byte_budget = request_bytes + 64;
3917 let client = EventsSplitClient::new_with_limits_and_delivery_timeout(
3918 socket,
3919 8,
3920 byte_budget,
3921 Duration::from_secs(30),
3922 )
3923 .unwrap();
3924 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3925 store.append_events(first.clone()).await.unwrap();
3926 assert_eq!(client.metrics().queued_bytes, request_bytes);
3927
3928 store.append_events(first).await.unwrap();
3929 let metrics = client.metrics();
3930 assert_eq!(metrics.dropped_batches, 1);
3931 assert_eq!(metrics.dropped_events, 2);
3932 assert_eq!(metrics.queued_bytes, request_bytes);
3933 assert!(metrics.queued_bytes <= byte_budget);
3934
3935 let oversized = test_event("local")
3938 .with_payload(serde_json::json!({"data": "x".repeat(crate::daemon::MAX_FRAME_BYTES)}));
3939 store.append_event(oversized).await.unwrap();
3940 let metrics = client.metrics();
3941 assert_eq!(metrics.dropped_batches, 2);
3942 assert_eq!(metrics.dropped_events, 3);
3943 assert_eq!(metrics.queued_bytes, request_bytes);
3944 server.abort();
3945 }
3946
3947 #[cfg(unix)]
3950 #[tokio::test]
3951 async fn protocol_version_skew_is_a_typed_refusal() {
3952 let dir = tempfile::tempdir().expect("tempdir");
3953 let (_db, socket) = boot_daemon(&dir).await;
3954
3955 let mut stream = UnixStream::connect(&socket).await.expect("connect");
3956 let request = EventsRequest::CountEvents {
3957 protocol_version: EVENTS_PROTOCOL_VERSION + 1,
3958 namespace: "local".into(),
3959 filter: EventFilter::default(),
3960 };
3961 let payload = serde_json::to_vec(&request).expect("serialize");
3962 write_frame(&mut stream, &payload).await.expect("write");
3963 let bytes = read_frame(&mut stream).await.expect("read");
3964 let response: EventsResponse = serde_json::from_slice(&bytes).expect("parse");
3965 match response {
3966 EventsResponse::Error {
3967 message, retryable, ..
3968 } => {
3969 assert!(!retryable, "version skew is not retryable");
3970 assert!(message.contains("protocol version"), "message: {message}");
3971 }
3972 other => panic!("expected a typed refusal, got {other:?}"),
3973 }
3974 }
3975
3976 #[cfg(unix)]
3980 #[test]
3981 fn events_daemon_guard_is_exclusive_then_reusable() {
3982 let dir = tempfile::tempdir().expect("tempdir");
3983 let socket = dir.path().join("events.sock");
3984 let first = match acquire_events_daemon_guard_outcome(&socket) {
3985 EventsDaemonGuardAcquisition::Held(guard) => guard,
3986 other => panic!("first acquire must succeed: {other:?}"),
3987 };
3988 let second = acquire_events_daemon_guard_outcome(&socket);
3989 assert!(
3990 matches!(second, EventsDaemonGuardAcquisition::Contended),
3991 "second acquire must report contention while the first guard is held: {second:?}"
3992 );
3993 drop(first);
3994 const MAX_ATTEMPTS: usize = 100;
3998 for attempt in 1..=MAX_ATTEMPTS {
3999 match acquire_events_daemon_guard_outcome(&socket) {
4000 EventsDaemonGuardAcquisition::Held(_) => return,
4001 EventsDaemonGuardAcquisition::Contended if attempt < MAX_ATTEMPTS => {
4002 std::thread::sleep(Duration::from_millis(10));
4003 }
4004 other => panic!(
4005 "acquire must succeed after release within {MAX_ATTEMPTS} attempts; \
4006 attempt {attempt}: {other:?}"
4007 ),
4008 }
4009 }
4010 unreachable!("the final acquisition attempt returns or reports its refusal");
4011 }
4012
4013 #[cfg(unix)]
4014 #[test]
4015 fn events_daemon_guard_open_failure_is_not_contention() {
4016 let dir = tempfile::tempdir().expect("tempdir");
4017 let clean = dir.path().join("clean.sock");
4018 assert!(matches!(
4019 acquire_events_daemon_guard_outcome(&clean),
4020 EventsDaemonGuardAcquisition::Held(_)
4021 ));
4022
4023 assert!(
4024 try_acquire_events_daemon_guard(&clean).is_some(),
4025 "the public Option entrance must preserve a successful acquisition"
4026 );
4027
4028 let socket = dir.path().join("blocked.sock");
4031 let lock_path = socket.with_extension("lock");
4032 std::fs::create_dir(&lock_path).expect("create directory at lock entry");
4033 let outcome = acquire_events_daemon_guard_outcome(&socket);
4034 assert!(
4035 matches!(outcome, EventsDaemonGuardAcquisition::OpenFailed(_)),
4036 "an actual lock-file open failure must not report contention: {outcome:?}"
4037 );
4038 assert!(
4039 try_acquire_events_daemon_guard(&socket).is_none(),
4040 "the public Option entrance must refuse an actual open failure"
4041 );
4042 assert!(
4043 lock_path.is_dir(),
4044 "the refused entry must remain a directory"
4045 );
4046 }
4047
4048 #[cfg(unix)]
4055 #[tokio::test]
4056 async fn read_only_runtime_never_creates_an_events_db() {
4057 use crate::{KhiveRuntime, Namespace, RuntimeConfig};
4058
4059 let dir = tempfile::tempdir().expect("tempdir");
4060 let _registry_guard = TestRegistryGuard::new(dir.path());
4061 let main_db = dir.path().join("main.db");
4062 drop(
4064 KhiveRuntime::new_for_test(RuntimeConfig {
4065 db_path: Some(main_db.clone()),
4066 ..RuntimeConfig::no_embeddings()
4067 })
4068 .expect("create main db"),
4069 );
4070 let events_db = events_db_path_beside(&main_db);
4071 assert!(!events_db.exists(), "precondition: no events db yet");
4072
4073 khive_storage::test_support::freeze_snapshot_sidecars(&main_db);
4076
4077 let split_config = |db: PathBuf| RuntimeConfig {
4078 db_path: Some(main_db.clone()),
4079 events_split: Some(EventsSplitConfig {
4080 db_path: db,
4081 socket_path: None,
4082 }),
4083 ..RuntimeConfig::no_embeddings()
4084 };
4085
4086 let ro = KhiveRuntime::new_readonly_for_test(split_config(events_db.clone()))
4089 .expect("read-only runtime");
4090 let token = ro.authorize(Namespace::local()).expect("token");
4091 let store = ro.events(&token).expect("events store");
4092 let count = store
4093 .count_events(EventFilter::default())
4094 .await
4095 .expect("count through legacy-only plane");
4096 assert_eq!(count, 0);
4097 assert!(
4098 !events_db.exists(),
4099 "a read-only runtime must not mint the events database"
4100 );
4101
4102 {
4109 let lane_backend = StorageBackend::sqlite_for_test(&events_db).expect("writable lane");
4110 lane_backend
4111 .events_for_namespace("local")
4112 .expect("lane store")
4113 .append_event(test_event("local"))
4114 .await
4115 .expect("seed lane row");
4116 }
4117 khive_storage::test_support::freeze_snapshot_sidecars(&events_db);
4118 let store = ro.events(&token).expect("events store with lane present");
4119 let count = store
4120 .count_events(EventFilter::default())
4121 .await
4122 .expect("merged count");
4123 assert_eq!(count, 1, "the pre-existing lane row must merge into reads");
4124 }
4125
4126 #[tokio::test]
4129 async fn direct_mode_appends_and_reads_without_a_daemon() {
4130 let dir = tempfile::tempdir().expect("tempdir");
4131 let _registry_guard = TestRegistryGuard::new(dir.path());
4132 let db = dir.path().join("events.db");
4133 let backend = direct_backend_for(&db).expect("direct backend");
4134 let store = backend.events_for_namespace("local").expect("store");
4135 store
4136 .append_event(test_event("local"))
4137 .await
4138 .expect("direct append");
4139 let count = store
4140 .count_events(EventFilter::default())
4141 .await
4142 .expect("count");
4143 assert_eq!(count, 1);
4144 }
4145}