use std::path::{Path, PathBuf};
#[cfg(unix)]
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::Arc;
#[cfg(unix)]
use std::time::Duration;
use async_trait::async_trait;
use khive_db::{StorageBackend, WalCeilingPolicy};
use khive_storage::event::{EventPageQuery, EventPageWindow, IdempotentEventBatchResult};
use khive_storage::{
BatchWriteSummary, Event, EventFilter, EventStore, Page, PageRequest, StorageError,
StorageResult,
};
use serde::{Deserialize, Serialize};
#[cfg(unix)]
use tokio::net::{UnixListener, UnixStream};
use uuid::Uuid;
#[cfg(unix)]
use crate::daemon::{read_frame, write_frame};
pub const EVENTS_PROTOCOL_VERSION: u32 = 4;
#[cfg(unix)]
pub const DEFAULT_APPEND_QUEUE_BATCHES: usize = 4096;
#[cfg(unix)]
pub const DEFAULT_APPEND_QUEUE_BYTES: usize = 32 * 1024 * 1024;
#[cfg(unix)]
const REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
#[cfg(unix)]
const FORWARDER_BACKOFF: Duration = Duration::from_secs(2);
#[cfg(unix)]
const MAX_EVENTS_CONNECTIONS: usize = 128;
#[cfg(unix)]
const CONN_IO_TIMEOUT: Duration = Duration::from_secs(60);
#[cfg(unix)]
const MAX_INFLIGHT_REQUEST_BYTES: usize = 64 * 1024 * 1024;
#[cfg(unix)]
const MAX_CACHED_NAMESPACE_STORES: usize = 1024;
const EVENTS_SYMLINK_HOP_BOUND: u32 = 40;
pub const MAX_QUERY_EVENTS_PAGE_ROWS: u32 = 4096;
pub fn events_db_path_beside(main_db: &Path) -> PathBuf {
let resolved = std::fs::canonicalize(main_db)
.unwrap_or_else(|_| resolve_dangling_final_component(main_db));
let mut name = resolved
.file_name()
.map(std::ffi::OsStr::to_os_string)
.unwrap_or_else(|| std::ffi::OsString::from("khive.db"));
name.push(".events.db");
let path = match resolved.parent().filter(|dir| !dir.as_os_str().is_empty()) {
Some(dir) => std::fs::canonicalize(dir)
.unwrap_or_else(|_| dir.to_path_buf())
.join(&name),
None => PathBuf::from(name),
};
absolutize(&path)
}
fn resolve_dangling_final_component(path: &Path) -> PathBuf {
let mut current = path.to_path_buf();
for _ in 0..EVENTS_SYMLINK_HOP_BOUND {
match std::fs::read_link(¤t) {
Ok(target) => {
current = if target.is_absolute() {
target
} else {
match current.parent().filter(|dir| !dir.as_os_str().is_empty()) {
Some(dir) => dir.join(target),
None => target,
}
};
}
Err(_) => break,
}
}
current
}
fn absolutize(path: &Path) -> PathBuf {
if path.is_absolute() {
return path.to_path_buf();
}
std::env::current_dir()
.map(|cwd| cwd.join(path))
.unwrap_or_else(|_| path.to_path_buf())
}
pub fn events_socket_path_beside(events_db: &Path) -> PathBuf {
events_db.with_extension("sock")
}
#[cfg(unix)]
#[path = "events_split_socket_path.rs"]
pub(crate) mod socket_path;
#[cfg(unix)]
pub use socket_path::validate_events_socket_path;
#[derive(Debug, Clone)]
pub struct EventsSplitConfig {
pub db_path: PathBuf,
pub socket_path: Option<PathBuf>,
}
#[cfg(unix)]
type ClientMap = std::collections::HashMap<PathBuf, Arc<EventsSplitClient>>;
type BackendMap = std::collections::HashMap<PathBuf, (bool, Arc<StorageBackend>)>;
#[cfg(unix)]
fn client_registry() -> &'static std::sync::Mutex<ClientMap> {
static REGISTRY: std::sync::OnceLock<std::sync::Mutex<ClientMap>> = std::sync::OnceLock::new();
REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
}
fn direct_backend_registry() -> &'static std::sync::Mutex<BackendMap> {
static REGISTRY: std::sync::OnceLock<std::sync::Mutex<BackendMap>> = std::sync::OnceLock::new();
REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
}
#[cfg(any(test, feature = "test-internals"))]
#[doc(hidden)]
pub struct TestRegistryGuard<'a> {
root: &'a Path,
canonical_root: PathBuf,
}
#[cfg(any(test, feature = "test-internals"))]
impl<'a> TestRegistryGuard<'a> {
pub fn new(root: &'a Path) -> Self {
Self {
root,
canonical_root: std::fs::canonicalize(root).expect("existing unique fixture root"),
}
}
fn remove_entries<T>(
&self,
registry: &std::sync::Mutex<std::collections::HashMap<PathBuf, T>>,
) -> Vec<T> {
let mut entries = registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let keys: Vec<_> = entries
.keys()
.filter(|path| path.starts_with(self.root) || path.starts_with(&self.canonical_root))
.cloned()
.collect();
keys.into_iter()
.filter_map(|key| entries.remove(&key))
.collect()
}
}
#[cfg(any(test, feature = "test-internals"))]
impl Drop for TestRegistryGuard<'_> {
fn drop(&mut self) {
let backends = self.remove_entries(direct_backend_registry());
#[cfg(unix)]
let clients = self.remove_entries(client_registry());
drop(backends);
#[cfg(unix)]
drop(clients);
}
}
#[cfg(unix)]
pub fn client_for(socket_path: &Path) -> crate::error::RuntimeResult<Arc<EventsSplitClient>> {
let mut registry = client_registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(existing) = registry.get(socket_path) {
return Ok(Arc::clone(existing));
}
let client = EventsSplitClient::new(socket_path.to_path_buf())?;
registry.insert(socket_path.to_path_buf(), Arc::clone(&client));
Ok(client)
}
pub fn direct_backend_for(db_path: &Path) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
direct_backend_with_max_readers(db_path, false, None)
}
pub fn direct_backend_read_only_for(
db_path: &Path,
) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
direct_backend_with_max_readers(db_path, true, None)
}
pub(crate) fn direct_backend_with_max_readers(
db_path: &Path,
read_only: bool,
max_readers: Option<usize>,
) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
let mut config = crate::RuntimeConfig {
db_path: Some(db_path.to_path_buf()),
..crate::RuntimeConfig::no_embeddings()
};
let wal_ceiling = config.resolve_wal_ceiling_policy(read_only)?;
let disk_guard = config.resolve_disk_guard_policy(read_only)?;
direct_backend_with_policies(
db_path,
read_only,
max_readers,
wal_ceiling,
disk_guard,
config.volume_lock_dir,
)
}
pub(crate) fn direct_backend_with_policies(
db_path: &Path,
read_only: bool,
max_readers: Option<usize>,
wal_ceiling: WalCeilingPolicy,
disk_guard: Option<khive_db::EffectiveDiskGuardConfig>,
volume_lock_dir: Option<PathBuf>,
) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
wal_ceiling.validate_static(true, true, read_only)?;
if !read_only {
disk_guard
.ok_or_else(|| {
crate::error::RuntimeError::Internal("missing events disk policy".into())
})?
.validate()?;
khive_db::require_volume_lock_dir(volume_lock_dir.clone())?;
}
let mut registry = direct_backend_registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let absolute_path = absolutize(db_path);
let db_path = absolute_path.as_path();
refuse_events_db_symlinks(db_path)
.map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
#[cfg(unix)]
if !read_only {
if let Some(parent) = db_path.parent().filter(|p| !p.as_os_str().is_empty()) {
let _ = std::fs::create_dir_all(parent);
}
}
#[cfg(unix)]
ensure_events_db_parent_trusted(db_path)
.map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
let key = std::fs::canonicalize(db_path)
.or_else(|error| {
if error.kind() != std::io::ErrorKind::NotFound {
return Err(error);
}
let parent = db_path.parent().ok_or(error)?;
let name = db_path.file_name().ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"events path has no file name",
)
})?;
std::fs::canonicalize(parent).map(|parent| parent.join(name))
})
.map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
if let Some((existing_read_only, existing)) = registry.get(&key) {
if *existing_read_only != read_only {
let existing_mode = if *existing_read_only {
"read-only"
} else {
"writable"
};
let requested_mode = if read_only { "read-only" } else { "writable" };
return Err(crate::error::RuntimeError::InvalidInput(format!(
"events database {} is already open {existing_mode} in this process; cannot \
open it {requested_mode}; read-only events access requires a separate frozen snapshot",
key.display()
)));
}
let numbers = |p: Option<khive_db::EffectiveDiskGuardConfig>| {
p.map(|p| (p.reserve_bytes, p.guard_deadline_ms))
};
if numbers(existing.pool().effective_disk_guard_config()) != numbers(disk_guard) {
return Err(crate::error::RuntimeError::Internal(
"events database is already open with a different disk reserve/deadline policy"
.into(),
));
}
if !read_only && existing.pool().config().volume_lock_dir != volume_lock_dir {
return Err(crate::error::RuntimeError::Internal(
"events database is already open with a different volume-lock directory".into(),
));
}
let existing_bytes = existing.pool().config().wal_ceiling.bytes;
if existing_bytes != wal_ceiling.bytes {
return Err(khive_db::SqliteError::InvalidConfig(format!(
"events database {} is already open with wal_ceiling_bytes={existing_bytes}; \
requested {}; drain and restart before changing the WAL ceiling",
key.display(),
wal_ceiling.bytes
))
.into());
}
return Ok(Arc::clone(existing));
}
let db_path = key.as_path();
#[cfg(unix)]
if !read_only {
use std::os::unix::fs::OpenOptionsExt;
let _ = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(db_path);
let _before_open = harden_events_db_sidecars(db_path)
.map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
}
let backend = Arc::new(if read_only {
StorageBackend::sqlite_read_only_with_max_readers_and_wal_ceiling(
db_path,
max_readers,
wal_ceiling,
)?
} else {
StorageBackend::sqlite_with_max_readers_and_policies(
db_path,
max_readers,
wal_ceiling,
disk_guard.ok_or_else(|| {
crate::error::RuntimeError::Internal("missing events disk policy".into())
})?,
khive_db::require_volume_lock_dir(volume_lock_dir)?,
)?
});
registry.insert(key, (read_only, Arc::clone(&backend)));
Ok(backend)
}
#[cfg(unix)]
pub fn forwarding_metrics(socket_path: &Path) -> Option<EventsForwardingMetrics> {
let registry = client_registry()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
registry.get(socket_path).map(|client| client.metrics())
}
#[derive(Debug, Serialize, Deserialize)]
#[serde(tag = "op", rename_all = "snake_case")]
pub enum EventsRequest {
AppendEvents {
protocol_version: u32,
namespace: String,
events: Vec<Event>,
},
AppendEventsIdempotent {
protocol_version: u32,
namespace: String,
events: Vec<Event>,
},
GetEvent {
protocol_version: u32,
namespace: String,
id: Uuid,
},
QueryEvents {
protocol_version: u32,
namespace: String,
filter: EventFilter,
page: PageRequest,
},
QueryEventPage {
protocol_version: u32,
namespace: String,
query: EventPageQuery,
},
CountEvents {
protocol_version: u32,
namespace: String,
filter: EventFilter,
},
}
impl EventsRequest {
fn protocol_version(&self) -> u32 {
match self {
Self::AppendEvents {
protocol_version, ..
}
| Self::AppendEventsIdempotent {
protocol_version, ..
}
| Self::GetEvent {
protocol_version, ..
}
| Self::QueryEvents {
protocol_version, ..
}
| Self::QueryEventPage {
protocol_version, ..
}
| Self::CountEvents {
protocol_version, ..
} => *protocol_version,
}
}
fn namespace(&self) -> &str {
match self {
Self::AppendEvents { namespace, .. }
| Self::AppendEventsIdempotent { namespace, .. }
| Self::GetEvent { namespace, .. }
| Self::QueryEvents { namespace, .. }
| Self::QueryEventPage { namespace, .. }
| Self::CountEvents { namespace, .. } => namespace,
}
}
}
#[derive(Debug, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum EventsResponse {
Appended {
summary: BatchWriteSummary,
},
Idempotent {
result: IdempotentEventBatchResult,
},
Event {
event: Option<Event>,
},
Pageful {
page: Page<Event>,
},
EventPageWindow {
window: EventPageWindow,
},
Count {
count: u64,
},
Error {
message: String,
retryable: bool,
#[serde(default)]
writer_task_failure: Option<WireWriterTaskFailure>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum WireWriterTaskFailure {
RequestFailed { request_state: WireWriterTaskState },
TaskTerminated { request_state: WireWriterTaskState },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WireWriterTaskState {
NotStarted,
TransactionRolledBack,
SideEffectsUnknown,
}
impl From<khive_storage::WriterTaskRequestState> for WireWriterTaskState {
fn from(state: khive_storage::WriterTaskRequestState) -> Self {
use khive_storage::WriterTaskRequestState as S;
match state {
S::NotStarted => Self::NotStarted,
S::TransactionRolledBack => Self::TransactionRolledBack,
S::SideEffectsUnknown => Self::SideEffectsUnknown,
}
}
}
impl From<WireWriterTaskState> for khive_storage::WriterTaskRequestState {
fn from(state: WireWriterTaskState) -> Self {
use khive_storage::WriterTaskRequestState as S;
match state {
WireWriterTaskState::NotStarted => S::NotStarted,
WireWriterTaskState::TransactionRolledBack => S::TransactionRolledBack,
WireWriterTaskState::SideEffectsUnknown => S::SideEffectsUnknown,
}
}
}
#[cfg(unix)]
pub struct EventsDaemonGuard {
_file: std::fs::File,
}
#[cfg(unix)]
pub fn try_acquire_events_daemon_guard(socket_path: &Path) -> Option<EventsDaemonGuard> {
match acquire_events_daemon_guard_outcome(socket_path) {
EventsDaemonGuardAcquisition::Held(guard) => Some(guard),
_ => None,
}
}
#[cfg(unix)]
enum EventsDaemonGuardAcquisition {
Held(EventsDaemonGuard),
Contended,
OpenFailed(std::io::Error),
HardeningRefused(std::io::Error),
}
#[cfg(unix)]
impl std::fmt::Debug for EventsDaemonGuardAcquisition {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Held(_) => f.write_str("Held"),
Self::Contended => f.write_str("Contended"),
Self::OpenFailed(error) => f.debug_tuple("OpenFailed").field(error).finish(),
Self::HardeningRefused(error) => {
f.debug_tuple("HardeningRefused").field(error).finish()
}
}
}
}
#[cfg(unix)]
fn acquire_events_daemon_guard_outcome(socket_path: &Path) -> EventsDaemonGuardAcquisition {
use EventsDaemonGuardAcquisition::{Contended, HardeningRefused, Held, OpenFailed};
let lock_path = socket_path.with_extension("lock");
if let Some(parent) = lock_path.parent() {
if let Err(error) = std::fs::create_dir_all(parent) {
return OpenFailed(error);
}
}
use std::os::unix::fs::OpenOptionsExt;
let file = match std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.write(true)
.mode(0o600)
.custom_flags(libc::O_NOFOLLOW)
.open(&lock_path)
{
Ok(file) => file,
Err(error) if error.raw_os_error() == Some(libc::ELOOP) => {
return HardeningRefused(error);
}
Err(error) => return OpenFailed(error),
};
if let Err(error) = file.set_permissions(
<std::fs::Permissions as std::os::unix::fs::PermissionsExt>::from_mode(0o600),
) {
return HardeningRefused(error);
}
use std::os::fd::AsRawFd;
let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
if rc == 0 {
Held(EventsDaemonGuard { _file: file })
} else {
let error = std::io::Error::last_os_error();
if error.raw_os_error() == Some(libc::EWOULDBLOCK)
|| error.raw_os_error() == Some(libc::EAGAIN)
{
Contended
} else {
HardeningRefused(error)
}
}
}
fn events_db_targets(db_path: &Path) -> [PathBuf; 3] {
["", "-wal", "-shm"].map(|suffix| {
let mut name = db_path.as_os_str().to_os_string();
name.push(suffix);
PathBuf::from(name)
})
}
fn refuse_events_db_symlinks(db_path: &Path) -> anyhow::Result<()> {
for path in events_db_targets(db_path) {
match std::fs::symlink_metadata(&path) {
Ok(meta) if meta.file_type().is_symlink() => anyhow::bail!(
"refusing to serve events: {} is a symlink; the events database and its \
sidecars must be regular files",
path.display()
),
Ok(_) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => anyhow::bail!(
"refusing to serve events: cannot inspect {}: {e}",
path.display()
),
}
}
Ok(())
}
#[cfg(unix)]
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct EventsFileIdentity {
device: u64,
inode: u64,
}
#[cfg(unix)]
impl EventsFileIdentity {
fn from_metadata(metadata: &std::fs::Metadata) -> Self {
use std::os::unix::fs::MetadataExt;
Self {
device: metadata.dev(),
inode: metadata.ino(),
}
}
}
#[cfg(unix)]
type EventsDbIdentities = [Option<EventsFileIdentity>; 3];
#[cfg(unix)]
fn ensure_events_db_parent_trusted(db_path: &Path) -> anyhow::Result<()> {
let parent = absolutize(db_path);
let parent = parent.parent().filter(|p| !p.as_os_str().is_empty());
match parent {
Some(dir) => crate::daemon::ensure_socket_dir_is_trusted(dir).map_err(|e| {
anyhow::anyhow!("refusing to serve events from an untrusted directory: {e}")
}),
None => Ok(()),
}
}
#[cfg(unix)]
fn ensure_events_db_owner_only(db_path: &Path) -> anyhow::Result<()> {
use std::os::unix::fs::OpenOptionsExt;
refuse_events_db_symlinks(db_path)?;
if let Some(parent) = db_path.parent() {
std::fs::create_dir_all(parent)?;
}
ensure_events_db_parent_trusted(db_path)?;
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(db_path)
{
Ok(_created) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
Err(e) => Err(anyhow::anyhow!(
"refusing to serve events: cannot create {} owner-only: {e}",
db_path.display()
)),
}
}
#[cfg(unix)]
fn harden_events_db_sidecars(db_path: &Path) -> anyhow::Result<EventsDbIdentities> {
use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
let mut identities = [None; 3];
for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
let file = match std::fs::OpenOptions::new()
.read(true)
.custom_flags(libc::O_NOFOLLOW)
.open(&path)
{
Ok(file) => file,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
Err(e) => {
anyhow::bail!(
"refusing to serve events: cannot open {} without following symlinks: \
{e}. The events database and its sidecars must be regular files.",
path.display()
)
}
};
let metadata = file.metadata()?;
if !metadata.file_type().is_file() {
anyhow::bail!(
"refusing to serve events: {} is not a regular file. The events \
database and its sidecars must be regular files.",
path.display()
);
}
file.set_permissions(std::fs::Permissions::from_mode(0o600))
.map_err(|e| {
anyhow::anyhow!(
"refusing to serve events: cannot chmod 0600 {}: {e}. The events \
database and its sidecars must be owner-only.",
path.display()
)
})?;
identities[index] = Some(EventsFileIdentity::from_metadata(&metadata));
}
Ok(identities)
}
#[cfg(unix)]
fn verify_events_db_owner_only_unopened(
db_path: &Path,
before: &EventsDbIdentities,
) -> anyhow::Result<()> {
use std::os::unix::fs::PermissionsExt;
for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
let metadata = match std::fs::symlink_metadata(&path) {
Ok(metadata) => metadata,
Err(e) if e.kind() == std::io::ErrorKind::NotFound && index != 0 => continue,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => anyhow::bail!(
"refusing to serve events: main database {} disappeared after SQLite open; \
refusing to serve an unlinked database",
path.display()
),
Err(e) => {
anyhow::bail!(
"refusing to serve events: cannot stat {}: {e}",
path.display()
)
}
};
if !metadata.file_type().is_file() {
anyhow::bail!(
"refusing to serve events: {} is not a regular file. The events database \
and its sidecars must be regular files.",
path.display()
);
}
let mode = metadata.permissions().mode() & 0o777;
if mode & 0o077 != 0 {
anyhow::bail!(
"refusing to serve events: {} is mode {mode:03o}, not owner-only. The events \
database and its sidecars must be owner-only.",
path.display()
);
}
if before[index]
.is_some_and(|identity| identity != EventsFileIdentity::from_metadata(&metadata))
{
anyhow::bail!(
"refusing to serve events: {} changed identity between pre-open hardening \
and post-open verification",
path.display()
);
}
}
Ok(())
}
#[cfg(unix)]
pub async fn supervise_events_daemon(db_path: PathBuf, socket_path: PathBuf) {
match standalone_daemon_policies(&db_path) {
Ok((wal_ceiling, disk_guard, volume_lock_dir)) => {
supervise_events_daemon_with_policies(
db_path,
socket_path,
wal_ceiling,
disk_guard,
volume_lock_dir,
)
.await;
}
Err(error) => {
tracing::warn!(%error, "invalid events daemon policy; events supervisor not started");
}
}
}
#[cfg(unix)]
pub async fn supervise_events_daemon_with_policies(
db_path: PathBuf,
socket_path: PathBuf,
wal_ceiling: WalCeilingPolicy,
disk_guard: khive_db::EffectiveDiskGuardConfig,
volume_lock_dir: PathBuf,
) {
if let Err(error) = disk_guard.validate() {
tracing::warn!(%error, "invalid disk policy; events supervisor not started");
return;
}
if let Err(error) = wal_ceiling.validate_static(true, true, false) {
tracing::warn!(error = %error, "invalid WAL ceiling; events supervisor not started");
return;
}
const PROBE_INTERVAL: Duration = Duration::from_secs(15);
let shutdown = crate::daemon::daemon_shutdown_token();
let mut child: Option<std::process::Child> = None;
let mut respawns: u64 = 0;
loop {
if let Some(c) = child.as_mut() {
match c.try_wait() {
Ok(Some(status)) => {
tracing::info!(%status, "events daemon child exited");
child = None;
}
Ok(None) => {}
Err(error) => {
tracing::warn!(error = %error, "cannot poll events daemon child; dropping handle");
child = None;
}
}
}
let reachable = match connect_verified(&socket_path).await {
Ok(_stream) => true,
Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
tracing::warn!(
socket = %socket_path.display(),
error = %error,
"events socket answered by a foreign uid; treating as unreachable"
);
false
}
Err(_) => false,
};
if !reachable && child.is_none() {
match std::env::current_exe() {
Ok(exe) => {
let spawned = events_daemon_command(
&exe,
&db_path,
&socket_path,
wal_ceiling,
disk_guard,
&volume_lock_dir,
)
.spawn();
match spawned {
Ok(spawned_child) => {
respawns += 1;
tracing::info!(
pid = spawned_child.id(),
respawns,
socket = %socket_path.display(),
"spawned events daemon"
);
child = Some(spawned_child);
}
Err(error) => {
tracing::warn!(error = %error, "failed to spawn events daemon");
}
}
}
Err(error) => {
tracing::warn!(error = %error, "cannot resolve current executable for events daemon spawn");
}
}
}
tokio::select! {
_ = shutdown.cancelled() => break,
_ = tokio::time::sleep(PROBE_INTERVAL) => {}
}
}
if let Some(mut c) = child.take() {
let _ = c.kill();
let _ = c.wait();
tracing::info!("events daemon child stopped with supervisor shutdown");
}
}
#[cfg(unix)]
fn events_daemon_command(
executable: &Path,
db_path: &Path,
socket_path: &Path,
wal_ceiling: WalCeilingPolicy,
disk_guard: khive_db::EffectiveDiskGuardConfig,
volume_lock_dir: &Path,
) -> std::process::Command {
let source = match wal_ceiling.source {
khive_db::WalCeilingSource::BackendField => "backend_field",
khive_db::WalCeilingSource::Environment => "environment",
khive_db::WalCeilingSource::Default => "default",
};
let mut command = std::process::Command::new(executable);
command
.arg("events-daemon")
.arg("--db")
.arg(db_path)
.arg("--socket")
.arg(socket_path)
.arg("--wal-ceiling-bytes")
.arg(wal_ceiling.bytes.to_string())
.arg("--wal-ceiling-source")
.arg(source)
.arg("--disk-reserve-bytes")
.arg(disk_guard.reserve_bytes.to_string())
.arg("--disk-guard-deadline-ms")
.arg(disk_guard.guard_deadline_ms.to_string())
.arg("--disk-reserve-source")
.arg(disk_guard.reserve_source.as_str())
.arg("--disk-deadline-source")
.arg(disk_guard.deadline_source.as_str())
.arg("--disk-legacy-environment-present")
.arg(disk_guard.legacy_environment_present.to_string())
.arg("--volume-lock-dir")
.arg(volume_lock_dir)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null());
command
}
#[cfg(unix)]
pub async fn run_events_daemon(db_path: &Path, socket_path: &Path) -> anyhow::Result<()> {
let (wal_ceiling, disk_guard, volume_lock_dir) = standalone_daemon_policies(db_path)?;
run_events_daemon_with_policies(
db_path,
socket_path,
wal_ceiling,
disk_guard,
volume_lock_dir,
)
.await
}
#[cfg(unix)]
pub async fn run_events_daemon_with_policies(
db_path: &Path,
socket_path: &Path,
wal_ceiling: WalCeilingPolicy,
disk_guard: khive_db::EffectiveDiskGuardConfig,
volume_lock_dir: PathBuf,
) -> anyhow::Result<()> {
disk_guard.validate()?;
wal_ceiling
.validate_static(true, true, false)
.map_err(crate::error::RuntimeError::from)?;
let db_path = &absolutize(db_path);
let socket_path = &absolutize(socket_path);
validate_events_socket_path(socket_path)?;
if let Some(parent) = socket_path.parent() {
std::fs::create_dir_all(parent)?;
crate::daemon::ensure_socket_dir_is_trusted(parent)?;
}
let _guard = match acquire_events_daemon_guard_outcome(socket_path) {
EventsDaemonGuardAcquisition::Held(guard) => guard,
refusal => {
tracing::info!(
socket = %socket_path.display(),
reason = ?refusal,
"events daemon lock unavailable; exiting"
);
return Ok(());
}
};
ensure_events_db_owner_only(db_path)?;
let before_open = harden_events_db_sidecars(db_path)?;
let backend = Arc::new(
StorageBackend::sqlite_with_max_readers_and_policies(
db_path,
None,
wal_ceiling,
disk_guard,
volume_lock_dir,
)
.map_err(crate::error::RuntimeError::from)?,
);
backend.events()?;
verify_events_db_owner_only_unopened(db_path, &before_open)?;
if socket_path.exists() {
std::fs::remove_file(socket_path)?;
}
let listener = UnixListener::bind(socket_path)?;
{
use std::os::unix::fs::PermissionsExt;
if let Err(e) =
std::fs::set_permissions(socket_path, std::fs::Permissions::from_mode(0o600))
{
drop(listener);
let _ = std::fs::remove_file(socket_path);
anyhow::bail!(
"refusing to serve events: cannot chmod 0600 {}: {e}. The events socket must \
be owner-only.",
socket_path.display()
);
}
}
let daemon_euid = unsafe { libc::geteuid() };
let stores: NamespaceStores = Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
let connections = Arc::new(tokio::sync::Semaphore::new(MAX_EVENTS_CONNECTIONS));
let frame_budget = Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES));
tracing::info!(
socket = %socket_path.display(),
db = %db_path.display(),
wal_ceiling_configured_bytes = wal_ceiling.bytes,
wal_ceiling_effective_bytes = wal_ceiling.effective_bytes(backend.is_read_only()),
wal_ceiling_source = ?wal_ceiling.source,
"events daemon listening"
);
loop {
let (stream, _addr) = match listener.accept().await {
Ok(pair) => pair,
Err(error) => {
tracing::warn!(error = %error, "events daemon accept failed");
continue;
}
};
match crate::daemon::peer_uid(&stream) {
Ok(uid) if crate::daemon::uid_is_permitted(uid, daemon_euid) => {}
Ok(uid) => {
tracing::warn!(peer_uid = uid, "events daemon rejected foreign-uid peer");
continue;
}
Err(error) => {
tracing::warn!(error = %error, "events daemon could not read peer credentials");
continue;
}
}
let permit = match Arc::clone(&connections).try_acquire_owned() {
Ok(permit) => permit,
Err(_) => {
tracing::warn!(
cap = MAX_EVENTS_CONNECTIONS,
"events daemon at connection cap; dropping new connection"
);
continue;
}
};
let backend = Arc::clone(&backend);
let stores = Arc::clone(&stores);
let frame_budget = Arc::clone(&frame_budget);
crate::daemon::spawn_named_tracked_task("events_connection", async move {
serve_events_conn(stream, backend, stores, frame_budget).await;
drop(permit);
});
}
}
#[cfg(unix)]
#[path = "events_split_policy_entries.rs"]
mod policy_entries;
#[cfg(unix)]
use policy_entries::standalone_daemon_policies;
#[cfg(unix)]
pub use policy_entries::{
run_events_daemon_with_wal_ceiling, supervise_events_daemon_with_wal_ceiling,
};
#[cfg(all(test, unix))]
#[path = "events_wal_policy_tests.rs"]
mod wal_policy_tests;
#[cfg(unix)]
async fn read_frame_budgeted(
stream: &mut UnixStream,
budget: &Arc<tokio::sync::Semaphore>,
) -> std::io::Result<(Vec<u8>, tokio::sync::OwnedSemaphorePermit)> {
use tokio::io::AsyncReadExt;
let mut len_buf = [0u8; 4];
stream.read_exact(&mut len_buf).await?;
let len = u32::from_be_bytes(len_buf) as usize;
if len > crate::daemon::MAX_FRAME_BYTES {
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"daemon frame of {len} bytes exceeds {} cap",
crate::daemon::MAX_FRAME_BYTES
),
));
}
let permit = Arc::clone(budget)
.acquire_many_owned(len as u32)
.await
.map_err(|_| std::io::Error::other("frame budget closed"))?;
let mut buf = vec![0u8; len];
stream.read_exact(&mut buf).await?;
Ok((buf, permit))
}
#[cfg(unix)]
type NamespaceStores =
Arc<std::sync::Mutex<std::collections::HashMap<String, Arc<dyn EventStore>>>>;
#[cfg(unix)]
fn namespace_store(
backend: &StorageBackend,
stores: &NamespaceStores,
namespace: &str,
) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
namespace_store_with_cap(backend, stores, namespace, MAX_CACHED_NAMESPACE_STORES)
}
#[cfg(unix)]
fn namespace_store_with_cap(
backend: &StorageBackend,
stores: &NamespaceStores,
namespace: &str,
cap: usize,
) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
let key = namespace.trim();
if let Some(store) = stores
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(key)
{
return Ok(Arc::clone(store));
}
let store = backend.events_for_namespace(namespace)?;
let mut map = stores
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if map.len() >= cap && !map.contains_key(key) {
if let Some(victim) = map.keys().next().cloned() {
map.remove(&victim);
}
}
map.insert(key.to_string(), Arc::clone(&store));
Ok(store)
}
#[cfg(unix)]
async fn serve_events_conn(
mut stream: UnixStream,
backend: Arc<StorageBackend>,
stores: NamespaceStores,
frame_budget: Arc<tokio::sync::Semaphore>,
) {
loop {
let (payload, budget_permit) = match tokio::time::timeout(
CONN_IO_TIMEOUT,
read_frame_budgeted(&mut stream, &frame_budget),
)
.await
{
Ok(Ok(pair)) => pair,
Ok(Err(_)) | Err(_) => return,
};
let response = match serde_json::from_slice::<EventsRequest>(&payload) {
Ok(request) => dispatch_events_request(request, &backend, &stores).await,
Err(error) => EventsResponse::Error {
message: format!("events daemon could not parse request frame: {error}"),
retryable: false,
writer_task_failure: None,
},
};
let bytes = match serde_json::to_vec(&response) {
Ok(bytes) => bytes,
Err(error) => {
tracing::error!(error = %error, "events daemon response serialization failed");
return;
}
};
let bytes = if bytes.len() > crate::daemon::MAX_FRAME_BYTES {
let refusal = EventsResponse::Error {
message: format!(
"response_frame_size_limit: {} response bytes exceed the {}-byte events IPC frame cap; request a narrower page",
bytes.len(),
crate::daemon::MAX_FRAME_BYTES
),
retryable: false,
writer_task_failure: None,
};
match serde_json::to_vec(&refusal) {
Ok(bytes) => bytes,
Err(error) => {
tracing::error!(error = %error, "events daemon frame-size refusal serialization failed");
return;
}
}
} else {
bytes
};
match tokio::time::timeout(CONN_IO_TIMEOUT, write_frame(&mut stream, &bytes)).await {
Ok(Ok(())) => {}
Ok(Err(_)) | Err(_) => return,
}
drop(budget_permit);
}
}
#[cfg(unix)]
async fn dispatch_events_request(
request: EventsRequest,
backend: &StorageBackend,
stores: &NamespaceStores,
) -> EventsResponse {
if request.protocol_version() != EVENTS_PROTOCOL_VERSION {
return EventsResponse::Error {
message: format!(
"events protocol version mismatch: daemon speaks {}, client sent {}",
EVENTS_PROTOCOL_VERSION,
request.protocol_version()
),
retryable: false,
writer_task_failure: None,
};
}
if let Err(error) = khive_types::Namespace::parse(request.namespace().trim()) {
return EventsResponse::Error {
message: format!("events request namespace rejected: {error}"),
retryable: false,
writer_task_failure: None,
};
}
let store = match namespace_store(backend, stores, request.namespace()) {
Ok(store) => store,
Err(error) => {
return EventsResponse::Error {
message: format!("events store unavailable: {error}"),
retryable: true,
writer_task_failure: None,
};
}
};
match request {
EventsRequest::AppendEvents { events, .. } => match store.append_events(events).await {
Ok(summary) => EventsResponse::Appended { summary },
Err(error) => storage_error_response(&error),
},
EventsRequest::AppendEventsIdempotent { events, .. } => {
match store.append_events_idempotent(events).await {
Ok(result) => EventsResponse::Idempotent { result },
Err(error) => storage_error_response(&error),
}
}
EventsRequest::GetEvent { id, .. } => match store.get_event(id).await {
Ok(event) => EventsResponse::Event { event },
Err(error) => storage_error_response(&error),
},
EventsRequest::QueryEvents { filter, page, .. } => {
if page.limit > MAX_QUERY_EVENTS_PAGE_ROWS {
return EventsResponse::Error {
message: format!(
"events query page limit {} exceeds the daemon cap of {} rows; \
request narrower pages",
page.limit, MAX_QUERY_EVENTS_PAGE_ROWS
),
retryable: false,
writer_task_failure: None,
};
}
match store.query_events(filter, page).await {
Ok(page) => EventsResponse::Pageful { page },
Err(error) => storage_error_response(&error),
}
}
EventsRequest::QueryEventPage { query, .. } => {
if let Err(error) = crate::event_page::validate_page_query(&query) {
return storage_error_response(&error);
}
match store.query_event_page(query).await {
Ok(window) => EventsResponse::EventPageWindow { window },
Err(error) => storage_error_response(&error),
}
}
EventsRequest::CountEvents { filter, .. } => match store.count_events(filter).await {
Ok(count) => EventsResponse::Count { count },
Err(error) => storage_error_response(&error),
},
}
}
#[cfg(unix)]
fn storage_error_response(error: &StorageError) -> EventsResponse {
let (message, writer_task_failure) = match error {
StorageError::WriterTaskRequestFailed {
request_state,
source,
} => (
source.to_string(),
Some(WireWriterTaskFailure::RequestFailed {
request_state: (*request_state).into(),
}),
),
StorageError::WriterTaskTerminated { request_state } => (
error.to_string(),
Some(WireWriterTaskFailure::TaskTerminated {
request_state: (*request_state).into(),
}),
),
_ => (error.to_string(), None),
};
EventsResponse::Error {
message,
retryable: error.is_retryable(),
writer_task_failure,
}
}
#[cfg(unix)]
async fn connect_verified(socket_path: &Path) -> std::io::Result<UnixStream> {
let stream = UnixStream::connect(socket_path).await?;
let own_euid = unsafe { libc::geteuid() } as u32;
let peer = crate::daemon::peer_uid(&stream)?;
if !crate::daemon::uid_is_permitted(peer, own_euid) {
return Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
format!(
"events socket {} answered by uid {peer}, not this process's uid {own_euid}; \
refusing to exchange event frames with an unowned daemon",
socket_path.display()
),
));
}
Ok(stream)
}
#[cfg(unix)]
#[derive(Debug, Clone, Copy, Serialize)]
pub struct EventsForwardingMetrics {
pub forwarded_batches: u64,
pub forwarded_events: u64,
pub dropped_batches: u64,
pub dropped_events: u64,
pub queued_bytes: usize,
}
#[cfg(unix)]
#[derive(Debug, Default)]
struct ForwardingCounters {
forwarded_batches: AtomicU64,
forwarded_events: AtomicU64,
dropped_batches: AtomicU64,
dropped_events: AtomicU64,
queued_bytes: AtomicUsize,
}
#[cfg(unix)]
#[derive(Debug)]
struct AppendByteReservation {
counters: Arc<ForwardingCounters>,
bytes: usize,
}
#[cfg(unix)]
impl Drop for AppendByteReservation {
fn drop(&mut self) {
self.counters
.queued_bytes
.fetch_sub(self.bytes, Ordering::AcqRel);
}
}
#[cfg(unix)]
#[derive(Debug)]
struct QueuedAppend {
request: EventsRequest,
count: u64,
_reservation: AppendByteReservation,
}
#[cfg(all(test, unix))]
impl QueuedAppend {
fn unmetered(namespace: &str, events: Vec<Event>, counters: Arc<ForwardingCounters>) -> Self {
let count = events.len() as u64;
Self {
request: EventsRequest::AppendEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: namespace.to_string(),
events,
},
count,
_reservation: AppendByteReservation { counters, bytes: 0 },
}
}
}
#[cfg(unix)]
#[derive(Default)]
struct BoundedFrameCounter {
bytes: usize,
exceeded: bool,
}
#[cfg(unix)]
impl std::io::Write for BoundedFrameCounter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let next = self.bytes.saturating_add(buf.len());
if next > crate::daemon::MAX_FRAME_BYTES {
self.exceeded = true;
return Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"events append request exceeds IPC frame cap",
));
}
self.bytes = next;
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[cfg(unix)]
fn reserve_append_bytes(
counters: &Arc<ForwardingCounters>,
bytes: usize,
budget: usize,
) -> Option<AppendByteReservation> {
let mut used = counters.queued_bytes.load(Ordering::Acquire);
loop {
let next = used.checked_add(bytes)?;
if next > budget {
return None;
}
match counters.queued_bytes.compare_exchange_weak(
used,
next,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
return Some(AppendByteReservation {
counters: Arc::clone(counters),
bytes,
});
}
Err(observed) => used = observed,
}
}
}
#[cfg(unix)]
pub struct EventsSplitClient {
socket_path: PathBuf,
append_tx: tokio::sync::mpsc::Sender<QueuedAppend>,
append_queue_byte_budget: usize,
counters: Arc<ForwardingCounters>,
outage_logged: Arc<AtomicBool>,
preflight_store: Arc<dyn EventStore>,
}
#[cfg(unix)]
impl std::fmt::Debug for EventsSplitClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("EventsSplitClient")
.field("socket_path", &self.socket_path)
.finish_non_exhaustive()
}
}
#[cfg(unix)]
impl EventsSplitClient {
pub fn new(socket_path: PathBuf) -> crate::error::RuntimeResult<Arc<Self>> {
Self::new_with_queue_depth(socket_path, DEFAULT_APPEND_QUEUE_BATCHES)
}
pub fn new_with_queue_depth(
socket_path: PathBuf,
queue_depth: usize,
) -> crate::error::RuntimeResult<Arc<Self>> {
Self::new_with_queue_depth_and_delivery_timeout(socket_path, queue_depth, REQUEST_TIMEOUT)
}
fn new_with_queue_depth_and_delivery_timeout(
socket_path: PathBuf,
queue_depth: usize,
delivery_timeout: Duration,
) -> crate::error::RuntimeResult<Arc<Self>> {
Self::new_with_limits_and_delivery_timeout(
socket_path,
queue_depth,
DEFAULT_APPEND_QUEUE_BYTES,
delivery_timeout,
)
}
fn new_with_limits_and_delivery_timeout(
socket_path: PathBuf,
queue_depth: usize,
byte_budget: usize,
delivery_timeout: Duration,
) -> crate::error::RuntimeResult<Arc<Self>> {
validate_events_socket_path(&socket_path)?;
let preflight_backend = StorageBackend::memory()?;
let preflight_store = preflight_backend.events()?;
let (append_tx, append_rx) = tokio::sync::mpsc::channel::<QueuedAppend>(queue_depth.max(1));
let counters = Arc::new(ForwardingCounters::default());
let outage_logged = Arc::new(AtomicBool::new(false));
let client = Arc::new(Self {
socket_path: socket_path.clone(),
append_tx,
append_queue_byte_budget: byte_budget,
counters: Arc::clone(&counters),
outage_logged: Arc::clone(&outage_logged),
preflight_store,
});
crate::daemon::spawn_named_tracked_task(
"events_forwarder",
run_forwarder(
socket_path,
append_rx,
counters,
outage_logged,
delivery_timeout,
crate::daemon::daemon_shutdown_token(),
),
);
Ok(client)
}
pub fn metrics(&self) -> EventsForwardingMetrics {
EventsForwardingMetrics {
forwarded_batches: self.counters.forwarded_batches.load(Ordering::Relaxed),
forwarded_events: self.counters.forwarded_events.load(Ordering::Relaxed),
dropped_batches: self.counters.dropped_batches.load(Ordering::Relaxed),
dropped_events: self.counters.dropped_events.load(Ordering::Relaxed),
queued_bytes: self.counters.queued_bytes.load(Ordering::Acquire),
}
}
fn enqueue(&self, namespace: &str, events: Vec<Event>) {
let count = events.len() as u64;
let request = EventsRequest::AppendEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: namespace.to_string(),
events,
};
let mut frame_size = BoundedFrameCounter::default();
if let Err(error) = serde_json::to_writer(&mut frame_size, &request) {
self.count_append_drop(
count,
if frame_size.exceeded {
"frame cap"
} else {
"serialization"
},
);
tracing::warn!(error = %error, "events append batch could not fit a wire frame");
return;
}
let Some(reservation) = reserve_append_bytes(
&self.counters,
frame_size.bytes,
self.append_queue_byte_budget,
) else {
self.count_append_drop(count, "queue byte budget");
return;
};
let batch = QueuedAppend {
request,
count,
_reservation: reservation,
};
match self.append_tx.try_send(batch) {
Ok(()) => {}
Err(_) => {
self.count_append_drop(count, "queue batch limit");
}
}
}
fn count_append_drop(&self, count: u64, reason: &'static str) {
self.counters
.dropped_batches
.fetch_add(1, Ordering::Relaxed);
self.counters
.dropped_events
.fetch_add(count, Ordering::Relaxed);
if !self.outage_logged.swap(true, Ordering::Relaxed) {
tracing::warn!(
dropped_events = count,
reason,
"events append queue rejected loss-tolerant batch"
);
}
}
async fn round_trip(&self, request: &EventsRequest) -> StorageResult<EventsResponse> {
let op = "events daemon round-trip";
let payload = serde_json::to_vec(request).map_err(|error| StorageError::Serialization {
capability: khive_storage::StorageCapability::Events,
message: format!("events request serialization failed: {error}"),
})?;
let attempt = async {
let mut stream = connect_verified(&self.socket_path)
.await
.map_err(|error| ("connect", error))?;
write_frame(&mut stream, &payload)
.await
.map_err(|error| ("write", error))?;
read_frame(&mut stream)
.await
.map_err(|error| ("read", error))
};
let bytes = match tokio::time::timeout(REQUEST_TIMEOUT, attempt).await {
Ok(Ok(bytes)) => bytes,
Ok(Err(("read", error))) => {
return Err(StorageError::Serialization {
capability: khive_storage::StorageCapability::Events,
message: format!(
"events daemon closed or broke the response after connection at {}: {error}",
self.socket_path.display()
),
});
}
Ok(Err((stage, error))) => {
return Err(StorageError::Pool {
operation: op.into(),
message: format!(
"events daemon {stage} failed at {}: {error}",
self.socket_path.display()
),
});
}
Err(_elapsed) => {
return Err(StorageError::Timeout {
operation: op.into(),
});
}
};
serde_json::from_slice::<EventsResponse>(&bytes).map_err(|error| {
StorageError::Serialization {
capability: khive_storage::StorageCapability::Events,
message: format!("events response deserialization failed: {error}"),
}
})
}
}
#[cfg(unix)]
fn drain_dropped_queue(
rx: &mut tokio::sync::mpsc::Receiver<QueuedAppend>,
counters: &ForwardingCounters,
) {
let mut dropped_batches = 0u64;
let mut dropped_events = 0u64;
while let Ok(batch) = rx.try_recv() {
dropped_batches += 1;
dropped_events += batch.count;
}
if dropped_batches > 0 {
counters
.dropped_batches
.fetch_add(dropped_batches, Ordering::Relaxed);
counters
.dropped_events
.fetch_add(dropped_events, Ordering::Relaxed);
tracing::warn!(
dropped_batches,
dropped_events,
"events forwarder shutting down; dropping queued loss-tolerant batches"
);
}
}
#[cfg(unix)]
async fn run_forwarder(
socket_path: PathBuf,
mut rx: tokio::sync::mpsc::Receiver<QueuedAppend>,
counters: Arc<ForwardingCounters>,
outage_logged: Arc<AtomicBool>,
delivery_timeout: Duration,
shutdown: tokio_util::sync::CancellationToken,
) {
let mut conn: Option<UnixStream> = None;
loop {
let batch = tokio::select! {
_ = shutdown.cancelled() => {
drain_dropped_queue(&mut rx, &counters);
break;
},
received = rx.recv() => match received {
Some(batch) => batch,
None => break,
},
};
let count = batch.count;
let payload = match serde_json::to_vec(&batch.request) {
Ok(payload) => payload,
Err(error) => {
counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
counters.dropped_events.fetch_add(count, Ordering::Relaxed);
tracing::error!(error = %error, "events forwarder serialization failed; batch dropped");
continue;
}
};
let delivered = tokio::select! {
_ = shutdown.cancelled() => {
counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
counters.dropped_events.fetch_add(count, Ordering::Relaxed);
tracing::warn!(
dropped_events = count,
"events forwarder shutting down mid-delivery; in-flight batch dropped"
);
drain_dropped_queue(&mut rx, &counters);
break;
},
outcome = tokio::time::timeout(
delivery_timeout,
deliver_batch(&socket_path, &mut conn, &payload),
) => match outcome {
Ok(delivered) => delivered,
Err(_elapsed) => {
conn = None;
false
}
},
};
if delivered {
counters.forwarded_batches.fetch_add(1, Ordering::Relaxed);
counters
.forwarded_events
.fetch_add(count, Ordering::Relaxed);
if outage_logged.swap(false, Ordering::Relaxed) {
tracing::info!("events daemon reachable again; forwarding resumed");
}
} else {
counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
counters.dropped_events.fetch_add(count, Ordering::Relaxed);
if !outage_logged.swap(true, Ordering::Relaxed) {
tracing::warn!(
socket = %socket_path.display(),
"events daemon unreachable; dropping loss-tolerant events until it returns"
);
}
tokio::select! {
_ = shutdown.cancelled() => {
drain_dropped_queue(&mut rx, &counters);
break;
},
_ = tokio::time::sleep(FORWARDER_BACKOFF) => {}
}
}
}
}
#[cfg(unix)]
async fn deliver_batch(socket_path: &Path, conn: &mut Option<UnixStream>, payload: &[u8]) -> bool {
for _attempt in 0..2u8 {
if conn.is_none() {
match connect_verified(socket_path).await {
Ok(stream) => *conn = Some(stream),
Err(_) => return false,
}
}
let stream = conn.as_mut().expect("connection populated above");
let ok = async {
write_frame(stream, payload).await?;
let bytes = read_frame(stream).await?;
std::io::Result::Ok(bytes)
}
.await;
match ok {
Ok(bytes) => {
return !matches!(
serde_json::from_slice::<EventsResponse>(&bytes),
Ok(EventsResponse::Error { .. }) | Err(_)
);
}
Err(_) => {
*conn = None;
}
}
}
false
}
#[cfg(unix)]
#[derive(Debug)]
pub struct ForwardingEventStore {
namespace: String,
client: Arc<EventsSplitClient>,
}
#[cfg(unix)]
impl ForwardingEventStore {
pub fn new(namespace: impl Into<String>, client: Arc<EventsSplitClient>) -> Self {
Self {
namespace: namespace.into(),
client,
}
}
fn unexpected(&self, op: &'static str, response: EventsResponse) -> StorageError {
match response {
EventsResponse::Error {
message,
retryable,
writer_task_failure,
} => {
if let Some(failure) = writer_task_failure {
match failure {
WireWriterTaskFailure::RequestFailed { request_state } => {
let source = if retryable {
StorageError::Pool {
operation: op.into(),
message,
}
} else {
StorageError::InvalidInput {
capability: khive_storage::StorageCapability::Events,
operation: op.into(),
message,
}
};
StorageError::WriterTaskRequestFailed {
request_state: request_state.into(),
source: Box::new(source),
}
}
WireWriterTaskFailure::TaskTerminated { request_state } => {
StorageError::WriterTaskTerminated {
request_state: request_state.into(),
}
}
}
} else if retryable {
StorageError::Pool {
operation: op.into(),
message,
}
} else {
StorageError::InvalidInput {
capability: khive_storage::StorageCapability::Events,
operation: op.into(),
message,
}
}
}
other => StorageError::Serialization {
capability: khive_storage::StorageCapability::Events,
message: format!("events daemon returned mismatched response for {op}: {other:?}"),
},
}
}
}
#[cfg(unix)]
#[async_trait]
impl EventStore for ForwardingEventStore {
async fn append_event(&self, event: Event) -> StorageResult<()> {
self.client.enqueue(&self.namespace, vec![event]);
Ok(())
}
async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
let attempted = events.len() as u64;
self.client.enqueue(&self.namespace, events);
Ok(BatchWriteSummary {
attempted,
affected: attempted,
..BatchWriteSummary::default()
})
}
async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
let request = EventsRequest::GetEvent {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: self.namespace.clone(),
id,
};
match self.client.round_trip(&request).await? {
EventsResponse::Event { event } => Ok(event),
other => Err(self.unexpected("get_event", other)),
}
}
async fn query_events(
&self,
filter: EventFilter,
page: PageRequest,
) -> StorageResult<Page<Event>> {
let request = EventsRequest::QueryEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: self.namespace.clone(),
filter,
page,
};
match self.client.round_trip(&request).await? {
EventsResponse::Pageful { page } => Ok(page),
other => Err(self.unexpected("query_events", other)),
}
}
async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
let request = EventsRequest::CountEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: self.namespace.clone(),
filter,
};
match self.client.round_trip(&request).await? {
EventsResponse::Count { count } => Ok(count),
other => Err(self.unexpected("count_events", other)),
}
}
async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
crate::event_page::validate_page_query(&query)?;
let request = EventsRequest::QueryEventPage {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: self.namespace.clone(),
query: query.clone(),
};
match self.client.round_trip(&request).await? {
EventsResponse::EventPageWindow { window } => {
crate::event_page::validate_window(&query, Some(&self.namespace), &window)?;
Ok(window)
}
error @ EventsResponse::Error { .. } => Err(self.unexpected("query_event_page", error)),
_ => Err(crate::event_page::page_error(
"events page response kind mismatch",
)),
}
}
fn preflight_event(&self, event: &Event) -> StorageResult<()> {
self.client.preflight_store.preflight_event(event)
}
async fn append_events_idempotent(
&self,
events: Vec<Event>,
) -> StorageResult<IdempotentEventBatchResult> {
let request = EventsRequest::AppendEventsIdempotent {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: self.namespace.clone(),
events,
};
match self.client.round_trip(&request).await? {
EventsResponse::Idempotent { result } => Ok(result),
other => Err(self.unexpected("append_events_idempotent", other)),
}
}
fn supports_idempotent_audit_batch(&self) -> bool {
true
}
}
pub struct SplitEventStore {
legacy: Arc<dyn EventStore>,
lane: Arc<dyn EventStore>,
}
impl std::fmt::Debug for SplitEventStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SplitEventStore").finish_non_exhaustive()
}
}
impl SplitEventStore {
pub const MAX_MERGED_WINDOW_ROWS: u64 = 100_000;
pub fn new(legacy: Arc<dyn EventStore>, lane: Arc<dyn EventStore>) -> Self {
Self { legacy, lane }
}
}
#[async_trait]
impl EventStore for SplitEventStore {
async fn append_event(&self, event: Event) -> StorageResult<()> {
self.legacy.append_event(event).await
}
async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
self.legacy.append_events(events).await
}
async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
if let Some(event) = self.legacy.get_event(id).await? {
return Ok(Some(event));
}
self.lane.get_event(id).await
}
async fn query_events(
&self,
filter: EventFilter,
page: PageRequest,
) -> StorageResult<Page<Event>> {
let window = page.offset.saturating_add(u64::from(page.limit));
if window > Self::MAX_MERGED_WINDOW_ROWS {
return Err(StorageError::InvalidInput {
capability: khive_storage::StorageCapability::Events,
operation: "query_events".into(),
message: format!(
"offset+limit ({window}) exceeds the merged event plane's window bound \
of {}; page deep windows with a `before` cursor at offset 0, or narrow \
the filter",
Self::MAX_MERGED_WINDOW_ROWS
),
});
}
let prefix = PageRequest {
offset: 0,
limit: window.min(u64::from(u32::MAX)) as u32,
};
let legacy = self
.legacy
.query_events(filter.clone(), prefix.clone())
.await?;
let lane = self.lane.query_events(filter, prefix).await?;
let total = match (legacy.total, lane.total) {
(Some(a), Some(b)) => Some(a + b),
_ => None,
};
let mut items = legacy.items;
items.extend(lane.items);
items.sort_by(|a, b| {
b.created_at
.cmp(&a.created_at)
.then_with(|| b.id.cmp(&a.id))
});
let items = items
.into_iter()
.skip(page.offset as usize)
.take(page.limit as usize)
.collect();
Ok(Page { items, total })
}
async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
let legacy = self.legacy.count_events(filter.clone()).await?;
let lane = self.lane.count_events(filter).await?;
Ok(legacy + lane)
}
async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
crate::event_page::split_page(self.legacy.as_ref(), self.lane.as_ref(), query).await
}
fn preflight_event(&self, event: &Event) -> StorageResult<()> {
self.lane.preflight_event(event)
}
async fn append_events_idempotent(
&self,
events: Vec<Event>,
) -> StorageResult<IdempotentEventBatchResult> {
if events.is_empty() {
return self.lane.append_events_idempotent(events).await;
}
let ids: Vec<Uuid> = events.iter().map(|event| event.id).collect();
let probe_limit = u32::try_from(ids.len()).map_err(|_| StorageError::InvalidInput {
capability: khive_storage::StorageCapability::Events,
operation: "append_events_idempotent".into(),
message: format!(
"batch of {} rows exceeds the legacy probe window",
ids.len()
),
})?;
let existing = self
.legacy
.query_events(
EventFilter {
ids,
..EventFilter::default()
},
PageRequest {
offset: 0,
limit: probe_limit,
},
)
.await?;
let legacy_resident: std::collections::HashSet<Uuid> =
existing.items.iter().map(|event| event.id).collect();
if legacy_resident.is_empty() {
return self.lane.append_events_idempotent(events).await;
}
let mut legacy_rows = Vec::new();
let mut lane_rows = Vec::new();
let mut routed_to_legacy = Vec::with_capacity(events.len());
for event in events {
if legacy_resident.contains(&event.id) {
routed_to_legacy.push(true);
legacy_rows.push(event);
} else {
routed_to_legacy.push(false);
lane_rows.push(event);
}
}
let legacy_expected = legacy_rows.len();
let lane_expected = lane_rows.len();
let legacy_result = self.legacy.append_events_idempotent(legacy_rows).await?;
let lane_result = if lane_expected == 0 {
IdempotentEventBatchResult { rows: Vec::new() }
} else {
self.lane.append_events_idempotent(lane_rows).await?
};
if legacy_result.rows.len() != legacy_expected || lane_result.rows.len() != lane_expected {
return Err(StorageError::Driver {
capability: khive_storage::StorageCapability::Events,
operation: "append_events_idempotent".into(),
source: format!(
"idempotent sub-batch result length mismatch: legacy {}/{}, lane {}/{}",
legacy_result.rows.len(),
legacy_expected,
lane_result.rows.len(),
lane_expected
)
.into(),
});
}
let mut legacy_iter = legacy_result.rows.into_iter();
let mut lane_iter = lane_result.rows.into_iter();
let rows = routed_to_legacy
.into_iter()
.map(|to_legacy| {
if to_legacy {
legacy_iter.next().expect("length checked above")
} else {
lane_iter.next().expect("length checked above")
}
})
.collect();
Ok(IdempotentEventBatchResult { rows })
}
fn supports_idempotent_audit_batch(&self) -> bool {
self.lane.supports_idempotent_audit_batch() && self.legacy.supports_idempotent_audit_batch()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn fixture_registry_cleanup_preserves_other_roots() {
let live_dir = tempfile::tempdir().unwrap();
let _live_guard = TestRegistryGuard::new(live_dir.path());
let live_path = live_dir.path().join("events.db");
let live_backend = direct_backend_for(&live_path).unwrap();
let finished_dir = tempfile::tempdir().unwrap();
let finished_backend = {
let _finished_guard = TestRegistryGuard::new(finished_dir.path());
let backend = direct_backend_for(&finished_dir.path().join("events.db")).unwrap();
let weak = Arc::downgrade(&backend);
drop(backend);
assert!(
weak.upgrade().is_some(),
"registry retains the live fixture"
);
weak
};
assert!(
finished_backend.upgrade().is_none(),
"FD_FIXTURE_REGISTRY_RELEASED"
);
let same_live_backend = direct_backend_for(&live_path).unwrap();
assert!(
Arc::ptr_eq(&live_backend, &same_live_backend),
"FD_FIXTURE_OTHER_ROOT_RETAINED"
);
}
#[cfg(unix)]
mod target_filter_tests {
include!("events_split_target_tests.rs");
}
include!("events_sidecar_path_tests.rs");
#[cfg(unix)]
#[test]
fn deep_symlink_chains_within_the_kernel_bound_derive_the_target_sidecar() {
let dir = tempfile::tempdir().unwrap();
let real = dir.path().join("real.db");
let mut prev = real.clone();
for i in 0..35 {
let link = dir.path().join(format!("hop{i}.db"));
std::os::unix::fs::symlink(&prev, &link).unwrap();
prev = link;
}
assert_eq!(
events_db_path_beside(&prev),
events_db_path_beside(&real),
"a 35-hop dangling chain must derive the target's sidecar"
);
assert_ne!(
events_db_path_beside(&dir.path().join("unrelated.db")),
events_db_path_beside(&real)
);
}
include!("events_sidecar_hardening_tests.rs");
#[cfg(unix)]
fn write_owner_only_test_file(path: &Path) {
use std::os::unix::fs::PermissionsExt;
std::fs::write(path, b"").unwrap();
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).unwrap();
}
#[cfg(unix)]
#[test]
fn unopened_check_refuses_same_mode_regular_file_replacement() {
for target_index in 0..3 {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("events.db");
let targets = events_db_targets(&db);
for path in &targets {
write_owner_only_test_file(path);
}
let before_open = harden_events_db_sidecars(&db).unwrap();
let replacement = dir.path().join("replacement");
write_owner_only_test_file(&replacement);
let replacement_id = EventsFileIdentity::from_metadata(
&std::fs::symlink_metadata(&replacement).unwrap(),
);
assert_ne!(before_open[target_index], Some(replacement_id));
let _connection = rusqlite::Connection::open(&db).unwrap();
verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
std::fs::rename(&replacement, &targets[target_index]).unwrap();
let error = verify_events_db_owner_only_unopened(&db, &before_open)
.unwrap_err()
.to_string();
assert!(error.contains("changed identity"), "{error}");
assert!(
error.contains(&targets[target_index].display().to_string()),
"{error}"
);
}
}
#[cfg(unix)]
#[test]
fn unopened_check_accepts_sidecars_missing_at_check_time() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("events.db");
let targets = events_db_targets(&db);
for path in &targets {
write_owner_only_test_file(path);
}
let before_open = harden_events_db_sidecars(&db).unwrap();
let _connection = rusqlite::Connection::open(&db).unwrap();
for path in targets.iter().skip(1) {
std::fs::remove_file(path).unwrap();
verify_events_db_owner_only_unopened(&db, &before_open)
.expect("SQLite sidecars may disappear without losing the main database");
}
}
#[cfg(unix)]
#[test]
fn unopened_check_refuses_main_database_removed_after_sqlite_open() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("events.db");
ensure_events_db_owner_only(&db).unwrap();
let before_open = harden_events_db_sidecars(&db).unwrap();
let backend = StorageBackend::sqlite_for_test(&db).unwrap();
backend.events().unwrap();
verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
std::fs::remove_file(&db).unwrap();
let error = verify_events_db_owner_only_unopened(&db, &before_open)
.unwrap_err()
.to_string();
assert!(
error.contains("main database") && error.contains("disappeared"),
"{error}"
);
}
#[cfg(unix)]
#[test]
fn unopened_check_accepts_fresh_sqlite_sidecars_and_unchanged_main() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("new.events.db");
assert!(!db.exists());
ensure_events_db_owner_only(&db).unwrap();
let before_open = harden_events_db_sidecars(&db).unwrap();
assert!(before_open[0].is_some());
assert_eq!(&before_open[1..], &[None, None]);
let backend = StorageBackend::sqlite_for_test(&db).unwrap();
backend.events().unwrap();
for path in events_db_targets(&db).iter().skip(1) {
assert!(path.exists(), "SQLite creates {}", path.display());
}
verify_events_db_owner_only_unopened(&db, &before_open)
.expect("new SQLite sidecars have no pre-open identity to contradict");
}
#[test]
fn socket_derives_beside_the_sidecar() {
let db = events_db_path_beside(Path::new("/data/khive.db"));
let socket = events_socket_path_beside(&db);
assert!(socket.ends_with("khive.db.events.sock"), "got {socket:?}");
}
#[test]
fn bare_relative_main_db_yields_an_absolute_sidecar() {
let path = events_db_path_beside(Path::new("khive.db"));
assert!(path.is_absolute(), "got relative {path:?}");
assert!(path.ends_with("khive.db.events.db"), "got {path:?}");
}
#[cfg(unix)]
#[tokio::test]
async fn over_cap_query_page_limit_is_refused_before_materialization() {
let backend = Arc::new(StorageBackend::memory().unwrap());
let stores: NamespaceStores =
Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
let request = |limit: u32| EventsRequest::QueryEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: "local".to_string(),
filter: EventFilter::default(),
page: PageRequest { offset: 0, limit },
};
let refused =
dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS + 1), &backend, &stores)
.await;
match refused {
EventsResponse::Error {
message,
retryable,
writer_task_failure,
} => {
assert!(
message.contains(&MAX_QUERY_EVENTS_PAGE_ROWS.to_string()),
"refusal must name the cap: {message}"
);
assert!(!retryable, "an over-cap page is not transient");
assert!(writer_task_failure.is_none());
}
other => panic!("over-cap query must be refused, got {other:?}"),
}
let at_cap =
dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS), &backend, &stores).await;
assert!(
matches!(at_cap, EventsResponse::Pageful { .. }),
"at-cap query must reach the store, got {at_cap:?}"
);
}
#[cfg(unix)]
#[test]
fn side_effects_unknown_state_crosses_the_wire_as_terminal_failure() {
let response = storage_error_response(&StorageError::WriterTaskTerminated {
request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
});
assert!(
matches!(
response,
EventsResponse::Error {
writer_task_failure: Some(WireWriterTaskFailure::TaskTerminated {
request_state: WireWriterTaskState::SideEffectsUnknown,
}),
..
}
),
"an unknown-commit-state termination must carry its state on the wire"
);
let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
assert!(matches!(
busy,
EventsResponse::Error {
writer_task_failure: None,
..
}
));
}
#[cfg(unix)]
#[test]
fn proven_rollback_crosses_the_wire_without_claiming_task_termination() {
let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
request_state: khive_storage::WriterTaskRequestState::TransactionRolledBack,
source: Box::new(StorageError::Pool {
operation: "writer_task_commit".into(),
message: "commit refused".into(),
}),
});
assert!(matches!(
response,
EventsResponse::Error {
writer_task_failure: Some(WireWriterTaskFailure::RequestFailed {
request_state: WireWriterTaskState::TransactionRolledBack,
}),
..
}
));
}
#[test]
fn error_frames_without_writer_failure_still_parse() {
let bytes = br#"{"kind":"error","message":"boom","retryable":true}"#;
let parsed: EventsResponse = serde_json::from_slice(bytes).expect("stateless frame parses");
assert!(matches!(
parsed,
EventsResponse::Error {
retryable: true,
writer_task_failure: None,
..
}
));
}
fn split_retry_event(verb: &str) -> Event {
Event::new(
"test",
verb,
khive_types::EventKind::RecallExecuted,
khive_types::SubstrateKind::Note,
"agent:test",
)
}
fn store_pair(dir: &Path) -> (Arc<dyn EventStore>, Arc<dyn EventStore>) {
let legacy = direct_backend_for(&dir.join("legacy.db"))
.expect("legacy backend")
.events_for_namespace("test")
.expect("legacy store");
let lane = direct_backend_for(&dir.join("lane.db"))
.expect("lane backend")
.events_for_namespace("test")
.expect("lane store");
(legacy, lane)
}
#[tokio::test]
async fn idempotent_retry_of_legacy_resident_rows_does_not_duplicate() {
use khive_storage::event::EventAppendDisposition;
let dir = tempfile::tempdir().unwrap();
let _registry_guard = TestRegistryGuard::new(dir.path());
let (legacy, lane) = store_pair(dir.path());
let resident = split_retry_event("recall");
legacy
.append_events_idempotent(vec![resident.clone()])
.await
.expect("pre-cutover landing");
let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
let fresh = split_retry_event("search");
let result = split
.append_events_idempotent(vec![resident.clone(), fresh.clone()])
.await
.expect("mixed retry batch");
assert_eq!(
result.rows,
vec![
EventAppendDisposition::AlreadyPresentIdentical,
EventAppendDisposition::Inserted,
],
"input order must be preserved across the two sub-batches"
);
assert_eq!(
lane.count_events(EventFilter::default()).await.unwrap(),
1,
"the legacy-resident row must not reach the lane"
);
assert_eq!(
split.count_events(EventFilter::default()).await.unwrap(),
2,
"merged count must not double-count the retried id"
);
let mut mutated = resident.clone();
mutated.verb = "other".to_string();
let conflict = split
.append_events_idempotent(vec![mutated])
.await
.expect("conflicting retry");
assert_eq!(
conflict.rows,
vec![EventAppendDisposition::IdentityConflict]
);
assert_eq!(split.count_events(EventFilter::default()).await.unwrap(), 2);
}
#[tokio::test]
async fn merged_offset_window_is_bounded() {
let dir = tempfile::tempdir().unwrap();
let _registry_guard = TestRegistryGuard::new(dir.path());
let (legacy, lane) = store_pair(dir.path());
let split = SplitEventStore::new(legacy, lane);
let err = split
.query_events(
EventFilter::default(),
PageRequest {
offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS,
limit: 1,
},
)
.await
.expect_err("a window past the bound must be refused");
assert!(
matches!(err, StorageError::InvalidInput { .. }),
"got {err:?}"
);
assert!(
err.to_string().contains("before"),
"the refusal must name the cursor remedy: {err}"
);
let at_bound = split
.query_events(
EventFilter::default(),
PageRequest {
offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS - 1,
limit: 1,
},
)
.await
.expect("a window at the bound is admitted");
assert!(at_bound.items.is_empty());
}
#[cfg(unix)]
#[tokio::test]
async fn stateless_non_retryable_error_maps_terminal_not_writer_terminated() {
let client = EventsSplitClient::new(std::path::PathBuf::from("/tmp/never-bound.sock"))
.expect("client builds");
let store = ForwardingEventStore::new("test", client);
let bytes = br#"{"kind":"error","message":"refused","retryable":false}"#;
let parsed: EventsResponse = serde_json::from_slice(bytes).expect("frame parses");
assert!(matches!(
store.unexpected("append", parsed),
StorageError::InvalidInput { .. }
));
let with_state = br#"{"kind":"error","message":"died","retryable":false,"writer_task_failure":{"kind":"task_terminated","request_state":"side_effects_unknown"}}"#;
let parsed: EventsResponse =
serde_json::from_slice(with_state).expect("stateful frame parses");
assert!(matches!(
store.unexpected("append", parsed),
StorageError::WriterTaskTerminated { .. }
));
}
#[cfg(unix)]
#[tokio::test]
async fn writer_task_states_cross_the_wire_verbatim() {
use khive_storage::WriterTaskRequestState as S;
let dir = tempfile::tempdir().unwrap();
let client =
EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
let store = ForwardingEventStore::new("test", client);
for state in [
S::NotStarted,
S::TransactionRolledBack,
S::SideEffectsUnknown,
] {
let response = storage_error_response(&StorageError::WriterTaskTerminated {
request_state: state,
});
let bytes = serde_json::to_vec(&response).unwrap();
let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
let err = store.unexpected("append_events_idempotent", parsed);
assert!(
matches!(
err,
StorageError::WriterTaskTerminated { request_state } if request_state == state
),
"state {state:?} did not survive the socket: got {err:?}"
);
}
let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
assert!(matches!(
busy,
EventsResponse::Error {
writer_task_failure: None,
..
}
));
}
#[cfg(unix)]
#[tokio::test]
async fn proven_rollback_round_trips_as_non_terminal_request_failure() {
use khive_storage::WriterTaskRequestState as S;
let dir = tempfile::tempdir().unwrap();
let client =
EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
let store = ForwardingEventStore::new("test", client);
let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
request_state: S::TransactionRolledBack,
source: Box::new(StorageError::Pool {
operation: "writer_task_commit".into(),
message: "commit refused".into(),
}),
});
let bytes = serde_json::to_vec(&response).unwrap();
let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
let err = store.unexpected("append_events_idempotent", parsed);
assert!(
matches!(
err,
StorageError::WriterTaskRequestFailed {
request_state: S::TransactionRolledBack,
..
}
),
"proven rollback must not reconstruct as a terminal writer: {err:?}"
);
}
#[cfg(unix)]
#[test]
fn direct_backend_hardens_preexisting_db_and_sidecars() {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().unwrap();
let _registry_guard = TestRegistryGuard::new(dir.path());
let db = dir.path().join("pre-existing.events.db");
let wal = dir.path().join("pre-existing.events.db-wal");
std::fs::write(&db, b"").unwrap();
std::fs::write(&wal, b"").unwrap();
for path in [&db, &wal] {
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o644)).unwrap();
}
assert_eq!(
std::fs::metadata(&db).unwrap().permissions().mode() & 0o777,
0o644
);
direct_backend_for(&db).expect("writable open succeeds");
for path in [&db, &wal] {
assert_eq!(
std::fs::metadata(path).unwrap().permissions().mode() & 0o777,
0o600,
"pre-existing {} must be tightened to owner-only",
path.display()
);
}
}
#[cfg(unix)]
#[test]
fn namespace_store_cache_is_bounded_and_trim_normalized() {
let backend = StorageBackend::memory().unwrap();
let stores: NamespaceStores =
Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
namespace_store_with_cap(&backend, &stores, "alpha", 1).expect("first store");
namespace_store_with_cap(&backend, &stores, "beta", 1).expect("second store evicts");
{
let map = stores.lock().unwrap();
assert_eq!(map.len(), 1, "cache must not grow past its cap");
assert!(
map.contains_key("beta"),
"the newest namespace must be admitted at the cap"
);
}
namespace_store_with_cap(&backend, &stores, "beta", 1).expect("cached beta");
namespace_store_with_cap(&backend, &stores, " beta ", 1).expect("trimmed spelling");
{
let map = stores.lock().unwrap();
assert_eq!(
map.len(),
1,
"spellings of one namespace must share one entry"
);
assert!(map.contains_key("beta"), "trimmed key is the cache key");
}
}
#[cfg(unix)]
#[tokio::test]
async fn frame_budget_admits_before_allocating_and_releases_after() {
let (mut client, mut server) = UnixStream::pair().expect("socketpair");
let budget = Arc::new(tokio::sync::Semaphore::new(1024));
let payload = vec![7u8; 100];
write_frame(&mut client, &payload).await.expect("write");
let (bytes, permit) = read_frame_budgeted(&mut server, &budget)
.await
.expect("budgeted read");
assert_eq!(bytes, payload);
assert_eq!(budget.available_permits(), 1024 - 100);
drop(permit);
assert_eq!(budget.available_permits(), 1024);
let mut oversized = Vec::from(((crate::daemon::MAX_FRAME_BYTES + 1) as u32).to_be_bytes());
oversized.extend_from_slice(&[0u8; 8]);
use tokio::io::AsyncWriteExt;
client.write_all(&oversized).await.expect("raw prefix");
let err = read_frame_budgeted(&mut server, &budget)
.await
.expect_err("oversized declaration must refuse");
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
assert_eq!(
budget.available_permits(),
1024,
"refusal must not consume budget"
);
}
#[cfg(unix)]
#[tokio::test]
async fn dispatch_rejects_invalid_wire_namespace() {
let backend = Arc::new(StorageBackend::memory().unwrap());
let stores: NamespaceStores =
Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
let huge = "n".repeat(64 * 1024);
let response = dispatch_events_request(
EventsRequest::CountEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: huge,
filter: Default::default(),
},
&backend,
&stores,
)
.await;
assert!(
matches!(
&response,
EventsResponse::Error {
retryable: false,
writer_task_failure: None,
..
}
),
"oversized namespace must be a typed refusal, got {response:?}"
);
assert_eq!(
stores.lock().unwrap().len(),
0,
"a rejected namespace must never enter the cache"
);
let ok = dispatch_events_request(
EventsRequest::CountEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: "local".to_string(),
filter: Default::default(),
},
&backend,
&stores,
)
.await;
assert!(
matches!(ok, EventsResponse::Count { .. }),
"valid namespace must dispatch, got {ok:?}"
);
}
include!("events_split_shutdown_tests.rs");
#[cfg(unix)]
#[tokio::test]
async fn forwarder_abandons_hung_delivery() {
let dir = tempfile::tempdir().unwrap();
let socket = dir.path().join("hung.sock");
let listener = tokio::net::UnixListener::bind(&socket).unwrap();
let _server = tokio::spawn(async move {
let mut held = Vec::new();
loop {
if let Ok((stream, _)) = listener.accept().await {
held.push(stream);
}
}
});
let client = EventsSplitClient::new_with_queue_depth_and_delivery_timeout(
socket.clone(),
4,
Duration::from_millis(100),
)
.expect("client builds");
let event = Event::new(
"test",
"noop",
khive_types::EventKind::Audit,
khive_types::SubstrateKind::Event,
"tester",
);
client.enqueue("test", vec![event]);
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
loop {
if client.metrics().dropped_batches >= 1 {
break;
}
assert!(
tokio::time::Instant::now() < deadline,
"forwarder never abandoned the hung delivery: {:?}",
client.metrics()
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
#[cfg(unix)]
#[test]
fn wire_retryability_follows_the_storage_classifier() {
let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
assert!(
matches!(
busy,
EventsResponse::Error {
retryable: true,
..
}
),
"transient writer contention must stay retryable across the socket"
);
let terminated = storage_error_response(&StorageError::WriterTaskTerminated {
request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
});
assert!(
matches!(
terminated,
EventsResponse::Error {
retryable: false,
..
}
),
"a terminated writer must not be reported transient"
);
}
use khive_storage::event::EventAppendDisposition;
use khive_types::{EventKind, SubstrateKind};
fn test_event(namespace: &str) -> Event {
Event::new(
namespace,
"test.verb",
EventKind::Audit,
SubstrateKind::Event,
"actor:test",
)
}
#[cfg(unix)]
async fn boot_daemon(dir: &tempfile::TempDir) -> (PathBuf, PathBuf) {
let db = dir.path().join("events.db");
let socket = dir.path().join("events.sock");
let (db_clone, socket_clone) = (db.clone(), socket.clone());
tokio::spawn(async move {
let _ = run_events_daemon(&db_clone, &socket_clone).await;
});
for _ in 0..100 {
if UnixStream::connect(&socket).await.is_ok() {
return (db, socket);
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
panic!("events daemon did not come up on {}", socket.display());
}
#[tokio::test]
async fn split_store_routes_plain_to_legacy_idempotent_to_lane_and_merges_reads() {
let dir = tempfile::tempdir().expect("tempdir");
let _registry_guard = TestRegistryGuard::new(dir.path());
let legacy_backend =
direct_backend_for(&dir.path().join("legacy.db")).expect("legacy backend");
let lane_backend = direct_backend_for(&dir.path().join("lane.db")).expect("lane backend");
let legacy = legacy_backend
.events_for_namespace("local")
.expect("legacy store");
let lane = lane_backend
.events_for_namespace("local")
.expect("lane store");
let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
let plain = test_event("local");
let plain_id = plain.id;
split.append_event(plain).await.expect("plain append");
let audit = test_event("local");
let audit_id = audit.id;
let result = split
.append_events_idempotent(vec![audit])
.await
.expect("idempotent append");
assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
assert!(legacy
.get_event(plain_id)
.await
.expect("legacy get")
.is_some());
assert!(lane.get_event(plain_id).await.expect("lane get").is_none());
assert!(legacy
.get_event(audit_id)
.await
.expect("legacy get")
.is_none());
assert!(lane.get_event(audit_id).await.expect("lane get").is_some());
assert!(split
.get_event(plain_id)
.await
.expect("split get")
.is_some());
assert!(split
.get_event(audit_id)
.await
.expect("split get")
.is_some());
assert_eq!(
split
.count_events(EventFilter::default())
.await
.expect("split count"),
2
);
let page = split
.query_events(
EventFilter::default(),
PageRequest {
offset: 0,
limit: 10,
},
)
.await
.expect("split query");
let ids: Vec<Uuid> = page.items.iter().map(|e| e.id).collect();
assert!(ids.contains(&plain_id) && ids.contains(&audit_id));
let mut seen = Vec::new();
for offset in 0..2 {
let page = split
.query_events(EventFilter::default(), PageRequest { offset, limit: 1 })
.await
.expect("windowed query");
assert_eq!(page.items.len(), 1);
seen.push(page.items[0].id);
}
seen.sort();
let mut expected = vec![plain_id, audit_id];
expected.sort();
assert_eq!(seen, expected);
}
#[tokio::test]
async fn raw_sql_consumers_of_the_legacy_events_table_still_see_plain_appends() {
use khive_storage::types::{SqlStatement, SqlValue};
let dir = tempfile::tempdir().expect("tempdir");
let _registry_guard = TestRegistryGuard::new(dir.path());
let legacy_backend =
direct_backend_for(&dir.path().join("legacy-guard.db")).expect("legacy backend");
let lane_backend =
direct_backend_for(&dir.path().join("lane-guard.db")).expect("lane backend");
let split = SplitEventStore::new(
legacy_backend
.events_for_namespace("local")
.expect("legacy store"),
lane_backend
.events_for_namespace("local")
.expect("lane store"),
);
let provenance = test_event("local");
let provenance_id = provenance.id;
split.append_event(provenance).await.expect("plain append");
let audit = test_event("local");
let audit_id = audit.id;
split
.append_events_idempotent(vec![audit])
.await
.expect("idempotent append");
let count_by_raw_sql = |backend: Arc<StorageBackend>, id: Uuid| async move {
let mut reader = backend.sql().reader().await.expect("sql reader");
let rows = reader
.query_all(SqlStatement {
sql: "SELECT actor FROM events WHERE id = ?1".to_string(),
params: vec![SqlValue::Text(id.to_string())],
label: None,
})
.await
.expect("raw events query");
rows.len()
};
assert_eq!(
count_by_raw_sql(Arc::clone(&legacy_backend), provenance_id).await,
1,
"plain appends must stay visible to raw-SQL consumers of the legacy events table"
);
assert_eq!(
count_by_raw_sql(Arc::clone(&legacy_backend), audit_id).await,
0,
"audit-lane rows must not land in the legacy events table"
);
assert_eq!(
count_by_raw_sql(Arc::clone(&lane_backend), audit_id).await,
1,
"the raw query must prove it can find rows where they actually live"
);
}
#[cfg(unix)]
#[tokio::test]
async fn idempotent_append_round_trips_through_the_daemon() {
let dir = tempfile::tempdir().expect("tempdir");
let (_db, socket) = boot_daemon(&dir).await;
let client = EventsSplitClient::new(socket).expect("client");
let store = ForwardingEventStore::new("local", client);
let event = test_event("local");
let id = event.id;
let result = store
.append_events_idempotent(vec![event.clone()])
.await
.expect("idempotent append over socket");
assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
let retry = store
.append_events_idempotent(vec![event])
.await
.expect("idempotent retry");
assert_eq!(
retry.rows,
vec![EventAppendDisposition::AlreadyPresentIdentical]
);
let fetched = store.get_event(id).await.expect("get over socket");
assert_eq!(fetched.map(|e| e.id), Some(id));
let count = store
.count_events(EventFilter::default())
.await
.expect("count over socket");
assert_eq!(count, 1);
}
#[cfg(unix)]
#[tokio::test]
async fn fire_and_forget_append_lands_in_the_daemon_store() {
let dir = tempfile::tempdir().expect("tempdir");
let (_db, socket) = boot_daemon(&dir).await;
let client = EventsSplitClient::new(socket).expect("client");
let store = ForwardingEventStore::new("local", Arc::clone(&client));
store
.append_event(test_event("local"))
.await
.expect("append_event is fire-and-forget");
let (count, metrics) = tokio::time::timeout(Duration::from_secs(30), async {
loop {
let count = store
.count_events(EventFilter::default())
.await
.expect("count over socket");
let metrics = client.metrics();
if (count == 1 && metrics.forwarded_events >= 1) || metrics.dropped_events > 0 {
break (count, metrics);
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
})
.await
.expect("forwarder must deliver or report a drop within 30 seconds");
assert_eq!(
metrics.dropped_events, 0,
"forwarder dropped an event before it reached the daemon store"
);
assert_eq!(count, 1, "forwarded event must land in the daemon store");
assert!(metrics.forwarded_events >= 1, "delivery must be counted");
}
#[cfg(unix)]
#[tokio::test]
async fn queue_overflow_drops_and_counts_but_never_errors() {
let dir = tempfile::tempdir().expect("tempdir");
let dead_socket = dir.path().join("nobody-home.sock");
let client = EventsSplitClient::new_with_queue_depth(dead_socket, 2).expect("client");
let store = ForwardingEventStore::new("local", Arc::clone(&client));
for _ in 0..20 {
store
.append_event(test_event("local"))
.await
.expect("append_event must not error under overflow");
}
let metrics = client.metrics();
assert!(
metrics.dropped_events > 0,
"overflow must register in the drop counter, got {metrics:?}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn dead_socket_reads_fail_typed_and_preflight_stays_local() {
let dir = tempfile::tempdir().expect("tempdir");
let dead_socket = dir.path().join("nobody-home.sock");
let client = EventsSplitClient::new(dead_socket).expect("client");
let store = ForwardingEventStore::new("local", client);
let error = store
.get_event(Uuid::new_v4())
.await
.expect_err("read against a dead socket must fail");
assert!(
matches!(error, StorageError::Pool { .. }),
"expected the typed unreachable error, got {error:?}"
);
assert!(store.supports_idempotent_audit_batch());
store
.preflight_event(&test_event("local"))
.expect("offline preflight validates a well-formed event");
}
#[cfg(unix)]
#[tokio::test]
async fn connected_events_peer_closing_without_response_is_not_unreachable() {
let dir = tempfile::tempdir().unwrap();
let socket = dir.path().join("closes-after-request.sock");
let listener = UnixListener::bind(&socket).unwrap();
let server = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await.unwrap();
read_frame(&mut stream)
.await
.expect("complete request arrives");
});
let client = EventsSplitClient::new(socket).unwrap();
let store = ForwardingEventStore::new("local", client);
let error = store.get_event(Uuid::new_v4()).await.unwrap_err();
assert!(
matches!(error, StorageError::Serialization { .. }),
"a post-connect response failure must not be Pool/unreachable: {error:?}"
);
server.await.unwrap();
}
#[cfg(unix)]
#[tokio::test]
async fn query_page_over_frame_cap_returns_typed_size_refusal() {
let dir = tempfile::tempdir().unwrap();
let socket = dir.path().join("large-query.sock");
let backend = Arc::new(StorageBackend::memory().unwrap());
let event_store = backend.events_for_namespace("local").unwrap();
let events: Vec<_> = (0..MAX_QUERY_EVENTS_PAGE_ROWS)
.map(|_| {
test_event("local").with_payload(serde_json::json!({"data": "x".repeat(3 * 1024)}))
})
.collect();
event_store
.append_events(events)
.await
.expect("seed large page");
let listener = UnixListener::bind(&socket).unwrap();
let stores: NamespaceStores =
Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
serve_events_conn(
stream,
backend,
stores,
Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES)),
)
.await;
});
let client = EventsSplitClient::new(socket).unwrap();
let store = ForwardingEventStore::new("local", client);
let error = store
.query_events(
EventFilter::default(),
PageRequest {
offset: 0,
limit: MAX_QUERY_EVENTS_PAGE_ROWS,
},
)
.await
.expect_err("the full page cannot fit one frame");
assert!(
matches!(error, StorageError::InvalidInput { .. }),
"expected non-retryable frame-size refusal, got {error:?}"
);
assert!(
error.to_string().contains("response_frame_size_limit"),
"{error}"
);
assert!(
error.to_string().contains("request a narrower page"),
"{error}"
);
server.abort();
}
#[cfg(unix)]
#[tokio::test]
async fn forwarding_queue_enforces_serialized_byte_and_frame_limits_with_drop_metrics() {
let dir = tempfile::tempdir().unwrap();
let socket = dir.path().join("stalled-forwarder.sock");
let listener = UnixListener::bind(&socket).unwrap();
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.unwrap();
let _hold_open = stream;
std::future::pending::<()>().await;
});
let event = test_event("local").with_payload(serde_json::json!({"data": "x".repeat(4096)}));
let first = vec![event.clone(), event.clone()];
let request_bytes = serde_json::to_vec(&EventsRequest::AppendEvents {
protocol_version: EVENTS_PROTOCOL_VERSION,
namespace: "local".into(),
events: first.clone(),
})
.unwrap()
.len();
let byte_budget = request_bytes + 64;
let client = EventsSplitClient::new_with_limits_and_delivery_timeout(
socket,
8,
byte_budget,
Duration::from_secs(30),
)
.unwrap();
let store = ForwardingEventStore::new("local", Arc::clone(&client));
store.append_events(first.clone()).await.unwrap();
assert_eq!(client.metrics().queued_bytes, request_bytes);
store.append_events(first).await.unwrap();
let metrics = client.metrics();
assert_eq!(metrics.dropped_batches, 1);
assert_eq!(metrics.dropped_events, 2);
assert_eq!(metrics.queued_bytes, request_bytes);
assert!(metrics.queued_bytes <= byte_budget);
let oversized = test_event("local")
.with_payload(serde_json::json!({"data": "x".repeat(crate::daemon::MAX_FRAME_BYTES)}));
store.append_event(oversized).await.unwrap();
let metrics = client.metrics();
assert_eq!(metrics.dropped_batches, 2);
assert_eq!(metrics.dropped_events, 3);
assert_eq!(metrics.queued_bytes, request_bytes);
server.abort();
}
#[cfg(unix)]
#[tokio::test]
async fn protocol_version_skew_is_a_typed_refusal() {
let dir = tempfile::tempdir().expect("tempdir");
let (_db, socket) = boot_daemon(&dir).await;
let mut stream = UnixStream::connect(&socket).await.expect("connect");
let request = EventsRequest::CountEvents {
protocol_version: EVENTS_PROTOCOL_VERSION + 1,
namespace: "local".into(),
filter: EventFilter::default(),
};
let payload = serde_json::to_vec(&request).expect("serialize");
write_frame(&mut stream, &payload).await.expect("write");
let bytes = read_frame(&mut stream).await.expect("read");
let response: EventsResponse = serde_json::from_slice(&bytes).expect("parse");
match response {
EventsResponse::Error {
message, retryable, ..
} => {
assert!(!retryable, "version skew is not retryable");
assert!(message.contains("protocol version"), "message: {message}");
}
other => panic!("expected a typed refusal, got {other:?}"),
}
}
#[cfg(unix)]
#[test]
fn events_daemon_guard_is_exclusive_then_reusable() {
let dir = tempfile::tempdir().expect("tempdir");
let socket = dir.path().join("events.sock");
let first = match acquire_events_daemon_guard_outcome(&socket) {
EventsDaemonGuardAcquisition::Held(guard) => guard,
other => panic!("first acquire must succeed: {other:?}"),
};
let second = acquire_events_daemon_guard_outcome(&socket);
assert!(
matches!(second, EventsDaemonGuardAcquisition::Contended),
"second acquire must report contention while the first guard is held: {second:?}"
);
drop(first);
const MAX_ATTEMPTS: usize = 100;
for attempt in 1..=MAX_ATTEMPTS {
match acquire_events_daemon_guard_outcome(&socket) {
EventsDaemonGuardAcquisition::Held(_) => return,
EventsDaemonGuardAcquisition::Contended if attempt < MAX_ATTEMPTS => {
std::thread::sleep(Duration::from_millis(10));
}
other => panic!(
"acquire must succeed after release within {MAX_ATTEMPTS} attempts; \
attempt {attempt}: {other:?}"
),
}
}
unreachable!("the final acquisition attempt returns or reports its refusal");
}
#[cfg(unix)]
#[test]
fn events_daemon_guard_open_failure_is_not_contention() {
let dir = tempfile::tempdir().expect("tempdir");
let clean = dir.path().join("clean.sock");
assert!(matches!(
acquire_events_daemon_guard_outcome(&clean),
EventsDaemonGuardAcquisition::Held(_)
));
assert!(
try_acquire_events_daemon_guard(&clean).is_some(),
"the public Option entrance must preserve a successful acquisition"
);
let socket = dir.path().join("blocked.sock");
let lock_path = socket.with_extension("lock");
std::fs::create_dir(&lock_path).expect("create directory at lock entry");
let outcome = acquire_events_daemon_guard_outcome(&socket);
assert!(
matches!(outcome, EventsDaemonGuardAcquisition::OpenFailed(_)),
"an actual lock-file open failure must not report contention: {outcome:?}"
);
assert!(
try_acquire_events_daemon_guard(&socket).is_none(),
"the public Option entrance must refuse an actual open failure"
);
assert!(
lock_path.is_dir(),
"the refused entry must remain a directory"
);
}
#[cfg(unix)]
#[tokio::test]
async fn read_only_runtime_never_creates_an_events_db() {
use crate::{KhiveRuntime, Namespace, RuntimeConfig};
let dir = tempfile::tempdir().expect("tempdir");
let _registry_guard = TestRegistryGuard::new(dir.path());
let main_db = dir.path().join("main.db");
drop(
KhiveRuntime::new_for_test(RuntimeConfig {
db_path: Some(main_db.clone()),
..RuntimeConfig::no_embeddings()
})
.expect("create main db"),
);
let events_db = events_db_path_beside(&main_db);
assert!(!events_db.exists(), "precondition: no events db yet");
khive_storage::test_support::freeze_snapshot_sidecars(&main_db);
let split_config = |db: PathBuf| RuntimeConfig {
db_path: Some(main_db.clone()),
events_split: Some(EventsSplitConfig {
db_path: db,
socket_path: None,
}),
..RuntimeConfig::no_embeddings()
};
let ro = KhiveRuntime::new_readonly_for_test(split_config(events_db.clone()))
.expect("read-only runtime");
let token = ro.authorize(Namespace::local()).expect("token");
let store = ro.events(&token).expect("events store");
let count = store
.count_events(EventFilter::default())
.await
.expect("count through legacy-only plane");
assert_eq!(count, 0);
assert!(
!events_db.exists(),
"a read-only runtime must not mint the events database"
);
{
let lane_backend = StorageBackend::sqlite_for_test(&events_db).expect("writable lane");
lane_backend
.events_for_namespace("local")
.expect("lane store")
.append_event(test_event("local"))
.await
.expect("seed lane row");
}
khive_storage::test_support::freeze_snapshot_sidecars(&events_db);
let store = ro.events(&token).expect("events store with lane present");
let count = store
.count_events(EventFilter::default())
.await
.expect("merged count");
assert_eq!(count, 1, "the pre-existing lane row must merge into reads");
}
#[tokio::test]
async fn direct_mode_appends_and_reads_without_a_daemon() {
let dir = tempfile::tempdir().expect("tempdir");
let _registry_guard = TestRegistryGuard::new(dir.path());
let db = dir.path().join("events.db");
let backend = direct_backend_for(&db).expect("direct backend");
let store = backend.events_for_namespace("local").expect("store");
store
.append_event(test_event("local"))
.await
.expect("direct append");
let count = store
.count_events(EventFilter::default())
.await
.expect("count");
assert_eq!(count, 1);
}
}