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 { request_state: WireWriterTaskState },
685 TaskTerminated { request_state: WireWriterTaskState },
686}
687
688#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
692#[serde(rename_all = "snake_case")]
693pub enum WireWriterTaskState {
694 NotStarted,
695 TransactionRolledBack,
696 SideEffectsUnknown,
697}
698
699impl From<khive_storage::WriterTaskRequestState> for WireWriterTaskState {
700 fn from(state: khive_storage::WriterTaskRequestState) -> Self {
701 use khive_storage::WriterTaskRequestState as S;
702 match state {
703 S::NotStarted => Self::NotStarted,
704 S::TransactionRolledBack => Self::TransactionRolledBack,
705 S::SideEffectsUnknown => Self::SideEffectsUnknown,
706 }
707 }
708}
709
710impl From<WireWriterTaskState> for khive_storage::WriterTaskRequestState {
711 fn from(state: WireWriterTaskState) -> Self {
712 use khive_storage::WriterTaskRequestState as S;
713 match state {
714 WireWriterTaskState::NotStarted => S::NotStarted,
715 WireWriterTaskState::TransactionRolledBack => S::TransactionRolledBack,
716 WireWriterTaskState::SideEffectsUnknown => S::SideEffectsUnknown,
717 }
718 }
719}
720
721#[cfg(unix)]
729pub struct EventsDaemonGuard {
730 _file: std::fs::File,
731}
732
733#[cfg(unix)]
737pub fn try_acquire_events_daemon_guard(socket_path: &Path) -> Option<EventsDaemonGuard> {
738 match acquire_events_daemon_guard_outcome(socket_path) {
739 EventsDaemonGuardAcquisition::Held(guard) => Some(guard),
740 _ => None,
741 }
742}
743
744#[cfg(unix)]
747enum EventsDaemonGuardAcquisition {
748 Held(EventsDaemonGuard),
749 Contended,
750 OpenFailed(std::io::Error),
751 HardeningRefused(std::io::Error),
752}
753
754#[cfg(unix)]
755impl std::fmt::Debug for EventsDaemonGuardAcquisition {
756 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
757 match self {
758 Self::Held(_) => f.write_str("Held"),
759 Self::Contended => f.write_str("Contended"),
760 Self::OpenFailed(error) => f.debug_tuple("OpenFailed").field(error).finish(),
761 Self::HardeningRefused(error) => {
762 f.debug_tuple("HardeningRefused").field(error).finish()
763 }
764 }
765 }
766}
767
768#[cfg(unix)]
774fn acquire_events_daemon_guard_outcome(socket_path: &Path) -> EventsDaemonGuardAcquisition {
775 use EventsDaemonGuardAcquisition::{Contended, HardeningRefused, Held, OpenFailed};
776
777 let lock_path = socket_path.with_extension("lock");
778 if let Some(parent) = lock_path.parent() {
779 if let Err(error) = std::fs::create_dir_all(parent) {
780 return OpenFailed(error);
781 }
782 }
783 use std::os::unix::fs::OpenOptionsExt;
784 let file = match std::fs::OpenOptions::new()
788 .create(true)
789 .truncate(false)
790 .write(true)
791 .mode(0o600)
792 .custom_flags(libc::O_NOFOLLOW)
793 .open(&lock_path)
794 {
795 Ok(file) => file,
796 Err(error) if error.raw_os_error() == Some(libc::ELOOP) => {
797 return HardeningRefused(error);
798 }
799 Err(error) => return OpenFailed(error),
800 };
801 if let Err(error) = file.set_permissions(
807 <std::fs::Permissions as std::os::unix::fs::PermissionsExt>::from_mode(0o600),
808 ) {
809 return HardeningRefused(error);
810 }
811 use std::os::fd::AsRawFd;
812 let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
815 if rc == 0 {
816 Held(EventsDaemonGuard { _file: file })
817 } else {
818 let error = std::io::Error::last_os_error();
819 if error.raw_os_error() == Some(libc::EWOULDBLOCK)
820 || error.raw_os_error() == Some(libc::EAGAIN)
821 {
822 Contended
823 } else {
824 HardeningRefused(error)
825 }
826 }
827}
828
829fn events_db_targets(db_path: &Path) -> [PathBuf; 3] {
830 ["", "-wal", "-shm"].map(|suffix| {
831 let mut name = db_path.as_os_str().to_os_string();
832 name.push(suffix);
833 PathBuf::from(name)
834 })
835}
836
837fn refuse_events_db_symlinks(db_path: &Path) -> anyhow::Result<()> {
847 for path in events_db_targets(db_path) {
848 match std::fs::symlink_metadata(&path) {
849 Ok(meta) if meta.file_type().is_symlink() => anyhow::bail!(
850 "refusing to serve events: {} is a symlink; the events database and its \
851 sidecars must be regular files",
852 path.display()
853 ),
854 Ok(_) => {}
855 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
856 Err(e) => anyhow::bail!(
857 "refusing to serve events: cannot inspect {}: {e}",
858 path.display()
859 ),
860 }
861 }
862 Ok(())
863}
864
865#[cfg(unix)]
866#[derive(Clone, Copy, Debug, PartialEq, Eq)]
867struct EventsFileIdentity {
868 device: u64,
869 inode: u64,
870}
871
872#[cfg(unix)]
873impl EventsFileIdentity {
874 fn from_metadata(metadata: &std::fs::Metadata) -> Self {
875 use std::os::unix::fs::MetadataExt;
876 Self {
877 device: metadata.dev(),
878 inode: metadata.ino(),
879 }
880 }
881}
882
883#[cfg(unix)]
886type EventsDbIdentities = [Option<EventsFileIdentity>; 3];
887
888#[cfg(unix)]
897fn ensure_events_db_parent_trusted(db_path: &Path) -> anyhow::Result<()> {
898 let parent = absolutize(db_path);
899 let parent = parent.parent().filter(|p| !p.as_os_str().is_empty());
900 match parent {
901 Some(dir) => crate::daemon::ensure_socket_dir_is_trusted(dir).map_err(|e| {
902 anyhow::anyhow!("refusing to serve events from an untrusted directory: {e}")
903 }),
904 None => Ok(()),
905 }
906}
907
908#[cfg(unix)]
914fn ensure_events_db_owner_only(db_path: &Path) -> anyhow::Result<()> {
915 use std::os::unix::fs::OpenOptionsExt;
916 refuse_events_db_symlinks(db_path)?;
917 if let Some(parent) = db_path.parent() {
918 std::fs::create_dir_all(parent)?;
919 }
920 ensure_events_db_parent_trusted(db_path)?;
923 match std::fs::OpenOptions::new()
924 .write(true)
925 .create_new(true)
926 .mode(0o600)
927 .open(db_path)
928 {
929 Ok(_created) => Ok(()),
931 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
932 Err(e) => Err(anyhow::anyhow!(
933 "refusing to serve events: cannot create {} owner-only: {e}",
934 db_path.display()
935 )),
936 }
937}
938
939#[cfg(unix)]
952fn harden_events_db_sidecars(db_path: &Path) -> anyhow::Result<EventsDbIdentities> {
953 use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
954 let mut identities = [None; 3];
955 for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
956 let file = match std::fs::OpenOptions::new()
963 .read(true)
964 .custom_flags(libc::O_NOFOLLOW)
965 .open(&path)
966 {
967 Ok(file) => file,
968 Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
969 Err(e) => {
970 anyhow::bail!(
971 "refusing to serve events: cannot open {} without following symlinks: \
972 {e}. The events database and its sidecars must be regular files.",
973 path.display()
974 )
975 }
976 };
977 let metadata = file.metadata()?;
980 if !metadata.file_type().is_file() {
981 anyhow::bail!(
982 "refusing to serve events: {} is not a regular file. The events \
983 database and its sidecars must be regular files.",
984 path.display()
985 );
986 }
987 file.set_permissions(std::fs::Permissions::from_mode(0o600))
988 .map_err(|e| {
989 anyhow::anyhow!(
990 "refusing to serve events: cannot chmod 0600 {}: {e}. The events \
991 database and its sidecars must be owner-only.",
992 path.display()
993 )
994 })?;
995 identities[index] = Some(EventsFileIdentity::from_metadata(&metadata));
998 }
999 Ok(identities)
1000}
1001
1002#[cfg(unix)]
1020fn verify_events_db_owner_only_unopened(
1021 db_path: &Path,
1022 before: &EventsDbIdentities,
1023) -> anyhow::Result<()> {
1024 use std::os::unix::fs::PermissionsExt;
1025 for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
1026 let metadata = match std::fs::symlink_metadata(&path) {
1027 Ok(metadata) => metadata,
1028 Err(e) if e.kind() == std::io::ErrorKind::NotFound && index != 0 => continue,
1029 Err(e) if e.kind() == std::io::ErrorKind::NotFound => anyhow::bail!(
1030 "refusing to serve events: main database {} disappeared after SQLite open; \
1031 refusing to serve an unlinked database",
1032 path.display()
1033 ),
1034 Err(e) => {
1035 anyhow::bail!(
1036 "refusing to serve events: cannot stat {}: {e}",
1037 path.display()
1038 )
1039 }
1040 };
1041 if !metadata.file_type().is_file() {
1042 anyhow::bail!(
1043 "refusing to serve events: {} is not a regular file. The events database \
1044 and its sidecars must be regular files.",
1045 path.display()
1046 );
1047 }
1048 let mode = metadata.permissions().mode() & 0o777;
1049 if mode & 0o077 != 0 {
1050 anyhow::bail!(
1051 "refusing to serve events: {} is mode {mode:03o}, not owner-only. The events \
1052 database and its sidecars must be owner-only.",
1053 path.display()
1054 );
1055 }
1056 if before[index]
1057 .is_some_and(|identity| identity != EventsFileIdentity::from_metadata(&metadata))
1058 {
1059 anyhow::bail!(
1060 "refusing to serve events: {} changed identity between pre-open hardening \
1061 and post-open verification",
1062 path.display()
1063 );
1064 }
1065 }
1066 Ok(())
1067}
1068
1069#[cfg(unix)]
1095pub async fn supervise_events_daemon(db_path: PathBuf, socket_path: PathBuf) {
1096 match standalone_daemon_policies(&db_path) {
1097 Ok((wal_ceiling, disk_guard, volume_lock_dir)) => {
1098 supervise_events_daemon_with_policies(
1099 db_path,
1100 socket_path,
1101 wal_ceiling,
1102 disk_guard,
1103 volume_lock_dir,
1104 )
1105 .await;
1106 }
1107 Err(error) => {
1108 tracing::warn!(%error, "invalid events daemon policy; events supervisor not started");
1109 }
1110 }
1111}
1112
1113#[cfg(unix)]
1115pub async fn supervise_events_daemon_with_policies(
1116 db_path: PathBuf,
1117 socket_path: PathBuf,
1118 wal_ceiling: WalCeilingPolicy,
1119 disk_guard: khive_db::EffectiveDiskGuardConfig,
1120 volume_lock_dir: PathBuf,
1121) {
1122 if let Err(error) = disk_guard.validate() {
1123 tracing::warn!(%error, "invalid disk policy; events supervisor not started");
1124 return;
1125 }
1126 if let Err(error) = wal_ceiling.validate_static(true, true, false) {
1127 tracing::warn!(error = %error, "invalid WAL ceiling; events supervisor not started");
1128 return;
1129 }
1130 const PROBE_INTERVAL: Duration = Duration::from_secs(15);
1131 let shutdown = crate::daemon::daemon_shutdown_token();
1132 let mut child: Option<std::process::Child> = None;
1133 let mut respawns: u64 = 0;
1134 loop {
1135 if let Some(c) = child.as_mut() {
1139 match c.try_wait() {
1140 Ok(Some(status)) => {
1141 tracing::info!(%status, "events daemon child exited");
1142 child = None;
1143 }
1144 Ok(None) => {}
1145 Err(error) => {
1146 tracing::warn!(error = %error, "cannot poll events daemon child; dropping handle");
1147 child = None;
1148 }
1149 }
1150 }
1151
1152 let reachable = match connect_verified(&socket_path).await {
1153 Ok(_stream) => true,
1154 Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
1155 tracing::warn!(
1156 socket = %socket_path.display(),
1157 error = %error,
1158 "events socket answered by a foreign uid; treating as unreachable"
1159 );
1160 false
1161 }
1162 Err(_) => false,
1163 };
1164 if !reachable && child.is_none() {
1165 match std::env::current_exe() {
1166 Ok(exe) => {
1167 let spawned = events_daemon_command(
1168 &exe,
1169 &db_path,
1170 &socket_path,
1171 wal_ceiling,
1172 disk_guard,
1173 &volume_lock_dir,
1174 )
1175 .spawn();
1176 match spawned {
1177 Ok(spawned_child) => {
1178 respawns += 1;
1179 tracing::info!(
1180 pid = spawned_child.id(),
1181 respawns,
1182 socket = %socket_path.display(),
1183 "spawned events daemon"
1184 );
1185 child = Some(spawned_child);
1186 }
1187 Err(error) => {
1188 tracing::warn!(error = %error, "failed to spawn events daemon");
1189 }
1190 }
1191 }
1192 Err(error) => {
1193 tracing::warn!(error = %error, "cannot resolve current executable for events daemon spawn");
1194 }
1195 }
1196 }
1197 tokio::select! {
1198 _ = shutdown.cancelled() => break,
1199 _ = tokio::time::sleep(PROBE_INTERVAL) => {}
1200 }
1201 }
1202 if let Some(mut c) = child.take() {
1203 let _ = c.kill();
1207 let _ = c.wait();
1208 tracing::info!("events daemon child stopped with supervisor shutdown");
1209 }
1210}
1211
1212#[cfg(unix)]
1213fn events_daemon_command(
1214 executable: &Path,
1215 db_path: &Path,
1216 socket_path: &Path,
1217 wal_ceiling: WalCeilingPolicy,
1218 disk_guard: khive_db::EffectiveDiskGuardConfig,
1219 volume_lock_dir: &Path,
1220) -> std::process::Command {
1221 let source = match wal_ceiling.source {
1222 khive_db::WalCeilingSource::BackendField => "backend_field",
1223 khive_db::WalCeilingSource::Environment => "environment",
1224 khive_db::WalCeilingSource::Default => "default",
1225 };
1226 let mut command = std::process::Command::new(executable);
1227 command
1228 .arg("events-daemon")
1229 .arg("--db")
1230 .arg(db_path)
1231 .arg("--socket")
1232 .arg(socket_path)
1233 .arg("--wal-ceiling-bytes")
1234 .arg(wal_ceiling.bytes.to_string())
1235 .arg("--wal-ceiling-source")
1236 .arg(source)
1237 .arg("--disk-reserve-bytes")
1238 .arg(disk_guard.reserve_bytes.to_string())
1239 .arg("--disk-guard-deadline-ms")
1240 .arg(disk_guard.guard_deadline_ms.to_string())
1241 .arg("--disk-reserve-source")
1242 .arg(disk_guard.reserve_source.as_str())
1243 .arg("--disk-deadline-source")
1244 .arg(disk_guard.deadline_source.as_str())
1245 .arg("--disk-legacy-environment-present")
1246 .arg(disk_guard.legacy_environment_present.to_string())
1247 .arg("--volume-lock-dir")
1248 .arg(volume_lock_dir)
1249 .stdin(std::process::Stdio::null())
1250 .stdout(std::process::Stdio::null())
1251 .stderr(std::process::Stdio::null());
1252 command
1253}
1254
1255#[cfg(unix)]
1274pub async fn run_events_daemon(db_path: &Path, socket_path: &Path) -> anyhow::Result<()> {
1275 let (wal_ceiling, disk_guard, volume_lock_dir) = standalone_daemon_policies(db_path)?;
1276 run_events_daemon_with_policies(
1277 db_path,
1278 socket_path,
1279 wal_ceiling,
1280 disk_guard,
1281 volume_lock_dir,
1282 )
1283 .await
1284}
1285
1286#[cfg(unix)]
1288pub async fn run_events_daemon_with_policies(
1289 db_path: &Path,
1290 socket_path: &Path,
1291 wal_ceiling: WalCeilingPolicy,
1292 disk_guard: khive_db::EffectiveDiskGuardConfig,
1293 volume_lock_dir: PathBuf,
1294) -> anyhow::Result<()> {
1295 disk_guard.validate()?;
1296 wal_ceiling
1297 .validate_static(true, true, false)
1298 .map_err(crate::error::RuntimeError::from)?;
1299 let db_path = &absolutize(db_path);
1302 let socket_path = &absolutize(socket_path);
1303 validate_events_socket_path(socket_path)?;
1304 if let Some(parent) = socket_path.parent() {
1308 std::fs::create_dir_all(parent)?;
1309 crate::daemon::ensure_socket_dir_is_trusted(parent)?;
1310 }
1311 let _guard = match acquire_events_daemon_guard_outcome(socket_path) {
1312 EventsDaemonGuardAcquisition::Held(guard) => guard,
1313 refusal => {
1314 tracing::info!(
1315 socket = %socket_path.display(),
1316 reason = ?refusal,
1317 "events daemon lock unavailable; exiting"
1318 );
1319 return Ok(());
1320 }
1321 };
1322 ensure_events_db_owner_only(db_path)?;
1323 let before_open = harden_events_db_sidecars(db_path)?;
1336 let backend = Arc::new(
1337 StorageBackend::sqlite_with_max_readers_and_policies(
1338 db_path,
1339 None,
1340 wal_ceiling,
1341 disk_guard,
1342 volume_lock_dir,
1343 )
1344 .map_err(crate::error::RuntimeError::from)?,
1345 );
1346 backend.events()?;
1348 verify_events_db_owner_only_unopened(db_path, &before_open)?;
1349
1350 if socket_path.exists() {
1351 std::fs::remove_file(socket_path)?;
1352 }
1353 let listener = UnixListener::bind(socket_path)?;
1354 {
1355 use std::os::unix::fs::PermissionsExt;
1356 if let Err(e) =
1360 std::fs::set_permissions(socket_path, std::fs::Permissions::from_mode(0o600))
1361 {
1362 drop(listener);
1363 let _ = std::fs::remove_file(socket_path);
1364 anyhow::bail!(
1365 "refusing to serve events: cannot chmod 0600 {}: {e}. The events socket must \
1366 be owner-only.",
1367 socket_path.display()
1368 );
1369 }
1370 }
1371 let daemon_euid = unsafe { libc::geteuid() };
1372 let stores: NamespaceStores = Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
1378 let connections = Arc::new(tokio::sync::Semaphore::new(MAX_EVENTS_CONNECTIONS));
1379 let frame_budget = Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES));
1380 tracing::info!(
1381 socket = %socket_path.display(),
1382 db = %db_path.display(),
1383 wal_ceiling_configured_bytes = wal_ceiling.bytes,
1384 wal_ceiling_effective_bytes = wal_ceiling.effective_bytes(backend.is_read_only()),
1385 wal_ceiling_source = ?wal_ceiling.source,
1386 "events daemon listening"
1387 );
1388
1389 loop {
1390 let (stream, _addr) = match listener.accept().await {
1391 Ok(pair) => pair,
1392 Err(error) => {
1393 tracing::warn!(error = %error, "events daemon accept failed");
1394 continue;
1395 }
1396 };
1397 match crate::daemon::peer_uid(&stream) {
1398 Ok(uid) if crate::daemon::uid_is_permitted(uid, daemon_euid) => {}
1399 Ok(uid) => {
1400 tracing::warn!(peer_uid = uid, "events daemon rejected foreign-uid peer");
1401 continue;
1402 }
1403 Err(error) => {
1404 tracing::warn!(error = %error, "events daemon could not read peer credentials");
1405 continue;
1406 }
1407 }
1408 let permit = match Arc::clone(&connections).try_acquire_owned() {
1409 Ok(permit) => permit,
1410 Err(_) => {
1411 tracing::warn!(
1414 cap = MAX_EVENTS_CONNECTIONS,
1415 "events daemon at connection cap; dropping new connection"
1416 );
1417 continue;
1418 }
1419 };
1420 let backend = Arc::clone(&backend);
1421 let stores = Arc::clone(&stores);
1422 let frame_budget = Arc::clone(&frame_budget);
1423 crate::daemon::spawn_named_tracked_task("events_connection", async move {
1424 serve_events_conn(stream, backend, stores, frame_budget).await;
1425 drop(permit);
1426 });
1427 }
1428}
1429
1430#[cfg(unix)]
1431#[path = "events_split_policy_entries.rs"]
1432mod policy_entries;
1433#[cfg(unix)]
1434use policy_entries::standalone_daemon_policies;
1435#[cfg(unix)]
1436pub use policy_entries::{
1437 run_events_daemon_with_wal_ceiling, supervise_events_daemon_with_wal_ceiling,
1438};
1439
1440#[cfg(all(test, unix))]
1441#[path = "events_wal_policy_tests.rs"]
1442mod wal_policy_tests;
1443
1444#[cfg(unix)]
1452async fn read_frame_budgeted(
1453 stream: &mut UnixStream,
1454 budget: &Arc<tokio::sync::Semaphore>,
1455) -> std::io::Result<(Vec<u8>, tokio::sync::OwnedSemaphorePermit)> {
1456 use tokio::io::AsyncReadExt;
1457 let mut len_buf = [0u8; 4];
1458 stream.read_exact(&mut len_buf).await?;
1459 let len = u32::from_be_bytes(len_buf) as usize;
1460 if len > crate::daemon::MAX_FRAME_BYTES {
1461 return Err(std::io::Error::new(
1462 std::io::ErrorKind::InvalidData,
1463 format!(
1464 "daemon frame of {len} bytes exceeds {} cap",
1465 crate::daemon::MAX_FRAME_BYTES
1466 ),
1467 ));
1468 }
1469 let permit = Arc::clone(budget)
1470 .acquire_many_owned(len as u32)
1471 .await
1472 .map_err(|_| std::io::Error::other("frame budget closed"))?;
1473 let mut buf = vec![0u8; len];
1474 stream.read_exact(&mut buf).await?;
1475 Ok((buf, permit))
1476}
1477
1478#[cfg(unix)]
1479type NamespaceStores =
1480 Arc<std::sync::Mutex<std::collections::HashMap<String, Arc<dyn EventStore>>>>;
1481
1482#[cfg(unix)]
1494fn namespace_store(
1495 backend: &StorageBackend,
1496 stores: &NamespaceStores,
1497 namespace: &str,
1498) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1499 namespace_store_with_cap(backend, stores, namespace, MAX_CACHED_NAMESPACE_STORES)
1500}
1501
1502#[cfg(unix)]
1503fn namespace_store_with_cap(
1504 backend: &StorageBackend,
1505 stores: &NamespaceStores,
1506 namespace: &str,
1507 cap: usize,
1508) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1509 let key = namespace.trim();
1510 if let Some(store) = stores
1511 .lock()
1512 .unwrap_or_else(std::sync::PoisonError::into_inner)
1513 .get(key)
1514 {
1515 return Ok(Arc::clone(store));
1516 }
1517 let store = backend.events_for_namespace(namespace)?;
1518 let mut map = stores
1519 .lock()
1520 .unwrap_or_else(std::sync::PoisonError::into_inner);
1521 if map.len() >= cap && !map.contains_key(key) {
1522 if let Some(victim) = map.keys().next().cloned() {
1530 map.remove(&victim);
1531 }
1532 }
1533 map.insert(key.to_string(), Arc::clone(&store));
1534 Ok(store)
1535}
1536
1537#[cfg(unix)]
1538async fn serve_events_conn(
1539 mut stream: UnixStream,
1540 backend: Arc<StorageBackend>,
1541 stores: NamespaceStores,
1542 frame_budget: Arc<tokio::sync::Semaphore>,
1543) {
1544 loop {
1545 let (payload, budget_permit) = match tokio::time::timeout(
1551 CONN_IO_TIMEOUT,
1552 read_frame_budgeted(&mut stream, &frame_budget),
1553 )
1554 .await
1555 {
1556 Ok(Ok(pair)) => pair,
1557 Ok(Err(_)) | Err(_) => return,
1559 };
1560 let response = match serde_json::from_slice::<EventsRequest>(&payload) {
1561 Ok(request) => dispatch_events_request(request, &backend, &stores).await,
1562 Err(error) => EventsResponse::Error {
1563 message: format!("events daemon could not parse request frame: {error}"),
1564 retryable: false,
1565 writer_task_failure: None,
1566 },
1567 };
1568 let bytes = match serde_json::to_vec(&response) {
1569 Ok(bytes) => bytes,
1570 Err(error) => {
1571 tracing::error!(error = %error, "events daemon response serialization failed");
1572 return;
1573 }
1574 };
1575 let bytes = if bytes.len() > crate::daemon::MAX_FRAME_BYTES {
1576 let refusal = EventsResponse::Error {
1580 message: format!(
1581 "response_frame_size_limit: {} response bytes exceed the {}-byte events IPC frame cap; request a narrower page",
1582 bytes.len(),
1583 crate::daemon::MAX_FRAME_BYTES
1584 ),
1585 retryable: false,
1586 writer_task_failure: None,
1587 };
1588 match serde_json::to_vec(&refusal) {
1589 Ok(bytes) => bytes,
1590 Err(error) => {
1591 tracing::error!(error = %error, "events daemon frame-size refusal serialization failed");
1592 return;
1593 }
1594 }
1595 } else {
1596 bytes
1597 };
1598 match tokio::time::timeout(CONN_IO_TIMEOUT, write_frame(&mut stream, &bytes)).await {
1601 Ok(Ok(())) => {}
1602 Ok(Err(_)) | Err(_) => return,
1603 }
1604 drop(budget_permit);
1608 }
1609}
1610
1611#[cfg(unix)]
1612async fn dispatch_events_request(
1613 request: EventsRequest,
1614 backend: &StorageBackend,
1615 stores: &NamespaceStores,
1616) -> EventsResponse {
1617 if request.protocol_version() != EVENTS_PROTOCOL_VERSION {
1618 return EventsResponse::Error {
1619 message: format!(
1620 "events protocol version mismatch: daemon speaks {}, client sent {}",
1621 EVENTS_PROTOCOL_VERSION,
1622 request.protocol_version()
1623 ),
1624 retryable: false,
1625 writer_task_failure: None,
1626 };
1627 }
1628 if let Err(error) = khive_types::Namespace::parse(request.namespace().trim()) {
1634 return EventsResponse::Error {
1635 message: format!("events request namespace rejected: {error}"),
1636 retryable: false,
1637 writer_task_failure: None,
1638 };
1639 }
1640 let store = match namespace_store(backend, stores, request.namespace()) {
1641 Ok(store) => store,
1642 Err(error) => {
1643 return EventsResponse::Error {
1644 message: format!("events store unavailable: {error}"),
1645 retryable: true,
1646 writer_task_failure: None,
1647 };
1648 }
1649 };
1650 match request {
1651 EventsRequest::AppendEvents { events, .. } => match store.append_events(events).await {
1652 Ok(summary) => EventsResponse::Appended { summary },
1653 Err(error) => storage_error_response(&error),
1654 },
1655 EventsRequest::AppendEventsIdempotent { events, .. } => {
1656 match store.append_events_idempotent(events).await {
1657 Ok(result) => EventsResponse::Idempotent { result },
1658 Err(error) => storage_error_response(&error),
1659 }
1660 }
1661 EventsRequest::GetEvent { id, .. } => match store.get_event(id).await {
1662 Ok(event) => EventsResponse::Event { event },
1663 Err(error) => storage_error_response(&error),
1664 },
1665 EventsRequest::QueryEvents { filter, page, .. } => {
1666 if page.limit > MAX_QUERY_EVENTS_PAGE_ROWS {
1667 return EventsResponse::Error {
1668 message: format!(
1669 "events query page limit {} exceeds the daemon cap of {} rows; \
1670 request narrower pages",
1671 page.limit, MAX_QUERY_EVENTS_PAGE_ROWS
1672 ),
1673 retryable: false,
1674 writer_task_failure: None,
1675 };
1676 }
1677 match store.query_events(filter, page).await {
1678 Ok(page) => EventsResponse::Pageful { page },
1679 Err(error) => storage_error_response(&error),
1680 }
1681 }
1682 EventsRequest::QueryEventPage { query, .. } => {
1683 if let Err(error) = crate::event_page::validate_page_query(&query) {
1684 return storage_error_response(&error);
1685 }
1686 match store.query_event_page(query).await {
1687 Ok(window) => EventsResponse::EventPageWindow { window },
1688 Err(error) => storage_error_response(&error),
1689 }
1690 }
1691 EventsRequest::CountEvents { filter, .. } => match store.count_events(filter).await {
1692 Ok(count) => EventsResponse::Count { count },
1693 Err(error) => storage_error_response(&error),
1694 },
1695 }
1696}
1697
1698#[cfg(unix)]
1699fn storage_error_response(error: &StorageError) -> EventsResponse {
1700 let (message, writer_task_failure) = match error {
1701 StorageError::WriterTaskRequestFailed {
1702 request_state,
1703 source,
1704 } => (
1705 source.to_string(),
1706 Some(WireWriterTaskFailure::RequestFailed {
1707 request_state: (*request_state).into(),
1708 }),
1709 ),
1710 StorageError::WriterTaskTerminated { request_state } => (
1711 error.to_string(),
1712 Some(WireWriterTaskFailure::TaskTerminated {
1713 request_state: (*request_state).into(),
1714 }),
1715 ),
1716 _ => (error.to_string(), None),
1717 };
1718 EventsResponse::Error {
1719 message,
1720 retryable: error.is_retryable(),
1725 writer_task_failure,
1726 }
1727}
1728
1729#[cfg(unix)]
1740async fn connect_verified(socket_path: &Path) -> std::io::Result<UnixStream> {
1741 let stream = UnixStream::connect(socket_path).await?;
1742 let own_euid = unsafe { libc::geteuid() } as u32;
1744 let peer = crate::daemon::peer_uid(&stream)?;
1745 if !crate::daemon::uid_is_permitted(peer, own_euid) {
1746 return Err(std::io::Error::new(
1747 std::io::ErrorKind::PermissionDenied,
1748 format!(
1749 "events socket {} answered by uid {peer}, not this process's uid {own_euid}; \
1750 refusing to exchange event frames with an unowned daemon",
1751 socket_path.display()
1752 ),
1753 ));
1754 }
1755 Ok(stream)
1756}
1757
1758#[cfg(unix)]
1762#[derive(Debug, Clone, Copy, Serialize)]
1763pub struct EventsForwardingMetrics {
1764 pub forwarded_batches: u64,
1765 pub forwarded_events: u64,
1766 pub dropped_batches: u64,
1767 pub dropped_events: u64,
1768 pub queued_bytes: usize,
1770}
1771
1772#[cfg(unix)]
1773#[derive(Debug, Default)]
1774struct ForwardingCounters {
1775 forwarded_batches: AtomicU64,
1776 forwarded_events: AtomicU64,
1777 dropped_batches: AtomicU64,
1778 dropped_events: AtomicU64,
1779 queued_bytes: AtomicUsize,
1780}
1781
1782#[cfg(unix)]
1785#[derive(Debug)]
1786struct AppendByteReservation {
1787 counters: Arc<ForwardingCounters>,
1788 bytes: usize,
1789}
1790
1791#[cfg(unix)]
1792impl Drop for AppendByteReservation {
1793 fn drop(&mut self) {
1794 self.counters
1795 .queued_bytes
1796 .fetch_sub(self.bytes, Ordering::AcqRel);
1797 }
1798}
1799
1800#[cfg(unix)]
1801#[derive(Debug)]
1802struct QueuedAppend {
1803 request: EventsRequest,
1804 count: u64,
1805 _reservation: AppendByteReservation,
1806}
1807
1808#[cfg(all(test, unix))]
1809impl QueuedAppend {
1810 fn unmetered(namespace: &str, events: Vec<Event>, counters: Arc<ForwardingCounters>) -> Self {
1811 let count = events.len() as u64;
1812 Self {
1813 request: EventsRequest::AppendEvents {
1814 protocol_version: EVENTS_PROTOCOL_VERSION,
1815 namespace: namespace.to_string(),
1816 events,
1817 },
1818 count,
1819 _reservation: AppendByteReservation { counters, bytes: 0 },
1820 }
1821 }
1822}
1823
1824#[cfg(unix)]
1827#[derive(Default)]
1828struct BoundedFrameCounter {
1829 bytes: usize,
1830 exceeded: bool,
1831}
1832
1833#[cfg(unix)]
1834impl std::io::Write for BoundedFrameCounter {
1835 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
1836 let next = self.bytes.saturating_add(buf.len());
1837 if next > crate::daemon::MAX_FRAME_BYTES {
1838 self.exceeded = true;
1839 return Err(std::io::Error::new(
1840 std::io::ErrorKind::InvalidData,
1841 "events append request exceeds IPC frame cap",
1842 ));
1843 }
1844 self.bytes = next;
1845 Ok(buf.len())
1846 }
1847
1848 fn flush(&mut self) -> std::io::Result<()> {
1849 Ok(())
1850 }
1851}
1852
1853#[cfg(unix)]
1854fn reserve_append_bytes(
1855 counters: &Arc<ForwardingCounters>,
1856 bytes: usize,
1857 budget: usize,
1858) -> Option<AppendByteReservation> {
1859 let mut used = counters.queued_bytes.load(Ordering::Acquire);
1860 loop {
1861 let next = used.checked_add(bytes)?;
1862 if next > budget {
1863 return None;
1864 }
1865 match counters.queued_bytes.compare_exchange_weak(
1866 used,
1867 next,
1868 Ordering::AcqRel,
1869 Ordering::Acquire,
1870 ) {
1871 Ok(_) => {
1872 return Some(AppendByteReservation {
1873 counters: Arc::clone(counters),
1874 bytes,
1875 });
1876 }
1877 Err(observed) => used = observed,
1878 }
1879 }
1880}
1881
1882#[cfg(unix)]
1890pub struct EventsSplitClient {
1891 socket_path: PathBuf,
1892 append_tx: tokio::sync::mpsc::Sender<QueuedAppend>,
1893 append_queue_byte_budget: usize,
1894 counters: Arc<ForwardingCounters>,
1895 outage_logged: Arc<AtomicBool>,
1898 preflight_store: Arc<dyn EventStore>,
1900}
1901
1902#[cfg(unix)]
1903impl std::fmt::Debug for EventsSplitClient {
1904 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1905 f.debug_struct("EventsSplitClient")
1906 .field("socket_path", &self.socket_path)
1907 .finish_non_exhaustive()
1908 }
1909}
1910
1911#[cfg(unix)]
1912impl EventsSplitClient {
1913 pub fn new(socket_path: PathBuf) -> crate::error::RuntimeResult<Arc<Self>> {
1915 Self::new_with_queue_depth(socket_path, DEFAULT_APPEND_QUEUE_BATCHES)
1916 }
1917
1918 pub fn new_with_queue_depth(
1921 socket_path: PathBuf,
1922 queue_depth: usize,
1923 ) -> crate::error::RuntimeResult<Arc<Self>> {
1924 Self::new_with_queue_depth_and_delivery_timeout(socket_path, queue_depth, REQUEST_TIMEOUT)
1925 }
1926
1927 fn new_with_queue_depth_and_delivery_timeout(
1931 socket_path: PathBuf,
1932 queue_depth: usize,
1933 delivery_timeout: Duration,
1934 ) -> crate::error::RuntimeResult<Arc<Self>> {
1935 Self::new_with_limits_and_delivery_timeout(
1936 socket_path,
1937 queue_depth,
1938 DEFAULT_APPEND_QUEUE_BYTES,
1939 delivery_timeout,
1940 )
1941 }
1942
1943 fn new_with_limits_and_delivery_timeout(
1944 socket_path: PathBuf,
1945 queue_depth: usize,
1946 byte_budget: usize,
1947 delivery_timeout: Duration,
1948 ) -> crate::error::RuntimeResult<Arc<Self>> {
1949 validate_events_socket_path(&socket_path)?;
1950 let preflight_backend = StorageBackend::memory()?;
1951 let preflight_store = preflight_backend.events()?;
1952 let (append_tx, append_rx) = tokio::sync::mpsc::channel::<QueuedAppend>(queue_depth.max(1));
1956 let counters = Arc::new(ForwardingCounters::default());
1957 let outage_logged = Arc::new(AtomicBool::new(false));
1958
1959 let client = Arc::new(Self {
1960 socket_path: socket_path.clone(),
1961 append_tx,
1962 append_queue_byte_budget: byte_budget,
1963 counters: Arc::clone(&counters),
1964 outage_logged: Arc::clone(&outage_logged),
1965 preflight_store,
1966 });
1967
1968 crate::daemon::spawn_named_tracked_task(
1969 "events_forwarder",
1970 run_forwarder(
1971 socket_path,
1972 append_rx,
1973 counters,
1974 outage_logged,
1975 delivery_timeout,
1976 crate::daemon::daemon_shutdown_token(),
1977 ),
1978 );
1979 Ok(client)
1980 }
1981
1982 pub fn metrics(&self) -> EventsForwardingMetrics {
1984 EventsForwardingMetrics {
1985 forwarded_batches: self.counters.forwarded_batches.load(Ordering::Relaxed),
1986 forwarded_events: self.counters.forwarded_events.load(Ordering::Relaxed),
1987 dropped_batches: self.counters.dropped_batches.load(Ordering::Relaxed),
1988 dropped_events: self.counters.dropped_events.load(Ordering::Relaxed),
1989 queued_bytes: self.counters.queued_bytes.load(Ordering::Acquire),
1990 }
1991 }
1992
1993 fn enqueue(&self, namespace: &str, events: Vec<Event>) {
1996 let count = events.len() as u64;
1997 let request = EventsRequest::AppendEvents {
1998 protocol_version: EVENTS_PROTOCOL_VERSION,
1999 namespace: namespace.to_string(),
2000 events,
2001 };
2002 let mut frame_size = BoundedFrameCounter::default();
2003 if let Err(error) = serde_json::to_writer(&mut frame_size, &request) {
2004 self.count_append_drop(
2005 count,
2006 if frame_size.exceeded {
2007 "frame cap"
2008 } else {
2009 "serialization"
2010 },
2011 );
2012 tracing::warn!(error = %error, "events append batch could not fit a wire frame");
2013 return;
2014 }
2015 let Some(reservation) = reserve_append_bytes(
2016 &self.counters,
2017 frame_size.bytes,
2018 self.append_queue_byte_budget,
2019 ) else {
2020 self.count_append_drop(count, "queue byte budget");
2021 return;
2022 };
2023 let batch = QueuedAppend {
2024 request,
2025 count,
2026 _reservation: reservation,
2027 };
2028 match self.append_tx.try_send(batch) {
2029 Ok(()) => {}
2030 Err(_) => {
2031 self.count_append_drop(count, "queue batch limit");
2033 }
2034 }
2035 }
2036
2037 fn count_append_drop(&self, count: u64, reason: &'static str) {
2038 self.counters
2039 .dropped_batches
2040 .fetch_add(1, Ordering::Relaxed);
2041 self.counters
2042 .dropped_events
2043 .fetch_add(count, Ordering::Relaxed);
2044 if !self.outage_logged.swap(true, Ordering::Relaxed) {
2045 tracing::warn!(
2046 dropped_events = count,
2047 reason,
2048 "events append queue rejected loss-tolerant batch"
2049 );
2050 }
2051 }
2052
2053 async fn round_trip(&self, request: &EventsRequest) -> StorageResult<EventsResponse> {
2058 let op = "events daemon round-trip";
2059 let payload = serde_json::to_vec(request).map_err(|error| StorageError::Serialization {
2060 capability: khive_storage::StorageCapability::Events,
2061 message: format!("events request serialization failed: {error}"),
2062 })?;
2063
2064 let attempt = async {
2065 let mut stream = connect_verified(&self.socket_path)
2066 .await
2067 .map_err(|error| ("connect", error))?;
2068 write_frame(&mut stream, &payload)
2069 .await
2070 .map_err(|error| ("write", error))?;
2071 read_frame(&mut stream)
2072 .await
2073 .map_err(|error| ("read", error))
2074 };
2075 let bytes = match tokio::time::timeout(REQUEST_TIMEOUT, attempt).await {
2076 Ok(Ok(bytes)) => bytes,
2077 Ok(Err(("read", error))) => {
2078 return Err(StorageError::Serialization {
2079 capability: khive_storage::StorageCapability::Events,
2080 message: format!(
2081 "events daemon closed or broke the response after connection at {}: {error}",
2082 self.socket_path.display()
2083 ),
2084 });
2085 }
2086 Ok(Err((stage, error))) => {
2087 return Err(StorageError::Pool {
2088 operation: op.into(),
2089 message: format!(
2090 "events daemon {stage} failed at {}: {error}",
2091 self.socket_path.display()
2092 ),
2093 });
2094 }
2095 Err(_elapsed) => {
2096 return Err(StorageError::Timeout {
2097 operation: op.into(),
2098 });
2099 }
2100 };
2101
2102 serde_json::from_slice::<EventsResponse>(&bytes).map_err(|error| {
2103 StorageError::Serialization {
2104 capability: khive_storage::StorageCapability::Events,
2105 message: format!("events response deserialization failed: {error}"),
2106 }
2107 })
2108 }
2109}
2110
2111#[cfg(unix)]
2118fn drain_dropped_queue(
2119 rx: &mut tokio::sync::mpsc::Receiver<QueuedAppend>,
2120 counters: &ForwardingCounters,
2121) {
2122 let mut dropped_batches = 0u64;
2123 let mut dropped_events = 0u64;
2124 while let Ok(batch) = rx.try_recv() {
2125 dropped_batches += 1;
2126 dropped_events += batch.count;
2127 }
2129 if dropped_batches > 0 {
2130 counters
2131 .dropped_batches
2132 .fetch_add(dropped_batches, Ordering::Relaxed);
2133 counters
2134 .dropped_events
2135 .fetch_add(dropped_events, Ordering::Relaxed);
2136 tracing::warn!(
2137 dropped_batches,
2138 dropped_events,
2139 "events forwarder shutting down; dropping queued loss-tolerant batches"
2140 );
2141 }
2142}
2143
2144#[cfg(unix)]
2149async fn run_forwarder(
2150 socket_path: PathBuf,
2151 mut rx: tokio::sync::mpsc::Receiver<QueuedAppend>,
2152 counters: Arc<ForwardingCounters>,
2153 outage_logged: Arc<AtomicBool>,
2154 delivery_timeout: Duration,
2155 shutdown: tokio_util::sync::CancellationToken,
2156) {
2157 let mut conn: Option<UnixStream> = None;
2165 loop {
2166 let batch = tokio::select! {
2167 _ = shutdown.cancelled() => {
2168 drain_dropped_queue(&mut rx, &counters);
2169 break;
2170 },
2171 received = rx.recv() => match received {
2172 Some(batch) => batch,
2173 None => break,
2174 },
2175 };
2176 let count = batch.count;
2177 let payload = match serde_json::to_vec(&batch.request) {
2178 Ok(payload) => payload,
2179 Err(error) => {
2180 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2181 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2182 tracing::error!(error = %error, "events forwarder serialization failed; batch dropped");
2183 continue;
2184 }
2185 };
2186
2187 let delivered = tokio::select! {
2195 _ = shutdown.cancelled() => {
2196 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2197 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2198 tracing::warn!(
2199 dropped_events = count,
2200 "events forwarder shutting down mid-delivery; in-flight batch dropped"
2201 );
2202 drain_dropped_queue(&mut rx, &counters);
2203 break;
2204 },
2205 outcome = tokio::time::timeout(
2206 delivery_timeout,
2207 deliver_batch(&socket_path, &mut conn, &payload),
2208 ) => match outcome {
2209 Ok(delivered) => delivered,
2210 Err(_elapsed) => {
2211 conn = None;
2212 false
2213 }
2214 },
2215 };
2216 if delivered {
2217 counters.forwarded_batches.fetch_add(1, Ordering::Relaxed);
2218 counters
2219 .forwarded_events
2220 .fetch_add(count, Ordering::Relaxed);
2221 if outage_logged.swap(false, Ordering::Relaxed) {
2222 tracing::info!("events daemon reachable again; forwarding resumed");
2223 }
2224 } else {
2225 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2226 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2227 if !outage_logged.swap(true, Ordering::Relaxed) {
2228 tracing::warn!(
2229 socket = %socket_path.display(),
2230 "events daemon unreachable; dropping loss-tolerant events until it returns"
2231 );
2232 }
2233 tokio::select! {
2234 _ = shutdown.cancelled() => {
2235 drain_dropped_queue(&mut rx, &counters);
2236 break;
2237 },
2238 _ = tokio::time::sleep(FORWARDER_BACKOFF) => {}
2239 }
2240 }
2241 }
2242}
2243
2244#[cfg(unix)]
2249async fn deliver_batch(socket_path: &Path, conn: &mut Option<UnixStream>, payload: &[u8]) -> bool {
2250 for _attempt in 0..2u8 {
2251 if conn.is_none() {
2252 match connect_verified(socket_path).await {
2253 Ok(stream) => *conn = Some(stream),
2254 Err(_) => return false,
2255 }
2256 }
2257 let stream = conn.as_mut().expect("connection populated above");
2258 let ok = async {
2259 write_frame(stream, payload).await?;
2260 let bytes = read_frame(stream).await?;
2261 std::io::Result::Ok(bytes)
2262 }
2263 .await;
2264 match ok {
2265 Ok(bytes) => {
2266 return !matches!(
2267 serde_json::from_slice::<EventsResponse>(&bytes),
2268 Ok(EventsResponse::Error { .. }) | Err(_)
2269 );
2270 }
2271 Err(_) => {
2272 *conn = None;
2275 }
2276 }
2277 }
2278 false
2279}
2280
2281#[cfg(unix)]
2290#[derive(Debug)]
2291pub struct ForwardingEventStore {
2292 namespace: String,
2293 client: Arc<EventsSplitClient>,
2294}
2295
2296#[cfg(unix)]
2297impl ForwardingEventStore {
2298 pub fn new(namespace: impl Into<String>, client: Arc<EventsSplitClient>) -> Self {
2299 Self {
2300 namespace: namespace.into(),
2301 client,
2302 }
2303 }
2304
2305 fn unexpected(&self, op: &'static str, response: EventsResponse) -> StorageError {
2306 match response {
2307 EventsResponse::Error {
2308 message,
2309 retryable,
2310 writer_task_failure,
2311 } => {
2312 if let Some(failure) = writer_task_failure {
2313 match failure {
2314 WireWriterTaskFailure::RequestFailed { request_state } => {
2315 let source = if retryable {
2316 StorageError::Pool {
2317 operation: op.into(),
2318 message,
2319 }
2320 } else {
2321 StorageError::InvalidInput {
2322 capability: khive_storage::StorageCapability::Events,
2323 operation: op.into(),
2324 message,
2325 }
2326 };
2327 StorageError::WriterTaskRequestFailed {
2328 request_state: request_state.into(),
2329 source: Box::new(source),
2330 }
2331 }
2332 WireWriterTaskFailure::TaskTerminated { request_state } => {
2333 StorageError::WriterTaskTerminated {
2334 request_state: request_state.into(),
2335 }
2336 }
2337 }
2338 } else if retryable {
2339 StorageError::Pool {
2340 operation: op.into(),
2341 message,
2342 }
2343 } else {
2344 StorageError::InvalidInput {
2345 capability: khive_storage::StorageCapability::Events,
2346 operation: op.into(),
2347 message,
2348 }
2349 }
2350 }
2351 other => StorageError::Serialization {
2352 capability: khive_storage::StorageCapability::Events,
2353 message: format!("events daemon returned mismatched response for {op}: {other:?}"),
2354 },
2355 }
2356 }
2357}
2358
2359#[cfg(unix)]
2360#[async_trait]
2361impl EventStore for ForwardingEventStore {
2362 async fn append_event(&self, event: Event) -> StorageResult<()> {
2363 self.client.enqueue(&self.namespace, vec![event]);
2364 Ok(())
2365 }
2366
2367 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2368 let attempted = events.len() as u64;
2369 self.client.enqueue(&self.namespace, events);
2370 Ok(BatchWriteSummary {
2373 attempted,
2374 affected: attempted,
2375 ..BatchWriteSummary::default()
2376 })
2377 }
2378
2379 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2380 let request = EventsRequest::GetEvent {
2381 protocol_version: EVENTS_PROTOCOL_VERSION,
2382 namespace: self.namespace.clone(),
2383 id,
2384 };
2385 match self.client.round_trip(&request).await? {
2386 EventsResponse::Event { event } => Ok(event),
2387 other => Err(self.unexpected("get_event", other)),
2388 }
2389 }
2390
2391 async fn query_events(
2392 &self,
2393 filter: EventFilter,
2394 page: PageRequest,
2395 ) -> StorageResult<Page<Event>> {
2396 let request = EventsRequest::QueryEvents {
2397 protocol_version: EVENTS_PROTOCOL_VERSION,
2398 namespace: self.namespace.clone(),
2399 filter,
2400 page,
2401 };
2402 match self.client.round_trip(&request).await? {
2403 EventsResponse::Pageful { page } => Ok(page),
2404 other => Err(self.unexpected("query_events", other)),
2405 }
2406 }
2407
2408 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2409 let request = EventsRequest::CountEvents {
2410 protocol_version: EVENTS_PROTOCOL_VERSION,
2411 namespace: self.namespace.clone(),
2412 filter,
2413 };
2414 match self.client.round_trip(&request).await? {
2415 EventsResponse::Count { count } => Ok(count),
2416 other => Err(self.unexpected("count_events", other)),
2417 }
2418 }
2419
2420 async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
2421 crate::event_page::validate_page_query(&query)?;
2422 let request = EventsRequest::QueryEventPage {
2423 protocol_version: EVENTS_PROTOCOL_VERSION,
2424 namespace: self.namespace.clone(),
2425 query: query.clone(),
2426 };
2427 match self.client.round_trip(&request).await? {
2428 EventsResponse::EventPageWindow { window } => {
2429 crate::event_page::validate_window(&query, Some(&self.namespace), &window)?;
2430 Ok(window)
2431 }
2432 error @ EventsResponse::Error { .. } => Err(self.unexpected("query_event_page", error)),
2433 _ => Err(crate::event_page::page_error(
2434 "events page response kind mismatch",
2435 )),
2436 }
2437 }
2438
2439 fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2440 self.client.preflight_store.preflight_event(event)
2442 }
2443
2444 async fn append_events_idempotent(
2445 &self,
2446 events: Vec<Event>,
2447 ) -> StorageResult<IdempotentEventBatchResult> {
2448 let request = EventsRequest::AppendEventsIdempotent {
2449 protocol_version: EVENTS_PROTOCOL_VERSION,
2450 namespace: self.namespace.clone(),
2451 events,
2452 };
2453 match self.client.round_trip(&request).await? {
2454 EventsResponse::Idempotent { result } => Ok(result),
2455 other => Err(self.unexpected("append_events_idempotent", other)),
2456 }
2457 }
2458
2459 fn supports_idempotent_audit_batch(&self) -> bool {
2460 true
2461 }
2462}
2463
2464pub struct SplitEventStore {
2481 legacy: Arc<dyn EventStore>,
2482 lane: Arc<dyn EventStore>,
2483}
2484
2485impl std::fmt::Debug for SplitEventStore {
2486 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2487 f.debug_struct("SplitEventStore").finish_non_exhaustive()
2488 }
2489}
2490
2491impl SplitEventStore {
2492 pub const MAX_MERGED_WINDOW_ROWS: u64 = 100_000;
2500
2501 pub fn new(legacy: Arc<dyn EventStore>, lane: Arc<dyn EventStore>) -> Self {
2502 Self { legacy, lane }
2503 }
2504}
2505
2506#[async_trait]
2507impl EventStore for SplitEventStore {
2508 async fn append_event(&self, event: Event) -> StorageResult<()> {
2509 self.legacy.append_event(event).await
2510 }
2511
2512 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2513 self.legacy.append_events(events).await
2514 }
2515
2516 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2517 if let Some(event) = self.legacy.get_event(id).await? {
2520 return Ok(Some(event));
2521 }
2522 self.lane.get_event(id).await
2523 }
2524
2525 async fn query_events(
2526 &self,
2527 filter: EventFilter,
2528 page: PageRequest,
2529 ) -> StorageResult<Page<Event>> {
2530 let window = page.offset.saturating_add(u64::from(page.limit));
2544 if window > Self::MAX_MERGED_WINDOW_ROWS {
2545 return Err(StorageError::InvalidInput {
2546 capability: khive_storage::StorageCapability::Events,
2547 operation: "query_events".into(),
2548 message: format!(
2549 "offset+limit ({window}) exceeds the merged event plane's window bound \
2550 of {}; page deep windows with a `before` cursor at offset 0, or narrow \
2551 the filter",
2552 Self::MAX_MERGED_WINDOW_ROWS
2553 ),
2554 });
2555 }
2556 let prefix = PageRequest {
2557 offset: 0,
2558 limit: window.min(u64::from(u32::MAX)) as u32,
2559 };
2560 let legacy = self
2561 .legacy
2562 .query_events(filter.clone(), prefix.clone())
2563 .await?;
2564 let lane = self.lane.query_events(filter, prefix).await?;
2565 let total = match (legacy.total, lane.total) {
2566 (Some(a), Some(b)) => Some(a + b),
2567 _ => None,
2568 };
2569 let mut items = legacy.items;
2570 items.extend(lane.items);
2571 items.sort_by(|a, b| {
2572 b.created_at
2573 .cmp(&a.created_at)
2574 .then_with(|| b.id.cmp(&a.id))
2575 });
2576 let items = items
2577 .into_iter()
2578 .skip(page.offset as usize)
2579 .take(page.limit as usize)
2580 .collect();
2581 Ok(Page { items, total })
2582 }
2583
2584 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2585 let legacy = self.legacy.count_events(filter.clone()).await?;
2586 let lane = self.lane.count_events(filter).await?;
2587 Ok(legacy + lane)
2588 }
2589
2590 async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
2591 crate::event_page::split_page(self.legacy.as_ref(), self.lane.as_ref(), query).await
2592 }
2593
2594 fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2595 self.lane.preflight_event(event)
2598 }
2599
2600 async fn append_events_idempotent(
2601 &self,
2602 events: Vec<Event>,
2603 ) -> StorageResult<IdempotentEventBatchResult> {
2604 if events.is_empty() {
2614 return self.lane.append_events_idempotent(events).await;
2615 }
2616 let ids: Vec<Uuid> = events.iter().map(|event| event.id).collect();
2617 let probe_limit = u32::try_from(ids.len()).map_err(|_| StorageError::InvalidInput {
2618 capability: khive_storage::StorageCapability::Events,
2619 operation: "append_events_idempotent".into(),
2620 message: format!(
2621 "batch of {} rows exceeds the legacy probe window",
2622 ids.len()
2623 ),
2624 })?;
2625 let existing = self
2626 .legacy
2627 .query_events(
2628 EventFilter {
2629 ids,
2630 ..EventFilter::default()
2631 },
2632 PageRequest {
2633 offset: 0,
2634 limit: probe_limit,
2635 },
2636 )
2637 .await?;
2638 let legacy_resident: std::collections::HashSet<Uuid> =
2639 existing.items.iter().map(|event| event.id).collect();
2640 if legacy_resident.is_empty() {
2641 return self.lane.append_events_idempotent(events).await;
2642 }
2643 let mut legacy_rows = Vec::new();
2644 let mut lane_rows = Vec::new();
2645 let mut routed_to_legacy = Vec::with_capacity(events.len());
2646 for event in events {
2647 if legacy_resident.contains(&event.id) {
2648 routed_to_legacy.push(true);
2649 legacy_rows.push(event);
2650 } else {
2651 routed_to_legacy.push(false);
2652 lane_rows.push(event);
2653 }
2654 }
2655 let legacy_expected = legacy_rows.len();
2656 let lane_expected = lane_rows.len();
2657 let legacy_result = self.legacy.append_events_idempotent(legacy_rows).await?;
2658 let lane_result = if lane_expected == 0 {
2659 IdempotentEventBatchResult { rows: Vec::new() }
2660 } else {
2661 self.lane.append_events_idempotent(lane_rows).await?
2662 };
2663 if legacy_result.rows.len() != legacy_expected || lane_result.rows.len() != lane_expected {
2664 return Err(StorageError::Driver {
2665 capability: khive_storage::StorageCapability::Events,
2666 operation: "append_events_idempotent".into(),
2667 source: format!(
2668 "idempotent sub-batch result length mismatch: legacy {}/{}, lane {}/{}",
2669 legacy_result.rows.len(),
2670 legacy_expected,
2671 lane_result.rows.len(),
2672 lane_expected
2673 )
2674 .into(),
2675 });
2676 }
2677 let mut legacy_iter = legacy_result.rows.into_iter();
2678 let mut lane_iter = lane_result.rows.into_iter();
2679 let rows = routed_to_legacy
2680 .into_iter()
2681 .map(|to_legacy| {
2682 if to_legacy {
2683 legacy_iter.next().expect("length checked above")
2684 } else {
2685 lane_iter.next().expect("length checked above")
2686 }
2687 })
2688 .collect();
2689 Ok(IdempotentEventBatchResult { rows })
2690 }
2691
2692 fn supports_idempotent_audit_batch(&self) -> bool {
2693 self.lane.supports_idempotent_audit_batch() && self.legacy.supports_idempotent_audit_batch()
2697 }
2698}
2699
2700#[cfg(test)]
2701mod tests {
2702 use super::*;
2703
2704 #[test]
2705 fn fixture_registry_cleanup_preserves_other_roots() {
2706 let live_dir = tempfile::tempdir().unwrap();
2707 let _live_guard = TestRegistryGuard::new(live_dir.path());
2708 let live_path = live_dir.path().join("events.db");
2709 let live_backend = direct_backend_for(&live_path).unwrap();
2710
2711 let finished_dir = tempfile::tempdir().unwrap();
2712 let finished_backend = {
2713 let _finished_guard = TestRegistryGuard::new(finished_dir.path());
2714 let backend = direct_backend_for(&finished_dir.path().join("events.db")).unwrap();
2715 let weak = Arc::downgrade(&backend);
2716 drop(backend);
2717 assert!(
2718 weak.upgrade().is_some(),
2719 "registry retains the live fixture"
2720 );
2721 weak
2722 };
2723
2724 assert!(
2725 finished_backend.upgrade().is_none(),
2726 "FD_FIXTURE_REGISTRY_RELEASED"
2727 );
2728 let same_live_backend = direct_backend_for(&live_path).unwrap();
2729 assert!(
2730 Arc::ptr_eq(&live_backend, &same_live_backend),
2731 "FD_FIXTURE_OTHER_ROOT_RETAINED"
2732 );
2733 }
2734
2735 #[cfg(unix)]
2736 mod target_filter_tests {
2737 include!("events_split_target_tests.rs");
2738 }
2739
2740 include!("events_sidecar_path_tests.rs");
2745
2746 #[cfg(unix)]
2747 #[test]
2748 fn deep_symlink_chains_within_the_kernel_bound_derive_the_target_sidecar() {
2749 let dir = tempfile::tempdir().unwrap();
2754 let real = dir.path().join("real.db");
2755 let mut prev = real.clone();
2756 for i in 0..35 {
2757 let link = dir.path().join(format!("hop{i}.db"));
2758 std::os::unix::fs::symlink(&prev, &link).unwrap();
2759 prev = link;
2760 }
2761 assert_eq!(
2762 events_db_path_beside(&prev),
2763 events_db_path_beside(&real),
2764 "a 35-hop dangling chain must derive the target's sidecar"
2765 );
2766 assert_ne!(
2768 events_db_path_beside(&dir.path().join("unrelated.db")),
2769 events_db_path_beside(&real)
2770 );
2771 }
2772
2773 include!("events_sidecar_hardening_tests.rs");
2774
2775 #[cfg(unix)]
2776 fn write_owner_only_test_file(path: &Path) {
2777 use std::os::unix::fs::PermissionsExt;
2778 std::fs::write(path, b"").unwrap();
2779 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).unwrap();
2780 }
2781
2782 #[cfg(unix)]
2783 #[test]
2784 fn unopened_check_refuses_same_mode_regular_file_replacement() {
2785 for target_index in 0..3 {
2789 let dir = tempfile::tempdir().unwrap();
2790 let db = dir.path().join("events.db");
2791 let targets = events_db_targets(&db);
2792 for path in &targets {
2793 write_owner_only_test_file(path);
2794 }
2795 let before_open = harden_events_db_sidecars(&db).unwrap();
2796 let replacement = dir.path().join("replacement");
2797 write_owner_only_test_file(&replacement);
2798 let replacement_id = EventsFileIdentity::from_metadata(
2799 &std::fs::symlink_metadata(&replacement).unwrap(),
2800 );
2801 assert_ne!(before_open[target_index], Some(replacement_id));
2802 let _connection = rusqlite::Connection::open(&db).unwrap();
2805 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2806 std::fs::rename(&replacement, &targets[target_index]).unwrap();
2807 let error = verify_events_db_owner_only_unopened(&db, &before_open)
2808 .unwrap_err()
2809 .to_string();
2810 assert!(error.contains("changed identity"), "{error}");
2811 assert!(
2812 error.contains(&targets[target_index].display().to_string()),
2813 "{error}"
2814 );
2815 }
2816 }
2817
2818 #[cfg(unix)]
2819 #[test]
2820 fn unopened_check_accepts_sidecars_missing_at_check_time() {
2821 let dir = tempfile::tempdir().unwrap();
2822 let db = dir.path().join("events.db");
2823 let targets = events_db_targets(&db);
2824 for path in &targets {
2825 write_owner_only_test_file(path);
2826 }
2827 let before_open = harden_events_db_sidecars(&db).unwrap();
2828 let _connection = rusqlite::Connection::open(&db).unwrap();
2829 for path in targets.iter().skip(1) {
2830 std::fs::remove_file(path).unwrap();
2831 verify_events_db_owner_only_unopened(&db, &before_open)
2832 .expect("SQLite sidecars may disappear without losing the main database");
2833 }
2834 }
2835
2836 #[cfg(unix)]
2837 #[test]
2838 fn unopened_check_refuses_main_database_removed_after_sqlite_open() {
2839 let dir = tempfile::tempdir().unwrap();
2840 let db = dir.path().join("events.db");
2841 ensure_events_db_owner_only(&db).unwrap();
2842 let before_open = harden_events_db_sidecars(&db).unwrap();
2843 let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2844 backend.events().unwrap();
2845 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2846 std::fs::remove_file(&db).unwrap();
2847 let error = verify_events_db_owner_only_unopened(&db, &before_open)
2850 .unwrap_err()
2851 .to_string();
2852 assert!(
2853 error.contains("main database") && error.contains("disappeared"),
2854 "{error}"
2855 );
2856 }
2857
2858 #[cfg(unix)]
2859 #[test]
2860 fn unopened_check_accepts_fresh_sqlite_sidecars_and_unchanged_main() {
2861 let dir = tempfile::tempdir().unwrap();
2862 let db = dir.path().join("new.events.db");
2863 assert!(!db.exists());
2864 ensure_events_db_owner_only(&db).unwrap();
2865 let before_open = harden_events_db_sidecars(&db).unwrap();
2866 assert!(before_open[0].is_some());
2867 assert_eq!(&before_open[1..], &[None, None]);
2868 let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2869 backend.events().unwrap();
2870 for path in events_db_targets(&db).iter().skip(1) {
2871 assert!(path.exists(), "SQLite creates {}", path.display());
2872 }
2873 verify_events_db_owner_only_unopened(&db, &before_open)
2874 .expect("new SQLite sidecars have no pre-open identity to contradict");
2875 }
2876
2877 #[test]
2878 fn socket_derives_beside_the_sidecar() {
2879 let db = events_db_path_beside(Path::new("/data/khive.db"));
2880 let socket = events_socket_path_beside(&db);
2881 assert!(socket.ends_with("khive.db.events.sock"), "got {socket:?}");
2882 }
2883
2884 #[test]
2885 fn bare_relative_main_db_yields_an_absolute_sidecar() {
2886 let path = events_db_path_beside(Path::new("khive.db"));
2890 assert!(path.is_absolute(), "got relative {path:?}");
2891 assert!(path.ends_with("khive.db.events.db"), "got {path:?}");
2892 }
2893
2894 #[cfg(unix)]
2895 #[tokio::test]
2896 async fn over_cap_query_page_limit_is_refused_before_materialization() {
2897 let backend = Arc::new(StorageBackend::memory().unwrap());
2902 let stores: NamespaceStores =
2903 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
2904 let request = |limit: u32| EventsRequest::QueryEvents {
2905 protocol_version: EVENTS_PROTOCOL_VERSION,
2906 namespace: "local".to_string(),
2907 filter: EventFilter::default(),
2908 page: PageRequest { offset: 0, limit },
2909 };
2910 let refused =
2911 dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS + 1), &backend, &stores)
2912 .await;
2913 match refused {
2914 EventsResponse::Error {
2915 message,
2916 retryable,
2917 writer_task_failure,
2918 } => {
2919 assert!(
2920 message.contains(&MAX_QUERY_EVENTS_PAGE_ROWS.to_string()),
2921 "refusal must name the cap: {message}"
2922 );
2923 assert!(!retryable, "an over-cap page is not transient");
2924 assert!(writer_task_failure.is_none());
2925 }
2926 other => panic!("over-cap query must be refused, got {other:?}"),
2927 }
2928 let at_cap =
2931 dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS), &backend, &stores).await;
2932 assert!(
2933 matches!(at_cap, EventsResponse::Pageful { .. }),
2934 "at-cap query must reach the store, got {at_cap:?}"
2935 );
2936 }
2937
2938 #[cfg(unix)]
2939 #[test]
2940 fn side_effects_unknown_state_crosses_the_wire_as_terminal_failure() {
2941 let response = storage_error_response(&StorageError::WriterTaskTerminated {
2942 request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
2943 });
2944 assert!(
2945 matches!(
2946 response,
2947 EventsResponse::Error {
2948 writer_task_failure: Some(WireWriterTaskFailure::TaskTerminated {
2949 request_state: WireWriterTaskState::SideEffectsUnknown,
2950 }),
2951 ..
2952 }
2953 ),
2954 "an unknown-commit-state termination must carry its state on the wire"
2955 );
2956 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
2958 assert!(matches!(
2959 busy,
2960 EventsResponse::Error {
2961 writer_task_failure: None,
2962 ..
2963 }
2964 ));
2965 }
2966
2967 #[cfg(unix)]
2968 #[test]
2969 fn proven_rollback_crosses_the_wire_without_claiming_task_termination() {
2970 let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
2971 request_state: khive_storage::WriterTaskRequestState::TransactionRolledBack,
2972 source: Box::new(StorageError::Pool {
2973 operation: "writer_task_commit".into(),
2974 message: "commit refused".into(),
2975 }),
2976 });
2977 assert!(matches!(
2978 response,
2979 EventsResponse::Error {
2980 writer_task_failure: Some(WireWriterTaskFailure::RequestFailed {
2981 request_state: WireWriterTaskState::TransactionRolledBack,
2982 }),
2983 ..
2984 }
2985 ));
2986 }
2987
2988 #[test]
2989 fn error_frames_without_writer_failure_still_parse() {
2990 let bytes = br#"{"kind":"error","message":"boom","retryable":true}"#;
2994 let parsed: EventsResponse = serde_json::from_slice(bytes).expect("stateless frame parses");
2995 assert!(matches!(
2996 parsed,
2997 EventsResponse::Error {
2998 retryable: true,
2999 writer_task_failure: None,
3000 ..
3001 }
3002 ));
3003 }
3004
3005 fn split_retry_event(verb: &str) -> Event {
3006 Event::new(
3007 "test",
3008 verb,
3009 khive_types::EventKind::RecallExecuted,
3010 khive_types::SubstrateKind::Note,
3011 "agent:test",
3012 )
3013 }
3014
3015 fn store_pair(dir: &Path) -> (Arc<dyn EventStore>, Arc<dyn EventStore>) {
3016 let legacy = direct_backend_for(&dir.join("legacy.db"))
3017 .expect("legacy backend")
3018 .events_for_namespace("test")
3019 .expect("legacy store");
3020 let lane = direct_backend_for(&dir.join("lane.db"))
3021 .expect("lane backend")
3022 .events_for_namespace("test")
3023 .expect("lane store");
3024 (legacy, lane)
3025 }
3026
3027 #[tokio::test]
3032 async fn idempotent_retry_of_legacy_resident_rows_does_not_duplicate() {
3033 use khive_storage::event::EventAppendDisposition;
3034
3035 let dir = tempfile::tempdir().unwrap();
3036 let _registry_guard = TestRegistryGuard::new(dir.path());
3037 let (legacy, lane) = store_pair(dir.path());
3038
3039 let resident = split_retry_event("recall");
3040 legacy
3041 .append_events_idempotent(vec![resident.clone()])
3042 .await
3043 .expect("pre-cutover landing");
3044
3045 let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3046 let fresh = split_retry_event("search");
3047 let result = split
3048 .append_events_idempotent(vec![resident.clone(), fresh.clone()])
3049 .await
3050 .expect("mixed retry batch");
3051 assert_eq!(
3052 result.rows,
3053 vec![
3054 EventAppendDisposition::AlreadyPresentIdentical,
3055 EventAppendDisposition::Inserted,
3056 ],
3057 "input order must be preserved across the two sub-batches"
3058 );
3059 assert_eq!(
3060 lane.count_events(EventFilter::default()).await.unwrap(),
3061 1,
3062 "the legacy-resident row must not reach the lane"
3063 );
3064 assert_eq!(
3065 split.count_events(EventFilter::default()).await.unwrap(),
3066 2,
3067 "merged count must not double-count the retried id"
3068 );
3069
3070 let mut mutated = resident.clone();
3073 mutated.verb = "other".to_string();
3074 let conflict = split
3075 .append_events_idempotent(vec![mutated])
3076 .await
3077 .expect("conflicting retry");
3078 assert_eq!(
3079 conflict.rows,
3080 vec![EventAppendDisposition::IdentityConflict]
3081 );
3082 assert_eq!(split.count_events(EventFilter::default()).await.unwrap(), 2);
3083 }
3084
3085 #[tokio::test]
3088 async fn merged_offset_window_is_bounded() {
3089 let dir = tempfile::tempdir().unwrap();
3090 let _registry_guard = TestRegistryGuard::new(dir.path());
3091 let (legacy, lane) = store_pair(dir.path());
3092 let split = SplitEventStore::new(legacy, lane);
3093
3094 let err = split
3095 .query_events(
3096 EventFilter::default(),
3097 PageRequest {
3098 offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS,
3099 limit: 1,
3100 },
3101 )
3102 .await
3103 .expect_err("a window past the bound must be refused");
3104 assert!(
3105 matches!(err, StorageError::InvalidInput { .. }),
3106 "got {err:?}"
3107 );
3108 assert!(
3109 err.to_string().contains("before"),
3110 "the refusal must name the cursor remedy: {err}"
3111 );
3112
3113 let at_bound = split
3114 .query_events(
3115 EventFilter::default(),
3116 PageRequest {
3117 offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS - 1,
3118 limit: 1,
3119 },
3120 )
3121 .await
3122 .expect("a window at the bound is admitted");
3123 assert!(at_bound.items.is_empty());
3124 }
3125
3126 #[cfg(unix)]
3127 #[tokio::test]
3128 async fn stateless_non_retryable_error_maps_terminal_not_writer_terminated() {
3129 let client = EventsSplitClient::new(std::path::PathBuf::from("/tmp/never-bound.sock"))
3136 .expect("client builds");
3137 let store = ForwardingEventStore::new("test", client);
3138 let bytes = br#"{"kind":"error","message":"refused","retryable":false}"#;
3139 let parsed: EventsResponse = serde_json::from_slice(bytes).expect("frame parses");
3140 assert!(matches!(
3141 store.unexpected("append", parsed),
3142 StorageError::InvalidInput { .. }
3143 ));
3144 let with_state = br#"{"kind":"error","message":"died","retryable":false,"writer_task_failure":{"kind":"task_terminated","request_state":"side_effects_unknown"}}"#;
3148 let parsed: EventsResponse =
3149 serde_json::from_slice(with_state).expect("stateful frame parses");
3150 assert!(matches!(
3151 store.unexpected("append", parsed),
3152 StorageError::WriterTaskTerminated { .. }
3153 ));
3154 }
3155
3156 #[cfg(unix)]
3157 #[tokio::test]
3158 async fn writer_task_states_cross_the_wire_verbatim() {
3159 use khive_storage::WriterTaskRequestState as S;
3165 let dir = tempfile::tempdir().unwrap();
3166 let client =
3167 EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3168 let store = ForwardingEventStore::new("test", client);
3169 for state in [
3170 S::NotStarted,
3171 S::TransactionRolledBack,
3172 S::SideEffectsUnknown,
3173 ] {
3174 let response = storage_error_response(&StorageError::WriterTaskTerminated {
3175 request_state: state,
3176 });
3177 let bytes = serde_json::to_vec(&response).unwrap();
3179 let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3180 let err = store.unexpected("append_events_idempotent", parsed);
3181 assert!(
3182 matches!(
3183 err,
3184 StorageError::WriterTaskTerminated { request_state } if request_state == state
3185 ),
3186 "state {state:?} did not survive the socket: got {err:?}"
3187 );
3188 }
3189 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3191 assert!(matches!(
3192 busy,
3193 EventsResponse::Error {
3194 writer_task_failure: None,
3195 ..
3196 }
3197 ));
3198 }
3199
3200 #[cfg(unix)]
3201 #[tokio::test]
3202 async fn proven_rollback_round_trips_as_non_terminal_request_failure() {
3203 use khive_storage::WriterTaskRequestState as S;
3204
3205 let dir = tempfile::tempdir().unwrap();
3206 let client =
3207 EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3208 let store = ForwardingEventStore::new("test", client);
3209 let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
3210 request_state: S::TransactionRolledBack,
3211 source: Box::new(StorageError::Pool {
3212 operation: "writer_task_commit".into(),
3213 message: "commit refused".into(),
3214 }),
3215 });
3216 let bytes = serde_json::to_vec(&response).unwrap();
3217 let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3218 let err = store.unexpected("append_events_idempotent", parsed);
3219 assert!(
3220 matches!(
3221 err,
3222 StorageError::WriterTaskRequestFailed {
3223 request_state: S::TransactionRolledBack,
3224 ..
3225 }
3226 ),
3227 "proven rollback must not reconstruct as a terminal writer: {err:?}"
3228 );
3229 }
3230
3231 #[cfg(unix)]
3232 #[test]
3233 fn direct_backend_hardens_preexisting_db_and_sidecars() {
3234 use std::os::unix::fs::PermissionsExt;
3235 let dir = tempfile::tempdir().unwrap();
3236 let _registry_guard = TestRegistryGuard::new(dir.path());
3237 let db = dir.path().join("pre-existing.events.db");
3238 let wal = dir.path().join("pre-existing.events.db-wal");
3239 std::fs::write(&db, b"").unwrap();
3240 std::fs::write(&wal, b"").unwrap();
3241 for path in [&db, &wal] {
3242 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o644)).unwrap();
3243 }
3244 assert_eq!(
3246 std::fs::metadata(&db).unwrap().permissions().mode() & 0o777,
3247 0o644
3248 );
3249 direct_backend_for(&db).expect("writable open succeeds");
3250 for path in [&db, &wal] {
3251 assert_eq!(
3252 std::fs::metadata(path).unwrap().permissions().mode() & 0o777,
3253 0o600,
3254 "pre-existing {} must be tightened to owner-only",
3255 path.display()
3256 );
3257 }
3258 }
3259
3260 #[cfg(unix)]
3261 #[test]
3262 fn namespace_store_cache_is_bounded_and_trim_normalized() {
3263 let backend = StorageBackend::memory().unwrap();
3264 let stores: NamespaceStores =
3265 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3266 namespace_store_with_cap(&backend, &stores, "alpha", 1).expect("first store");
3269 namespace_store_with_cap(&backend, &stores, "beta", 1).expect("second store evicts");
3270 {
3271 let map = stores.lock().unwrap();
3272 assert_eq!(map.len(), 1, "cache must not grow past its cap");
3273 assert!(
3274 map.contains_key("beta"),
3275 "the newest namespace must be admitted at the cap"
3276 );
3277 }
3278 namespace_store_with_cap(&backend, &stores, "beta", 1).expect("cached beta");
3281 namespace_store_with_cap(&backend, &stores, " beta ", 1).expect("trimmed spelling");
3282 {
3283 let map = stores.lock().unwrap();
3284 assert_eq!(
3285 map.len(),
3286 1,
3287 "spellings of one namespace must share one entry"
3288 );
3289 assert!(map.contains_key("beta"), "trimmed key is the cache key");
3290 }
3291 }
3292
3293 #[cfg(unix)]
3294 #[tokio::test]
3295 async fn frame_budget_admits_before_allocating_and_releases_after() {
3296 let (mut client, mut server) = UnixStream::pair().expect("socketpair");
3297 let budget = Arc::new(tokio::sync::Semaphore::new(1024));
3298 let payload = vec![7u8; 100];
3299 write_frame(&mut client, &payload).await.expect("write");
3300 let (bytes, permit) = read_frame_budgeted(&mut server, &budget)
3301 .await
3302 .expect("budgeted read");
3303 assert_eq!(bytes, payload);
3304 assert_eq!(budget.available_permits(), 1024 - 100);
3306 drop(permit);
3308 assert_eq!(budget.available_permits(), 1024);
3309 let mut oversized = Vec::from(((crate::daemon::MAX_FRAME_BYTES + 1) as u32).to_be_bytes());
3312 oversized.extend_from_slice(&[0u8; 8]);
3313 use tokio::io::AsyncWriteExt;
3314 client.write_all(&oversized).await.expect("raw prefix");
3315 let err = read_frame_budgeted(&mut server, &budget)
3316 .await
3317 .expect_err("oversized declaration must refuse");
3318 assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
3319 assert_eq!(
3320 budget.available_permits(),
3321 1024,
3322 "refusal must not consume budget"
3323 );
3324 }
3325
3326 #[cfg(unix)]
3327 #[tokio::test]
3328 async fn dispatch_rejects_invalid_wire_namespace() {
3329 let backend = Arc::new(StorageBackend::memory().unwrap());
3330 let stores: NamespaceStores =
3331 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3332 let huge = "n".repeat(64 * 1024);
3336 let response = dispatch_events_request(
3337 EventsRequest::CountEvents {
3338 protocol_version: EVENTS_PROTOCOL_VERSION,
3339 namespace: huge,
3340 filter: Default::default(),
3341 },
3342 &backend,
3343 &stores,
3344 )
3345 .await;
3346 assert!(
3347 matches!(
3348 &response,
3349 EventsResponse::Error {
3350 retryable: false,
3351 writer_task_failure: None,
3352 ..
3353 }
3354 ),
3355 "oversized namespace must be a typed refusal, got {response:?}"
3356 );
3357 assert_eq!(
3358 stores.lock().unwrap().len(),
3359 0,
3360 "a rejected namespace must never enter the cache"
3361 );
3362 let ok = dispatch_events_request(
3364 EventsRequest::CountEvents {
3365 protocol_version: EVENTS_PROTOCOL_VERSION,
3366 namespace: "local".to_string(),
3367 filter: Default::default(),
3368 },
3369 &backend,
3370 &stores,
3371 )
3372 .await;
3373 assert!(
3374 matches!(ok, EventsResponse::Count { .. }),
3375 "valid namespace must dispatch, got {ok:?}"
3376 );
3377 }
3378
3379 include!("events_split_shutdown_tests.rs");
3380
3381 #[cfg(unix)]
3382 #[tokio::test]
3383 async fn forwarder_abandons_hung_delivery() {
3384 let dir = tempfile::tempdir().unwrap();
3388 let socket = dir.path().join("hung.sock");
3389 let listener = tokio::net::UnixListener::bind(&socket).unwrap();
3390 let _server = tokio::spawn(async move {
3392 let mut held = Vec::new();
3393 loop {
3394 if let Ok((stream, _)) = listener.accept().await {
3395 held.push(stream);
3396 }
3397 }
3398 });
3399 let client = EventsSplitClient::new_with_queue_depth_and_delivery_timeout(
3400 socket.clone(),
3401 4,
3402 Duration::from_millis(100),
3403 )
3404 .expect("client builds");
3405 let event = Event::new(
3406 "test",
3407 "noop",
3408 khive_types::EventKind::Audit,
3409 khive_types::SubstrateKind::Event,
3410 "tester",
3411 );
3412 client.enqueue("test", vec![event]);
3413 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
3414 loop {
3415 if client.metrics().dropped_batches >= 1 {
3416 break;
3417 }
3418 assert!(
3419 tokio::time::Instant::now() < deadline,
3420 "forwarder never abandoned the hung delivery: {:?}",
3421 client.metrics()
3422 );
3423 tokio::time::sleep(Duration::from_millis(20)).await;
3424 }
3425 }
3426
3427 #[cfg(unix)]
3428 #[test]
3429 fn wire_retryability_follows_the_storage_classifier() {
3430 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3431 assert!(
3432 matches!(
3433 busy,
3434 EventsResponse::Error {
3435 retryable: true,
3436 ..
3437 }
3438 ),
3439 "transient writer contention must stay retryable across the socket"
3440 );
3441 let terminated = storage_error_response(&StorageError::WriterTaskTerminated {
3444 request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3445 });
3446 assert!(
3447 matches!(
3448 terminated,
3449 EventsResponse::Error {
3450 retryable: false,
3451 ..
3452 }
3453 ),
3454 "a terminated writer must not be reported transient"
3455 );
3456 }
3457
3458 use khive_storage::event::EventAppendDisposition;
3459 use khive_types::{EventKind, SubstrateKind};
3460
3461 fn test_event(namespace: &str) -> Event {
3462 Event::new(
3463 namespace,
3464 "test.verb",
3465 EventKind::Audit,
3466 SubstrateKind::Event,
3467 "actor:test",
3468 )
3469 }
3470
3471 #[cfg(unix)]
3473 async fn boot_daemon(dir: &tempfile::TempDir) -> (PathBuf, PathBuf) {
3474 let db = dir.path().join("events.db");
3475 let socket = dir.path().join("events.sock");
3476 let (db_clone, socket_clone) = (db.clone(), socket.clone());
3477 tokio::spawn(async move {
3478 let _ = run_events_daemon(&db_clone, &socket_clone).await;
3479 });
3480 for _ in 0..100 {
3481 if UnixStream::connect(&socket).await.is_ok() {
3482 return (db, socket);
3483 }
3484 tokio::time::sleep(Duration::from_millis(20)).await;
3485 }
3486 panic!("events daemon did not come up on {}", socket.display());
3487 }
3488
3489 #[tokio::test]
3494 async fn split_store_routes_plain_to_legacy_idempotent_to_lane_and_merges_reads() {
3495 let dir = tempfile::tempdir().expect("tempdir");
3496 let _registry_guard = TestRegistryGuard::new(dir.path());
3497 let legacy_backend =
3498 direct_backend_for(&dir.path().join("legacy.db")).expect("legacy backend");
3499 let lane_backend = direct_backend_for(&dir.path().join("lane.db")).expect("lane backend");
3500 let legacy = legacy_backend
3501 .events_for_namespace("local")
3502 .expect("legacy store");
3503 let lane = lane_backend
3504 .events_for_namespace("local")
3505 .expect("lane store");
3506 let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3507
3508 let plain = test_event("local");
3509 let plain_id = plain.id;
3510 split.append_event(plain).await.expect("plain append");
3511
3512 let audit = test_event("local");
3513 let audit_id = audit.id;
3514 let result = split
3515 .append_events_idempotent(vec![audit])
3516 .await
3517 .expect("idempotent append");
3518 assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3519
3520 assert!(legacy
3522 .get_event(plain_id)
3523 .await
3524 .expect("legacy get")
3525 .is_some());
3526 assert!(lane.get_event(plain_id).await.expect("lane get").is_none());
3527 assert!(legacy
3528 .get_event(audit_id)
3529 .await
3530 .expect("legacy get")
3531 .is_none());
3532 assert!(lane.get_event(audit_id).await.expect("lane get").is_some());
3533
3534 assert!(split
3536 .get_event(plain_id)
3537 .await
3538 .expect("split get")
3539 .is_some());
3540 assert!(split
3541 .get_event(audit_id)
3542 .await
3543 .expect("split get")
3544 .is_some());
3545 assert_eq!(
3546 split
3547 .count_events(EventFilter::default())
3548 .await
3549 .expect("split count"),
3550 2
3551 );
3552 let page = split
3553 .query_events(
3554 EventFilter::default(),
3555 PageRequest {
3556 offset: 0,
3557 limit: 10,
3558 },
3559 )
3560 .await
3561 .expect("split query");
3562 let ids: Vec<Uuid> = page.items.iter().map(|e| e.id).collect();
3563 assert!(ids.contains(&plain_id) && ids.contains(&audit_id));
3564
3565 let mut seen = Vec::new();
3568 for offset in 0..2 {
3569 let page = split
3570 .query_events(EventFilter::default(), PageRequest { offset, limit: 1 })
3571 .await
3572 .expect("windowed query");
3573 assert_eq!(page.items.len(), 1);
3574 seen.push(page.items[0].id);
3575 }
3576 seen.sort();
3577 let mut expected = vec![plain_id, audit_id];
3578 expected.sort();
3579 assert_eq!(seen, expected);
3580 }
3581
3582 #[tokio::test]
3595 async fn raw_sql_consumers_of_the_legacy_events_table_still_see_plain_appends() {
3596 use khive_storage::types::{SqlStatement, SqlValue};
3597
3598 let dir = tempfile::tempdir().expect("tempdir");
3599 let _registry_guard = TestRegistryGuard::new(dir.path());
3600 let legacy_backend =
3601 direct_backend_for(&dir.path().join("legacy-guard.db")).expect("legacy backend");
3602 let lane_backend =
3603 direct_backend_for(&dir.path().join("lane-guard.db")).expect("lane backend");
3604 let split = SplitEventStore::new(
3605 legacy_backend
3606 .events_for_namespace("local")
3607 .expect("legacy store"),
3608 lane_backend
3609 .events_for_namespace("local")
3610 .expect("lane store"),
3611 );
3612
3613 let provenance = test_event("local");
3616 let provenance_id = provenance.id;
3617 split.append_event(provenance).await.expect("plain append");
3618 let audit = test_event("local");
3619 let audit_id = audit.id;
3620 split
3621 .append_events_idempotent(vec![audit])
3622 .await
3623 .expect("idempotent append");
3624
3625 let count_by_raw_sql = |backend: Arc<StorageBackend>, id: Uuid| async move {
3626 let mut reader = backend.sql().reader().await.expect("sql reader");
3627 let rows = reader
3628 .query_all(SqlStatement {
3629 sql: "SELECT actor FROM events WHERE id = ?1".to_string(),
3630 params: vec![SqlValue::Text(id.to_string())],
3631 label: None,
3632 })
3633 .await
3634 .expect("raw events query");
3635 rows.len()
3636 };
3637
3638 assert_eq!(
3641 count_by_raw_sql(Arc::clone(&legacy_backend), provenance_id).await,
3642 1,
3643 "plain appends must stay visible to raw-SQL consumers of the legacy events table"
3644 );
3645 assert_eq!(
3648 count_by_raw_sql(Arc::clone(&legacy_backend), audit_id).await,
3649 0,
3650 "audit-lane rows must not land in the legacy events table"
3651 );
3652 assert_eq!(
3656 count_by_raw_sql(Arc::clone(&lane_backend), audit_id).await,
3657 1,
3658 "the raw query must prove it can find rows where they actually live"
3659 );
3660 }
3661
3662 #[cfg(unix)]
3665 #[tokio::test]
3666 async fn idempotent_append_round_trips_through_the_daemon() {
3667 let dir = tempfile::tempdir().expect("tempdir");
3668 let (_db, socket) = boot_daemon(&dir).await;
3669 let client = EventsSplitClient::new(socket).expect("client");
3670 let store = ForwardingEventStore::new("local", client);
3671
3672 let event = test_event("local");
3673 let id = event.id;
3674 let result = store
3675 .append_events_idempotent(vec![event.clone()])
3676 .await
3677 .expect("idempotent append over socket");
3678 assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3679
3680 let retry = store
3683 .append_events_idempotent(vec![event])
3684 .await
3685 .expect("idempotent retry");
3686 assert_eq!(
3687 retry.rows,
3688 vec![EventAppendDisposition::AlreadyPresentIdentical]
3689 );
3690
3691 let fetched = store.get_event(id).await.expect("get over socket");
3692 assert_eq!(fetched.map(|e| e.id), Some(id));
3693 let count = store
3694 .count_events(EventFilter::default())
3695 .await
3696 .expect("count over socket");
3697 assert_eq!(count, 1);
3698 }
3699
3700 #[cfg(unix)]
3703 #[tokio::test]
3704 async fn fire_and_forget_append_lands_in_the_daemon_store() {
3705 let dir = tempfile::tempdir().expect("tempdir");
3706 let (_db, socket) = boot_daemon(&dir).await;
3707 let client = EventsSplitClient::new(socket).expect("client");
3708 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3709
3710 store
3711 .append_event(test_event("local"))
3712 .await
3713 .expect("append_event is fire-and-forget");
3714
3715 let (count, metrics) = tokio::time::timeout(Duration::from_secs(30), async {
3719 loop {
3720 let count = store
3721 .count_events(EventFilter::default())
3722 .await
3723 .expect("count over socket");
3724 let metrics = client.metrics();
3725 if (count == 1 && metrics.forwarded_events >= 1) || metrics.dropped_events > 0 {
3726 break (count, metrics);
3727 }
3728 tokio::time::sleep(Duration::from_millis(20)).await;
3729 }
3730 })
3731 .await
3732 .expect("forwarder must deliver or report a drop within 30 seconds");
3733 assert_eq!(
3734 metrics.dropped_events, 0,
3735 "forwarder dropped an event before it reached the daemon store"
3736 );
3737 assert_eq!(count, 1, "forwarded event must land in the daemon store");
3738 assert!(metrics.forwarded_events >= 1, "delivery must be counted");
3739 }
3740
3741 #[cfg(unix)]
3745 #[tokio::test]
3746 async fn queue_overflow_drops_and_counts_but_never_errors() {
3747 let dir = tempfile::tempdir().expect("tempdir");
3748 let dead_socket = dir.path().join("nobody-home.sock");
3749 let client = EventsSplitClient::new_with_queue_depth(dead_socket, 2).expect("client");
3750 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3751
3752 for _ in 0..20 {
3753 store
3754 .append_event(test_event("local"))
3755 .await
3756 .expect("append_event must not error under overflow");
3757 }
3758 let metrics = client.metrics();
3759 assert!(
3760 metrics.dropped_events > 0,
3761 "overflow must register in the drop counter, got {metrics:?}"
3762 );
3763 }
3764
3765 #[cfg(unix)]
3769 #[tokio::test]
3770 async fn dead_socket_reads_fail_typed_and_preflight_stays_local() {
3771 let dir = tempfile::tempdir().expect("tempdir");
3772 let dead_socket = dir.path().join("nobody-home.sock");
3773 let client = EventsSplitClient::new(dead_socket).expect("client");
3774 let store = ForwardingEventStore::new("local", client);
3775
3776 let error = store
3777 .get_event(Uuid::new_v4())
3778 .await
3779 .expect_err("read against a dead socket must fail");
3780 assert!(
3781 matches!(error, StorageError::Pool { .. }),
3782 "expected the typed unreachable error, got {error:?}"
3783 );
3784
3785 assert!(store.supports_idempotent_audit_batch());
3788 store
3789 .preflight_event(&test_event("local"))
3790 .expect("offline preflight validates a well-formed event");
3791 }
3792
3793 #[cfg(unix)]
3796 #[tokio::test]
3797 async fn connected_events_peer_closing_without_response_is_not_unreachable() {
3798 let dir = tempfile::tempdir().unwrap();
3799 let socket = dir.path().join("closes-after-request.sock");
3800 let listener = UnixListener::bind(&socket).unwrap();
3801 let server = tokio::spawn(async move {
3802 let (mut stream, _) = listener.accept().await.unwrap();
3803 read_frame(&mut stream)
3804 .await
3805 .expect("complete request arrives");
3806 });
3808 let client = EventsSplitClient::new(socket).unwrap();
3809 let store = ForwardingEventStore::new("local", client);
3810 let error = store.get_event(Uuid::new_v4()).await.unwrap_err();
3811 assert!(
3812 matches!(error, StorageError::Serialization { .. }),
3813 "a post-connect response failure must not be Pool/unreachable: {error:?}"
3814 );
3815 server.await.unwrap();
3816 }
3817
3818 #[cfg(unix)]
3822 #[tokio::test]
3823 async fn query_page_over_frame_cap_returns_typed_size_refusal() {
3824 let dir = tempfile::tempdir().unwrap();
3825 let socket = dir.path().join("large-query.sock");
3826 let backend = Arc::new(StorageBackend::memory().unwrap());
3827 let event_store = backend.events_for_namespace("local").unwrap();
3828 let events: Vec<_> = (0..MAX_QUERY_EVENTS_PAGE_ROWS)
3829 .map(|_| {
3830 test_event("local").with_payload(serde_json::json!({"data": "x".repeat(3 * 1024)}))
3831 })
3832 .collect();
3833 event_store
3834 .append_events(events)
3835 .await
3836 .expect("seed large page");
3837
3838 let listener = UnixListener::bind(&socket).unwrap();
3839 let stores: NamespaceStores =
3840 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3841 let server = tokio::spawn(async move {
3842 let (stream, _) = listener.accept().await.unwrap();
3843 serve_events_conn(
3844 stream,
3845 backend,
3846 stores,
3847 Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES)),
3848 )
3849 .await;
3850 });
3851 let client = EventsSplitClient::new(socket).unwrap();
3852 let store = ForwardingEventStore::new("local", client);
3853 let error = store
3854 .query_events(
3855 EventFilter::default(),
3856 PageRequest {
3857 offset: 0,
3858 limit: MAX_QUERY_EVENTS_PAGE_ROWS,
3859 },
3860 )
3861 .await
3862 .expect_err("the full page cannot fit one frame");
3863 assert!(
3864 matches!(error, StorageError::InvalidInput { .. }),
3865 "expected non-retryable frame-size refusal, got {error:?}"
3866 );
3867 assert!(
3868 error.to_string().contains("response_frame_size_limit"),
3869 "{error}"
3870 );
3871 assert!(
3872 error.to_string().contains("request a narrower page"),
3873 "{error}"
3874 );
3875 server.abort();
3876 }
3877
3878 #[cfg(unix)]
3882 #[tokio::test]
3883 async fn forwarding_queue_enforces_serialized_byte_and_frame_limits_with_drop_metrics() {
3884 let dir = tempfile::tempdir().unwrap();
3885 let socket = dir.path().join("stalled-forwarder.sock");
3886 let listener = UnixListener::bind(&socket).unwrap();
3887 let server = tokio::spawn(async move {
3888 let (stream, _) = listener.accept().await.unwrap();
3889 let _hold_open = stream;
3890 std::future::pending::<()>().await;
3891 });
3892 let event = test_event("local").with_payload(serde_json::json!({"data": "x".repeat(4096)}));
3893 let first = vec![event.clone(), event.clone()];
3894 let request_bytes = serde_json::to_vec(&EventsRequest::AppendEvents {
3895 protocol_version: EVENTS_PROTOCOL_VERSION,
3896 namespace: "local".into(),
3897 events: first.clone(),
3898 })
3899 .unwrap()
3900 .len();
3901 let byte_budget = request_bytes + 64;
3902 let client = EventsSplitClient::new_with_limits_and_delivery_timeout(
3903 socket,
3904 8,
3905 byte_budget,
3906 Duration::from_secs(30),
3907 )
3908 .unwrap();
3909 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3910 store.append_events(first.clone()).await.unwrap();
3911 assert_eq!(client.metrics().queued_bytes, request_bytes);
3912
3913 store.append_events(first).await.unwrap();
3914 let metrics = client.metrics();
3915 assert_eq!(metrics.dropped_batches, 1);
3916 assert_eq!(metrics.dropped_events, 2);
3917 assert_eq!(metrics.queued_bytes, request_bytes);
3918 assert!(metrics.queued_bytes <= byte_budget);
3919
3920 let oversized = test_event("local")
3923 .with_payload(serde_json::json!({"data": "x".repeat(crate::daemon::MAX_FRAME_BYTES)}));
3924 store.append_event(oversized).await.unwrap();
3925 let metrics = client.metrics();
3926 assert_eq!(metrics.dropped_batches, 2);
3927 assert_eq!(metrics.dropped_events, 3);
3928 assert_eq!(metrics.queued_bytes, request_bytes);
3929 server.abort();
3930 }
3931
3932 #[cfg(unix)]
3935 #[tokio::test]
3936 async fn protocol_version_skew_is_a_typed_refusal() {
3937 let dir = tempfile::tempdir().expect("tempdir");
3938 let (_db, socket) = boot_daemon(&dir).await;
3939
3940 let mut stream = UnixStream::connect(&socket).await.expect("connect");
3941 let request = EventsRequest::CountEvents {
3942 protocol_version: EVENTS_PROTOCOL_VERSION + 1,
3943 namespace: "local".into(),
3944 filter: EventFilter::default(),
3945 };
3946 let payload = serde_json::to_vec(&request).expect("serialize");
3947 write_frame(&mut stream, &payload).await.expect("write");
3948 let bytes = read_frame(&mut stream).await.expect("read");
3949 let response: EventsResponse = serde_json::from_slice(&bytes).expect("parse");
3950 match response {
3951 EventsResponse::Error {
3952 message, retryable, ..
3953 } => {
3954 assert!(!retryable, "version skew is not retryable");
3955 assert!(message.contains("protocol version"), "message: {message}");
3956 }
3957 other => panic!("expected a typed refusal, got {other:?}"),
3958 }
3959 }
3960
3961 #[cfg(unix)]
3965 #[test]
3966 fn events_daemon_guard_is_exclusive_then_reusable() {
3967 let dir = tempfile::tempdir().expect("tempdir");
3968 let socket = dir.path().join("events.sock");
3969 let first = match acquire_events_daemon_guard_outcome(&socket) {
3970 EventsDaemonGuardAcquisition::Held(guard) => guard,
3971 other => panic!("first acquire must succeed: {other:?}"),
3972 };
3973 let second = acquire_events_daemon_guard_outcome(&socket);
3974 assert!(
3975 matches!(second, EventsDaemonGuardAcquisition::Contended),
3976 "second acquire must report contention while the first guard is held: {second:?}"
3977 );
3978 drop(first);
3979 const MAX_ATTEMPTS: usize = 100;
3983 for attempt in 1..=MAX_ATTEMPTS {
3984 match acquire_events_daemon_guard_outcome(&socket) {
3985 EventsDaemonGuardAcquisition::Held(_) => return,
3986 EventsDaemonGuardAcquisition::Contended if attempt < MAX_ATTEMPTS => {
3987 std::thread::sleep(Duration::from_millis(10));
3988 }
3989 other => panic!(
3990 "acquire must succeed after release within {MAX_ATTEMPTS} attempts; \
3991 attempt {attempt}: {other:?}"
3992 ),
3993 }
3994 }
3995 unreachable!("the final acquisition attempt returns or reports its refusal");
3996 }
3997
3998 #[cfg(unix)]
3999 #[test]
4000 fn events_daemon_guard_open_failure_is_not_contention() {
4001 let dir = tempfile::tempdir().expect("tempdir");
4002 let clean = dir.path().join("clean.sock");
4003 assert!(matches!(
4004 acquire_events_daemon_guard_outcome(&clean),
4005 EventsDaemonGuardAcquisition::Held(_)
4006 ));
4007
4008 assert!(
4009 try_acquire_events_daemon_guard(&clean).is_some(),
4010 "the public Option entrance must preserve a successful acquisition"
4011 );
4012
4013 let socket = dir.path().join("blocked.sock");
4016 let lock_path = socket.with_extension("lock");
4017 std::fs::create_dir(&lock_path).expect("create directory at lock entry");
4018 let outcome = acquire_events_daemon_guard_outcome(&socket);
4019 assert!(
4020 matches!(outcome, EventsDaemonGuardAcquisition::OpenFailed(_)),
4021 "an actual lock-file open failure must not report contention: {outcome:?}"
4022 );
4023 assert!(
4024 try_acquire_events_daemon_guard(&socket).is_none(),
4025 "the public Option entrance must refuse an actual open failure"
4026 );
4027 assert!(
4028 lock_path.is_dir(),
4029 "the refused entry must remain a directory"
4030 );
4031 }
4032
4033 #[cfg(unix)]
4040 #[tokio::test]
4041 async fn read_only_runtime_never_creates_an_events_db() {
4042 use crate::{KhiveRuntime, Namespace, RuntimeConfig};
4043
4044 let dir = tempfile::tempdir().expect("tempdir");
4045 let _registry_guard = TestRegistryGuard::new(dir.path());
4046 let main_db = dir.path().join("main.db");
4047 drop(
4049 KhiveRuntime::new_for_test(RuntimeConfig {
4050 db_path: Some(main_db.clone()),
4051 ..RuntimeConfig::no_embeddings()
4052 })
4053 .expect("create main db"),
4054 );
4055 let events_db = events_db_path_beside(&main_db);
4056 assert!(!events_db.exists(), "precondition: no events db yet");
4057
4058 khive_storage::test_support::freeze_snapshot_sidecars(&main_db);
4061
4062 let split_config = |db: PathBuf| RuntimeConfig {
4063 db_path: Some(main_db.clone()),
4064 events_split: Some(EventsSplitConfig {
4065 db_path: db,
4066 socket_path: None,
4067 }),
4068 ..RuntimeConfig::no_embeddings()
4069 };
4070
4071 let ro = KhiveRuntime::new_readonly_for_test(split_config(events_db.clone()))
4074 .expect("read-only runtime");
4075 let token = ro.authorize(Namespace::local()).expect("token");
4076 let store = ro.events(&token).expect("events store");
4077 let count = store
4078 .count_events(EventFilter::default())
4079 .await
4080 .expect("count through legacy-only plane");
4081 assert_eq!(count, 0);
4082 assert!(
4083 !events_db.exists(),
4084 "a read-only runtime must not mint the events database"
4085 );
4086
4087 {
4094 let lane_backend = StorageBackend::sqlite_for_test(&events_db).expect("writable lane");
4095 lane_backend
4096 .events_for_namespace("local")
4097 .expect("lane store")
4098 .append_event(test_event("local"))
4099 .await
4100 .expect("seed lane row");
4101 }
4102 khive_storage::test_support::freeze_snapshot_sidecars(&events_db);
4103 let store = ro.events(&token).expect("events store with lane present");
4104 let count = store
4105 .count_events(EventFilter::default())
4106 .await
4107 .expect("merged count");
4108 assert_eq!(count, 1, "the pre-existing lane row must merge into reads");
4109 }
4110
4111 #[tokio::test]
4114 async fn direct_mode_appends_and_reads_without_a_daemon() {
4115 let dir = tempfile::tempdir().expect("tempdir");
4116 let _registry_guard = TestRegistryGuard::new(dir.path());
4117 let db = dir.path().join("events.db");
4118 let backend = direct_backend_for(&db).expect("direct backend");
4119 let store = backend.events_for_namespace("local").expect("store");
4120 store
4121 .append_event(test_event("local"))
4122 .await
4123 .expect("direct append");
4124 let count = store
4125 .count_events(EventFilter::default())
4126 .await
4127 .expect("count");
4128 assert_eq!(count, 1);
4129 }
4130}