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;
49use khive_storage::event::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#[derive(Debug, Clone)]
262pub struct EventsSplitConfig {
263 pub db_path: PathBuf,
266 pub socket_path: Option<PathBuf>,
270}
271
272#[cfg(unix)]
282type ClientMap = std::collections::HashMap<PathBuf, Arc<EventsSplitClient>>;
283type BackendMap = std::collections::HashMap<PathBuf, (bool, Arc<StorageBackend>)>;
287
288#[cfg(unix)]
289fn client_registry() -> &'static std::sync::Mutex<ClientMap> {
290 static REGISTRY: std::sync::OnceLock<std::sync::Mutex<ClientMap>> = std::sync::OnceLock::new();
291 REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
292}
293
294fn direct_backend_registry() -> &'static std::sync::Mutex<BackendMap> {
295 static REGISTRY: std::sync::OnceLock<std::sync::Mutex<BackendMap>> = std::sync::OnceLock::new();
296 REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
297}
298
299#[cfg(any(test, feature = "test-internals"))]
306#[doc(hidden)]
307pub struct TestRegistryGuard<'a> {
308 root: &'a Path,
309 canonical_root: PathBuf,
310}
311
312#[cfg(any(test, feature = "test-internals"))]
313impl<'a> TestRegistryGuard<'a> {
314 pub fn new(root: &'a Path) -> Self {
316 Self {
317 root,
318 canonical_root: std::fs::canonicalize(root).expect("existing unique fixture root"),
319 }
320 }
321
322 fn remove_entries<T>(
323 &self,
324 registry: &std::sync::Mutex<std::collections::HashMap<PathBuf, T>>,
325 ) -> Vec<T> {
326 let mut entries = registry
327 .lock()
328 .unwrap_or_else(std::sync::PoisonError::into_inner);
329 let keys: Vec<_> = entries
330 .keys()
331 .filter(|path| path.starts_with(self.root) || path.starts_with(&self.canonical_root))
332 .cloned()
333 .collect();
334 keys.into_iter()
335 .filter_map(|key| entries.remove(&key))
336 .collect()
337 }
338}
339
340#[cfg(any(test, feature = "test-internals"))]
341impl Drop for TestRegistryGuard<'_> {
342 fn drop(&mut self) {
343 let backends = self.remove_entries(direct_backend_registry());
346 #[cfg(unix)]
347 let clients = self.remove_entries(client_registry());
348 drop(backends);
349 #[cfg(unix)]
350 drop(clients);
351 }
352}
353
354#[cfg(unix)]
357pub fn client_for(socket_path: &Path) -> crate::error::RuntimeResult<Arc<EventsSplitClient>> {
358 let mut registry = client_registry()
359 .lock()
360 .unwrap_or_else(std::sync::PoisonError::into_inner);
361 if let Some(existing) = registry.get(socket_path) {
362 return Ok(Arc::clone(existing));
363 }
364 let client = EventsSplitClient::new(socket_path.to_path_buf())?;
365 registry.insert(socket_path.to_path_buf(), Arc::clone(&client));
366 Ok(client)
367}
368
369pub fn direct_backend_for(db_path: &Path) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
372 direct_backend_with_max_readers(db_path, false, None)
373}
374
375pub fn direct_backend_read_only_for(
380 db_path: &Path,
381) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
382 direct_backend_with_max_readers(db_path, true, None)
383}
384
385pub(crate) fn direct_backend_with_max_readers(
386 db_path: &Path,
387 read_only: bool,
388 max_readers: Option<usize>,
389) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
390 let mut registry = direct_backend_registry()
391 .lock()
392 .unwrap_or_else(std::sync::PoisonError::into_inner);
393 let absolute_path = absolutize(db_path);
394 let db_path = absolute_path.as_path();
395 refuse_events_db_symlinks(db_path)
396 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
397 #[cfg(unix)]
398 if !read_only {
399 if let Some(parent) = db_path.parent().filter(|p| !p.as_os_str().is_empty()) {
400 let _ = std::fs::create_dir_all(parent);
401 }
402 }
403 #[cfg(unix)]
409 ensure_events_db_parent_trusted(db_path)
410 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
411 let key = std::fs::canonicalize(db_path)
412 .or_else(|error| {
413 if error.kind() != std::io::ErrorKind::NotFound {
414 return Err(error);
415 }
416 let parent = db_path.parent().ok_or(error)?;
417 let name = db_path.file_name().ok_or_else(|| {
418 std::io::Error::new(
419 std::io::ErrorKind::InvalidInput,
420 "events path has no file name",
421 )
422 })?;
423 std::fs::canonicalize(parent).map(|parent| parent.join(name))
424 })
425 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
426 if let Some((existing_read_only, existing)) = registry.get(&key) {
427 if *existing_read_only != read_only {
428 let existing_mode = if *existing_read_only {
429 "read-only"
430 } else {
431 "writable"
432 };
433 let requested_mode = if read_only { "read-only" } else { "writable" };
434 return Err(crate::error::RuntimeError::InvalidInput(format!(
435 "events database {} is already open {existing_mode} in this process; cannot \
436 open it {requested_mode}; read-only events access requires a separate frozen snapshot",
437 key.display()
438 )));
439 }
440 return Ok(Arc::clone(existing));
441 }
442 let db_path = key.as_path();
443 #[cfg(unix)]
449 if !read_only {
450 use std::os::unix::fs::OpenOptionsExt;
451 let _ = std::fs::OpenOptions::new()
452 .write(true)
453 .create_new(true)
454 .mode(0o600)
455 .open(db_path);
456 let _before_open = harden_events_db_sidecars(db_path)
459 .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
460 }
461 let backend = Arc::new(if read_only {
462 StorageBackend::sqlite_read_only_with_max_readers(db_path, max_readers)?
463 } else {
464 StorageBackend::sqlite_with_max_readers(db_path, max_readers)?
465 });
466 registry.insert(key, (read_only, Arc::clone(&backend)));
467 Ok(backend)
468}
469
470#[cfg(unix)]
473pub fn forwarding_metrics(socket_path: &Path) -> Option<EventsForwardingMetrics> {
474 let registry = client_registry()
475 .lock()
476 .unwrap_or_else(std::sync::PoisonError::into_inner);
477 registry.get(socket_path).map(|client| client.metrics())
478}
479
480#[derive(Debug, Serialize, Deserialize)]
486#[serde(tag = "op", rename_all = "snake_case")]
487pub enum EventsRequest {
488 AppendEvents {
492 protocol_version: u32,
493 namespace: String,
494 events: Vec<Event>,
495 },
496 AppendEventsIdempotent {
498 protocol_version: u32,
499 namespace: String,
500 events: Vec<Event>,
501 },
502 GetEvent {
503 protocol_version: u32,
504 namespace: String,
505 id: Uuid,
506 },
507 QueryEvents {
508 protocol_version: u32,
509 namespace: String,
510 filter: EventFilter,
511 page: PageRequest,
512 },
513 CountEvents {
514 protocol_version: u32,
515 namespace: String,
516 filter: EventFilter,
517 },
518}
519
520impl EventsRequest {
521 fn protocol_version(&self) -> u32 {
522 match self {
523 Self::AppendEvents {
524 protocol_version, ..
525 }
526 | Self::AppendEventsIdempotent {
527 protocol_version, ..
528 }
529 | Self::GetEvent {
530 protocol_version, ..
531 }
532 | Self::QueryEvents {
533 protocol_version, ..
534 }
535 | Self::CountEvents {
536 protocol_version, ..
537 } => *protocol_version,
538 }
539 }
540
541 fn namespace(&self) -> &str {
542 match self {
543 Self::AppendEvents { namespace, .. }
544 | Self::AppendEventsIdempotent { namespace, .. }
545 | Self::GetEvent { namespace, .. }
546 | Self::QueryEvents { namespace, .. }
547 | Self::CountEvents { namespace, .. } => namespace,
548 }
549 }
550}
551
552#[derive(Debug, Serialize, Deserialize)]
554#[serde(tag = "kind", rename_all = "snake_case")]
555pub enum EventsResponse {
556 Appended {
557 summary: BatchWriteSummary,
558 },
559 Idempotent {
560 result: IdempotentEventBatchResult,
561 },
562 Event {
563 event: Option<Event>,
564 },
565 Pageful {
566 page: Page<Event>,
567 },
568 Count {
569 count: u64,
570 },
571 Error {
581 message: String,
582 retryable: bool,
583 #[serde(default)]
584 writer_task_failure: Option<WireWriterTaskFailure>,
585 },
586}
587
588#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
591#[serde(tag = "kind", rename_all = "snake_case")]
592pub enum WireWriterTaskFailure {
593 RequestFailed { request_state: WireWriterTaskState },
594 TaskTerminated { request_state: WireWriterTaskState },
595}
596
597#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
601#[serde(rename_all = "snake_case")]
602pub enum WireWriterTaskState {
603 NotStarted,
604 TransactionRolledBack,
605 SideEffectsUnknown,
606}
607
608impl From<khive_storage::WriterTaskRequestState> for WireWriterTaskState {
609 fn from(state: khive_storage::WriterTaskRequestState) -> Self {
610 use khive_storage::WriterTaskRequestState as S;
611 match state {
612 S::NotStarted => Self::NotStarted,
613 S::TransactionRolledBack => Self::TransactionRolledBack,
614 S::SideEffectsUnknown => Self::SideEffectsUnknown,
615 }
616 }
617}
618
619impl From<WireWriterTaskState> for khive_storage::WriterTaskRequestState {
620 fn from(state: WireWriterTaskState) -> Self {
621 use khive_storage::WriterTaskRequestState as S;
622 match state {
623 WireWriterTaskState::NotStarted => S::NotStarted,
624 WireWriterTaskState::TransactionRolledBack => S::TransactionRolledBack,
625 WireWriterTaskState::SideEffectsUnknown => S::SideEffectsUnknown,
626 }
627 }
628}
629
630#[cfg(unix)]
638pub struct EventsDaemonGuard {
639 _file: std::fs::File,
640}
641
642#[cfg(unix)]
650pub fn try_acquire_events_daemon_guard(socket_path: &Path) -> Option<EventsDaemonGuard> {
651 let lock_path = socket_path.with_extension("lock");
652 if let Some(parent) = lock_path.parent() {
653 std::fs::create_dir_all(parent).ok()?;
654 }
655 use std::os::unix::fs::OpenOptionsExt;
656 let file = std::fs::OpenOptions::new()
660 .create(true)
661 .truncate(false)
662 .write(true)
663 .mode(0o600)
664 .custom_flags(libc::O_NOFOLLOW)
665 .open(&lock_path)
666 .ok()?;
667 file.set_permissions(
673 <std::fs::Permissions as std::os::unix::fs::PermissionsExt>::from_mode(0o600),
674 )
675 .ok()?;
676 use std::os::fd::AsRawFd;
677 let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
680 if rc == 0 {
681 Some(EventsDaemonGuard { _file: file })
682 } else {
683 None
684 }
685}
686
687fn events_db_targets(db_path: &Path) -> [PathBuf; 3] {
688 ["", "-wal", "-shm"].map(|suffix| {
689 let mut name = db_path.as_os_str().to_os_string();
690 name.push(suffix);
691 PathBuf::from(name)
692 })
693}
694
695fn refuse_events_db_symlinks(db_path: &Path) -> anyhow::Result<()> {
705 for path in events_db_targets(db_path) {
706 match std::fs::symlink_metadata(&path) {
707 Ok(meta) if meta.file_type().is_symlink() => anyhow::bail!(
708 "refusing to serve events: {} is a symlink; the events database and its \
709 sidecars must be regular files",
710 path.display()
711 ),
712 Ok(_) => {}
713 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
714 Err(e) => anyhow::bail!(
715 "refusing to serve events: cannot inspect {}: {e}",
716 path.display()
717 ),
718 }
719 }
720 Ok(())
721}
722
723#[cfg(unix)]
724#[derive(Clone, Copy, Debug, PartialEq, Eq)]
725struct EventsFileIdentity {
726 device: u64,
727 inode: u64,
728}
729
730#[cfg(unix)]
731impl EventsFileIdentity {
732 fn from_metadata(metadata: &std::fs::Metadata) -> Self {
733 use std::os::unix::fs::MetadataExt;
734 Self {
735 device: metadata.dev(),
736 inode: metadata.ino(),
737 }
738 }
739}
740
741#[cfg(unix)]
744type EventsDbIdentities = [Option<EventsFileIdentity>; 3];
745
746#[cfg(unix)]
755fn ensure_events_db_parent_trusted(db_path: &Path) -> anyhow::Result<()> {
756 let parent = absolutize(db_path);
757 let parent = parent.parent().filter(|p| !p.as_os_str().is_empty());
758 match parent {
759 Some(dir) => crate::daemon::ensure_socket_dir_is_trusted(dir).map_err(|e| {
760 anyhow::anyhow!("refusing to serve events from an untrusted directory: {e}")
761 }),
762 None => Ok(()),
763 }
764}
765
766#[cfg(unix)]
772fn ensure_events_db_owner_only(db_path: &Path) -> anyhow::Result<()> {
773 use std::os::unix::fs::OpenOptionsExt;
774 refuse_events_db_symlinks(db_path)?;
775 if let Some(parent) = db_path.parent() {
776 std::fs::create_dir_all(parent)?;
777 }
778 ensure_events_db_parent_trusted(db_path)?;
781 match std::fs::OpenOptions::new()
782 .write(true)
783 .create_new(true)
784 .mode(0o600)
785 .open(db_path)
786 {
787 Ok(_created) => Ok(()),
789 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
790 Err(e) => Err(anyhow::anyhow!(
791 "refusing to serve events: cannot create {} owner-only: {e}",
792 db_path.display()
793 )),
794 }
795}
796
797#[cfg(unix)]
810fn harden_events_db_sidecars(db_path: &Path) -> anyhow::Result<EventsDbIdentities> {
811 use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
812 let mut identities = [None; 3];
813 for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
814 let file = match std::fs::OpenOptions::new()
821 .read(true)
822 .custom_flags(libc::O_NOFOLLOW)
823 .open(&path)
824 {
825 Ok(file) => file,
826 Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
827 Err(e) => {
828 anyhow::bail!(
829 "refusing to serve events: cannot open {} without following symlinks: \
830 {e}. The events database and its sidecars must be regular files.",
831 path.display()
832 )
833 }
834 };
835 let metadata = file.metadata()?;
838 if !metadata.file_type().is_file() {
839 anyhow::bail!(
840 "refusing to serve events: {} is not a regular file. The events \
841 database and its sidecars must be regular files.",
842 path.display()
843 );
844 }
845 file.set_permissions(std::fs::Permissions::from_mode(0o600))
846 .map_err(|e| {
847 anyhow::anyhow!(
848 "refusing to serve events: cannot chmod 0600 {}: {e}. The events \
849 database and its sidecars must be owner-only.",
850 path.display()
851 )
852 })?;
853 identities[index] = Some(EventsFileIdentity::from_metadata(&metadata));
856 }
857 Ok(identities)
858}
859
860#[cfg(unix)]
878fn verify_events_db_owner_only_unopened(
879 db_path: &Path,
880 before: &EventsDbIdentities,
881) -> anyhow::Result<()> {
882 use std::os::unix::fs::PermissionsExt;
883 for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
884 let metadata = match std::fs::symlink_metadata(&path) {
885 Ok(metadata) => metadata,
886 Err(e) if e.kind() == std::io::ErrorKind::NotFound && index != 0 => continue,
887 Err(e) if e.kind() == std::io::ErrorKind::NotFound => anyhow::bail!(
888 "refusing to serve events: main database {} disappeared after SQLite open; \
889 refusing to serve an unlinked database",
890 path.display()
891 ),
892 Err(e) => {
893 anyhow::bail!(
894 "refusing to serve events: cannot stat {}: {e}",
895 path.display()
896 )
897 }
898 };
899 if !metadata.file_type().is_file() {
900 anyhow::bail!(
901 "refusing to serve events: {} is not a regular file. The events database \
902 and its sidecars must be regular files.",
903 path.display()
904 );
905 }
906 let mode = metadata.permissions().mode() & 0o777;
907 if mode & 0o077 != 0 {
908 anyhow::bail!(
909 "refusing to serve events: {} is mode {mode:03o}, not owner-only. The events \
910 database and its sidecars must be owner-only.",
911 path.display()
912 );
913 }
914 if before[index]
915 .is_some_and(|identity| identity != EventsFileIdentity::from_metadata(&metadata))
916 {
917 anyhow::bail!(
918 "refusing to serve events: {} changed identity between pre-open hardening \
919 and post-open verification",
920 path.display()
921 );
922 }
923 }
924 Ok(())
925}
926
927#[cfg(unix)]
948pub async fn supervise_events_daemon(db_path: PathBuf, socket_path: PathBuf) {
949 const PROBE_INTERVAL: Duration = Duration::from_secs(15);
950 let shutdown = crate::daemon::daemon_shutdown_token();
951 let mut child: Option<std::process::Child> = None;
952 let mut respawns: u64 = 0;
953 loop {
954 if let Some(c) = child.as_mut() {
958 match c.try_wait() {
959 Ok(Some(status)) => {
960 tracing::info!(%status, "events daemon child exited");
961 child = None;
962 }
963 Ok(None) => {}
964 Err(error) => {
965 tracing::warn!(error = %error, "cannot poll events daemon child; dropping handle");
966 child = None;
967 }
968 }
969 }
970
971 let reachable = match connect_verified(&socket_path).await {
972 Ok(_stream) => true,
973 Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
974 tracing::warn!(
975 socket = %socket_path.display(),
976 error = %error,
977 "events socket answered by a foreign uid; treating as unreachable"
978 );
979 false
980 }
981 Err(_) => false,
982 };
983 if !reachable && child.is_none() {
984 match std::env::current_exe() {
985 Ok(exe) => {
986 let spawned = std::process::Command::new(exe)
987 .arg("events-daemon")
988 .arg("--db")
989 .arg(&db_path)
990 .arg("--socket")
991 .arg(&socket_path)
992 .stdin(std::process::Stdio::null())
993 .stdout(std::process::Stdio::null())
994 .stderr(std::process::Stdio::null())
995 .spawn();
996 match spawned {
997 Ok(spawned_child) => {
998 respawns += 1;
999 tracing::info!(
1000 pid = spawned_child.id(),
1001 respawns,
1002 socket = %socket_path.display(),
1003 "spawned events daemon"
1004 );
1005 child = Some(spawned_child);
1006 }
1007 Err(error) => {
1008 tracing::warn!(error = %error, "failed to spawn events daemon");
1009 }
1010 }
1011 }
1012 Err(error) => {
1013 tracing::warn!(error = %error, "cannot resolve current executable for events daemon spawn");
1014 }
1015 }
1016 }
1017 tokio::select! {
1018 _ = shutdown.cancelled() => break,
1019 _ = tokio::time::sleep(PROBE_INTERVAL) => {}
1020 }
1021 }
1022 if let Some(mut c) = child.take() {
1023 let _ = c.kill();
1027 let _ = c.wait();
1028 tracing::info!("events daemon child stopped with supervisor shutdown");
1029 }
1030}
1031
1032#[cfg(unix)]
1048pub async fn run_events_daemon(db_path: &Path, socket_path: &Path) -> anyhow::Result<()> {
1049 let db_path = &absolutize(db_path);
1052 let socket_path = &absolutize(socket_path);
1053 if let Some(parent) = socket_path.parent() {
1057 std::fs::create_dir_all(parent)?;
1058 crate::daemon::ensure_socket_dir_is_trusted(parent)?;
1059 }
1060 let Some(_guard) = try_acquire_events_daemon_guard(socket_path) else {
1061 tracing::info!(
1062 socket = %socket_path.display(),
1063 "events daemon lock unavailable (held by another daemon, or hardening refused); exiting"
1064 );
1065 return Ok(());
1066 };
1067 ensure_events_db_owner_only(db_path)?;
1068 let before_open = harden_events_db_sidecars(db_path)?;
1081 let backend = Arc::new(StorageBackend::sqlite(db_path)?);
1082 backend.events()?;
1084 verify_events_db_owner_only_unopened(db_path, &before_open)?;
1085
1086 if socket_path.exists() {
1087 std::fs::remove_file(socket_path)?;
1088 }
1089 let listener = UnixListener::bind(socket_path)?;
1090 {
1091 use std::os::unix::fs::PermissionsExt;
1092 if let Err(e) =
1096 std::fs::set_permissions(socket_path, std::fs::Permissions::from_mode(0o600))
1097 {
1098 drop(listener);
1099 let _ = std::fs::remove_file(socket_path);
1100 anyhow::bail!(
1101 "refusing to serve events: cannot chmod 0600 {}: {e}. The events socket must \
1102 be owner-only.",
1103 socket_path.display()
1104 );
1105 }
1106 }
1107 let daemon_euid = unsafe { libc::geteuid() };
1108 let stores: NamespaceStores = Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
1114 let connections = Arc::new(tokio::sync::Semaphore::new(MAX_EVENTS_CONNECTIONS));
1115 let frame_budget = Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES));
1116 tracing::info!(
1117 socket = %socket_path.display(),
1118 db = %db_path.display(),
1119 "events daemon listening"
1120 );
1121
1122 loop {
1123 let (stream, _addr) = match listener.accept().await {
1124 Ok(pair) => pair,
1125 Err(error) => {
1126 tracing::warn!(error = %error, "events daemon accept failed");
1127 continue;
1128 }
1129 };
1130 match crate::daemon::peer_uid(&stream) {
1131 Ok(uid) if crate::daemon::uid_is_permitted(uid, daemon_euid) => {}
1132 Ok(uid) => {
1133 tracing::warn!(peer_uid = uid, "events daemon rejected foreign-uid peer");
1134 continue;
1135 }
1136 Err(error) => {
1137 tracing::warn!(error = %error, "events daemon could not read peer credentials");
1138 continue;
1139 }
1140 }
1141 let permit = match Arc::clone(&connections).try_acquire_owned() {
1142 Ok(permit) => permit,
1143 Err(_) => {
1144 tracing::warn!(
1147 cap = MAX_EVENTS_CONNECTIONS,
1148 "events daemon at connection cap; dropping new connection"
1149 );
1150 continue;
1151 }
1152 };
1153 let backend = Arc::clone(&backend);
1154 let stores = Arc::clone(&stores);
1155 let frame_budget = Arc::clone(&frame_budget);
1156 crate::daemon::spawn_named_tracked_task("events_connection", async move {
1157 serve_events_conn(stream, backend, stores, frame_budget).await;
1158 drop(permit);
1159 });
1160 }
1161}
1162
1163#[cfg(unix)]
1171async fn read_frame_budgeted(
1172 stream: &mut UnixStream,
1173 budget: &Arc<tokio::sync::Semaphore>,
1174) -> std::io::Result<(Vec<u8>, tokio::sync::OwnedSemaphorePermit)> {
1175 use tokio::io::AsyncReadExt;
1176 let mut len_buf = [0u8; 4];
1177 stream.read_exact(&mut len_buf).await?;
1178 let len = u32::from_be_bytes(len_buf) as usize;
1179 if len > crate::daemon::MAX_FRAME_BYTES {
1180 return Err(std::io::Error::new(
1181 std::io::ErrorKind::InvalidData,
1182 format!(
1183 "daemon frame of {len} bytes exceeds {} cap",
1184 crate::daemon::MAX_FRAME_BYTES
1185 ),
1186 ));
1187 }
1188 let permit = Arc::clone(budget)
1189 .acquire_many_owned(len as u32)
1190 .await
1191 .map_err(|_| std::io::Error::other("frame budget closed"))?;
1192 let mut buf = vec![0u8; len];
1193 stream.read_exact(&mut buf).await?;
1194 Ok((buf, permit))
1195}
1196
1197#[cfg(unix)]
1198type NamespaceStores =
1199 Arc<std::sync::Mutex<std::collections::HashMap<String, Arc<dyn EventStore>>>>;
1200
1201#[cfg(unix)]
1213fn namespace_store(
1214 backend: &StorageBackend,
1215 stores: &NamespaceStores,
1216 namespace: &str,
1217) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1218 namespace_store_with_cap(backend, stores, namespace, MAX_CACHED_NAMESPACE_STORES)
1219}
1220
1221#[cfg(unix)]
1222fn namespace_store_with_cap(
1223 backend: &StorageBackend,
1224 stores: &NamespaceStores,
1225 namespace: &str,
1226 cap: usize,
1227) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1228 let key = namespace.trim();
1229 if let Some(store) = stores
1230 .lock()
1231 .unwrap_or_else(std::sync::PoisonError::into_inner)
1232 .get(key)
1233 {
1234 return Ok(Arc::clone(store));
1235 }
1236 let store = backend.events_for_namespace(namespace)?;
1237 let mut map = stores
1238 .lock()
1239 .unwrap_or_else(std::sync::PoisonError::into_inner);
1240 if map.len() >= cap && !map.contains_key(key) {
1241 if let Some(victim) = map.keys().next().cloned() {
1249 map.remove(&victim);
1250 }
1251 }
1252 map.insert(key.to_string(), Arc::clone(&store));
1253 Ok(store)
1254}
1255
1256#[cfg(unix)]
1257async fn serve_events_conn(
1258 mut stream: UnixStream,
1259 backend: Arc<StorageBackend>,
1260 stores: NamespaceStores,
1261 frame_budget: Arc<tokio::sync::Semaphore>,
1262) {
1263 loop {
1264 let (payload, budget_permit) = match tokio::time::timeout(
1270 CONN_IO_TIMEOUT,
1271 read_frame_budgeted(&mut stream, &frame_budget),
1272 )
1273 .await
1274 {
1275 Ok(Ok(pair)) => pair,
1276 Ok(Err(_)) | Err(_) => return,
1278 };
1279 let response = match serde_json::from_slice::<EventsRequest>(&payload) {
1280 Ok(request) => dispatch_events_request(request, &backend, &stores).await,
1281 Err(error) => EventsResponse::Error {
1282 message: format!("events daemon could not parse request frame: {error}"),
1283 retryable: false,
1284 writer_task_failure: None,
1285 },
1286 };
1287 let bytes = match serde_json::to_vec(&response) {
1288 Ok(bytes) => bytes,
1289 Err(error) => {
1290 tracing::error!(error = %error, "events daemon response serialization failed");
1291 return;
1292 }
1293 };
1294 let bytes = if bytes.len() > crate::daemon::MAX_FRAME_BYTES {
1295 let refusal = EventsResponse::Error {
1299 message: format!(
1300 "response_frame_size_limit: {} response bytes exceed the {}-byte events IPC frame cap; request a narrower page",
1301 bytes.len(),
1302 crate::daemon::MAX_FRAME_BYTES
1303 ),
1304 retryable: false,
1305 writer_task_failure: None,
1306 };
1307 match serde_json::to_vec(&refusal) {
1308 Ok(bytes) => bytes,
1309 Err(error) => {
1310 tracing::error!(error = %error, "events daemon frame-size refusal serialization failed");
1311 return;
1312 }
1313 }
1314 } else {
1315 bytes
1316 };
1317 match tokio::time::timeout(CONN_IO_TIMEOUT, write_frame(&mut stream, &bytes)).await {
1320 Ok(Ok(())) => {}
1321 Ok(Err(_)) | Err(_) => return,
1322 }
1323 drop(budget_permit);
1327 }
1328}
1329
1330#[cfg(unix)]
1331async fn dispatch_events_request(
1332 request: EventsRequest,
1333 backend: &StorageBackend,
1334 stores: &NamespaceStores,
1335) -> EventsResponse {
1336 if request.protocol_version() != EVENTS_PROTOCOL_VERSION {
1337 return EventsResponse::Error {
1338 message: format!(
1339 "events protocol version mismatch: daemon speaks {}, client sent {}",
1340 EVENTS_PROTOCOL_VERSION,
1341 request.protocol_version()
1342 ),
1343 retryable: false,
1344 writer_task_failure: None,
1345 };
1346 }
1347 if let Err(error) = khive_types::Namespace::parse(request.namespace().trim()) {
1353 return EventsResponse::Error {
1354 message: format!("events request namespace rejected: {error}"),
1355 retryable: false,
1356 writer_task_failure: None,
1357 };
1358 }
1359 let store = match namespace_store(backend, stores, request.namespace()) {
1360 Ok(store) => store,
1361 Err(error) => {
1362 return EventsResponse::Error {
1363 message: format!("events store unavailable: {error}"),
1364 retryable: true,
1365 writer_task_failure: None,
1366 };
1367 }
1368 };
1369 match request {
1370 EventsRequest::AppendEvents { events, .. } => match store.append_events(events).await {
1371 Ok(summary) => EventsResponse::Appended { summary },
1372 Err(error) => storage_error_response(&error),
1373 },
1374 EventsRequest::AppendEventsIdempotent { events, .. } => {
1375 match store.append_events_idempotent(events).await {
1376 Ok(result) => EventsResponse::Idempotent { result },
1377 Err(error) => storage_error_response(&error),
1378 }
1379 }
1380 EventsRequest::GetEvent { id, .. } => match store.get_event(id).await {
1381 Ok(event) => EventsResponse::Event { event },
1382 Err(error) => storage_error_response(&error),
1383 },
1384 EventsRequest::QueryEvents { filter, page, .. } => {
1385 if page.limit > MAX_QUERY_EVENTS_PAGE_ROWS {
1386 return EventsResponse::Error {
1387 message: format!(
1388 "events query page limit {} exceeds the daemon cap of {} rows; \
1389 request narrower pages",
1390 page.limit, MAX_QUERY_EVENTS_PAGE_ROWS
1391 ),
1392 retryable: false,
1393 writer_task_failure: None,
1394 };
1395 }
1396 match store.query_events(filter, page).await {
1397 Ok(page) => EventsResponse::Pageful { page },
1398 Err(error) => storage_error_response(&error),
1399 }
1400 }
1401 EventsRequest::CountEvents { filter, .. } => match store.count_events(filter).await {
1402 Ok(count) => EventsResponse::Count { count },
1403 Err(error) => storage_error_response(&error),
1404 },
1405 }
1406}
1407
1408#[cfg(unix)]
1409fn storage_error_response(error: &StorageError) -> EventsResponse {
1410 let (message, writer_task_failure) = match error {
1411 StorageError::WriterTaskRequestFailed {
1412 request_state,
1413 source,
1414 } => (
1415 source.to_string(),
1416 Some(WireWriterTaskFailure::RequestFailed {
1417 request_state: (*request_state).into(),
1418 }),
1419 ),
1420 StorageError::WriterTaskTerminated { request_state } => (
1421 error.to_string(),
1422 Some(WireWriterTaskFailure::TaskTerminated {
1423 request_state: (*request_state).into(),
1424 }),
1425 ),
1426 _ => (error.to_string(), None),
1427 };
1428 EventsResponse::Error {
1429 message,
1430 retryable: error.is_retryable(),
1435 writer_task_failure,
1436 }
1437}
1438
1439#[cfg(unix)]
1450async fn connect_verified(socket_path: &Path) -> std::io::Result<UnixStream> {
1451 let stream = UnixStream::connect(socket_path).await?;
1452 let own_euid = unsafe { libc::geteuid() } as u32;
1454 let peer = crate::daemon::peer_uid(&stream)?;
1455 if !crate::daemon::uid_is_permitted(peer, own_euid) {
1456 return Err(std::io::Error::new(
1457 std::io::ErrorKind::PermissionDenied,
1458 format!(
1459 "events socket {} answered by uid {peer}, not this process's uid {own_euid}; \
1460 refusing to exchange event frames with an unowned daemon",
1461 socket_path.display()
1462 ),
1463 ));
1464 }
1465 Ok(stream)
1466}
1467
1468#[cfg(unix)]
1472#[derive(Debug, Clone, Copy, Serialize)]
1473pub struct EventsForwardingMetrics {
1474 pub forwarded_batches: u64,
1475 pub forwarded_events: u64,
1476 pub dropped_batches: u64,
1477 pub dropped_events: u64,
1478 pub queued_bytes: usize,
1480}
1481
1482#[cfg(unix)]
1483#[derive(Debug, Default)]
1484struct ForwardingCounters {
1485 forwarded_batches: AtomicU64,
1486 forwarded_events: AtomicU64,
1487 dropped_batches: AtomicU64,
1488 dropped_events: AtomicU64,
1489 queued_bytes: AtomicUsize,
1490}
1491
1492#[cfg(unix)]
1495#[derive(Debug)]
1496struct AppendByteReservation {
1497 counters: Arc<ForwardingCounters>,
1498 bytes: usize,
1499}
1500
1501#[cfg(unix)]
1502impl Drop for AppendByteReservation {
1503 fn drop(&mut self) {
1504 self.counters
1505 .queued_bytes
1506 .fetch_sub(self.bytes, Ordering::AcqRel);
1507 }
1508}
1509
1510#[cfg(unix)]
1511#[derive(Debug)]
1512struct QueuedAppend {
1513 request: EventsRequest,
1514 count: u64,
1515 _reservation: AppendByteReservation,
1516}
1517
1518#[cfg(all(test, unix))]
1519impl QueuedAppend {
1520 fn unmetered(namespace: &str, events: Vec<Event>, counters: Arc<ForwardingCounters>) -> Self {
1521 let count = events.len() as u64;
1522 Self {
1523 request: EventsRequest::AppendEvents {
1524 protocol_version: EVENTS_PROTOCOL_VERSION,
1525 namespace: namespace.to_string(),
1526 events,
1527 },
1528 count,
1529 _reservation: AppendByteReservation { counters, bytes: 0 },
1530 }
1531 }
1532}
1533
1534#[cfg(unix)]
1537#[derive(Default)]
1538struct BoundedFrameCounter {
1539 bytes: usize,
1540 exceeded: bool,
1541}
1542
1543#[cfg(unix)]
1544impl std::io::Write for BoundedFrameCounter {
1545 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
1546 let next = self.bytes.saturating_add(buf.len());
1547 if next > crate::daemon::MAX_FRAME_BYTES {
1548 self.exceeded = true;
1549 return Err(std::io::Error::new(
1550 std::io::ErrorKind::InvalidData,
1551 "events append request exceeds IPC frame cap",
1552 ));
1553 }
1554 self.bytes = next;
1555 Ok(buf.len())
1556 }
1557
1558 fn flush(&mut self) -> std::io::Result<()> {
1559 Ok(())
1560 }
1561}
1562
1563#[cfg(unix)]
1564fn reserve_append_bytes(
1565 counters: &Arc<ForwardingCounters>,
1566 bytes: usize,
1567 budget: usize,
1568) -> Option<AppendByteReservation> {
1569 let mut used = counters.queued_bytes.load(Ordering::Acquire);
1570 loop {
1571 let next = used.checked_add(bytes)?;
1572 if next > budget {
1573 return None;
1574 }
1575 match counters.queued_bytes.compare_exchange_weak(
1576 used,
1577 next,
1578 Ordering::AcqRel,
1579 Ordering::Acquire,
1580 ) {
1581 Ok(_) => {
1582 return Some(AppendByteReservation {
1583 counters: Arc::clone(counters),
1584 bytes,
1585 });
1586 }
1587 Err(observed) => used = observed,
1588 }
1589 }
1590}
1591
1592#[cfg(unix)]
1600pub struct EventsSplitClient {
1601 socket_path: PathBuf,
1602 append_tx: tokio::sync::mpsc::Sender<QueuedAppend>,
1603 append_queue_byte_budget: usize,
1604 counters: Arc<ForwardingCounters>,
1605 outage_logged: Arc<AtomicBool>,
1608 preflight_store: Arc<dyn EventStore>,
1610}
1611
1612#[cfg(unix)]
1613impl std::fmt::Debug for EventsSplitClient {
1614 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1615 f.debug_struct("EventsSplitClient")
1616 .field("socket_path", &self.socket_path)
1617 .finish_non_exhaustive()
1618 }
1619}
1620
1621#[cfg(unix)]
1622impl EventsSplitClient {
1623 pub fn new(socket_path: PathBuf) -> crate::error::RuntimeResult<Arc<Self>> {
1625 Self::new_with_queue_depth(socket_path, DEFAULT_APPEND_QUEUE_BATCHES)
1626 }
1627
1628 pub fn new_with_queue_depth(
1631 socket_path: PathBuf,
1632 queue_depth: usize,
1633 ) -> crate::error::RuntimeResult<Arc<Self>> {
1634 Self::new_with_queue_depth_and_delivery_timeout(socket_path, queue_depth, REQUEST_TIMEOUT)
1635 }
1636
1637 fn new_with_queue_depth_and_delivery_timeout(
1641 socket_path: PathBuf,
1642 queue_depth: usize,
1643 delivery_timeout: Duration,
1644 ) -> crate::error::RuntimeResult<Arc<Self>> {
1645 Self::new_with_limits_and_delivery_timeout(
1646 socket_path,
1647 queue_depth,
1648 DEFAULT_APPEND_QUEUE_BYTES,
1649 delivery_timeout,
1650 )
1651 }
1652
1653 fn new_with_limits_and_delivery_timeout(
1654 socket_path: PathBuf,
1655 queue_depth: usize,
1656 byte_budget: usize,
1657 delivery_timeout: Duration,
1658 ) -> crate::error::RuntimeResult<Arc<Self>> {
1659 let preflight_backend = StorageBackend::memory()?;
1660 let preflight_store = preflight_backend.events()?;
1661 let (append_tx, append_rx) = tokio::sync::mpsc::channel::<QueuedAppend>(queue_depth.max(1));
1665 let counters = Arc::new(ForwardingCounters::default());
1666 let outage_logged = Arc::new(AtomicBool::new(false));
1667
1668 let client = Arc::new(Self {
1669 socket_path: socket_path.clone(),
1670 append_tx,
1671 append_queue_byte_budget: byte_budget,
1672 counters: Arc::clone(&counters),
1673 outage_logged: Arc::clone(&outage_logged),
1674 preflight_store,
1675 });
1676
1677 crate::daemon::spawn_named_tracked_task(
1678 "events_forwarder",
1679 run_forwarder(
1680 socket_path,
1681 append_rx,
1682 counters,
1683 outage_logged,
1684 delivery_timeout,
1685 crate::daemon::daemon_shutdown_token(),
1686 ),
1687 );
1688 Ok(client)
1689 }
1690
1691 pub fn metrics(&self) -> EventsForwardingMetrics {
1693 EventsForwardingMetrics {
1694 forwarded_batches: self.counters.forwarded_batches.load(Ordering::Relaxed),
1695 forwarded_events: self.counters.forwarded_events.load(Ordering::Relaxed),
1696 dropped_batches: self.counters.dropped_batches.load(Ordering::Relaxed),
1697 dropped_events: self.counters.dropped_events.load(Ordering::Relaxed),
1698 queued_bytes: self.counters.queued_bytes.load(Ordering::Acquire),
1699 }
1700 }
1701
1702 fn enqueue(&self, namespace: &str, events: Vec<Event>) {
1705 let count = events.len() as u64;
1706 let request = EventsRequest::AppendEvents {
1707 protocol_version: EVENTS_PROTOCOL_VERSION,
1708 namespace: namespace.to_string(),
1709 events,
1710 };
1711 let mut frame_size = BoundedFrameCounter::default();
1712 if let Err(error) = serde_json::to_writer(&mut frame_size, &request) {
1713 self.count_append_drop(
1714 count,
1715 if frame_size.exceeded {
1716 "frame cap"
1717 } else {
1718 "serialization"
1719 },
1720 );
1721 tracing::warn!(error = %error, "events append batch could not fit a wire frame");
1722 return;
1723 }
1724 let Some(reservation) = reserve_append_bytes(
1725 &self.counters,
1726 frame_size.bytes,
1727 self.append_queue_byte_budget,
1728 ) else {
1729 self.count_append_drop(count, "queue byte budget");
1730 return;
1731 };
1732 let batch = QueuedAppend {
1733 request,
1734 count,
1735 _reservation: reservation,
1736 };
1737 match self.append_tx.try_send(batch) {
1738 Ok(()) => {}
1739 Err(_) => {
1740 self.count_append_drop(count, "queue batch limit");
1742 }
1743 }
1744 }
1745
1746 fn count_append_drop(&self, count: u64, reason: &'static str) {
1747 self.counters
1748 .dropped_batches
1749 .fetch_add(1, Ordering::Relaxed);
1750 self.counters
1751 .dropped_events
1752 .fetch_add(count, Ordering::Relaxed);
1753 if !self.outage_logged.swap(true, Ordering::Relaxed) {
1754 tracing::warn!(
1755 dropped_events = count,
1756 reason,
1757 "events append queue rejected loss-tolerant batch"
1758 );
1759 }
1760 }
1761
1762 async fn round_trip(&self, request: &EventsRequest) -> StorageResult<EventsResponse> {
1767 let op = "events daemon round-trip";
1768 let payload = serde_json::to_vec(request).map_err(|error| StorageError::Serialization {
1769 capability: khive_storage::StorageCapability::Events,
1770 message: format!("events request serialization failed: {error}"),
1771 })?;
1772
1773 let attempt = async {
1774 let mut stream = connect_verified(&self.socket_path)
1775 .await
1776 .map_err(|error| ("connect", error))?;
1777 write_frame(&mut stream, &payload)
1778 .await
1779 .map_err(|error| ("write", error))?;
1780 read_frame(&mut stream)
1781 .await
1782 .map_err(|error| ("read", error))
1783 };
1784 let bytes = match tokio::time::timeout(REQUEST_TIMEOUT, attempt).await {
1785 Ok(Ok(bytes)) => bytes,
1786 Ok(Err(("read", error))) => {
1787 return Err(StorageError::Serialization {
1788 capability: khive_storage::StorageCapability::Events,
1789 message: format!(
1790 "events daemon closed or broke the response after connection at {}: {error}",
1791 self.socket_path.display()
1792 ),
1793 });
1794 }
1795 Ok(Err((stage, error))) => {
1796 return Err(StorageError::Pool {
1797 operation: op.into(),
1798 message: format!(
1799 "events daemon {stage} failed at {}: {error}",
1800 self.socket_path.display()
1801 ),
1802 });
1803 }
1804 Err(_elapsed) => {
1805 return Err(StorageError::Timeout {
1806 operation: op.into(),
1807 });
1808 }
1809 };
1810
1811 serde_json::from_slice::<EventsResponse>(&bytes).map_err(|error| {
1812 StorageError::Serialization {
1813 capability: khive_storage::StorageCapability::Events,
1814 message: format!("events response deserialization failed: {error}"),
1815 }
1816 })
1817 }
1818}
1819
1820#[cfg(unix)]
1827fn drain_dropped_queue(
1828 rx: &mut tokio::sync::mpsc::Receiver<QueuedAppend>,
1829 counters: &ForwardingCounters,
1830) {
1831 let mut dropped_batches = 0u64;
1832 let mut dropped_events = 0u64;
1833 while let Ok(batch) = rx.try_recv() {
1834 dropped_batches += 1;
1835 dropped_events += batch.count;
1836 }
1838 if dropped_batches > 0 {
1839 counters
1840 .dropped_batches
1841 .fetch_add(dropped_batches, Ordering::Relaxed);
1842 counters
1843 .dropped_events
1844 .fetch_add(dropped_events, Ordering::Relaxed);
1845 tracing::warn!(
1846 dropped_batches,
1847 dropped_events,
1848 "events forwarder shutting down; dropping queued loss-tolerant batches"
1849 );
1850 }
1851}
1852
1853#[cfg(unix)]
1858async fn run_forwarder(
1859 socket_path: PathBuf,
1860 mut rx: tokio::sync::mpsc::Receiver<QueuedAppend>,
1861 counters: Arc<ForwardingCounters>,
1862 outage_logged: Arc<AtomicBool>,
1863 delivery_timeout: Duration,
1864 shutdown: tokio_util::sync::CancellationToken,
1865) {
1866 let mut conn: Option<UnixStream> = None;
1874 loop {
1875 let batch = tokio::select! {
1876 _ = shutdown.cancelled() => {
1877 drain_dropped_queue(&mut rx, &counters);
1878 break;
1879 },
1880 received = rx.recv() => match received {
1881 Some(batch) => batch,
1882 None => break,
1883 },
1884 };
1885 let count = batch.count;
1886 let payload = match serde_json::to_vec(&batch.request) {
1887 Ok(payload) => payload,
1888 Err(error) => {
1889 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
1890 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
1891 tracing::error!(error = %error, "events forwarder serialization failed; batch dropped");
1892 continue;
1893 }
1894 };
1895
1896 let delivered = tokio::select! {
1904 _ = shutdown.cancelled() => {
1905 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
1906 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
1907 tracing::warn!(
1908 dropped_events = count,
1909 "events forwarder shutting down mid-delivery; in-flight batch dropped"
1910 );
1911 drain_dropped_queue(&mut rx, &counters);
1912 break;
1913 },
1914 outcome = tokio::time::timeout(
1915 delivery_timeout,
1916 deliver_batch(&socket_path, &mut conn, &payload),
1917 ) => match outcome {
1918 Ok(delivered) => delivered,
1919 Err(_elapsed) => {
1920 conn = None;
1921 false
1922 }
1923 },
1924 };
1925 if delivered {
1926 counters.forwarded_batches.fetch_add(1, Ordering::Relaxed);
1927 counters
1928 .forwarded_events
1929 .fetch_add(count, Ordering::Relaxed);
1930 if outage_logged.swap(false, Ordering::Relaxed) {
1931 tracing::info!("events daemon reachable again; forwarding resumed");
1932 }
1933 } else {
1934 counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
1935 counters.dropped_events.fetch_add(count, Ordering::Relaxed);
1936 if !outage_logged.swap(true, Ordering::Relaxed) {
1937 tracing::warn!(
1938 socket = %socket_path.display(),
1939 "events daemon unreachable; dropping loss-tolerant events until it returns"
1940 );
1941 }
1942 tokio::select! {
1943 _ = shutdown.cancelled() => {
1944 drain_dropped_queue(&mut rx, &counters);
1945 break;
1946 },
1947 _ = tokio::time::sleep(FORWARDER_BACKOFF) => {}
1948 }
1949 }
1950 }
1951}
1952
1953#[cfg(unix)]
1958async fn deliver_batch(socket_path: &Path, conn: &mut Option<UnixStream>, payload: &[u8]) -> bool {
1959 for _attempt in 0..2u8 {
1960 if conn.is_none() {
1961 match connect_verified(socket_path).await {
1962 Ok(stream) => *conn = Some(stream),
1963 Err(_) => return false,
1964 }
1965 }
1966 let stream = conn.as_mut().expect("connection populated above");
1967 let ok = async {
1968 write_frame(stream, payload).await?;
1969 let bytes = read_frame(stream).await?;
1970 std::io::Result::Ok(bytes)
1971 }
1972 .await;
1973 match ok {
1974 Ok(bytes) => {
1975 return !matches!(
1976 serde_json::from_slice::<EventsResponse>(&bytes),
1977 Ok(EventsResponse::Error { .. }) | Err(_)
1978 );
1979 }
1980 Err(_) => {
1981 *conn = None;
1984 }
1985 }
1986 }
1987 false
1988}
1989
1990#[cfg(unix)]
1999#[derive(Debug)]
2000pub struct ForwardingEventStore {
2001 namespace: String,
2002 client: Arc<EventsSplitClient>,
2003}
2004
2005#[cfg(unix)]
2006impl ForwardingEventStore {
2007 pub fn new(namespace: impl Into<String>, client: Arc<EventsSplitClient>) -> Self {
2008 Self {
2009 namespace: namespace.into(),
2010 client,
2011 }
2012 }
2013
2014 fn unexpected(&self, op: &'static str, response: EventsResponse) -> StorageError {
2015 match response {
2016 EventsResponse::Error {
2017 message,
2018 retryable,
2019 writer_task_failure,
2020 } => {
2021 if let Some(failure) = writer_task_failure {
2022 match failure {
2023 WireWriterTaskFailure::RequestFailed { request_state } => {
2024 let source = if retryable {
2025 StorageError::Pool {
2026 operation: op.into(),
2027 message,
2028 }
2029 } else {
2030 StorageError::InvalidInput {
2031 capability: khive_storage::StorageCapability::Events,
2032 operation: op.into(),
2033 message,
2034 }
2035 };
2036 StorageError::WriterTaskRequestFailed {
2037 request_state: request_state.into(),
2038 source: Box::new(source),
2039 }
2040 }
2041 WireWriterTaskFailure::TaskTerminated { request_state } => {
2042 StorageError::WriterTaskTerminated {
2043 request_state: request_state.into(),
2044 }
2045 }
2046 }
2047 } else if retryable {
2048 StorageError::Pool {
2049 operation: op.into(),
2050 message,
2051 }
2052 } else {
2053 StorageError::InvalidInput {
2054 capability: khive_storage::StorageCapability::Events,
2055 operation: op.into(),
2056 message,
2057 }
2058 }
2059 }
2060 other => StorageError::Serialization {
2061 capability: khive_storage::StorageCapability::Events,
2062 message: format!("events daemon returned mismatched response for {op}: {other:?}"),
2063 },
2064 }
2065 }
2066}
2067
2068#[cfg(unix)]
2069#[async_trait]
2070impl EventStore for ForwardingEventStore {
2071 async fn append_event(&self, event: Event) -> StorageResult<()> {
2072 self.client.enqueue(&self.namespace, vec![event]);
2073 Ok(())
2074 }
2075
2076 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2077 let attempted = events.len() as u64;
2078 self.client.enqueue(&self.namespace, events);
2079 Ok(BatchWriteSummary {
2082 attempted,
2083 affected: attempted,
2084 ..BatchWriteSummary::default()
2085 })
2086 }
2087
2088 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2089 let request = EventsRequest::GetEvent {
2090 protocol_version: EVENTS_PROTOCOL_VERSION,
2091 namespace: self.namespace.clone(),
2092 id,
2093 };
2094 match self.client.round_trip(&request).await? {
2095 EventsResponse::Event { event } => Ok(event),
2096 other => Err(self.unexpected("get_event", other)),
2097 }
2098 }
2099
2100 async fn query_events(
2101 &self,
2102 filter: EventFilter,
2103 page: PageRequest,
2104 ) -> StorageResult<Page<Event>> {
2105 let request = EventsRequest::QueryEvents {
2106 protocol_version: EVENTS_PROTOCOL_VERSION,
2107 namespace: self.namespace.clone(),
2108 filter,
2109 page,
2110 };
2111 match self.client.round_trip(&request).await? {
2112 EventsResponse::Pageful { page } => Ok(page),
2113 other => Err(self.unexpected("query_events", other)),
2114 }
2115 }
2116
2117 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2118 let request = EventsRequest::CountEvents {
2119 protocol_version: EVENTS_PROTOCOL_VERSION,
2120 namespace: self.namespace.clone(),
2121 filter,
2122 };
2123 match self.client.round_trip(&request).await? {
2124 EventsResponse::Count { count } => Ok(count),
2125 other => Err(self.unexpected("count_events", other)),
2126 }
2127 }
2128
2129 fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2130 self.client.preflight_store.preflight_event(event)
2132 }
2133
2134 async fn append_events_idempotent(
2135 &self,
2136 events: Vec<Event>,
2137 ) -> StorageResult<IdempotentEventBatchResult> {
2138 let request = EventsRequest::AppendEventsIdempotent {
2139 protocol_version: EVENTS_PROTOCOL_VERSION,
2140 namespace: self.namespace.clone(),
2141 events,
2142 };
2143 match self.client.round_trip(&request).await? {
2144 EventsResponse::Idempotent { result } => Ok(result),
2145 other => Err(self.unexpected("append_events_idempotent", other)),
2146 }
2147 }
2148
2149 fn supports_idempotent_audit_batch(&self) -> bool {
2150 true
2151 }
2152}
2153
2154pub struct SplitEventStore {
2171 legacy: Arc<dyn EventStore>,
2172 lane: Arc<dyn EventStore>,
2173}
2174
2175impl std::fmt::Debug for SplitEventStore {
2176 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2177 f.debug_struct("SplitEventStore").finish_non_exhaustive()
2178 }
2179}
2180
2181impl SplitEventStore {
2182 pub const MAX_MERGED_WINDOW_ROWS: u64 = 100_000;
2190
2191 pub fn new(legacy: Arc<dyn EventStore>, lane: Arc<dyn EventStore>) -> Self {
2192 Self { legacy, lane }
2193 }
2194}
2195
2196#[async_trait]
2197impl EventStore for SplitEventStore {
2198 async fn append_event(&self, event: Event) -> StorageResult<()> {
2199 self.legacy.append_event(event).await
2200 }
2201
2202 async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2203 self.legacy.append_events(events).await
2204 }
2205
2206 async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2207 if let Some(event) = self.legacy.get_event(id).await? {
2210 return Ok(Some(event));
2211 }
2212 self.lane.get_event(id).await
2213 }
2214
2215 async fn query_events(
2216 &self,
2217 filter: EventFilter,
2218 page: PageRequest,
2219 ) -> StorageResult<Page<Event>> {
2220 let window = page.offset.saturating_add(u64::from(page.limit));
2234 if window > Self::MAX_MERGED_WINDOW_ROWS {
2235 return Err(StorageError::InvalidInput {
2236 capability: khive_storage::StorageCapability::Events,
2237 operation: "query_events".into(),
2238 message: format!(
2239 "offset+limit ({window}) exceeds the merged event plane's window bound \
2240 of {}; page deep windows with a `before` cursor at offset 0, or narrow \
2241 the filter",
2242 Self::MAX_MERGED_WINDOW_ROWS
2243 ),
2244 });
2245 }
2246 let prefix = PageRequest {
2247 offset: 0,
2248 limit: window.min(u64::from(u32::MAX)) as u32,
2249 };
2250 let legacy = self
2251 .legacy
2252 .query_events(filter.clone(), prefix.clone())
2253 .await?;
2254 let lane = self.lane.query_events(filter, prefix).await?;
2255 let total = match (legacy.total, lane.total) {
2256 (Some(a), Some(b)) => Some(a + b),
2257 _ => None,
2258 };
2259 let mut items = legacy.items;
2260 items.extend(lane.items);
2261 items.sort_by(|a, b| {
2262 b.created_at
2263 .cmp(&a.created_at)
2264 .then_with(|| b.id.cmp(&a.id))
2265 });
2266 let items = items
2267 .into_iter()
2268 .skip(page.offset as usize)
2269 .take(page.limit as usize)
2270 .collect();
2271 Ok(Page { items, total })
2272 }
2273
2274 async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2275 let legacy = self.legacy.count_events(filter.clone()).await?;
2276 let lane = self.lane.count_events(filter).await?;
2277 Ok(legacy + lane)
2278 }
2279
2280 fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2281 self.lane.preflight_event(event)
2284 }
2285
2286 async fn append_events_idempotent(
2287 &self,
2288 events: Vec<Event>,
2289 ) -> StorageResult<IdempotentEventBatchResult> {
2290 if events.is_empty() {
2300 return self.lane.append_events_idempotent(events).await;
2301 }
2302 let ids: Vec<Uuid> = events.iter().map(|event| event.id).collect();
2303 let probe_limit = u32::try_from(ids.len()).map_err(|_| StorageError::InvalidInput {
2304 capability: khive_storage::StorageCapability::Events,
2305 operation: "append_events_idempotent".into(),
2306 message: format!(
2307 "batch of {} rows exceeds the legacy probe window",
2308 ids.len()
2309 ),
2310 })?;
2311 let existing = self
2312 .legacy
2313 .query_events(
2314 EventFilter {
2315 ids,
2316 ..EventFilter::default()
2317 },
2318 PageRequest {
2319 offset: 0,
2320 limit: probe_limit,
2321 },
2322 )
2323 .await?;
2324 let legacy_resident: std::collections::HashSet<Uuid> =
2325 existing.items.iter().map(|event| event.id).collect();
2326 if legacy_resident.is_empty() {
2327 return self.lane.append_events_idempotent(events).await;
2328 }
2329 let mut legacy_rows = Vec::new();
2330 let mut lane_rows = Vec::new();
2331 let mut routed_to_legacy = Vec::with_capacity(events.len());
2332 for event in events {
2333 if legacy_resident.contains(&event.id) {
2334 routed_to_legacy.push(true);
2335 legacy_rows.push(event);
2336 } else {
2337 routed_to_legacy.push(false);
2338 lane_rows.push(event);
2339 }
2340 }
2341 let legacy_expected = legacy_rows.len();
2342 let lane_expected = lane_rows.len();
2343 let legacy_result = self.legacy.append_events_idempotent(legacy_rows).await?;
2344 let lane_result = if lane_expected == 0 {
2345 IdempotentEventBatchResult { rows: Vec::new() }
2346 } else {
2347 self.lane.append_events_idempotent(lane_rows).await?
2348 };
2349 if legacy_result.rows.len() != legacy_expected || lane_result.rows.len() != lane_expected {
2350 return Err(StorageError::Driver {
2351 capability: khive_storage::StorageCapability::Events,
2352 operation: "append_events_idempotent".into(),
2353 source: format!(
2354 "idempotent sub-batch result length mismatch: legacy {}/{}, lane {}/{}",
2355 legacy_result.rows.len(),
2356 legacy_expected,
2357 lane_result.rows.len(),
2358 lane_expected
2359 )
2360 .into(),
2361 });
2362 }
2363 let mut legacy_iter = legacy_result.rows.into_iter();
2364 let mut lane_iter = lane_result.rows.into_iter();
2365 let rows = routed_to_legacy
2366 .into_iter()
2367 .map(|to_legacy| {
2368 if to_legacy {
2369 legacy_iter.next().expect("length checked above")
2370 } else {
2371 lane_iter.next().expect("length checked above")
2372 }
2373 })
2374 .collect();
2375 Ok(IdempotentEventBatchResult { rows })
2376 }
2377
2378 fn supports_idempotent_audit_batch(&self) -> bool {
2379 self.lane.supports_idempotent_audit_batch() && self.legacy.supports_idempotent_audit_batch()
2383 }
2384}
2385
2386#[cfg(test)]
2387mod tests {
2388 use super::*;
2389
2390 #[test]
2391 fn fixture_registry_cleanup_preserves_other_roots() {
2392 let live_dir = tempfile::tempdir().unwrap();
2393 let _live_guard = TestRegistryGuard::new(live_dir.path());
2394 let live_path = live_dir.path().join("events.db");
2395 let live_backend = direct_backend_for(&live_path).unwrap();
2396
2397 let finished_dir = tempfile::tempdir().unwrap();
2398 let finished_backend = {
2399 let _finished_guard = TestRegistryGuard::new(finished_dir.path());
2400 let backend = direct_backend_for(&finished_dir.path().join("events.db")).unwrap();
2401 let weak = Arc::downgrade(&backend);
2402 drop(backend);
2403 assert!(
2404 weak.upgrade().is_some(),
2405 "registry retains the live fixture"
2406 );
2407 weak
2408 };
2409
2410 assert!(
2411 finished_backend.upgrade().is_none(),
2412 "FD_FIXTURE_REGISTRY_RELEASED"
2413 );
2414 let same_live_backend = direct_backend_for(&live_path).unwrap();
2415 assert!(
2416 Arc::ptr_eq(&live_backend, &same_live_backend),
2417 "FD_FIXTURE_OTHER_ROOT_RETAINED"
2418 );
2419 }
2420
2421 #[cfg(unix)]
2422 mod target_filter_tests {
2423 include!("events_split_target_tests.rs");
2424 }
2425
2426 #[test]
2431 fn sidecar_name_derives_from_the_full_file_name() {
2432 let path = events_db_path_beside(Path::new("/data/khive.db"));
2433 assert!(path.ends_with("khive.db.events.db"), "got {path:?}");
2434 }
2435
2436 #[test]
2437 fn databases_sharing_a_stem_get_distinct_sidecars() {
2438 let a = events_db_path_beside(Path::new("/data/a.db"));
2439 let b = events_db_path_beside(Path::new("/data/a.sqlite"));
2440 assert_ne!(a, b, "a.db and a.sqlite must not share an event plane");
2441 assert!(a.ends_with("a.db.events.db"), "got {a:?}");
2442 assert!(b.ends_with("a.sqlite.events.db"), "got {b:?}");
2443 }
2444
2445 #[cfg(unix)]
2446 #[test]
2447 fn directory_aliases_resolve_to_one_sidecar() {
2448 let dir = tempfile::tempdir().unwrap();
2449 let real = dir.path().join("real");
2450 std::fs::create_dir(&real).unwrap();
2451 let alias = dir.path().join("alias");
2452 std::os::unix::fs::symlink(&real, &alias).unwrap();
2453 assert_eq!(
2454 events_db_path_beside(&real.join("khive.db")),
2455 events_db_path_beside(&alias.join("khive.db")),
2456 "a symlinked spelling of one directory must not mint a second sidecar"
2457 );
2458 }
2459
2460 #[cfg(unix)]
2461 #[test]
2462 fn final_component_aliases_of_one_database_share_one_sidecar() {
2463 let dir = tempfile::tempdir().unwrap();
2468 let real = dir.path().join("real.db");
2469 std::fs::write(&real, b"").unwrap();
2470 let alias = dir.path().join("alias.db");
2471 std::os::unix::fs::symlink(&real, &alias).unwrap();
2472 let real_sidecar = events_db_path_beside(&real);
2473 let alias_sidecar = events_db_path_beside(&alias);
2474 assert_eq!(
2475 real_sidecar, alias_sidecar,
2476 "a symlink alias of one database file must not mint a second event store"
2477 );
2478 assert_eq!(
2479 events_socket_path_beside(&real_sidecar),
2480 events_socket_path_beside(&alias_sidecar),
2481 );
2482 let other = dir.path().join("other.db");
2485 std::fs::write(&other, b"").unwrap();
2486 assert_ne!(events_db_path_beside(&other), real_sidecar);
2487 }
2488
2489 #[cfg(unix)]
2490 #[test]
2491 fn dangling_alias_derives_the_target_sidecar_on_cold_start() {
2492 let dir = tempfile::tempdir().unwrap();
2499 let real = dir.path().join("real.db");
2500 let alias = dir.path().join("alias.db");
2501 std::os::unix::fs::symlink(&real, &alias).unwrap();
2502 assert_eq!(
2504 events_db_path_beside(&alias),
2505 events_db_path_beside(&real),
2506 "a dangling alias must derive the same sidecar its target will use"
2507 );
2508 let chain = dir.path().join("chain.db");
2510 std::os::unix::fs::symlink(&alias, &chain).unwrap();
2511 assert_eq!(events_db_path_beside(&chain), events_db_path_beside(&real));
2512 assert_ne!(
2514 events_db_path_beside(&dir.path().join("unrelated.db")),
2515 events_db_path_beside(&real)
2516 );
2517 }
2518
2519 #[cfg(unix)]
2520 #[test]
2521 fn planted_symlink_at_the_sidecar_path_is_refused() {
2522 let dir = tempfile::tempdir().unwrap();
2526 let _registry_guard = TestRegistryGuard::new(dir.path());
2527 let victim = dir.path().join("victim.txt");
2528 std::fs::write(&victim, b"victim-bytes").unwrap();
2529 let mode_before = victim.metadata().unwrap().permissions();
2530 let sidecar = dir.path().join("khive.db.events.db");
2531 std::os::unix::fs::symlink(&victim, &sidecar).unwrap();
2532 let err = match direct_backend_for(&sidecar) {
2533 Ok(_) => panic!("a planted symlink at the sidecar path must be refused"),
2534 Err(e) => e,
2535 };
2536 assert!(err.to_string().contains("symlink"), "got: {err}");
2537 assert_eq!(std::fs::read(&victim).unwrap(), b"victim-bytes");
2539 assert_eq!(victim.metadata().unwrap().permissions(), mode_before);
2540 let ghost_target = dir.path().join("ghost.db");
2544 let dangling = dir.path().join("other.db.events.db");
2545 std::os::unix::fs::symlink(&ghost_target, &dangling).unwrap();
2546 assert!(direct_backend_for(&dangling).is_err());
2547 assert!(
2548 !ghost_target.exists(),
2549 "refusal must not create the redirect target"
2550 );
2551 let regular = dir.path().join("plain.db.events.db");
2553 direct_backend_for(®ular).expect("regular sidecar path must open");
2554 assert!(regular.exists());
2555 }
2556
2557 #[cfg(unix)]
2558 #[test]
2559 fn deep_symlink_chains_within_the_kernel_bound_derive_the_target_sidecar() {
2560 let dir = tempfile::tempdir().unwrap();
2565 let real = dir.path().join("real.db");
2566 let mut prev = real.clone();
2567 for i in 0..35 {
2568 let link = dir.path().join(format!("hop{i}.db"));
2569 std::os::unix::fs::symlink(&prev, &link).unwrap();
2570 prev = link;
2571 }
2572 assert_eq!(
2573 events_db_path_beside(&prev),
2574 events_db_path_beside(&real),
2575 "a 35-hop dangling chain must derive the target's sidecar"
2576 );
2577 assert_ne!(
2579 events_db_path_beside(&dir.path().join("unrelated.db")),
2580 events_db_path_beside(&real)
2581 );
2582 }
2583
2584 #[cfg(unix)]
2585 #[test]
2586 fn untrusted_parent_directory_is_refused() {
2587 use std::os::unix::fs::PermissionsExt;
2588 let dir = tempfile::tempdir().unwrap();
2593 let _registry_guard = TestRegistryGuard::new(dir.path());
2594 let open_dir = dir.path().join("shared");
2595 std::fs::create_dir(&open_dir).unwrap();
2596 std::fs::set_permissions(&open_dir, std::fs::Permissions::from_mode(0o777)).unwrap();
2597 let refused = direct_backend_for(&open_dir.join("khive.db.events.db"));
2598 match refused {
2599 Ok(_) => panic!("a world-writable events directory must be refused"),
2600 Err(e) => assert!(
2601 e.to_string().contains("untrusted directory"),
2602 "refusal must name the directory trust rule: {e}"
2603 ),
2604 }
2605 let safe_dir = dir.path().join("owned");
2607 std::fs::create_dir(&safe_dir).unwrap();
2608 std::fs::set_permissions(&safe_dir, std::fs::Permissions::from_mode(0o700)).unwrap();
2609 direct_backend_for(&safe_dir.join("khive.db.events.db"))
2610 .expect("an owner-only events directory must open");
2611 }
2612
2613 #[cfg(unix)]
2614 #[test]
2615 fn harden_refuses_a_sidecar_symlink_at_use_time() {
2616 use std::os::unix::fs::PermissionsExt;
2617 let dir = tempfile::tempdir().unwrap();
2622 let victim = dir.path().join("victim.txt");
2623 std::fs::write(&victim, b"v").unwrap();
2624 std::fs::set_permissions(&victim, std::fs::Permissions::from_mode(0o644)).unwrap();
2625 let db = dir.path().join("db.events.db");
2626 std::fs::write(&db, b"").unwrap();
2627 let mut wal = db.as_os_str().to_os_string();
2628 wal.push("-wal");
2629 std::os::unix::fs::symlink(&victim, PathBuf::from(&wal)).unwrap();
2630 let result = harden_events_db_sidecars(&db);
2631 assert!(
2632 result.is_err(),
2633 "a symlinked -wal must be refused at hardening time"
2634 );
2635 let mode = victim.metadata().unwrap().permissions().mode() & 0o7777;
2636 assert_eq!(mode, 0o644, "the link's target must keep its mode");
2637 }
2638
2639 #[cfg(unix)]
2640 #[test]
2641 fn lock_symlink_is_refused_and_target_untouched() {
2642 use std::os::unix::fs::PermissionsExt;
2643 let dir = tempfile::tempdir().unwrap();
2649 let victim = dir.path().join("victim.txt");
2650 std::fs::write(&victim, b"v").unwrap();
2651 std::fs::set_permissions(&victim, std::fs::Permissions::from_mode(0o644)).unwrap();
2652 let socket = dir.path().join("events.sock");
2653 std::os::unix::fs::symlink(&victim, socket.with_extension("lock")).unwrap();
2654 assert!(
2655 try_acquire_events_daemon_guard(&socket).is_none(),
2656 "a symlinked lock entry must refuse the guard"
2657 );
2658 let mode = victim.metadata().unwrap().permissions().mode() & 0o7777;
2659 assert_eq!(mode, 0o644, "the symlink's target must keep its mode");
2660 assert_eq!(std::fs::read(&victim).unwrap(), b"v");
2661 let clean = dir.path().join("clean.sock");
2662 assert!(
2663 try_acquire_events_daemon_guard(&clean).is_some(),
2664 "a plain lock path in the same directory must still acquire"
2665 );
2666 }
2667
2668 #[cfg(unix)]
2669 #[test]
2670 fn planted_wal_symlink_is_refused_before_open() {
2671 let dir = tempfile::tempdir().unwrap();
2675 let _registry_guard = TestRegistryGuard::new(dir.path());
2676 let victim = dir.path().join("victim.txt");
2677 std::fs::write(&victim, b"w").unwrap();
2678 let sidecar = dir.path().join("db.events.db");
2679 let mut wal = sidecar.as_os_str().to_os_string();
2680 wal.push("-wal");
2681 std::os::unix::fs::symlink(&victim, PathBuf::from(&wal)).unwrap();
2682 assert!(direct_backend_for(&sidecar).is_err());
2683 assert!(
2684 !sidecar.exists(),
2685 "refusal must precede creation of the events database"
2686 );
2687 }
2688
2689 #[cfg(unix)]
2690 #[test]
2691 fn hardening_refuses_a_directory_at_the_database_path_without_touching_it() {
2692 use std::os::unix::fs::PermissionsExt;
2693 let dir = tempfile::tempdir().unwrap();
2694 let db = dir.path().join("events.db");
2695 std::fs::create_dir(&db).unwrap();
2696 std::fs::set_permissions(&db, std::fs::Permissions::from_mode(0o700)).unwrap();
2697 let err = harden_events_db_sidecars(&db).unwrap_err().to_string();
2698 assert!(err.contains("not a regular file"), "{err}");
2699 let mode = std::fs::metadata(&db).unwrap().permissions().mode() & 0o777;
2700 assert_eq!(mode, 0o700, "the directory keeps its mode");
2701 }
2702
2703 #[cfg(unix)]
2704 #[test]
2705 fn unopened_check_refuses_a_loosened_or_replaced_sidecar() {
2706 use std::os::unix::fs::PermissionsExt;
2707 let dir = tempfile::tempdir().unwrap();
2708 std::fs::set_permissions(dir.path(), std::fs::Permissions::from_mode(0o700)).unwrap();
2709 let db = dir.path().join("events.db");
2710 let wal = PathBuf::from(format!("{}-wal", db.display()));
2711 let shm = PathBuf::from(format!("{}-shm", db.display()));
2712 for path in [&db, &wal] {
2713 std::fs::write(path, b"").unwrap();
2714 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).unwrap();
2715 }
2716 let before_open = harden_events_db_sidecars(&db).unwrap();
2717 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2719
2720 std::fs::set_permissions(&wal, std::fs::Permissions::from_mode(0o644)).unwrap();
2721 let err = verify_events_db_owner_only_unopened(&db, &before_open)
2722 .unwrap_err()
2723 .to_string();
2724 assert!(
2725 err.contains("events.db-wal") && err.contains("644"),
2726 "{err}"
2727 );
2728 std::fs::set_permissions(&wal, std::fs::Permissions::from_mode(0o600)).unwrap();
2729 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2730
2731 let victim = dir.path().join("victim");
2734 std::fs::write(&victim, b"").unwrap();
2735 std::fs::set_permissions(&victim, std::fs::Permissions::from_mode(0o644)).unwrap();
2736 std::os::unix::fs::symlink(&victim, &shm).unwrap();
2737 let err = verify_events_db_owner_only_unopened(&db, &before_open)
2738 .unwrap_err()
2739 .to_string();
2740 assert!(
2741 err.contains("events.db-shm") && err.contains("regular file"),
2742 "{err}"
2743 );
2744 assert_eq!(
2745 std::fs::metadata(&victim).unwrap().permissions().mode() & 0o777,
2746 0o644
2747 );
2748 }
2749
2750 #[cfg(unix)]
2751 fn write_owner_only_test_file(path: &Path) {
2752 use std::os::unix::fs::PermissionsExt;
2753 std::fs::write(path, b"").unwrap();
2754 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).unwrap();
2755 }
2756
2757 #[cfg(unix)]
2758 #[test]
2759 fn unopened_check_refuses_same_mode_regular_file_replacement() {
2760 for target_index in 0..3 {
2764 let dir = tempfile::tempdir().unwrap();
2765 let db = dir.path().join("events.db");
2766 let targets = events_db_targets(&db);
2767 for path in &targets {
2768 write_owner_only_test_file(path);
2769 }
2770 let before_open = harden_events_db_sidecars(&db).unwrap();
2771 let replacement = dir.path().join("replacement");
2772 write_owner_only_test_file(&replacement);
2773 let replacement_id = EventsFileIdentity::from_metadata(
2774 &std::fs::symlink_metadata(&replacement).unwrap(),
2775 );
2776 assert_ne!(before_open[target_index], Some(replacement_id));
2777 let _connection = rusqlite::Connection::open(&db).unwrap();
2780 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2781 std::fs::rename(&replacement, &targets[target_index]).unwrap();
2782 let error = verify_events_db_owner_only_unopened(&db, &before_open)
2783 .unwrap_err()
2784 .to_string();
2785 assert!(error.contains("changed identity"), "{error}");
2786 assert!(
2787 error.contains(&targets[target_index].display().to_string()),
2788 "{error}"
2789 );
2790 }
2791 }
2792
2793 #[cfg(unix)]
2794 #[test]
2795 fn unopened_check_accepts_sidecars_missing_at_check_time() {
2796 let dir = tempfile::tempdir().unwrap();
2797 let db = dir.path().join("events.db");
2798 let targets = events_db_targets(&db);
2799 for path in &targets {
2800 write_owner_only_test_file(path);
2801 }
2802 let before_open = harden_events_db_sidecars(&db).unwrap();
2803 let _connection = rusqlite::Connection::open(&db).unwrap();
2804 for path in targets.iter().skip(1) {
2805 std::fs::remove_file(path).unwrap();
2806 verify_events_db_owner_only_unopened(&db, &before_open)
2807 .expect("SQLite sidecars may disappear without losing the main database");
2808 }
2809 }
2810
2811 #[cfg(unix)]
2812 #[test]
2813 fn unopened_check_refuses_main_database_removed_after_sqlite_open() {
2814 let dir = tempfile::tempdir().unwrap();
2815 let db = dir.path().join("events.db");
2816 ensure_events_db_owner_only(&db).unwrap();
2817 let before_open = harden_events_db_sidecars(&db).unwrap();
2818 let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2819 backend.events().unwrap();
2820 verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2821 std::fs::remove_file(&db).unwrap();
2822 let error = verify_events_db_owner_only_unopened(&db, &before_open)
2825 .unwrap_err()
2826 .to_string();
2827 assert!(
2828 error.contains("main database") && error.contains("disappeared"),
2829 "{error}"
2830 );
2831 }
2832
2833 #[cfg(unix)]
2834 #[test]
2835 fn unopened_check_accepts_fresh_sqlite_sidecars_and_unchanged_main() {
2836 let dir = tempfile::tempdir().unwrap();
2837 let db = dir.path().join("new.events.db");
2838 assert!(!db.exists());
2839 ensure_events_db_owner_only(&db).unwrap();
2840 let before_open = harden_events_db_sidecars(&db).unwrap();
2841 assert!(before_open[0].is_some());
2842 assert_eq!(&before_open[1..], &[None, None]);
2843 let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2844 backend.events().unwrap();
2845 for path in events_db_targets(&db).iter().skip(1) {
2846 assert!(path.exists(), "SQLite creates {}", path.display());
2847 }
2848 verify_events_db_owner_only_unopened(&db, &before_open)
2849 .expect("new SQLite sidecars have no pre-open identity to contradict");
2850 }
2851
2852 #[test]
2853 fn socket_derives_beside_the_sidecar() {
2854 let db = events_db_path_beside(Path::new("/data/khive.db"));
2855 let socket = events_socket_path_beside(&db);
2856 assert!(socket.ends_with("khive.db.events.sock"), "got {socket:?}");
2857 }
2858
2859 #[test]
2860 fn bare_relative_main_db_yields_an_absolute_sidecar() {
2861 let path = events_db_path_beside(Path::new("khive.db"));
2865 assert!(path.is_absolute(), "got relative {path:?}");
2866 assert!(path.ends_with("khive.db.events.db"), "got {path:?}");
2867 }
2868
2869 #[cfg(unix)]
2870 #[tokio::test]
2871 async fn over_cap_query_page_limit_is_refused_before_materialization() {
2872 let backend = Arc::new(StorageBackend::memory().unwrap());
2877 let stores: NamespaceStores =
2878 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
2879 let request = |limit: u32| EventsRequest::QueryEvents {
2880 protocol_version: EVENTS_PROTOCOL_VERSION,
2881 namespace: "local".to_string(),
2882 filter: EventFilter::default(),
2883 page: PageRequest { offset: 0, limit },
2884 };
2885 let refused =
2886 dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS + 1), &backend, &stores)
2887 .await;
2888 match refused {
2889 EventsResponse::Error {
2890 message,
2891 retryable,
2892 writer_task_failure,
2893 } => {
2894 assert!(
2895 message.contains(&MAX_QUERY_EVENTS_PAGE_ROWS.to_string()),
2896 "refusal must name the cap: {message}"
2897 );
2898 assert!(!retryable, "an over-cap page is not transient");
2899 assert!(writer_task_failure.is_none());
2900 }
2901 other => panic!("over-cap query must be refused, got {other:?}"),
2902 }
2903 let at_cap =
2906 dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS), &backend, &stores).await;
2907 assert!(
2908 matches!(at_cap, EventsResponse::Pageful { .. }),
2909 "at-cap query must reach the store, got {at_cap:?}"
2910 );
2911 }
2912
2913 #[cfg(unix)]
2914 #[test]
2915 fn side_effects_unknown_state_crosses_the_wire_as_terminal_failure() {
2916 let response = storage_error_response(&StorageError::WriterTaskTerminated {
2917 request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
2918 });
2919 assert!(
2920 matches!(
2921 response,
2922 EventsResponse::Error {
2923 writer_task_failure: Some(WireWriterTaskFailure::TaskTerminated {
2924 request_state: WireWriterTaskState::SideEffectsUnknown,
2925 }),
2926 ..
2927 }
2928 ),
2929 "an unknown-commit-state termination must carry its state on the wire"
2930 );
2931 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
2933 assert!(matches!(
2934 busy,
2935 EventsResponse::Error {
2936 writer_task_failure: None,
2937 ..
2938 }
2939 ));
2940 }
2941
2942 #[cfg(unix)]
2943 #[test]
2944 fn proven_rollback_crosses_the_wire_without_claiming_task_termination() {
2945 let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
2946 request_state: khive_storage::WriterTaskRequestState::TransactionRolledBack,
2947 source: Box::new(StorageError::Pool {
2948 operation: "writer_task_commit".into(),
2949 message: "commit refused".into(),
2950 }),
2951 });
2952 assert!(matches!(
2953 response,
2954 EventsResponse::Error {
2955 writer_task_failure: Some(WireWriterTaskFailure::RequestFailed {
2956 request_state: WireWriterTaskState::TransactionRolledBack,
2957 }),
2958 ..
2959 }
2960 ));
2961 }
2962
2963 #[test]
2964 fn error_frames_without_writer_failure_still_parse() {
2965 let bytes = br#"{"kind":"error","message":"boom","retryable":true}"#;
2969 let parsed: EventsResponse = serde_json::from_slice(bytes).expect("stateless frame parses");
2970 assert!(matches!(
2971 parsed,
2972 EventsResponse::Error {
2973 retryable: true,
2974 writer_task_failure: None,
2975 ..
2976 }
2977 ));
2978 }
2979
2980 fn split_retry_event(verb: &str) -> Event {
2981 Event::new(
2982 "test",
2983 verb,
2984 khive_types::EventKind::RecallExecuted,
2985 khive_types::SubstrateKind::Note,
2986 "agent:test",
2987 )
2988 }
2989
2990 fn store_pair(dir: &Path) -> (Arc<dyn EventStore>, Arc<dyn EventStore>) {
2991 let legacy = direct_backend_for(&dir.join("legacy.db"))
2992 .expect("legacy backend")
2993 .events_for_namespace("test")
2994 .expect("legacy store");
2995 let lane = direct_backend_for(&dir.join("lane.db"))
2996 .expect("lane backend")
2997 .events_for_namespace("test")
2998 .expect("lane store");
2999 (legacy, lane)
3000 }
3001
3002 #[tokio::test]
3007 async fn idempotent_retry_of_legacy_resident_rows_does_not_duplicate() {
3008 use khive_storage::event::EventAppendDisposition;
3009
3010 let dir = tempfile::tempdir().unwrap();
3011 let _registry_guard = TestRegistryGuard::new(dir.path());
3012 let (legacy, lane) = store_pair(dir.path());
3013
3014 let resident = split_retry_event("recall");
3015 legacy
3016 .append_events_idempotent(vec![resident.clone()])
3017 .await
3018 .expect("pre-cutover landing");
3019
3020 let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3021 let fresh = split_retry_event("search");
3022 let result = split
3023 .append_events_idempotent(vec![resident.clone(), fresh.clone()])
3024 .await
3025 .expect("mixed retry batch");
3026 assert_eq!(
3027 result.rows,
3028 vec![
3029 EventAppendDisposition::AlreadyPresentIdentical,
3030 EventAppendDisposition::Inserted,
3031 ],
3032 "input order must be preserved across the two sub-batches"
3033 );
3034 assert_eq!(
3035 lane.count_events(EventFilter::default()).await.unwrap(),
3036 1,
3037 "the legacy-resident row must not reach the lane"
3038 );
3039 assert_eq!(
3040 split.count_events(EventFilter::default()).await.unwrap(),
3041 2,
3042 "merged count must not double-count the retried id"
3043 );
3044
3045 let mut mutated = resident.clone();
3048 mutated.verb = "other".to_string();
3049 let conflict = split
3050 .append_events_idempotent(vec![mutated])
3051 .await
3052 .expect("conflicting retry");
3053 assert_eq!(
3054 conflict.rows,
3055 vec![EventAppendDisposition::IdentityConflict]
3056 );
3057 assert_eq!(split.count_events(EventFilter::default()).await.unwrap(), 2);
3058 }
3059
3060 #[tokio::test]
3063 async fn merged_offset_window_is_bounded() {
3064 let dir = tempfile::tempdir().unwrap();
3065 let _registry_guard = TestRegistryGuard::new(dir.path());
3066 let (legacy, lane) = store_pair(dir.path());
3067 let split = SplitEventStore::new(legacy, lane);
3068
3069 let err = split
3070 .query_events(
3071 EventFilter::default(),
3072 PageRequest {
3073 offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS,
3074 limit: 1,
3075 },
3076 )
3077 .await
3078 .expect_err("a window past the bound must be refused");
3079 assert!(
3080 matches!(err, StorageError::InvalidInput { .. }),
3081 "got {err:?}"
3082 );
3083 assert!(
3084 err.to_string().contains("before"),
3085 "the refusal must name the cursor remedy: {err}"
3086 );
3087
3088 let at_bound = split
3089 .query_events(
3090 EventFilter::default(),
3091 PageRequest {
3092 offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS - 1,
3093 limit: 1,
3094 },
3095 )
3096 .await
3097 .expect("a window at the bound is admitted");
3098 assert!(at_bound.items.is_empty());
3099 }
3100
3101 #[cfg(unix)]
3102 #[tokio::test]
3103 async fn stateless_non_retryable_error_maps_terminal_not_writer_terminated() {
3104 let client = EventsSplitClient::new(std::path::PathBuf::from("/tmp/never-bound.sock"))
3111 .expect("client builds");
3112 let store = ForwardingEventStore::new("test", client);
3113 let bytes = br#"{"kind":"error","message":"refused","retryable":false}"#;
3114 let parsed: EventsResponse = serde_json::from_slice(bytes).expect("frame parses");
3115 assert!(matches!(
3116 store.unexpected("append", parsed),
3117 StorageError::InvalidInput { .. }
3118 ));
3119 let with_state = br#"{"kind":"error","message":"died","retryable":false,"writer_task_failure":{"kind":"task_terminated","request_state":"side_effects_unknown"}}"#;
3123 let parsed: EventsResponse =
3124 serde_json::from_slice(with_state).expect("stateful frame parses");
3125 assert!(matches!(
3126 store.unexpected("append", parsed),
3127 StorageError::WriterTaskTerminated { .. }
3128 ));
3129 }
3130
3131 #[cfg(unix)]
3132 #[tokio::test]
3133 async fn writer_task_states_cross_the_wire_verbatim() {
3134 use khive_storage::WriterTaskRequestState as S;
3140 let dir = tempfile::tempdir().unwrap();
3141 let client =
3142 EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3143 let store = ForwardingEventStore::new("test", client);
3144 for state in [
3145 S::NotStarted,
3146 S::TransactionRolledBack,
3147 S::SideEffectsUnknown,
3148 ] {
3149 let response = storage_error_response(&StorageError::WriterTaskTerminated {
3150 request_state: state,
3151 });
3152 let bytes = serde_json::to_vec(&response).unwrap();
3154 let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3155 let err = store.unexpected("append_events_idempotent", parsed);
3156 assert!(
3157 matches!(
3158 err,
3159 StorageError::WriterTaskTerminated { request_state } if request_state == state
3160 ),
3161 "state {state:?} did not survive the socket: got {err:?}"
3162 );
3163 }
3164 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3166 assert!(matches!(
3167 busy,
3168 EventsResponse::Error {
3169 writer_task_failure: None,
3170 ..
3171 }
3172 ));
3173 }
3174
3175 #[cfg(unix)]
3176 #[tokio::test]
3177 async fn proven_rollback_round_trips_as_non_terminal_request_failure() {
3178 use khive_storage::WriterTaskRequestState as S;
3179
3180 let dir = tempfile::tempdir().unwrap();
3181 let client =
3182 EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3183 let store = ForwardingEventStore::new("test", client);
3184 let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
3185 request_state: S::TransactionRolledBack,
3186 source: Box::new(StorageError::Pool {
3187 operation: "writer_task_commit".into(),
3188 message: "commit refused".into(),
3189 }),
3190 });
3191 let bytes = serde_json::to_vec(&response).unwrap();
3192 let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3193 let err = store.unexpected("append_events_idempotent", parsed);
3194 assert!(
3195 matches!(
3196 err,
3197 StorageError::WriterTaskRequestFailed {
3198 request_state: S::TransactionRolledBack,
3199 ..
3200 }
3201 ),
3202 "proven rollback must not reconstruct as a terminal writer: {err:?}"
3203 );
3204 }
3205
3206 #[cfg(unix)]
3207 #[test]
3208 fn direct_backend_hardens_preexisting_db_and_sidecars() {
3209 use std::os::unix::fs::PermissionsExt;
3210 let dir = tempfile::tempdir().unwrap();
3211 let _registry_guard = TestRegistryGuard::new(dir.path());
3212 let db = dir.path().join("pre-existing.events.db");
3213 let wal = dir.path().join("pre-existing.events.db-wal");
3214 std::fs::write(&db, b"").unwrap();
3215 std::fs::write(&wal, b"").unwrap();
3216 for path in [&db, &wal] {
3217 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o644)).unwrap();
3218 }
3219 assert_eq!(
3221 std::fs::metadata(&db).unwrap().permissions().mode() & 0o777,
3222 0o644
3223 );
3224 direct_backend_for(&db).expect("writable open succeeds");
3225 for path in [&db, &wal] {
3226 assert_eq!(
3227 std::fs::metadata(path).unwrap().permissions().mode() & 0o777,
3228 0o600,
3229 "pre-existing {} must be tightened to owner-only",
3230 path.display()
3231 );
3232 }
3233 }
3234
3235 #[cfg(unix)]
3236 #[test]
3237 fn namespace_store_cache_is_bounded_and_trim_normalized() {
3238 let backend = StorageBackend::memory().unwrap();
3239 let stores: NamespaceStores =
3240 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3241 namespace_store_with_cap(&backend, &stores, "alpha", 1).expect("first store");
3244 namespace_store_with_cap(&backend, &stores, "beta", 1).expect("second store evicts");
3245 {
3246 let map = stores.lock().unwrap();
3247 assert_eq!(map.len(), 1, "cache must not grow past its cap");
3248 assert!(
3249 map.contains_key("beta"),
3250 "the newest namespace must be admitted at the cap"
3251 );
3252 }
3253 namespace_store_with_cap(&backend, &stores, "beta", 1).expect("cached beta");
3256 namespace_store_with_cap(&backend, &stores, " beta ", 1).expect("trimmed spelling");
3257 {
3258 let map = stores.lock().unwrap();
3259 assert_eq!(
3260 map.len(),
3261 1,
3262 "spellings of one namespace must share one entry"
3263 );
3264 assert!(map.contains_key("beta"), "trimmed key is the cache key");
3265 }
3266 }
3267
3268 #[cfg(unix)]
3269 #[tokio::test]
3270 async fn frame_budget_admits_before_allocating_and_releases_after() {
3271 let (mut client, mut server) = UnixStream::pair().expect("socketpair");
3272 let budget = Arc::new(tokio::sync::Semaphore::new(1024));
3273 let payload = vec![7u8; 100];
3274 write_frame(&mut client, &payload).await.expect("write");
3275 let (bytes, permit) = read_frame_budgeted(&mut server, &budget)
3276 .await
3277 .expect("budgeted read");
3278 assert_eq!(bytes, payload);
3279 assert_eq!(budget.available_permits(), 1024 - 100);
3281 drop(permit);
3283 assert_eq!(budget.available_permits(), 1024);
3284 let mut oversized = Vec::from(((crate::daemon::MAX_FRAME_BYTES + 1) as u32).to_be_bytes());
3287 oversized.extend_from_slice(&[0u8; 8]);
3288 use tokio::io::AsyncWriteExt;
3289 client.write_all(&oversized).await.expect("raw prefix");
3290 let err = read_frame_budgeted(&mut server, &budget)
3291 .await
3292 .expect_err("oversized declaration must refuse");
3293 assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
3294 assert_eq!(
3295 budget.available_permits(),
3296 1024,
3297 "refusal must not consume budget"
3298 );
3299 }
3300
3301 #[cfg(unix)]
3302 #[tokio::test]
3303 async fn dispatch_rejects_invalid_wire_namespace() {
3304 let backend = Arc::new(StorageBackend::memory().unwrap());
3305 let stores: NamespaceStores =
3306 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3307 let huge = "n".repeat(64 * 1024);
3311 let response = dispatch_events_request(
3312 EventsRequest::CountEvents {
3313 protocol_version: EVENTS_PROTOCOL_VERSION,
3314 namespace: huge,
3315 filter: Default::default(),
3316 },
3317 &backend,
3318 &stores,
3319 )
3320 .await;
3321 assert!(
3322 matches!(
3323 &response,
3324 EventsResponse::Error {
3325 retryable: false,
3326 writer_task_failure: None,
3327 ..
3328 }
3329 ),
3330 "oversized namespace must be a typed refusal, got {response:?}"
3331 );
3332 assert_eq!(
3333 stores.lock().unwrap().len(),
3334 0,
3335 "a rejected namespace must never enter the cache"
3336 );
3337 let ok = dispatch_events_request(
3339 EventsRequest::CountEvents {
3340 protocol_version: EVENTS_PROTOCOL_VERSION,
3341 namespace: "local".to_string(),
3342 filter: Default::default(),
3343 },
3344 &backend,
3345 &stores,
3346 )
3347 .await;
3348 assert!(
3349 matches!(ok, EventsResponse::Count { .. }),
3350 "valid namespace must dispatch, got {ok:?}"
3351 );
3352 }
3353
3354 #[cfg(unix)]
3355 #[tokio::test]
3356 async fn shutdown_counts_and_logs_queued_and_in_flight_batches_as_dropped() {
3357 let dir = tempfile::tempdir().unwrap();
3366 let socket = dir.path().join("hung-shutdown.sock");
3367 let listener = tokio::net::UnixListener::bind(&socket).unwrap();
3368 let _server = tokio::spawn(async move {
3372 let mut held = Vec::new();
3373 loop {
3374 if let Ok((stream, _)) = listener.accept().await {
3375 held.push(stream);
3376 }
3377 }
3378 });
3379
3380 let (tx, rx) = tokio::sync::mpsc::channel::<QueuedAppend>(8);
3381 let counters = Arc::new(ForwardingCounters::default());
3382 let outage_logged = Arc::new(AtomicBool::new(false));
3383 let shutdown = tokio_util::sync::CancellationToken::new();
3384
3385 let forwarder = tokio::spawn(run_forwarder(
3386 socket,
3387 rx,
3388 Arc::clone(&counters),
3389 outage_logged,
3390 Duration::from_secs(30),
3391 shutdown.clone(),
3392 ));
3393
3394 fn probe_event(tag: &str) -> Event {
3395 Event::new(
3396 "test",
3397 tag,
3398 khive_types::EventKind::Audit,
3399 khive_types::SubstrateKind::Event,
3400 "tester",
3401 )
3402 }
3403
3404 tx.try_send(QueuedAppend::unmetered(
3405 "test",
3406 vec![probe_event("in-flight")],
3407 Arc::clone(&counters),
3408 ))
3409 .expect("queue has room for the in-flight batch");
3410 tokio::time::sleep(Duration::from_millis(200)).await;
3415
3416 tx.try_send(QueuedAppend::unmetered(
3417 "test",
3418 vec![probe_event("queued-1"), probe_event("queued-2")],
3419 Arc::clone(&counters),
3420 ))
3421 .expect("queue has room for the first queued batch");
3422 tx.try_send(QueuedAppend::unmetered(
3423 "test",
3424 vec![probe_event("queued-3")],
3425 Arc::clone(&counters),
3426 ))
3427 .expect("queue has room for the second queued batch");
3428
3429 shutdown.cancel();
3430
3431 tokio::time::timeout(Duration::from_secs(5), forwarder)
3432 .await
3433 .expect("forwarder must exit promptly on shutdown")
3434 .expect("forwarder task must not panic");
3435
3436 assert_eq!(
3437 counters.dropped_batches.load(Ordering::Relaxed),
3438 3,
3439 "the in-flight batch and both queued batches must all be counted as dropped"
3440 );
3441 assert_eq!(
3442 counters.dropped_events.load(Ordering::Relaxed),
3443 4,
3444 "1 in-flight + 2 + 1 queued events must all be counted as dropped"
3445 );
3446 assert_eq!(
3447 counters.forwarded_batches.load(Ordering::Relaxed),
3448 0,
3449 "a hung peer never acknowledges anything in this test"
3450 );
3451 }
3452
3453 #[cfg(unix)]
3454 #[tokio::test]
3455 async fn shutdown_during_backoff_drains_queued_batches() {
3456 let dir = tempfile::tempdir().unwrap();
3463 let socket = dir.path().join("no-listener.sock");
3464
3465 let (tx, rx) = tokio::sync::mpsc::channel::<QueuedAppend>(8);
3466 let counters = Arc::new(ForwardingCounters::default());
3467 let outage_logged = Arc::new(AtomicBool::new(false));
3468 let shutdown = tokio_util::sync::CancellationToken::new();
3469
3470 let forwarder = tokio::spawn(run_forwarder(
3471 socket,
3472 rx,
3473 Arc::clone(&counters),
3474 outage_logged,
3475 Duration::from_secs(30),
3476 shutdown.clone(),
3477 ));
3478
3479 fn probe_event(tag: &str) -> Event {
3480 Event::new(
3481 "test",
3482 tag,
3483 khive_types::EventKind::Audit,
3484 khive_types::SubstrateKind::Event,
3485 "tester",
3486 )
3487 }
3488
3489 tx.try_send(QueuedAppend::unmetered(
3490 "test",
3491 vec![probe_event("failed-delivery")],
3492 Arc::clone(&counters),
3493 ))
3494 .expect("queue has room for the first batch");
3495
3496 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
3499 loop {
3500 if counters.dropped_batches.load(Ordering::Relaxed) >= 1 {
3501 break;
3502 }
3503 assert!(
3504 tokio::time::Instant::now() < deadline,
3505 "forwarder never dropped the first (unreachable-daemon) batch"
3506 );
3507 tokio::time::sleep(Duration::from_millis(20)).await;
3508 }
3509
3510 tx.try_send(QueuedAppend::unmetered(
3511 "test",
3512 vec![
3513 probe_event("queued-during-backoff-1"),
3514 probe_event("queued-during-backoff-2"),
3515 ],
3516 Arc::clone(&counters),
3517 ))
3518 .expect("queue has room for the batch queued during backoff");
3519
3520 shutdown.cancel();
3521
3522 tokio::time::timeout(Duration::from_secs(5), forwarder)
3523 .await
3524 .expect("forwarder must exit promptly on shutdown, not wait out the full backoff")
3525 .expect("forwarder task must not panic");
3526
3527 assert_eq!(
3528 counters.dropped_batches.load(Ordering::Relaxed),
3529 2,
3530 "the failed-delivery batch and the batch queued during backoff must both be counted as dropped"
3531 );
3532 assert_eq!(
3533 counters.dropped_events.load(Ordering::Relaxed),
3534 3,
3535 "1 failed-delivery + 2 queued-during-backoff events must all be counted as dropped"
3536 );
3537 assert_eq!(
3538 counters.forwarded_batches.load(Ordering::Relaxed),
3539 0,
3540 "no listener is bound, so nothing can ever be acknowledged"
3541 );
3542 }
3543
3544 #[cfg(unix)]
3545 #[tokio::test]
3546 async fn forwarder_abandons_hung_delivery() {
3547 let dir = tempfile::tempdir().unwrap();
3551 let socket = dir.path().join("hung.sock");
3552 let listener = tokio::net::UnixListener::bind(&socket).unwrap();
3553 let _server = tokio::spawn(async move {
3555 let mut held = Vec::new();
3556 loop {
3557 if let Ok((stream, _)) = listener.accept().await {
3558 held.push(stream);
3559 }
3560 }
3561 });
3562 let client = EventsSplitClient::new_with_queue_depth_and_delivery_timeout(
3563 socket.clone(),
3564 4,
3565 Duration::from_millis(100),
3566 )
3567 .expect("client builds");
3568 let event = Event::new(
3569 "test",
3570 "noop",
3571 khive_types::EventKind::Audit,
3572 khive_types::SubstrateKind::Event,
3573 "tester",
3574 );
3575 client.enqueue("test", vec![event]);
3576 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
3577 loop {
3578 if client.metrics().dropped_batches >= 1 {
3579 break;
3580 }
3581 assert!(
3582 tokio::time::Instant::now() < deadline,
3583 "forwarder never abandoned the hung delivery: {:?}",
3584 client.metrics()
3585 );
3586 tokio::time::sleep(Duration::from_millis(20)).await;
3587 }
3588 }
3589
3590 #[cfg(unix)]
3591 #[test]
3592 fn wire_retryability_follows_the_storage_classifier() {
3593 let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3594 assert!(
3595 matches!(
3596 busy,
3597 EventsResponse::Error {
3598 retryable: true,
3599 ..
3600 }
3601 ),
3602 "transient writer contention must stay retryable across the socket"
3603 );
3604 let terminated = storage_error_response(&StorageError::WriterTaskTerminated {
3607 request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3608 });
3609 assert!(
3610 matches!(
3611 terminated,
3612 EventsResponse::Error {
3613 retryable: false,
3614 ..
3615 }
3616 ),
3617 "a terminated writer must not be reported transient"
3618 );
3619 }
3620
3621 use khive_storage::event::EventAppendDisposition;
3622 use khive_types::{EventKind, SubstrateKind};
3623
3624 fn test_event(namespace: &str) -> Event {
3625 Event::new(
3626 namespace,
3627 "test.verb",
3628 EventKind::Audit,
3629 SubstrateKind::Event,
3630 "actor:test",
3631 )
3632 }
3633
3634 #[cfg(unix)]
3636 async fn boot_daemon(dir: &tempfile::TempDir) -> (PathBuf, PathBuf) {
3637 let db = dir.path().join("events.db");
3638 let socket = dir.path().join("events.sock");
3639 let (db_clone, socket_clone) = (db.clone(), socket.clone());
3640 tokio::spawn(async move {
3641 let _ = run_events_daemon(&db_clone, &socket_clone).await;
3642 });
3643 for _ in 0..100 {
3644 if UnixStream::connect(&socket).await.is_ok() {
3645 return (db, socket);
3646 }
3647 tokio::time::sleep(Duration::from_millis(20)).await;
3648 }
3649 panic!("events daemon did not come up on {}", socket.display());
3650 }
3651
3652 #[tokio::test]
3657 async fn split_store_routes_plain_to_legacy_idempotent_to_lane_and_merges_reads() {
3658 let dir = tempfile::tempdir().expect("tempdir");
3659 let _registry_guard = TestRegistryGuard::new(dir.path());
3660 let legacy_backend =
3661 direct_backend_for(&dir.path().join("legacy.db")).expect("legacy backend");
3662 let lane_backend = direct_backend_for(&dir.path().join("lane.db")).expect("lane backend");
3663 let legacy = legacy_backend
3664 .events_for_namespace("local")
3665 .expect("legacy store");
3666 let lane = lane_backend
3667 .events_for_namespace("local")
3668 .expect("lane store");
3669 let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3670
3671 let plain = test_event("local");
3672 let plain_id = plain.id;
3673 split.append_event(plain).await.expect("plain append");
3674
3675 let audit = test_event("local");
3676 let audit_id = audit.id;
3677 let result = split
3678 .append_events_idempotent(vec![audit])
3679 .await
3680 .expect("idempotent append");
3681 assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3682
3683 assert!(legacy
3685 .get_event(plain_id)
3686 .await
3687 .expect("legacy get")
3688 .is_some());
3689 assert!(lane.get_event(plain_id).await.expect("lane get").is_none());
3690 assert!(legacy
3691 .get_event(audit_id)
3692 .await
3693 .expect("legacy get")
3694 .is_none());
3695 assert!(lane.get_event(audit_id).await.expect("lane get").is_some());
3696
3697 assert!(split
3699 .get_event(plain_id)
3700 .await
3701 .expect("split get")
3702 .is_some());
3703 assert!(split
3704 .get_event(audit_id)
3705 .await
3706 .expect("split get")
3707 .is_some());
3708 assert_eq!(
3709 split
3710 .count_events(EventFilter::default())
3711 .await
3712 .expect("split count"),
3713 2
3714 );
3715 let page = split
3716 .query_events(
3717 EventFilter::default(),
3718 PageRequest {
3719 offset: 0,
3720 limit: 10,
3721 },
3722 )
3723 .await
3724 .expect("split query");
3725 let ids: Vec<Uuid> = page.items.iter().map(|e| e.id).collect();
3726 assert!(ids.contains(&plain_id) && ids.contains(&audit_id));
3727
3728 let mut seen = Vec::new();
3731 for offset in 0..2 {
3732 let page = split
3733 .query_events(EventFilter::default(), PageRequest { offset, limit: 1 })
3734 .await
3735 .expect("windowed query");
3736 assert_eq!(page.items.len(), 1);
3737 seen.push(page.items[0].id);
3738 }
3739 seen.sort();
3740 let mut expected = vec![plain_id, audit_id];
3741 expected.sort();
3742 assert_eq!(seen, expected);
3743 }
3744
3745 #[tokio::test]
3758 async fn raw_sql_consumers_of_the_legacy_events_table_still_see_plain_appends() {
3759 use khive_storage::types::{SqlStatement, SqlValue};
3760
3761 let dir = tempfile::tempdir().expect("tempdir");
3762 let _registry_guard = TestRegistryGuard::new(dir.path());
3763 let legacy_backend =
3764 direct_backend_for(&dir.path().join("legacy-guard.db")).expect("legacy backend");
3765 let lane_backend =
3766 direct_backend_for(&dir.path().join("lane-guard.db")).expect("lane backend");
3767 let split = SplitEventStore::new(
3768 legacy_backend
3769 .events_for_namespace("local")
3770 .expect("legacy store"),
3771 lane_backend
3772 .events_for_namespace("local")
3773 .expect("lane store"),
3774 );
3775
3776 let provenance = test_event("local");
3779 let provenance_id = provenance.id;
3780 split.append_event(provenance).await.expect("plain append");
3781 let audit = test_event("local");
3782 let audit_id = audit.id;
3783 split
3784 .append_events_idempotent(vec![audit])
3785 .await
3786 .expect("idempotent append");
3787
3788 let count_by_raw_sql = |backend: Arc<StorageBackend>, id: Uuid| async move {
3789 let mut reader = backend.sql().reader().await.expect("sql reader");
3790 let rows = reader
3791 .query_all(SqlStatement {
3792 sql: "SELECT actor FROM events WHERE id = ?1".to_string(),
3793 params: vec![SqlValue::Text(id.to_string())],
3794 label: None,
3795 })
3796 .await
3797 .expect("raw events query");
3798 rows.len()
3799 };
3800
3801 assert_eq!(
3804 count_by_raw_sql(Arc::clone(&legacy_backend), provenance_id).await,
3805 1,
3806 "plain appends must stay visible to raw-SQL consumers of the legacy events table"
3807 );
3808 assert_eq!(
3811 count_by_raw_sql(Arc::clone(&legacy_backend), audit_id).await,
3812 0,
3813 "audit-lane rows must not land in the legacy events table"
3814 );
3815 assert_eq!(
3819 count_by_raw_sql(Arc::clone(&lane_backend), audit_id).await,
3820 1,
3821 "the raw query must prove it can find rows where they actually live"
3822 );
3823 }
3824
3825 #[cfg(unix)]
3828 #[tokio::test]
3829 async fn idempotent_append_round_trips_through_the_daemon() {
3830 let dir = tempfile::tempdir().expect("tempdir");
3831 let (_db, socket) = boot_daemon(&dir).await;
3832 let client = EventsSplitClient::new(socket).expect("client");
3833 let store = ForwardingEventStore::new("local", client);
3834
3835 let event = test_event("local");
3836 let id = event.id;
3837 let result = store
3838 .append_events_idempotent(vec![event.clone()])
3839 .await
3840 .expect("idempotent append over socket");
3841 assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3842
3843 let retry = store
3846 .append_events_idempotent(vec![event])
3847 .await
3848 .expect("idempotent retry");
3849 assert_eq!(
3850 retry.rows,
3851 vec![EventAppendDisposition::AlreadyPresentIdentical]
3852 );
3853
3854 let fetched = store.get_event(id).await.expect("get over socket");
3855 assert_eq!(fetched.map(|e| e.id), Some(id));
3856 let count = store
3857 .count_events(EventFilter::default())
3858 .await
3859 .expect("count over socket");
3860 assert_eq!(count, 1);
3861 }
3862
3863 #[cfg(unix)]
3866 #[tokio::test]
3867 async fn fire_and_forget_append_lands_in_the_daemon_store() {
3868 let dir = tempfile::tempdir().expect("tempdir");
3869 let (_db, socket) = boot_daemon(&dir).await;
3870 let client = EventsSplitClient::new(socket).expect("client");
3871 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3872
3873 store
3874 .append_event(test_event("local"))
3875 .await
3876 .expect("append_event is fire-and-forget");
3877
3878 let (count, metrics) = tokio::time::timeout(Duration::from_secs(30), async {
3882 loop {
3883 let count = store
3884 .count_events(EventFilter::default())
3885 .await
3886 .expect("count over socket");
3887 let metrics = client.metrics();
3888 if (count == 1 && metrics.forwarded_events >= 1) || metrics.dropped_events > 0 {
3889 break (count, metrics);
3890 }
3891 tokio::time::sleep(Duration::from_millis(20)).await;
3892 }
3893 })
3894 .await
3895 .expect("forwarder must deliver or report a drop within 30 seconds");
3896 assert_eq!(
3897 metrics.dropped_events, 0,
3898 "forwarder dropped an event before it reached the daemon store"
3899 );
3900 assert_eq!(count, 1, "forwarded event must land in the daemon store");
3901 assert!(metrics.forwarded_events >= 1, "delivery must be counted");
3902 }
3903
3904 #[cfg(unix)]
3908 #[tokio::test]
3909 async fn queue_overflow_drops_and_counts_but_never_errors() {
3910 let dir = tempfile::tempdir().expect("tempdir");
3911 let dead_socket = dir.path().join("nobody-home.sock");
3912 let client = EventsSplitClient::new_with_queue_depth(dead_socket, 2).expect("client");
3913 let store = ForwardingEventStore::new("local", Arc::clone(&client));
3914
3915 for _ in 0..20 {
3916 store
3917 .append_event(test_event("local"))
3918 .await
3919 .expect("append_event must not error under overflow");
3920 }
3921 let metrics = client.metrics();
3922 assert!(
3923 metrics.dropped_events > 0,
3924 "overflow must register in the drop counter, got {metrics:?}"
3925 );
3926 }
3927
3928 #[cfg(unix)]
3932 #[tokio::test]
3933 async fn dead_socket_reads_fail_typed_and_preflight_stays_local() {
3934 let dir = tempfile::tempdir().expect("tempdir");
3935 let dead_socket = dir.path().join("nobody-home.sock");
3936 let client = EventsSplitClient::new(dead_socket).expect("client");
3937 let store = ForwardingEventStore::new("local", client);
3938
3939 let error = store
3940 .get_event(Uuid::new_v4())
3941 .await
3942 .expect_err("read against a dead socket must fail");
3943 assert!(
3944 matches!(error, StorageError::Pool { .. }),
3945 "expected the typed unreachable error, got {error:?}"
3946 );
3947
3948 assert!(store.supports_idempotent_audit_batch());
3951 store
3952 .preflight_event(&test_event("local"))
3953 .expect("offline preflight validates a well-formed event");
3954 }
3955
3956 #[cfg(unix)]
3959 #[tokio::test]
3960 async fn connected_events_peer_closing_without_response_is_not_unreachable() {
3961 let dir = tempfile::tempdir().unwrap();
3962 let socket = dir.path().join("closes-after-request.sock");
3963 let listener = UnixListener::bind(&socket).unwrap();
3964 let server = tokio::spawn(async move {
3965 let (mut stream, _) = listener.accept().await.unwrap();
3966 read_frame(&mut stream)
3967 .await
3968 .expect("complete request arrives");
3969 });
3971 let client = EventsSplitClient::new(socket).unwrap();
3972 let store = ForwardingEventStore::new("local", client);
3973 let error = store.get_event(Uuid::new_v4()).await.unwrap_err();
3974 assert!(
3975 matches!(error, StorageError::Serialization { .. }),
3976 "a post-connect response failure must not be Pool/unreachable: {error:?}"
3977 );
3978 server.await.unwrap();
3979 }
3980
3981 #[cfg(unix)]
3985 #[tokio::test]
3986 async fn query_page_over_frame_cap_returns_typed_size_refusal() {
3987 let dir = tempfile::tempdir().unwrap();
3988 let socket = dir.path().join("large-query.sock");
3989 let backend = Arc::new(StorageBackend::memory().unwrap());
3990 let event_store = backend.events_for_namespace("local").unwrap();
3991 let events: Vec<_> = (0..MAX_QUERY_EVENTS_PAGE_ROWS)
3992 .map(|_| {
3993 test_event("local").with_payload(serde_json::json!({"data": "x".repeat(3 * 1024)}))
3994 })
3995 .collect();
3996 event_store
3997 .append_events(events)
3998 .await
3999 .expect("seed large page");
4000
4001 let listener = UnixListener::bind(&socket).unwrap();
4002 let stores: NamespaceStores =
4003 Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
4004 let server = tokio::spawn(async move {
4005 let (stream, _) = listener.accept().await.unwrap();
4006 serve_events_conn(
4007 stream,
4008 backend,
4009 stores,
4010 Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES)),
4011 )
4012 .await;
4013 });
4014 let client = EventsSplitClient::new(socket).unwrap();
4015 let store = ForwardingEventStore::new("local", client);
4016 let error = store
4017 .query_events(
4018 EventFilter::default(),
4019 PageRequest {
4020 offset: 0,
4021 limit: MAX_QUERY_EVENTS_PAGE_ROWS,
4022 },
4023 )
4024 .await
4025 .expect_err("the full page cannot fit one frame");
4026 assert!(
4027 matches!(error, StorageError::InvalidInput { .. }),
4028 "expected non-retryable frame-size refusal, got {error:?}"
4029 );
4030 assert!(
4031 error.to_string().contains("response_frame_size_limit"),
4032 "{error}"
4033 );
4034 assert!(
4035 error.to_string().contains("request a narrower page"),
4036 "{error}"
4037 );
4038 server.abort();
4039 }
4040
4041 #[cfg(unix)]
4045 #[tokio::test]
4046 async fn forwarding_queue_enforces_serialized_byte_and_frame_limits_with_drop_metrics() {
4047 let dir = tempfile::tempdir().unwrap();
4048 let socket = dir.path().join("stalled-forwarder.sock");
4049 let listener = UnixListener::bind(&socket).unwrap();
4050 let server = tokio::spawn(async move {
4051 let (stream, _) = listener.accept().await.unwrap();
4052 let _hold_open = stream;
4053 std::future::pending::<()>().await;
4054 });
4055 let event = test_event("local").with_payload(serde_json::json!({"data": "x".repeat(4096)}));
4056 let first = vec![event.clone(), event.clone()];
4057 let request_bytes = serde_json::to_vec(&EventsRequest::AppendEvents {
4058 protocol_version: EVENTS_PROTOCOL_VERSION,
4059 namespace: "local".into(),
4060 events: first.clone(),
4061 })
4062 .unwrap()
4063 .len();
4064 let byte_budget = request_bytes + 64;
4065 let client = EventsSplitClient::new_with_limits_and_delivery_timeout(
4066 socket,
4067 8,
4068 byte_budget,
4069 Duration::from_secs(30),
4070 )
4071 .unwrap();
4072 let store = ForwardingEventStore::new("local", Arc::clone(&client));
4073 store.append_events(first.clone()).await.unwrap();
4074 assert_eq!(client.metrics().queued_bytes, request_bytes);
4075
4076 store.append_events(first).await.unwrap();
4077 let metrics = client.metrics();
4078 assert_eq!(metrics.dropped_batches, 1);
4079 assert_eq!(metrics.dropped_events, 2);
4080 assert_eq!(metrics.queued_bytes, request_bytes);
4081 assert!(metrics.queued_bytes <= byte_budget);
4082
4083 let oversized = test_event("local")
4086 .with_payload(serde_json::json!({"data": "x".repeat(crate::daemon::MAX_FRAME_BYTES)}));
4087 store.append_event(oversized).await.unwrap();
4088 let metrics = client.metrics();
4089 assert_eq!(metrics.dropped_batches, 2);
4090 assert_eq!(metrics.dropped_events, 3);
4091 assert_eq!(metrics.queued_bytes, request_bytes);
4092 server.abort();
4093 }
4094
4095 #[cfg(unix)]
4098 #[tokio::test]
4099 async fn protocol_version_skew_is_a_typed_refusal() {
4100 let dir = tempfile::tempdir().expect("tempdir");
4101 let (_db, socket) = boot_daemon(&dir).await;
4102
4103 let mut stream = UnixStream::connect(&socket).await.expect("connect");
4104 let request = EventsRequest::CountEvents {
4105 protocol_version: EVENTS_PROTOCOL_VERSION + 1,
4106 namespace: "local".into(),
4107 filter: EventFilter::default(),
4108 };
4109 let payload = serde_json::to_vec(&request).expect("serialize");
4110 write_frame(&mut stream, &payload).await.expect("write");
4111 let bytes = read_frame(&mut stream).await.expect("read");
4112 let response: EventsResponse = serde_json::from_slice(&bytes).expect("parse");
4113 match response {
4114 EventsResponse::Error {
4115 message, retryable, ..
4116 } => {
4117 assert!(!retryable, "version skew is not retryable");
4118 assert!(message.contains("protocol version"), "message: {message}");
4119 }
4120 other => panic!("expected a typed refusal, got {other:?}"),
4121 }
4122 }
4123
4124 #[cfg(unix)]
4128 #[test]
4129 fn events_daemon_guard_is_exclusive_then_reusable() {
4130 let dir = tempfile::tempdir().expect("tempdir");
4131 let socket = dir.path().join("events.sock");
4132 let first = try_acquire_events_daemon_guard(&socket).expect("first acquire");
4133 assert!(
4134 try_acquire_events_daemon_guard(&socket).is_none(),
4135 "second acquire must fail while the first guard is held"
4136 );
4137 drop(first);
4138 assert!(
4139 try_acquire_events_daemon_guard(&socket).is_some(),
4140 "acquire must succeed again after the guard is released"
4141 );
4142 }
4143
4144 #[cfg(unix)]
4151 #[tokio::test]
4152 async fn read_only_runtime_never_creates_an_events_db() {
4153 use crate::{KhiveRuntime, Namespace, RuntimeConfig};
4154
4155 let dir = tempfile::tempdir().expect("tempdir");
4156 let _registry_guard = TestRegistryGuard::new(dir.path());
4157 let main_db = dir.path().join("main.db");
4158 drop(
4160 KhiveRuntime::new_for_test(RuntimeConfig {
4161 db_path: Some(main_db.clone()),
4162 ..RuntimeConfig::no_embeddings()
4163 })
4164 .expect("create main db"),
4165 );
4166 let events_db = events_db_path_beside(&main_db);
4167 assert!(!events_db.exists(), "precondition: no events db yet");
4168
4169 khive_storage::test_support::freeze_snapshot_sidecars(&main_db);
4172
4173 let split_config = |db: PathBuf| RuntimeConfig {
4174 db_path: Some(main_db.clone()),
4175 events_split: Some(EventsSplitConfig {
4176 db_path: db,
4177 socket_path: None,
4178 }),
4179 ..RuntimeConfig::no_embeddings()
4180 };
4181
4182 let ro = KhiveRuntime::new_readonly_for_test(split_config(events_db.clone()))
4185 .expect("read-only runtime");
4186 let token = ro.authorize(Namespace::local()).expect("token");
4187 let store = ro.events(&token).expect("events store");
4188 let count = store
4189 .count_events(EventFilter::default())
4190 .await
4191 .expect("count through legacy-only plane");
4192 assert_eq!(count, 0);
4193 assert!(
4194 !events_db.exists(),
4195 "a read-only runtime must not mint the events database"
4196 );
4197
4198 {
4205 let lane_backend = StorageBackend::sqlite_for_test(&events_db).expect("writable lane");
4206 lane_backend
4207 .events_for_namespace("local")
4208 .expect("lane store")
4209 .append_event(test_event("local"))
4210 .await
4211 .expect("seed lane row");
4212 }
4213 khive_storage::test_support::freeze_snapshot_sidecars(&events_db);
4214 let store = ro.events(&token).expect("events store with lane present");
4215 let count = store
4216 .count_events(EventFilter::default())
4217 .await
4218 .expect("merged count");
4219 assert_eq!(count, 1, "the pre-existing lane row must merge into reads");
4220 }
4221
4222 #[tokio::test]
4225 async fn direct_mode_appends_and_reads_without_a_daemon() {
4226 let dir = tempfile::tempdir().expect("tempdir");
4227 let _registry_guard = TestRegistryGuard::new(dir.path());
4228 let db = dir.path().join("events.db");
4229 let backend = direct_backend_for(&db).expect("direct backend");
4230 let store = backend.events_for_namespace("local").expect("store");
4231 store
4232 .append_event(test_event("local"))
4233 .await
4234 .expect("direct append");
4235 let count = store
4236 .count_events(EventFilter::default())
4237 .await
4238 .expect("count");
4239 assert_eq!(count, 1);
4240 }
4241}