Skip to main content

khive_runtime/
events_split.rs

1//! Events-daemon split (ADR-170): the audit lane leaves the domain store.
2//!
3//! The ADR-133 idempotent audit batch is the measured bulk of event write
4//! volume, and in a single-store deployment its rows queue on the same SQLite
5//! writer lane as domain mutations. This module moves that lane into a
6//! dedicated events daemon that owns the events database, reachable over a
7//! Unix socket with the same length-prefixed framing and peer-uid admission the
8//! main daemon socket uses. Plain event appends stay on the domain store —
9//! the legacy `events` table has raw-SQL consumers (schedule provenance, kg
10//! projection guards, graph-query substrate unions) whose correctness
11//! depends on finding those rows there.
12//!
13//! Cooperating pieces:
14//!
15//! - [`run_events_daemon`] — the server loop the `events-daemon` subcommand
16//!   runs: binds the events socket, owns the only resident writer of the
17//!   events database, and serves append/read requests through the ordinary
18//!   `SqlEventStore`.
19//! - [`EventsSplitClient`] — one per domain process. Plain appends ride a
20//!   bounded in-memory queue drained by a background forwarder
21//!   (fire-and-forget; overflow or a dead daemon drops the batch, counts it,
22//!   and logs — the loss-tolerant durability class made concrete; unused by
23//!   the default routing until telemetry producers opt in). Idempotent
24//!   audit-batch appends and reads are synchronous framed round-trips with
25//!   bounded timeouts, because their callers are background flushers or query
26//!   paths, never the dispatch hot path.
27//! - [`ForwardingEventStore`] — the lane-side [`EventStore`] over the socket.
28//!   Preflight validation delegates to an in-memory `SqlEventStore`, so the
29//!   ADR-133 audit-batch seam keeps its pre-enqueue shape check without any
30//!   I/O on the dispatch path.
31//! - [`SplitEventStore`] — the per-namespace handle
32//!   [`crate::runtime::KhiveRuntime::events`] returns when the split is
33//!   configured: routes the idempotent lane to the events store, plain
34//!   appends to the legacy store, and merges reads across both.
35//!
36//! Domain availability never depends on events-daemon liveness: every failure
37//! path here degrades (drop + count + log, or a typed storage error for the
38//! synchronous lanes) instead of blocking the caller.
39
40use std::path::{Path, PathBuf};
41#[cfg(unix)]
42use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
43use std::sync::Arc;
44#[cfg(unix)]
45use std::time::Duration;
46
47use async_trait::async_trait;
48use khive_db::{StorageBackend, WalCeilingPolicy};
49use khive_storage::event::{EventPageQuery, EventPageWindow, IdempotentEventBatchResult};
50use khive_storage::{
51    BatchWriteSummary, Event, EventFilter, EventStore, Page, PageRequest, StorageError,
52    StorageResult,
53};
54use serde::{Deserialize, Serialize};
55
56#[cfg(unix)]
57use tokio::net::{UnixListener, UnixStream};
58use uuid::Uuid;
59
60#[cfg(unix)]
61use crate::daemon::{read_frame, write_frame};
62
63/// Bump whenever the request or response frame shape changes incompatibly.
64/// The server rejects frames whose version it does not speak, so a skewed
65/// client gets a typed refusal instead of a deserialization panic.
66///
67/// Version 2 was the first version any release shipped: version 1 carried a
68/// `side_effects_unknown` error field that was replaced by the
69/// `writer_task_state` carrier before this module reached any released ref,
70/// so v1 speakers existed only on unreleased development heads. The bump
71/// exists so that even such a process gets the version refusal above rather
72/// than having its retryable writer states silently mapped to terminal
73/// `InvalidInput`. Version 3 replaces that state-only carrier with a failure
74/// disposition so a proven rollback can cross the socket without falsely
75/// claiming the remote writer task terminated. Version 4 adds exact event-target
76/// filtering and the refusal kind. Older readers would ignore the new filter
77/// field, silently broadening a query, so both peers must speak version 4.
78pub const EVENTS_PROTOCOL_VERSION: u32 = 4;
79
80/// Default bound on the fire-and-forget append queue, in batches. The byte
81/// bound below also applies, so a large batch cannot multiply this depth
82/// into unbounded retained event memory.
83#[cfg(unix)]
84pub const DEFAULT_APPEND_QUEUE_BATCHES: usize = 4096;
85
86/// Maximum serialized append-request bytes retained by the fire-and-forget
87/// queue and its one in-flight delivery. A single request must also fit the
88/// daemon's per-frame cap; the queue budget covers several such requests.
89#[cfg(unix)]
90pub const DEFAULT_APPEND_QUEUE_BYTES: usize = 32 * 1024 * 1024;
91
92/// Timeout for one synchronous round-trip (idempotent appends, reads).
93/// Covers the WHOLE attempt — connect, peer verification, write, read —
94/// because each round-trip owns its connection (no shared lock to wait on
95/// outside the clock).
96#[cfg(unix)]
97const REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
98
99/// Reconnect backoff for the background forwarder after a failed connect.
100#[cfg(unix)]
101const FORWARDER_BACKOFF: Duration = Duration::from_secs(2);
102
103/// Cap on concurrently served daemon connections. Admission control, not a
104/// correctness bound: every permitted-uid process on the machine shares the
105/// daemon, and without a cap each accepted connection holds a task and a file
106/// descriptor for as long as the peer keeps the socket open — a stalled or
107/// hostile same-uid peer could exhaust both. Excess connects are dropped
108/// (peer sees EOF) and retried by the clients' own reconnect paths.
109#[cfg(unix)]
110const MAX_EVENTS_CONNECTIONS: usize = 128;
111
112/// Per-frame I/O deadline on a served connection: the longest the daemon
113/// waits for one request frame to finish arriving, or one response frame to
114/// finish sending. A peer that sends a partial frame and stalls is closed
115/// instead of holding its task and descriptor forever. Idle well-behaved
116/// clients are closed too — that is fine by construction: round-trips open a
117/// fresh connection per call, and the forwarder's delivery path retries once
118/// on a stale connection.
119#[cfg(unix)]
120const CONN_IO_TIMEOUT: Duration = Duration::from_secs(60);
121
122/// Aggregate cap on request-frame bytes buffered across all served
123/// connections at once. The per-frame cap bounds one buffer and the
124/// connection cap bounds the task count, but their product is the real
125/// allocation exposure (128 connections × 8 MiB declared frames ≈ 1 GiB).
126/// Body buffers are admitted against this shared byte budget before they are
127/// allocated; a connection whose frame cannot be admitted waits inside its
128/// own I/O deadline, so exhaustion degrades into per-connection timeouts,
129/// never into daemon memory growth. Must be at least `MAX_FRAME_BYTES`, or a
130/// maximum-size frame could never be admitted.
131#[cfg(unix)]
132const MAX_INFLIGHT_REQUEST_BYTES: usize = 64 * 1024 * 1024;
133
134/// Cap on the per-namespace store cache. Stores are cheap pool handles, but
135/// the namespace string is client-supplied, so an unbounded map is
136/// attacker-controlled memory growth. Beyond the cap, stores are built
137/// per-request instead of cached — slower, never wrong.
138#[cfg(unix)]
139const MAX_CACHED_NAMESPACE_STORES: usize = 1024;
140
141/// Symlink-chain hop ceiling for events-path resolution. Matches the bound
142/// the daemon's socket-path traversal guard enforces, which is itself the
143/// kernel's total-links ceiling on Linux (40): both resolvers must agree,
144/// or a chain one accepts derives a different identity than the other.
145const EVENTS_SYMLINK_HOP_BOUND: u32 = 40;
146
147/// Cap on `QueryEvents` page size, in rows. The wire `PageRequest.limit` is
148/// client-supplied `u32`, and the daemon materializes the full page as a
149/// `Vec<Event>` and serializes it into one response frame — so an unbounded
150/// limit is attacker-controlled memory and serialization work in a process
151/// that lives for months. Over-cap requests get a typed refusal naming the
152/// cap, never a silently clamped page: a caller that asked for more rows
153/// than it got would otherwise read the short page as the end of the data.
154/// The split client's merged read requests a prefix of `offset + limit`
155/// rows, so this cap also bounds the deep-offset window a socket client can
156/// demand in one request.
157///
158/// Public (and defined on every platform) because in-tree consumers that
159/// read events through a possibly-split store must size their page requests
160/// under it — a deep read is a `before`-cursor walk at `offset: 0` in pages
161/// of at most this many rows, never one wide page. Even a page under this row
162/// cap can exceed the IPC frame cap; the daemon then returns a non-retryable
163/// `response_frame_size_limit` refusal so the caller can narrow the page.
164pub const MAX_QUERY_EVENTS_PAGE_ROWS: u32 = 4096;
165
166/// Default events database file, beside the main database file.
167///
168/// The name is derived from the main database's full file name (`khive.db` →
169/// `khive.db.events.db`), never from its stem and never a fixed name in the
170/// parent directory: two independent databases that happen to share a
171/// directory (`a.db`, `b.db`) — or a stem (`a.db`, `a.sqlite`) — must each
172/// get their own event plane, not silently share one. The whole path is
173/// canonicalized when the database file exists — final-component symlink
174/// aliases of one database must derive the same sidecar and socket as the
175/// target spelling, because backend identity already treats those aliases as
176/// one database. When the file does not exist yet, the parent directory alone
177/// is canonicalized when it resolves, so path aliases of a fresh database
178/// (relative spellings, symlinked directories) still map to one sidecar
179/// instead of minting a distinct events database per spelling.
180pub fn events_db_path_beside(main_db: &Path) -> PathBuf {
181    // When the database file does not exist yet, the path can still be a
182    // symlink — a dangling alias whose target the first open will create.
183    // `canonicalize` refuses dangling links, so resolve final-component
184    // links by hand on that arm: alias-first and target-first cold starts
185    // must derive the same sidecar, or the first process to open each
186    // spelling writes audit records the other never reads.
187    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
203/// Follow final-component symlinks whose eventual target need not exist.
204/// The identity being recovered is the target's NAME — the file the first
205/// open through the alias will actually create — so a chain of links is
206/// walked (relative targets anchored at each link's parent) up to the same
207/// bound kernels use; a cycle or over-long chain stops at the last spelling,
208/// which `open()` will refuse anyway.
209fn resolve_dangling_final_component(path: &Path) -> PathBuf {
210    let mut current = path.to_path_buf();
211    // Same hop bound as the daemon's socket-path traversal guard (and the
212    // kernel's own total-links ceiling on Linux, 40): a chain the kernel
213    // would resolve must derive the target's sidecar, or a 33-hop alias
214    // opens fine yet splits its audit records from the target spelling.
215    for _ in 0..EVENTS_SYMLINK_HOP_BOUND {
216        match std::fs::read_link(&current) {
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
233/// Anchor a relative path to the current working directory. A bare relative
234/// spelling (`khive.db`) yields `Some("")` from `Path::parent`, and every
235/// consumer downstream — lock-file parenting, socket-directory trust
236/// validation, the daemon spawn contract — needs a real directory to stat,
237/// not an empty string. Absolute paths pass through untouched.
238fn 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
247/// Events daemon socket path, beside the events database it serves.
248///
249/// Derived from the events db path (not a process-global location, not a
250/// fixed name in the parent directory) so every events database gets its own
251/// daemon and socket: a shared socket would route events from any second
252/// database (another seat, a test tempdir) to whichever daemon happens to
253/// own it, persisting them beside the wrong main store. `khive.db.events.db`
254/// yields `khive.db.events.sock`; the daemon's advisory lock is the same
255/// path with a `.lock` extension.
256pub fn events_socket_path_beside(events_db: &Path) -> PathBuf {
257    events_db.with_extension("sock")
258}
259
260#[cfg(unix)]
261#[path = "events_split_socket_path.rs"]
262pub(crate) mod socket_path;
263#[cfg(unix)]
264pub use socket_path::validate_events_socket_path;
265
266/// How a runtime reaches event storage when the split is configured.
267#[derive(Debug, Clone)]
268pub struct EventsSplitConfig {
269    /// The events database file. The events daemon is its only writer in
270    /// daemon deployments; embedded mode writes it directly.
271    pub db_path: PathBuf,
272    /// `Some(socket)` = forward appends to the events daemon at this socket
273    /// (daemon deployments). `None` = embedded mode: open `db_path` directly
274    /// in-process (one-shot CLI, tests).
275    pub socket_path: Option<PathBuf>,
276}
277
278// ---------------------------------------------------------------------------
279// Process-global handles
280// ---------------------------------------------------------------------------
281//
282// One events-split client (or one direct backend) per process, whatever the
283// number of runtimes and clones — the same shape as the daemon module's other
284// process-global state. Keyed by path so tests exercising two distinct paths
285// in one process stay isolated.
286
287#[cfg(unix)]
288type ClientMap = std::collections::HashMap<PathBuf, Arc<EventsSplitClient>>;
289/// One admission per canonical path. The access mode belongs to the entry,
290/// so another mode cannot run raw-file preflight over a live SQLite pool.
291/// Strong entries retain the existing process lifetime, including derived stores.
292type BackendMap = std::collections::HashMap<PathBuf, (bool, Arc<StorageBackend>)>;
293
294#[cfg(unix)]
295fn client_registry() -> &'static std::sync::Mutex<ClientMap> {
296    static REGISTRY: std::sync::OnceLock<std::sync::Mutex<ClientMap>> = std::sync::OnceLock::new();
297    REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
298}
299
300fn direct_backend_registry() -> &'static std::sync::Mutex<BackendMap> {
301    static REGISTRY: std::sync::OnceLock<std::sync::Mutex<BackendMap>> = std::sync::OnceLock::new();
302    REGISTRY.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
303}
304
305/// Release only one test fixture's process-registry entries when its scope ends.
306///
307/// Borrow a unique temporary directory immediately after creating it, before
308/// declaring its backend/store/client handles. Those handles then drop before
309/// this guard, and the borrowed directory outlives registry cleanup. Production
310/// callers retain the normal process-lifetime registry ownership.
311#[cfg(any(test, feature = "test-internals"))]
312#[doc(hidden)]
313pub struct TestRegistryGuard<'a> {
314    root: &'a Path,
315    canonical_root: PathBuf,
316}
317
318#[cfg(any(test, feature = "test-internals"))]
319impl<'a> TestRegistryGuard<'a> {
320    /// Keep `root` alive until this fixture's registry references are released.
321    pub fn new(root: &'a Path) -> Self {
322        Self {
323            root,
324            canonical_root: std::fs::canonicalize(root).expect("existing unique fixture root"),
325        }
326    }
327
328    fn remove_entries<T>(
329        &self,
330        registry: &std::sync::Mutex<std::collections::HashMap<PathBuf, T>>,
331    ) -> Vec<T> {
332        let mut entries = registry
333            .lock()
334            .unwrap_or_else(std::sync::PoisonError::into_inner);
335        let keys: Vec<_> = entries
336            .keys()
337            .filter(|path| path.starts_with(self.root) || path.starts_with(&self.canonical_root))
338            .cloned()
339            .collect();
340        keys.into_iter()
341            .filter_map(|key| entries.remove(&key))
342            .collect()
343    }
344}
345
346#[cfg(any(test, feature = "test-internals"))]
347impl Drop for TestRegistryGuard<'_> {
348    fn drop(&mut self) {
349        // Destructors may close SQLite connections or senders. Run them only
350        // after releasing each registry lock; unrelated fixtures keep serving.
351        let backends = self.remove_entries(direct_backend_registry());
352        #[cfg(unix)]
353        let clients = self.remove_entries(client_registry());
354        drop(backends);
355        #[cfg(unix)]
356        drop(clients);
357    }
358}
359
360/// The process-wide client for `socket_path`, created (and its forwarder
361/// spawned) on first use. Requires a tokio runtime context on first call.
362#[cfg(unix)]
363pub fn client_for(socket_path: &Path) -> crate::error::RuntimeResult<Arc<EventsSplitClient>> {
364    let mut registry = client_registry()
365        .lock()
366        .unwrap_or_else(std::sync::PoisonError::into_inner);
367    if let Some(existing) = registry.get(socket_path) {
368        return Ok(Arc::clone(existing));
369    }
370    let client = EventsSplitClient::new(socket_path.to_path_buf())?;
371    registry.insert(socket_path.to_path_buf(), Arc::clone(&client));
372    Ok(client)
373}
374
375/// The process-wide direct (embedded-mode) backend for `db_path`, opened
376/// read-write on first use.
377pub fn direct_backend_for(db_path: &Path) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
378    direct_backend_with_max_readers(db_path, false, None)
379}
380
381/// The process-wide direct backend for `db_path`, opened READ-ONLY on first
382/// use. The file must already exist — this constructor never creates or
383/// schema-initializes an events database, which is what a read-only runtime's
384/// no-DB-creation contract requires of its event lane.
385pub fn direct_backend_read_only_for(
386    db_path: &Path,
387) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
388    direct_backend_with_max_readers(db_path, true, None)
389}
390
391/// Standalone opener: resolve its environment once before opening the event lane.
392pub(crate) fn direct_backend_with_max_readers(
393    db_path: &Path,
394    read_only: bool,
395    max_readers: Option<usize>,
396) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
397    let mut config = crate::RuntimeConfig {
398        db_path: Some(db_path.to_path_buf()),
399        ..crate::RuntimeConfig::no_embeddings()
400    };
401    let wal_ceiling = config.resolve_wal_ceiling_policy(read_only)?;
402    let disk_guard = config.resolve_disk_guard_policy(read_only)?;
403    direct_backend_with_policies(
404        db_path,
405        read_only,
406        max_readers,
407        wal_ceiling,
408        disk_guard,
409        config.volume_lock_dir,
410    )
411}
412
413/// Open the direct event lane with the policies already resolved for its main
414/// backend. A writable open also needs the disk policy and the volume lock
415/// directory; a read-only open uses neither.
416pub(crate) fn direct_backend_with_policies(
417    db_path: &Path,
418    read_only: bool,
419    max_readers: Option<usize>,
420    wal_ceiling: WalCeilingPolicy,
421    disk_guard: Option<khive_db::EffectiveDiskGuardConfig>,
422    volume_lock_dir: Option<PathBuf>,
423) -> crate::error::RuntimeResult<Arc<StorageBackend>> {
424    wal_ceiling.validate_static(true, true, read_only)?;
425    if !read_only {
426        disk_guard
427            .ok_or_else(|| {
428                crate::error::RuntimeError::Internal("missing events disk policy".into())
429            })?
430            .validate()?;
431        khive_db::require_volume_lock_dir(volume_lock_dir.clone())?;
432    }
433    let mut registry = direct_backend_registry()
434        .lock()
435        .unwrap_or_else(std::sync::PoisonError::into_inner);
436    let absolute_path = absolutize(db_path);
437    let db_path = absolute_path.as_path();
438    refuse_events_db_symlinks(db_path)
439        .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
440    #[cfg(unix)]
441    if !read_only {
442        if let Some(parent) = db_path.parent().filter(|p| !p.as_os_str().is_empty()) {
443            let _ = std::fs::create_dir_all(parent);
444        }
445    }
446    // Directory trust comes after creation (a fresh parent in a trusted
447    // ancestor must be statable) and before any open: SQLite's open is
448    // path-based, so the only defense against a component swapped between
449    // validation and open is that no untrusted local user can write the
450    // directories the path traverses.
451    #[cfg(unix)]
452    ensure_events_db_parent_trusted(db_path)
453        .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
454    let key = std::fs::canonicalize(db_path)
455        .or_else(|error| {
456            if error.kind() != std::io::ErrorKind::NotFound {
457                return Err(error);
458            }
459            let parent = db_path.parent().ok_or(error)?;
460            let name = db_path.file_name().ok_or_else(|| {
461                std::io::Error::new(
462                    std::io::ErrorKind::InvalidInput,
463                    "events path has no file name",
464                )
465            })?;
466            std::fs::canonicalize(parent).map(|parent| parent.join(name))
467        })
468        .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
469    if let Some((existing_read_only, existing)) = registry.get(&key) {
470        if *existing_read_only != read_only {
471            let existing_mode = if *existing_read_only {
472                "read-only"
473            } else {
474                "writable"
475            };
476            let requested_mode = if read_only { "read-only" } else { "writable" };
477            return Err(crate::error::RuntimeError::InvalidInput(format!(
478                "events database {} is already open {existing_mode} in this process; cannot \
479                 open it {requested_mode}; read-only events access requires a separate frozen snapshot",
480                key.display()
481            )));
482        }
483        let numbers = |p: Option<khive_db::EffectiveDiskGuardConfig>| {
484            p.map(|p| (p.reserve_bytes, p.guard_deadline_ms))
485        };
486        if numbers(existing.pool().effective_disk_guard_config()) != numbers(disk_guard) {
487            return Err(crate::error::RuntimeError::Internal(
488                "events database is already open with a different disk reserve/deadline policy"
489                    .into(),
490            ));
491        }
492        if !read_only && existing.pool().config().volume_lock_dir != volume_lock_dir {
493            return Err(crate::error::RuntimeError::Internal(
494                "events database is already open with a different volume-lock directory".into(),
495            ));
496        }
497        let existing_bytes = existing.pool().config().wal_ceiling.bytes;
498        if existing_bytes != wal_ceiling.bytes {
499            return Err(khive_db::SqliteError::InvalidConfig(format!(
500                "events database {} is already open with wal_ceiling_bytes={existing_bytes}; \
501                 requested {}; drain and restart before changing the WAL ceiling",
502                key.display(),
503                wal_ceiling.bytes
504            ))
505            .into());
506        }
507        return Ok(Arc::clone(existing));
508    }
509    let db_path = key.as_path();
510    // Embedded writable mode holds the events database to the daemon's own
511    // contract: owner-only from the first byte, never at the process umask,
512    // and — because event rows carry the same audit payloads either way — a
513    // pre-existing database or -wal/-shm sidecar is tightened to 0600 too,
514    // fail-closed. A create race just means SQLite finds the file present.
515    #[cfg(unix)]
516    if !read_only {
517        use std::os::unix::fs::OpenOptionsExt;
518        let _ = std::fs::OpenOptions::new()
519            .write(true)
520            .create_new(true)
521            .mode(0o600)
522            .open(db_path);
523        // Before the open, and only before it: see the precondition on
524        // `harden_events_db_sidecars`.
525        let _before_open = harden_events_db_sidecars(db_path)
526            .map_err(|e| crate::error::RuntimeError::Internal(e.to_string()))?;
527    }
528    let backend = Arc::new(if read_only {
529        StorageBackend::sqlite_read_only_with_max_readers_and_wal_ceiling(
530            db_path,
531            max_readers,
532            wal_ceiling,
533        )?
534    } else {
535        StorageBackend::sqlite_with_max_readers_and_policies(
536            db_path,
537            max_readers,
538            wal_ceiling,
539            disk_guard.ok_or_else(|| {
540                crate::error::RuntimeError::Internal("missing events disk policy".into())
541            })?,
542            khive_db::require_volume_lock_dir(volume_lock_dir)?,
543        )?
544    });
545    registry.insert(key, (read_only, Arc::clone(&backend)));
546    Ok(backend)
547}
548
549/// Forwarding metrics for the process-wide client at `socket_path`, if one
550/// exists. `None` means the split never initialized in this process.
551#[cfg(unix)]
552pub fn forwarding_metrics(socket_path: &Path) -> Option<EventsForwardingMetrics> {
553    let registry = client_registry()
554        .lock()
555        .unwrap_or_else(std::sync::PoisonError::into_inner);
556    registry.get(socket_path).map(|client| client.metrics())
557}
558
559// ---------------------------------------------------------------------------
560// Wire protocol
561// ---------------------------------------------------------------------------
562
563/// Request frame sent from a domain process to the events daemon.
564#[derive(Debug, Serialize, Deserialize)]
565#[serde(tag = "op", rename_all = "snake_case")]
566pub enum EventsRequest {
567    /// Fire-and-forget append lane. The daemon replies with a summary, but the
568    /// forwarder treats a failed reply as a counted drop, never an error to
569    /// the original caller.
570    AppendEvents {
571        protocol_version: u32,
572        namespace: String,
573        events: Vec<Event>,
574    },
575    /// ADR-133 audit-batch lane: real dispositions come back.
576    AppendEventsIdempotent {
577        protocol_version: u32,
578        namespace: String,
579        events: Vec<Event>,
580    },
581    GetEvent {
582        protocol_version: u32,
583        namespace: String,
584        id: Uuid,
585    },
586    QueryEvents {
587        protocol_version: u32,
588        namespace: String,
589        filter: EventFilter,
590        page: PageRequest,
591    },
592    QueryEventPage {
593        protocol_version: u32,
594        namespace: String,
595        query: EventPageQuery,
596    },
597    CountEvents {
598        protocol_version: u32,
599        namespace: String,
600        filter: EventFilter,
601    },
602}
603
604impl EventsRequest {
605    fn protocol_version(&self) -> u32 {
606        match self {
607            Self::AppendEvents {
608                protocol_version, ..
609            }
610            | Self::AppendEventsIdempotent {
611                protocol_version, ..
612            }
613            | Self::GetEvent {
614                protocol_version, ..
615            }
616            | Self::QueryEvents {
617                protocol_version, ..
618            }
619            | Self::QueryEventPage {
620                protocol_version, ..
621            }
622            | Self::CountEvents {
623                protocol_version, ..
624            } => *protocol_version,
625        }
626    }
627
628    fn namespace(&self) -> &str {
629        match self {
630            Self::AppendEvents { namespace, .. }
631            | Self::AppendEventsIdempotent { namespace, .. }
632            | Self::GetEvent { namespace, .. }
633            | Self::QueryEvents { namespace, .. }
634            | Self::QueryEventPage { namespace, .. }
635            | Self::CountEvents { namespace, .. } => namespace,
636        }
637    }
638}
639
640/// Response frame from the events daemon.
641#[derive(Debug, Serialize, Deserialize)]
642#[serde(tag = "kind", rename_all = "snake_case")]
643pub enum EventsResponse {
644    Appended {
645        summary: BatchWriteSummary,
646    },
647    Idempotent {
648        result: IdempotentEventBatchResult,
649    },
650    Event {
651        event: Option<Event>,
652    },
653    Pageful {
654        page: Page<Event>,
655    },
656    EventPageWindow {
657        window: EventPageWindow,
658    },
659    Count {
660        count: u64,
661    },
662    /// Typed refusal. `retryable` distinguishes transient daemon-side
663    /// conditions from contract errors (bad frame, version skew).
664    /// `writer_task_failure` carries both request finality and whether the
665    /// remote writer seam terminated. The retryable bit cannot express that
666    /// distinction: an ordinary operation or COMMIT error may be transient or
667    /// permanent while its transaction is independently proven rolled back.
668    /// `None` means the error is not from a writer request. Cross-version
669    /// frames never reach this mapping because `dispatch_events_request`
670    /// refuses a mismatched version before any write executes.
671    Error {
672        message: String,
673        retryable: bool,
674        #[serde(default)]
675        writer_task_failure: Option<WireWriterTaskFailure>,
676    },
677}
678
679/// Whether a writer-state error belongs to one failed request on a still-live
680/// seam or to a permanently retired writer task.
681#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
682#[serde(tag = "kind", rename_all = "snake_case")]
683pub enum WireWriterTaskFailure {
684    RequestFailed { request_state: WireWriterTaskState },
685    TaskTerminated { request_state: WireWriterTaskState },
686}
687
688/// Wire mirror of [`khive_storage::WriterTaskRequestState`]. A separate type
689/// because the storage enum is not serializable and the wire shape must stay
690/// under this module's protocol-version control, not the storage crate's.
691#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
692#[serde(rename_all = "snake_case")]
693pub enum WireWriterTaskState {
694    NotStarted,
695    TransactionRolledBack,
696    SideEffectsUnknown,
697}
698
699impl From<khive_storage::WriterTaskRequestState> for WireWriterTaskState {
700    fn from(state: khive_storage::WriterTaskRequestState) -> Self {
701        use khive_storage::WriterTaskRequestState as S;
702        match state {
703            S::NotStarted => Self::NotStarted,
704            S::TransactionRolledBack => Self::TransactionRolledBack,
705            S::SideEffectsUnknown => Self::SideEffectsUnknown,
706        }
707    }
708}
709
710impl From<WireWriterTaskState> for khive_storage::WriterTaskRequestState {
711    fn from(state: WireWriterTaskState) -> Self {
712        use khive_storage::WriterTaskRequestState as S;
713        match state {
714            WireWriterTaskState::NotStarted => S::NotStarted,
715            WireWriterTaskState::TransactionRolledBack => S::TransactionRolledBack,
716            WireWriterTaskState::SideEffectsUnknown => S::SideEffectsUnknown,
717        }
718    }
719}
720
721// ---------------------------------------------------------------------------
722// Server side
723// ---------------------------------------------------------------------------
724
725/// Advisory lock guaranteeing at most one events daemon per socket path.
726/// Held for the daemon's lifetime; a second daemon exits instead of stealing
727/// the socket path from the live one.
728#[cfg(unix)]
729pub struct EventsDaemonGuard {
730    _file: std::fs::File,
731}
732
733/// Try to become the events daemon for `socket_path`, returning `None` for
734/// contention or any refusal. Callers must validate the parent with
735/// `ensure_socket_dir_is_trusted` before lock-path operations in that directory.
736#[cfg(unix)]
737pub fn try_acquire_events_daemon_guard(socket_path: &Path) -> Option<EventsDaemonGuard> {
738    match acquire_events_daemon_guard_outcome(socket_path) {
739        EventsDaemonGuardAcquisition::Held(guard) => Some(guard),
740        _ => None,
741    }
742}
743
744/// Why a non-blocking daemon-lock acquisition did or did not succeed.
745/// Refusals retain the I/O error instead of treating every failure as contention.
746#[cfg(unix)]
747enum EventsDaemonGuardAcquisition {
748    Held(EventsDaemonGuard),
749    Contended,
750    OpenFailed(std::io::Error),
751    HardeningRefused(std::io::Error),
752}
753
754#[cfg(unix)]
755impl std::fmt::Debug for EventsDaemonGuardAcquisition {
756    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
757        match self {
758            Self::Held(_) => f.write_str("Held"),
759            Self::Contended => f.write_str("Contended"),
760            Self::OpenFailed(error) => f.debug_tuple("OpenFailed").field(error).finish(),
761            Self::HardeningRefused(error) => {
762                f.debug_tuple("HardeningRefused").field(error).finish()
763            }
764        }
765    }
766}
767
768/// Try to become the events daemon for `socket_path`. Every refusal means
769/// the caller must not serve; only `Contended` identifies a held lock.
770/// Callers that need the socket directory validated must run
771/// `ensure_socket_dir_is_trusted` on the parent BEFORE calling this, so no
772/// lock-path operation happens in an untrusted directory.
773#[cfg(unix)]
774fn acquire_events_daemon_guard_outcome(socket_path: &Path) -> EventsDaemonGuardAcquisition {
775    use EventsDaemonGuardAcquisition::{Contended, HardeningRefused, Held, OpenFailed};
776
777    let lock_path = socket_path.with_extension("lock");
778    if let Some(parent) = lock_path.parent() {
779        if let Err(error) = std::fs::create_dir_all(parent) {
780            return OpenFailed(error);
781        }
782    }
783    use std::os::unix::fs::OpenOptionsExt;
784    // `O_NOFOLLOW` pins the open to the final component: a symlink planted
785    // at the lock name is refused instead of redirecting the open (and the
786    // chmod below) to an attacker-selected target.
787    let file = match std::fs::OpenOptions::new()
788        .create(true)
789        .truncate(false)
790        .write(true)
791        .mode(0o600)
792        .custom_flags(libc::O_NOFOLLOW)
793        .open(&lock_path)
794    {
795        Ok(file) => file,
796        Err(error) if error.raw_os_error() == Some(libc::ELOOP) => {
797            return HardeningRefused(error);
798        }
799        Err(error) => return OpenFailed(error),
800    };
801    // `mode` applies only at creation; tighten a pre-existing lock file too.
802    // Descriptor-based (`fchmod` on the handle just opened), never a second
803    // path lookup — and fail closed: with the inode pinned, a failed chmod
804    // is abnormal, and serving behind a lock file another user can open is
805    // exactly what the hardening exists to refuse.
806    if let Err(error) = file.set_permissions(
807        <std::fs::Permissions as std::os::unix::fs::PermissionsExt>::from_mode(0o600),
808    ) {
809        return HardeningRefused(error);
810    }
811    use std::os::fd::AsRawFd;
812    // SAFETY: `fd` is a live descriptor owned by `file` for the duration of
813    // the call; `flock` reads nothing else.
814    let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
815    if rc == 0 {
816        Held(EventsDaemonGuard { _file: file })
817    } else {
818        let error = std::io::Error::last_os_error();
819        if error.raw_os_error() == Some(libc::EWOULDBLOCK)
820            || error.raw_os_error() == Some(libc::EAGAIN)
821        {
822            Contended
823        } else {
824            HardeningRefused(error)
825        }
826    }
827}
828
829fn events_db_targets(db_path: &Path) -> [PathBuf; 3] {
830    ["", "-wal", "-shm"].map(|suffix| {
831        let mut name = db_path.as_os_str().to_os_string();
832        name.push(suffix);
833        PathBuf::from(name)
834    })
835}
836
837/// Refuse to serve an events database whose path — or whose `-wal`/`-shm`
838/// sidecar path — is a pre-existing symlink. These paths are derived, never
839/// user-chosen (`events_db_path_beside` canonicalizes the main database
840/// spelling first), so a link here is a planted redirect, not an alias:
841/// permission hardening and SQLite would otherwise follow it and tighten or
842/// write event rows through to whatever file the link's author chose
843/// (CWE-59). Runs before any open on both the embedded and daemon arms; a
844/// link planted after admission is bounded by the daemon's trusted-directory
845/// contract on the socket parent.
846fn refuse_events_db_symlinks(db_path: &Path) -> anyhow::Result<()> {
847    for path in events_db_targets(db_path) {
848        match std::fs::symlink_metadata(&path) {
849            Ok(meta) if meta.file_type().is_symlink() => anyhow::bail!(
850                "refusing to serve events: {} is a symlink; the events database and its \
851                 sidecars must be regular files",
852                path.display()
853            ),
854            Ok(_) => {}
855            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
856            Err(e) => anyhow::bail!(
857                "refusing to serve events: cannot inspect {}: {e}",
858                path.display()
859            ),
860        }
861    }
862    Ok(())
863}
864
865#[cfg(unix)]
866#[derive(Clone, Copy, Debug, PartialEq, Eq)]
867struct EventsFileIdentity {
868    device: u64,
869    inode: u64,
870}
871
872#[cfg(unix)]
873impl EventsFileIdentity {
874    fn from_metadata(metadata: &std::fs::Metadata) -> Self {
875        use std::os::unix::fs::MetadataExt;
876        Self {
877            device: metadata.dev(),
878            inode: metadata.ino(),
879        }
880    }
881}
882
883/// Main/WAL/SHM identities, in `events_db_targets` order. None means the
884/// target was absent when pre-open hardening tried to open it.
885#[cfg(unix)]
886type EventsDbIdentities = [Option<EventsFileIdentity>; 3];
887
888/// Require the directory holding the events database to be trusted before
889/// anything opens a path inside it. The symlink pre-checks and the fd-pinned
890/// chmod close the races they can see, but SQLite's own open is path-based
891/// and cannot be inode-pinned from here — so the remaining defense is the
892/// same one the daemon socket uses: no untrusted local user may be able to
893/// write (or swap components of) the directory the path traverses. Delegates
894/// to the socket-path walk, which validates every ancestor the kernel will
895/// visit, with the same ownership and sticky-bit rules.
896#[cfg(unix)]
897fn ensure_events_db_parent_trusted(db_path: &Path) -> anyhow::Result<()> {
898    let parent = absolutize(db_path);
899    let parent = parent.parent().filter(|p| !p.as_os_str().is_empty());
900    match parent {
901        Some(dir) => crate::daemon::ensure_socket_dir_is_trusted(dir).map_err(|e| {
902            anyhow::anyhow!("refusing to serve events from an untrusted directory: {e}")
903        }),
904        None => Ok(()),
905    }
906}
907
908/// Create the events database file owner-only if it does not exist yet.
909/// The 0600 socket and peer-uid admission bound who can *talk to* the
910/// daemon; they bound nothing if the database file itself is readable by
911/// other local users, so the daemon creates it 0600 before SQLite ever
912/// opens it (SQLite would otherwise create it at the process umask).
913#[cfg(unix)]
914fn ensure_events_db_owner_only(db_path: &Path) -> anyhow::Result<()> {
915    use std::os::unix::fs::OpenOptionsExt;
916    refuse_events_db_symlinks(db_path)?;
917    if let Some(parent) = db_path.parent() {
918        std::fs::create_dir_all(parent)?;
919    }
920    // After creation so a fresh parent can be validated; a directory this
921    // process just created in a trusted ancestor passes by construction.
922    ensure_events_db_parent_trusted(db_path)?;
923    match std::fs::OpenOptions::new()
924        .write(true)
925        .create_new(true)
926        .mode(0o600)
927        .open(db_path)
928    {
929        // A zero-byte file is a valid empty SQLite database.
930        Ok(_created) => Ok(()),
931        Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
932        Err(e) => Err(anyhow::anyhow!(
933            "refusing to serve events: cannot create {} owner-only: {e}",
934            db_path.display()
935        )),
936    }
937}
938
939/// Tighten the events database and its SQLite sidecars to owner-only.
940/// Fail closed, same contract as the socket chmod: a daemon that cannot
941/// keep its database owner-only must not serve it.
942///
943/// PRECONDITION: no SQLite connection to `db_path` may be open in this
944/// process. This function opens and closes a descriptor on the database and
945/// on each sidecar, and POSIX advisory locks are per process and per inode,
946/// released by the close of ANY descriptor for the inode (fcntl(2)): called
947/// on a live database it silently drops the SHARED lock SQLite keeps in WAL
948/// mode and the `-shm` DMS lock, while the connection believes it still
949/// holds them. Both callers run it before their open; keep it that way. The
950/// post-open check that opens nothing is `verify_events_db_owner_only_unopened`.
951#[cfg(unix)]
952fn harden_events_db_sidecars(db_path: &Path) -> anyhow::Result<EventsDbIdentities> {
953    use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
954    let mut identities = [None; 3];
955    for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
956        // Pin the inode before touching it: `O_NOFOLLOW` makes the open
957        // itself refuse a symlink at the final component, and the chmod is
958        // then issued on the returned handle (fchmod), so no path re-lookup
959        // exists between validation and the permission change for a swapped
960        // entry to exploit. A path-based lstat-then-chmod pair here would
961        // re-open the exact race it is defending against.
962        let file = match std::fs::OpenOptions::new()
963            .read(true)
964            .custom_flags(libc::O_NOFOLLOW)
965            .open(&path)
966        {
967            Ok(file) => file,
968            Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
969            Err(e) => {
970                anyhow::bail!(
971                    "refusing to serve events: cannot open {} without following symlinks: \
972                     {e}. The events database and its sidecars must be regular files.",
973                    path.display()
974                )
975            }
976        };
977        // The open succeeds on a directory too; fstat the handle (no path
978        // re-lookup) so a directory at the path is refused, never chmod'ed.
979        let metadata = file.metadata()?;
980        if !metadata.file_type().is_file() {
981            anyhow::bail!(
982                "refusing to serve events: {} is not a regular file. The events \
983                 database and its sidecars must be regular files.",
984                path.display()
985            );
986        }
987        file.set_permissions(std::fs::Permissions::from_mode(0o600))
988            .map_err(|e| {
989                anyhow::anyhow!(
990                    "refusing to serve events: cannot chmod 0600 {}: {e}. The events \
991                     database and its sidecars must be owner-only.",
992                    path.display()
993                )
994            })?;
995        // The same descriptor supplied validation and fchmod; record its
996        // identity before closing it, never by re-looking up the path.
997        identities[index] = Some(EventsFileIdentity::from_metadata(&metadata));
998    }
999    Ok(identities)
1000}
1001
1002/// Check, after SQLite has opened the database, that it and its sidecars are
1003/// owner-only regular files WITHOUT opening any of them: `lstat` takes no
1004/// descriptor, so it cannot release the advisory locks the open connections
1005/// hold (see the precondition on `harden_events_db_sidecars`). Sidecars
1006/// SQLite creates inherit the database file's mode, so a failure here means
1007/// something else changed the file set; refuse to serve rather than tighten
1008/// through a path-based chmod and its lookup race. Present paths must retain
1009/// their pre-open device/inode when one was observed. A missing main database
1010/// is refused: serving an unlinked database would detach writes from its path.
1011/// Sidecars may legitimately be absent or newly created by SQLite, so an absent
1012/// pre- or post-observation does not assert identity continuity for that sidecar.
1013///
1014/// This is a two-observation check, not a pin of SQLite's own descriptors. It
1015/// cannot detect swaps restored between observations, inode reuse, or changes
1016/// after the final lstat. Even a legitimate sidecar delete/recreate is refused
1017/// if both observations see different identities; a later startup can retry
1018/// with a fresh snapshot. Trusted parent directories remain required. Unix only.
1019#[cfg(unix)]
1020fn verify_events_db_owner_only_unopened(
1021    db_path: &Path,
1022    before: &EventsDbIdentities,
1023) -> anyhow::Result<()> {
1024    use std::os::unix::fs::PermissionsExt;
1025    for (index, path) in events_db_targets(db_path).into_iter().enumerate() {
1026        let metadata = match std::fs::symlink_metadata(&path) {
1027            Ok(metadata) => metadata,
1028            Err(e) if e.kind() == std::io::ErrorKind::NotFound && index != 0 => continue,
1029            Err(e) if e.kind() == std::io::ErrorKind::NotFound => anyhow::bail!(
1030                "refusing to serve events: main database {} disappeared after SQLite open; \
1031                 refusing to serve an unlinked database",
1032                path.display()
1033            ),
1034            Err(e) => {
1035                anyhow::bail!(
1036                    "refusing to serve events: cannot stat {}: {e}",
1037                    path.display()
1038                )
1039            }
1040        };
1041        if !metadata.file_type().is_file() {
1042            anyhow::bail!(
1043                "refusing to serve events: {} is not a regular file. The events database \
1044                 and its sidecars must be regular files.",
1045                path.display()
1046            );
1047        }
1048        let mode = metadata.permissions().mode() & 0o777;
1049        if mode & 0o077 != 0 {
1050            anyhow::bail!(
1051                "refusing to serve events: {} is mode {mode:03o}, not owner-only. The events \
1052                 database and its sidecars must be owner-only.",
1053                path.display()
1054            );
1055        }
1056        if before[index]
1057            .is_some_and(|identity| identity != EventsFileIdentity::from_metadata(&metadata))
1058        {
1059            anyhow::bail!(
1060                "refusing to serve events: {} changed identity between pre-open hardening \
1061                 and post-open verification",
1062                path.display()
1063            );
1064        }
1065    }
1066    Ok(())
1067}
1068
1069/// Supervise the events daemon from the main daemon process: probe the
1070/// socket periodically and (re)spawn the daemon subcommand when unreachable.
1071///
1072/// The spawned command contract is fixed here once: the current executable
1073/// re-invoked as `events-daemon --db <db> --socket <socket>` with the resolved
1074/// WAL ceiling bytes and source — the subcommand the kernel binary registers
1075/// for [`run_events_daemon`]. The child holds the per-socket advisory lock,
1076/// so a probe/spawn race resolves to one survivor.
1077///
1078/// Lifecycle: the loop observes the process-wide daemon shutdown token, so
1079/// `drain()` never waits on it forever, and it retains the handle of the
1080/// child it spawned — reaping it with `try_wait` on every probe (no zombie
1081/// accumulation during a persistent startup failure) and never stacking a
1082/// second spawn on a still-live child. On shutdown, a child this supervisor
1083/// spawned is killed and reaped; a pre-existing events daemon it never
1084/// spawned is left alone.
1085///
1086/// The reachability probe uses the peer-verified connect: a socket answered
1087/// by a foreign-uid process is treated as UNREACHABLE (and logged loudly),
1088/// so a pre-bound spoof socket triggers a real-daemon spawn instead of being
1089/// reported healthy.
1090///
1091/// This standalone entry resolves its environment once before supervision.
1092/// Hosts with an opened main backend pass its resolved policies through
1093/// [`supervise_events_daemon_with_policies`].
1094#[cfg(unix)]
1095pub async fn supervise_events_daemon(db_path: PathBuf, socket_path: PathBuf) {
1096    match standalone_daemon_policies(&db_path) {
1097        Ok((wal_ceiling, disk_guard, volume_lock_dir)) => {
1098            supervise_events_daemon_with_policies(
1099                db_path,
1100                socket_path,
1101                wal_ceiling,
1102                disk_guard,
1103                volume_lock_dir,
1104            )
1105            .await;
1106        }
1107        Err(error) => {
1108            tracing::warn!(%error, "invalid events daemon policy; events supervisor not started");
1109        }
1110    }
1111}
1112
1113/// Supervise using the main pool's captured writer policies and lock directory.
1114#[cfg(unix)]
1115pub async fn supervise_events_daemon_with_policies(
1116    db_path: PathBuf,
1117    socket_path: PathBuf,
1118    wal_ceiling: WalCeilingPolicy,
1119    disk_guard: khive_db::EffectiveDiskGuardConfig,
1120    volume_lock_dir: PathBuf,
1121) {
1122    if let Err(error) = disk_guard.validate() {
1123        tracing::warn!(%error, "invalid disk policy; events supervisor not started");
1124        return;
1125    }
1126    if let Err(error) = wal_ceiling.validate_static(true, true, false) {
1127        tracing::warn!(error = %error, "invalid WAL ceiling; events supervisor not started");
1128        return;
1129    }
1130    const PROBE_INTERVAL: Duration = Duration::from_secs(15);
1131    let shutdown = crate::daemon::daemon_shutdown_token();
1132    let mut child: Option<std::process::Child> = None;
1133    let mut respawns: u64 = 0;
1134    loop {
1135        // Reap first: a child that exited (crashed, lost the advisory lock
1136        // race, refused an untrusted directory) must not linger as a zombie,
1137        // and clearing the slot is what re-arms the spawn below.
1138        if let Some(c) = child.as_mut() {
1139            match c.try_wait() {
1140                Ok(Some(status)) => {
1141                    tracing::info!(%status, "events daemon child exited");
1142                    child = None;
1143                }
1144                Ok(None) => {}
1145                Err(error) => {
1146                    tracing::warn!(error = %error, "cannot poll events daemon child; dropping handle");
1147                    child = None;
1148                }
1149            }
1150        }
1151
1152        let reachable = match connect_verified(&socket_path).await {
1153            Ok(_stream) => true,
1154            Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
1155                tracing::warn!(
1156                    socket = %socket_path.display(),
1157                    error = %error,
1158                    "events socket answered by a foreign uid; treating as unreachable"
1159                );
1160                false
1161            }
1162            Err(_) => false,
1163        };
1164        if !reachable && child.is_none() {
1165            match std::env::current_exe() {
1166                Ok(exe) => {
1167                    let spawned = events_daemon_command(
1168                        &exe,
1169                        &db_path,
1170                        &socket_path,
1171                        wal_ceiling,
1172                        disk_guard,
1173                        &volume_lock_dir,
1174                    )
1175                    .spawn();
1176                    match spawned {
1177                        Ok(spawned_child) => {
1178                            respawns += 1;
1179                            tracing::info!(
1180                                pid = spawned_child.id(),
1181                                respawns,
1182                                socket = %socket_path.display(),
1183                                "spawned events daemon"
1184                            );
1185                            child = Some(spawned_child);
1186                        }
1187                        Err(error) => {
1188                            tracing::warn!(error = %error, "failed to spawn events daemon");
1189                        }
1190                    }
1191                }
1192                Err(error) => {
1193                    tracing::warn!(error = %error, "cannot resolve current executable for events daemon spawn");
1194                }
1195            }
1196        }
1197        tokio::select! {
1198            _ = shutdown.cancelled() => break,
1199            _ = tokio::time::sleep(PROBE_INTERVAL) => {}
1200        }
1201    }
1202    if let Some(mut c) = child.take() {
1203        // Our child, our cleanup: the events daemon holds no volatile queue
1204        // state (SQLite is the durability), and the next daemon host's
1205        // supervisor respawns it.
1206        let _ = c.kill();
1207        let _ = c.wait();
1208        tracing::info!("events daemon child stopped with supervisor shutdown");
1209    }
1210}
1211
1212#[cfg(unix)]
1213fn events_daemon_command(
1214    executable: &Path,
1215    db_path: &Path,
1216    socket_path: &Path,
1217    wal_ceiling: WalCeilingPolicy,
1218    disk_guard: khive_db::EffectiveDiskGuardConfig,
1219    volume_lock_dir: &Path,
1220) -> std::process::Command {
1221    let source = match wal_ceiling.source {
1222        khive_db::WalCeilingSource::BackendField => "backend_field",
1223        khive_db::WalCeilingSource::Environment => "environment",
1224        khive_db::WalCeilingSource::Default => "default",
1225    };
1226    let mut command = std::process::Command::new(executable);
1227    command
1228        .arg("events-daemon")
1229        .arg("--db")
1230        .arg(db_path)
1231        .arg("--socket")
1232        .arg(socket_path)
1233        .arg("--wal-ceiling-bytes")
1234        .arg(wal_ceiling.bytes.to_string())
1235        .arg("--wal-ceiling-source")
1236        .arg(source)
1237        .arg("--disk-reserve-bytes")
1238        .arg(disk_guard.reserve_bytes.to_string())
1239        .arg("--disk-guard-deadline-ms")
1240        .arg(disk_guard.guard_deadline_ms.to_string())
1241        .arg("--disk-reserve-source")
1242        .arg(disk_guard.reserve_source.as_str())
1243        .arg("--disk-deadline-source")
1244        .arg(disk_guard.deadline_source.as_str())
1245        .arg("--disk-legacy-environment-present")
1246        .arg(disk_guard.legacy_environment_present.to_string())
1247        .arg("--volume-lock-dir")
1248        .arg(volume_lock_dir)
1249        .stdin(std::process::Stdio::null())
1250        .stdout(std::process::Stdio::null())
1251        .stderr(std::process::Stdio::null());
1252    command
1253}
1254
1255/// Serve the events daemon loop on `socket_path`, owning `db_path`.
1256///
1257/// Binds the socket (removing a stale path first), then accepts connections
1258/// for the process lifetime. Each connection is served sequentially:
1259/// same-uid admission, then a read-frame → dispatch → write-frame loop until
1260/// the peer disconnects. All storage goes through `SqlEventStore` on a
1261/// backend opened read-write against `db_path`; the events schema is ensured
1262/// once at boot.
1263///
1264/// Bind-path trust mirrors the main daemon socket: the socket directory must
1265/// pass the same ownership/mode/swap-resistance validation, and the bound
1266/// socket entry is chmod'd 0600 fail-closed. Pathname reachability is not
1267/// identity — clients additionally verify the peer uid on every connect —
1268/// but a hardened bind path is what keeps the *bind* itself out of another
1269/// user's hands.
1270///
1271/// This standalone entry resolves its environment once before opening the
1272/// daemon. Hosts with resolved policies use [`run_events_daemon_with_policies`].
1273#[cfg(unix)]
1274pub async fn run_events_daemon(db_path: &Path, socket_path: &Path) -> anyhow::Result<()> {
1275    let (wal_ceiling, disk_guard, volume_lock_dir) = standalone_daemon_policies(db_path)?;
1276    run_events_daemon_with_policies(
1277        db_path,
1278        socket_path,
1279        wal_ceiling,
1280        disk_guard,
1281        volume_lock_dir,
1282    )
1283    .await
1284}
1285
1286/// Serve using the exact policies inherited from the opened main backend.
1287#[cfg(unix)]
1288pub async fn run_events_daemon_with_policies(
1289    db_path: &Path,
1290    socket_path: &Path,
1291    wal_ceiling: WalCeilingPolicy,
1292    disk_guard: khive_db::EffectiveDiskGuardConfig,
1293    volume_lock_dir: PathBuf,
1294) -> anyhow::Result<()> {
1295    disk_guard.validate()?;
1296    wal_ceiling
1297        .validate_static(true, true, false)
1298        .map_err(crate::error::RuntimeError::from)?;
1299    // The subcommand's `--db`/`--socket` arrive from argv and may be
1300    // relative; anchor them before anything derives a parent from them.
1301    let db_path = &absolutize(db_path);
1302    let socket_path = &absolutize(socket_path);
1303    validate_events_socket_path(socket_path)?;
1304    // Directory trust comes FIRST: the lock guard below opens and chmods a
1305    // path in this directory, and validating only before the later bind
1306    // would let those operations run in a directory another user controls.
1307    if let Some(parent) = socket_path.parent() {
1308        std::fs::create_dir_all(parent)?;
1309        crate::daemon::ensure_socket_dir_is_trusted(parent)?;
1310    }
1311    let _guard = match acquire_events_daemon_guard_outcome(socket_path) {
1312        EventsDaemonGuardAcquisition::Held(guard) => guard,
1313        refusal => {
1314            tracing::info!(
1315                socket = %socket_path.display(),
1316                reason = ?refusal,
1317                "events daemon lock unavailable; exiting"
1318            );
1319            return Ok(());
1320        }
1321    };
1322    ensure_events_db_owner_only(db_path)?;
1323    // Tighten a pre-existing database and any `-wal`/`-shm` an earlier
1324    // process left behind BEFORE SQLite opens the database. The order is
1325    // load-bearing: hardening opens and closes its own descriptor on each of
1326    // these files, and POSIX advisory locks are per process and per inode,
1327    // released by the close of ANY descriptor for the inode (fcntl(2)). Run
1328    // after the open, it silently dropped the SHARED lock SQLite keeps on the
1329    // database in WAL mode and the shared lock on the `-shm` DMS byte, so the
1330    // next external connection to close (a backup tool, an inspection shell)
1331    // took itself for the last connection, checkpointed, and unlinked the
1332    // sidecars underneath this daemon, which kept writing to the unlinked
1333    // inodes. Sidecars SQLite creates from here on inherit the database
1334    // file's mode; the check after the open below opens nothing.
1335    let before_open = harden_events_db_sidecars(db_path)?;
1336    let backend = Arc::new(
1337        StorageBackend::sqlite_with_max_readers_and_policies(
1338            db_path,
1339            None,
1340            wal_ceiling,
1341            disk_guard,
1342            volume_lock_dir,
1343        )
1344        .map_err(crate::error::RuntimeError::from)?,
1345    );
1346    // Ensure the schema once, loudly, before accepting traffic.
1347    backend.events()?;
1348    verify_events_db_owner_only_unopened(db_path, &before_open)?;
1349
1350    if socket_path.exists() {
1351        std::fs::remove_file(socket_path)?;
1352    }
1353    let listener = UnixListener::bind(socket_path)?;
1354    {
1355        use std::os::unix::fs::PermissionsExt;
1356        // Fail closed, same contract as the main daemon socket: if the entry
1357        // cannot be made owner-only, drop the listener and remove it rather
1358        // than serve a world-reachable socket the design never covered.
1359        if let Err(e) =
1360            std::fs::set_permissions(socket_path, std::fs::Permissions::from_mode(0o600))
1361        {
1362            drop(listener);
1363            let _ = std::fs::remove_file(socket_path);
1364            anyhow::bail!(
1365                "refusing to serve events: cannot chmod 0600 {}: {e}. The events socket must \
1366                 be owner-only.",
1367                socket_path.display()
1368            );
1369        }
1370    }
1371    let daemon_euid = unsafe { libc::geteuid() };
1372    // Per-namespace store cache: `events_for_namespace` takes a writer-lane
1373    // checkout and re-runs the schema DDL on every call, so paying it once
1374    // per namespace instead of once per request keeps the writer lane for
1375    // actual writes. The namespace is client-supplied, so the cache is
1376    // bounded (`MAX_CACHED_NAMESPACE_STORES`) rather than trusted small.
1377    let stores: NamespaceStores = Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
1378    let connections = Arc::new(tokio::sync::Semaphore::new(MAX_EVENTS_CONNECTIONS));
1379    let frame_budget = Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES));
1380    tracing::info!(
1381        socket = %socket_path.display(),
1382        db = %db_path.display(),
1383        wal_ceiling_configured_bytes = wal_ceiling.bytes,
1384        wal_ceiling_effective_bytes = wal_ceiling.effective_bytes(backend.is_read_only()),
1385        wal_ceiling_source = ?wal_ceiling.source,
1386        "events daemon listening"
1387    );
1388
1389    loop {
1390        let (stream, _addr) = match listener.accept().await {
1391            Ok(pair) => pair,
1392            Err(error) => {
1393                tracing::warn!(error = %error, "events daemon accept failed");
1394                continue;
1395            }
1396        };
1397        match crate::daemon::peer_uid(&stream) {
1398            Ok(uid) if crate::daemon::uid_is_permitted(uid, daemon_euid) => {}
1399            Ok(uid) => {
1400                tracing::warn!(peer_uid = uid, "events daemon rejected foreign-uid peer");
1401                continue;
1402            }
1403            Err(error) => {
1404                tracing::warn!(error = %error, "events daemon could not read peer credentials");
1405                continue;
1406            }
1407        }
1408        let permit = match Arc::clone(&connections).try_acquire_owned() {
1409            Ok(permit) => permit,
1410            Err(_) => {
1411                // At the cap: drop the stream (peer sees EOF and retries via
1412                // its own reconnect path) instead of queueing unbounded work.
1413                tracing::warn!(
1414                    cap = MAX_EVENTS_CONNECTIONS,
1415                    "events daemon at connection cap; dropping new connection"
1416                );
1417                continue;
1418            }
1419        };
1420        let backend = Arc::clone(&backend);
1421        let stores = Arc::clone(&stores);
1422        let frame_budget = Arc::clone(&frame_budget);
1423        crate::daemon::spawn_named_tracked_task("events_connection", async move {
1424            serve_events_conn(stream, backend, stores, frame_budget).await;
1425            drop(permit);
1426        });
1427    }
1428}
1429
1430#[cfg(unix)]
1431#[path = "events_split_policy_entries.rs"]
1432mod policy_entries;
1433#[cfg(unix)]
1434use policy_entries::standalone_daemon_policies;
1435#[cfg(unix)]
1436pub use policy_entries::{
1437    run_events_daemon_with_wal_ceiling, supervise_events_daemon_with_wal_ceiling,
1438};
1439
1440#[cfg(all(test, unix))]
1441#[path = "events_wal_policy_tests.rs"]
1442mod wal_policy_tests;
1443
1444/// Read one length-prefixed request frame, admitting the body buffer against
1445/// the shared byte budget before allocating it. The returned permit holds
1446/// the admitted bytes for as long as the buffer may be alive — the caller
1447/// drops it after the response is written. The per-frame cap is checked
1448/// before admission, so a single frame is always satisfiable against the
1449/// full budget (`MAX_INFLIGHT_REQUEST_BYTES >= MAX_FRAME_BYTES`) and a
1450/// waiter can never deadlock on an impossible request.
1451#[cfg(unix)]
1452async fn read_frame_budgeted(
1453    stream: &mut UnixStream,
1454    budget: &Arc<tokio::sync::Semaphore>,
1455) -> std::io::Result<(Vec<u8>, tokio::sync::OwnedSemaphorePermit)> {
1456    use tokio::io::AsyncReadExt;
1457    let mut len_buf = [0u8; 4];
1458    stream.read_exact(&mut len_buf).await?;
1459    let len = u32::from_be_bytes(len_buf) as usize;
1460    if len > crate::daemon::MAX_FRAME_BYTES {
1461        return Err(std::io::Error::new(
1462            std::io::ErrorKind::InvalidData,
1463            format!(
1464                "daemon frame of {len} bytes exceeds {} cap",
1465                crate::daemon::MAX_FRAME_BYTES
1466            ),
1467        ));
1468    }
1469    let permit = Arc::clone(budget)
1470        .acquire_many_owned(len as u32)
1471        .await
1472        .map_err(|_| std::io::Error::other("frame budget closed"))?;
1473    let mut buf = vec![0u8; len];
1474    stream.read_exact(&mut buf).await?;
1475    Ok((buf, permit))
1476}
1477
1478#[cfg(unix)]
1479type NamespaceStores =
1480    Arc<std::sync::Mutex<std::collections::HashMap<String, Arc<dyn EventStore>>>>;
1481
1482/// The cached per-namespace store, constructing (and caching) it on first
1483/// use. Construction failures are not cached — the next request retries.
1484///
1485/// The cache key is the trimmed namespace, matching the backend's own
1486/// normalization, so spellings that resolve to one store share one handle.
1487/// The map is bounded: at `MAX_CACHED_NAMESPACE_STORES` entries an arbitrary
1488/// existing entry is evicted to admit the new namespace, so the map never
1489/// grows past the cap and no namespace is permanently condemned to
1490/// per-request store rebuilds. Keys reaching this cache have already passed
1491/// `Namespace` validation at dispatch (charset + length bound), so cap ×
1492/// bounded key is the worst-case retained memory.
1493#[cfg(unix)]
1494fn namespace_store(
1495    backend: &StorageBackend,
1496    stores: &NamespaceStores,
1497    namespace: &str,
1498) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1499    namespace_store_with_cap(backend, stores, namespace, MAX_CACHED_NAMESPACE_STORES)
1500}
1501
1502#[cfg(unix)]
1503fn namespace_store_with_cap(
1504    backend: &StorageBackend,
1505    stores: &NamespaceStores,
1506    namespace: &str,
1507    cap: usize,
1508) -> Result<Arc<dyn EventStore>, khive_db::SqliteError> {
1509    let key = namespace.trim();
1510    if let Some(store) = stores
1511        .lock()
1512        .unwrap_or_else(std::sync::PoisonError::into_inner)
1513        .get(key)
1514    {
1515        return Ok(Arc::clone(store));
1516    }
1517    let store = backend.events_for_namespace(namespace)?;
1518    let mut map = stores
1519        .lock()
1520        .unwrap_or_else(std::sync::PoisonError::into_inner);
1521    if map.len() >= cap && !map.contains_key(key) {
1522        // At the cap, evict an arbitrary entry instead of refusing
1523        // admission: stores are cheap pool handles, so eviction costs one
1524        // rebuild on the victim's next request, while refusing admission
1525        // would make every request beyond the cap rebuild its store and
1526        // rerun schema init on the SQLite writer forever. Keys are
1527        // validated `Namespace` values (bounded length), so the map's
1528        // worst-case memory is cap × small key + cap handles.
1529        if let Some(victim) = map.keys().next().cloned() {
1530            map.remove(&victim);
1531        }
1532    }
1533    map.insert(key.to_string(), Arc::clone(&store));
1534    Ok(store)
1535}
1536
1537#[cfg(unix)]
1538async fn serve_events_conn(
1539    mut stream: UnixStream,
1540    backend: Arc<StorageBackend>,
1541    stores: NamespaceStores,
1542    frame_budget: Arc<tokio::sync::Semaphore>,
1543) {
1544    loop {
1545        // Deadline on the whole frame read — budget admission included: a
1546        // peer that opens a connection and stalls, or that cannot be
1547        // admitted because other connections hold the byte budget, is
1548        // closed instead of holding this task and its descriptor
1549        // indefinitely.
1550        let (payload, budget_permit) = match tokio::time::timeout(
1551            CONN_IO_TIMEOUT,
1552            read_frame_budgeted(&mut stream, &frame_budget),
1553        )
1554        .await
1555        {
1556            Ok(Ok(pair)) => pair,
1557            // Includes clean EOF on peer disconnect, and the expired deadline.
1558            Ok(Err(_)) | Err(_) => return,
1559        };
1560        let response = match serde_json::from_slice::<EventsRequest>(&payload) {
1561            Ok(request) => dispatch_events_request(request, &backend, &stores).await,
1562            Err(error) => EventsResponse::Error {
1563                message: format!("events daemon could not parse request frame: {error}"),
1564                retryable: false,
1565                writer_task_failure: None,
1566            },
1567        };
1568        let bytes = match serde_json::to_vec(&response) {
1569            Ok(bytes) => bytes,
1570            Err(error) => {
1571                tracing::error!(error = %error, "events daemon response serialization failed");
1572                return;
1573            }
1574        };
1575        let bytes = if bytes.len() > crate::daemon::MAX_FRAME_BYTES {
1576            // A page can obey the row cap yet exceed the framing cap. Reply
1577            // with a small, non-retryable refusal instead of closing the
1578            // socket and making a healthy daemon look unreachable.
1579            let refusal = EventsResponse::Error {
1580                message: format!(
1581                    "response_frame_size_limit: {} response bytes exceed the {}-byte events IPC frame cap; request a narrower page",
1582                    bytes.len(),
1583                    crate::daemon::MAX_FRAME_BYTES
1584                ),
1585                retryable: false,
1586                writer_task_failure: None,
1587            };
1588            match serde_json::to_vec(&refusal) {
1589                Ok(bytes) => bytes,
1590                Err(error) => {
1591                    tracing::error!(error = %error, "events daemon frame-size refusal serialization failed");
1592                    return;
1593                }
1594            }
1595        } else {
1596            bytes
1597        };
1598        // Same deadline on the response write: a peer that stops reading
1599        // would otherwise park this task in a full socket buffer.
1600        match tokio::time::timeout(CONN_IO_TIMEOUT, write_frame(&mut stream, &bytes)).await {
1601            Ok(Ok(())) => {}
1602            Ok(Err(_)) | Err(_) => return,
1603        }
1604        // The request buffer is long dropped; release its byte budget only
1605        // now, after the whole request lifecycle, so admission tracks live
1606        // work rather than just live buffers.
1607        drop(budget_permit);
1608    }
1609}
1610
1611#[cfg(unix)]
1612async fn dispatch_events_request(
1613    request: EventsRequest,
1614    backend: &StorageBackend,
1615    stores: &NamespaceStores,
1616) -> EventsResponse {
1617    if request.protocol_version() != EVENTS_PROTOCOL_VERSION {
1618        return EventsResponse::Error {
1619            message: format!(
1620                "events protocol version mismatch: daemon speaks {}, client sent {}",
1621                EVENTS_PROTOCOL_VERSION,
1622                request.protocol_version()
1623            ),
1624            retryable: false,
1625            writer_task_failure: None,
1626        };
1627    }
1628    // The wire namespace is an unrestricted client-supplied String until this
1629    // point. Validate it as a real `Namespace` (charset plus the 256-byte
1630    // length bound) before it can become a cache key, a store, or a database
1631    // row — without this, near-frame-sized strings are attacker-controlled
1632    // retained memory in a process that lives for months.
1633    if let Err(error) = khive_types::Namespace::parse(request.namespace().trim()) {
1634        return EventsResponse::Error {
1635            message: format!("events request namespace rejected: {error}"),
1636            retryable: false,
1637            writer_task_failure: None,
1638        };
1639    }
1640    let store = match namespace_store(backend, stores, request.namespace()) {
1641        Ok(store) => store,
1642        Err(error) => {
1643            return EventsResponse::Error {
1644                message: format!("events store unavailable: {error}"),
1645                retryable: true,
1646                writer_task_failure: None,
1647            };
1648        }
1649    };
1650    match request {
1651        EventsRequest::AppendEvents { events, .. } => match store.append_events(events).await {
1652            Ok(summary) => EventsResponse::Appended { summary },
1653            Err(error) => storage_error_response(&error),
1654        },
1655        EventsRequest::AppendEventsIdempotent { events, .. } => {
1656            match store.append_events_idempotent(events).await {
1657                Ok(result) => EventsResponse::Idempotent { result },
1658                Err(error) => storage_error_response(&error),
1659            }
1660        }
1661        EventsRequest::GetEvent { id, .. } => match store.get_event(id).await {
1662            Ok(event) => EventsResponse::Event { event },
1663            Err(error) => storage_error_response(&error),
1664        },
1665        EventsRequest::QueryEvents { filter, page, .. } => {
1666            if page.limit > MAX_QUERY_EVENTS_PAGE_ROWS {
1667                return EventsResponse::Error {
1668                    message: format!(
1669                        "events query page limit {} exceeds the daemon cap of {} rows; \
1670                         request narrower pages",
1671                        page.limit, MAX_QUERY_EVENTS_PAGE_ROWS
1672                    ),
1673                    retryable: false,
1674                    writer_task_failure: None,
1675                };
1676            }
1677            match store.query_events(filter, page).await {
1678                Ok(page) => EventsResponse::Pageful { page },
1679                Err(error) => storage_error_response(&error),
1680            }
1681        }
1682        EventsRequest::QueryEventPage { query, .. } => {
1683            if let Err(error) = crate::event_page::validate_page_query(&query) {
1684                return storage_error_response(&error);
1685            }
1686            match store.query_event_page(query).await {
1687                Ok(window) => EventsResponse::EventPageWindow { window },
1688                Err(error) => storage_error_response(&error),
1689            }
1690        }
1691        EventsRequest::CountEvents { filter, .. } => match store.count_events(filter).await {
1692            Ok(count) => EventsResponse::Count { count },
1693            Err(error) => storage_error_response(&error),
1694        },
1695    }
1696}
1697
1698#[cfg(unix)]
1699fn storage_error_response(error: &StorageError) -> EventsResponse {
1700    let (message, writer_task_failure) = match error {
1701        StorageError::WriterTaskRequestFailed {
1702            request_state,
1703            source,
1704        } => (
1705            source.to_string(),
1706            Some(WireWriterTaskFailure::RequestFailed {
1707                request_state: (*request_state).into(),
1708            }),
1709        ),
1710        StorageError::WriterTaskTerminated { request_state } => (
1711            error.to_string(),
1712            Some(WireWriterTaskFailure::TaskTerminated {
1713                request_state: (*request_state).into(),
1714            }),
1715        ),
1716        _ => (error.to_string(), None),
1717    };
1718    EventsResponse::Error {
1719        message,
1720        // Defer to the storage layer's own transience classifier rather than
1721        // re-enumerating variants here: a hand-rolled subset silently turns
1722        // transient writer contention (`WriterTaskBusy`, `Transaction`) into
1723        // a terminal error on the client side of the socket.
1724        retryable: error.is_retryable(),
1725        writer_task_failure,
1726    }
1727}
1728
1729// ---------------------------------------------------------------------------
1730// Client side
1731// ---------------------------------------------------------------------------
1732
1733/// Connect to the events socket and verify who answered. Pathname
1734/// reachability is not identity: any process that can write the directory can
1735/// pre-bind the path. The kernel-reported peer uid is the one identity the
1736/// far side cannot choose, so every client connect — round-trips, the
1737/// forwarder, the supervisor probe — refuses a socket answered by a foreign
1738/// uid before a single frame is written.
1739#[cfg(unix)]
1740async fn connect_verified(socket_path: &Path) -> std::io::Result<UnixStream> {
1741    let stream = UnixStream::connect(socket_path).await?;
1742    // SAFETY: `geteuid` is always successful and takes no arguments.
1743    let own_euid = unsafe { libc::geteuid() } as u32;
1744    let peer = crate::daemon::peer_uid(&stream)?;
1745    if !crate::daemon::uid_is_permitted(peer, own_euid) {
1746        return Err(std::io::Error::new(
1747            std::io::ErrorKind::PermissionDenied,
1748            format!(
1749                "events socket {} answered by uid {peer}, not this process's uid {own_euid}; \
1750                 refusing to exchange event frames with an unowned daemon",
1751                socket_path.display()
1752            ),
1753        ));
1754    }
1755    Ok(stream)
1756}
1757
1758/// Counters describing the fire-and-forget lane's degradation. Zero drops is
1759/// the healthy state; any non-zero `dropped_batches` means the loss-tolerant
1760/// contract was exercised and says so.
1761#[cfg(unix)]
1762#[derive(Debug, Clone, Copy, Serialize)]
1763pub struct EventsForwardingMetrics {
1764    pub forwarded_batches: u64,
1765    pub forwarded_events: u64,
1766    pub dropped_batches: u64,
1767    pub dropped_events: u64,
1768    /// Serialized request bytes reserved by queued and in-flight batches.
1769    pub queued_bytes: usize,
1770}
1771
1772#[cfg(unix)]
1773#[derive(Debug, Default)]
1774struct ForwardingCounters {
1775    forwarded_batches: AtomicU64,
1776    forwarded_events: AtomicU64,
1777    dropped_batches: AtomicU64,
1778    dropped_events: AtomicU64,
1779    queued_bytes: AtomicUsize,
1780}
1781
1782/// A reservation lives with the batch until delivery or a counted drop.
1783/// Dropping a forwarder mid-delivery also returns its byte budget.
1784#[cfg(unix)]
1785#[derive(Debug)]
1786struct AppendByteReservation {
1787    counters: Arc<ForwardingCounters>,
1788    bytes: usize,
1789}
1790
1791#[cfg(unix)]
1792impl Drop for AppendByteReservation {
1793    fn drop(&mut self) {
1794        self.counters
1795            .queued_bytes
1796            .fetch_sub(self.bytes, Ordering::AcqRel);
1797    }
1798}
1799
1800#[cfg(unix)]
1801#[derive(Debug)]
1802struct QueuedAppend {
1803    request: EventsRequest,
1804    count: u64,
1805    _reservation: AppendByteReservation,
1806}
1807
1808#[cfg(all(test, unix))]
1809impl QueuedAppend {
1810    fn unmetered(namespace: &str, events: Vec<Event>, counters: Arc<ForwardingCounters>) -> Self {
1811        let count = events.len() as u64;
1812        Self {
1813            request: EventsRequest::AppendEvents {
1814                protocol_version: EVENTS_PROTOCOL_VERSION,
1815                namespace: namespace.to_string(),
1816                events,
1817            },
1818            count,
1819            _reservation: AppendByteReservation { counters, bytes: 0 },
1820        }
1821    }
1822}
1823
1824/// Measure the actual wire request without retaining a second, potentially
1825/// huge serialized copy of the caller's batch. Serde stops at the frame cap.
1826#[cfg(unix)]
1827#[derive(Default)]
1828struct BoundedFrameCounter {
1829    bytes: usize,
1830    exceeded: bool,
1831}
1832
1833#[cfg(unix)]
1834impl std::io::Write for BoundedFrameCounter {
1835    fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
1836        let next = self.bytes.saturating_add(buf.len());
1837        if next > crate::daemon::MAX_FRAME_BYTES {
1838            self.exceeded = true;
1839            return Err(std::io::Error::new(
1840                std::io::ErrorKind::InvalidData,
1841                "events append request exceeds IPC frame cap",
1842            ));
1843        }
1844        self.bytes = next;
1845        Ok(buf.len())
1846    }
1847
1848    fn flush(&mut self) -> std::io::Result<()> {
1849        Ok(())
1850    }
1851}
1852
1853#[cfg(unix)]
1854fn reserve_append_bytes(
1855    counters: &Arc<ForwardingCounters>,
1856    bytes: usize,
1857    budget: usize,
1858) -> Option<AppendByteReservation> {
1859    let mut used = counters.queued_bytes.load(Ordering::Acquire);
1860    loop {
1861        let next = used.checked_add(bytes)?;
1862        if next > budget {
1863            return None;
1864        }
1865        match counters.queued_bytes.compare_exchange_weak(
1866            used,
1867            next,
1868            Ordering::AcqRel,
1869            Ordering::Acquire,
1870        ) {
1871            Ok(_) => {
1872                return Some(AppendByteReservation {
1873                    counters: Arc::clone(counters),
1874                    bytes,
1875                });
1876            }
1877            Err(observed) => used = observed,
1878        }
1879    }
1880}
1881
1882/// One per domain process: the connection to the events daemon plus the
1883/// bounded fire-and-forget append queue.
1884///
1885/// Synchronous round-trips each own a fresh connection: there is no shared
1886/// request connection and no lock in front of one, so concurrent reads and
1887/// audit flushes never queue behind a single stalled call, and the round-trip
1888/// timeout bounds each caller's whole attempt.
1889#[cfg(unix)]
1890pub struct EventsSplitClient {
1891    socket_path: PathBuf,
1892    append_tx: tokio::sync::mpsc::Sender<QueuedAppend>,
1893    append_queue_byte_budget: usize,
1894    counters: Arc<ForwardingCounters>,
1895    /// Flipped by the forwarder while the daemon is unreachable so the drop
1896    /// log fires once per outage, not once per batch.
1897    outage_logged: Arc<AtomicBool>,
1898    /// In-memory validator backing `preflight_event` without I/O.
1899    preflight_store: Arc<dyn EventStore>,
1900}
1901
1902#[cfg(unix)]
1903impl std::fmt::Debug for EventsSplitClient {
1904    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1905        f.debug_struct("EventsSplitClient")
1906            .field("socket_path", &self.socket_path)
1907            .finish_non_exhaustive()
1908    }
1909}
1910
1911#[cfg(unix)]
1912impl EventsSplitClient {
1913    /// Build the client and spawn its background forwarder.
1914    pub fn new(socket_path: PathBuf) -> crate::error::RuntimeResult<Arc<Self>> {
1915        Self::new_with_queue_depth(socket_path, DEFAULT_APPEND_QUEUE_BATCHES)
1916    }
1917
1918    /// [`Self::new`] with an explicit fire-and-forget queue bound. Tests use a
1919    /// tiny depth to exercise the overflow drop arm deterministically.
1920    pub fn new_with_queue_depth(
1921        socket_path: PathBuf,
1922        queue_depth: usize,
1923    ) -> crate::error::RuntimeResult<Arc<Self>> {
1924        Self::new_with_queue_depth_and_delivery_timeout(socket_path, queue_depth, REQUEST_TIMEOUT)
1925    }
1926
1927    /// [`Self::new_with_queue_depth`] with an explicit per-batch delivery
1928    /// timeout, so tests can prove the forwarder abandons a hung peer without
1929    /// waiting out the production clock.
1930    fn new_with_queue_depth_and_delivery_timeout(
1931        socket_path: PathBuf,
1932        queue_depth: usize,
1933        delivery_timeout: Duration,
1934    ) -> crate::error::RuntimeResult<Arc<Self>> {
1935        Self::new_with_limits_and_delivery_timeout(
1936            socket_path,
1937            queue_depth,
1938            DEFAULT_APPEND_QUEUE_BYTES,
1939            delivery_timeout,
1940        )
1941    }
1942
1943    fn new_with_limits_and_delivery_timeout(
1944        socket_path: PathBuf,
1945        queue_depth: usize,
1946        byte_budget: usize,
1947        delivery_timeout: Duration,
1948    ) -> crate::error::RuntimeResult<Arc<Self>> {
1949        validate_events_socket_path(&socket_path)?;
1950        let preflight_backend = StorageBackend::memory()?;
1951        let preflight_store = preflight_backend.events()?;
1952        // The in-memory backend must outlive the store handle; the store holds
1953        // the pool Arc internally, so dropping the backend wrapper here is fine.
1954
1955        let (append_tx, append_rx) = tokio::sync::mpsc::channel::<QueuedAppend>(queue_depth.max(1));
1956        let counters = Arc::new(ForwardingCounters::default());
1957        let outage_logged = Arc::new(AtomicBool::new(false));
1958
1959        let client = Arc::new(Self {
1960            socket_path: socket_path.clone(),
1961            append_tx,
1962            append_queue_byte_budget: byte_budget,
1963            counters: Arc::clone(&counters),
1964            outage_logged: Arc::clone(&outage_logged),
1965            preflight_store,
1966        });
1967
1968        crate::daemon::spawn_named_tracked_task(
1969            "events_forwarder",
1970            run_forwarder(
1971                socket_path,
1972                append_rx,
1973                counters,
1974                outage_logged,
1975                delivery_timeout,
1976                crate::daemon::daemon_shutdown_token(),
1977            ),
1978        );
1979        Ok(client)
1980    }
1981
1982    /// Snapshot of the fire-and-forget lane's health.
1983    pub fn metrics(&self) -> EventsForwardingMetrics {
1984        EventsForwardingMetrics {
1985            forwarded_batches: self.counters.forwarded_batches.load(Ordering::Relaxed),
1986            forwarded_events: self.counters.forwarded_events.load(Ordering::Relaxed),
1987            dropped_batches: self.counters.dropped_batches.load(Ordering::Relaxed),
1988            dropped_events: self.counters.dropped_events.load(Ordering::Relaxed),
1989            queued_bytes: self.counters.queued_bytes.load(Ordering::Acquire),
1990        }
1991    }
1992
1993    /// Enqueue a batch on the fire-and-forget lane. Never blocks; a full
1994    /// queue is a counted, logged drop.
1995    fn enqueue(&self, namespace: &str, events: Vec<Event>) {
1996        let count = events.len() as u64;
1997        let request = EventsRequest::AppendEvents {
1998            protocol_version: EVENTS_PROTOCOL_VERSION,
1999            namespace: namespace.to_string(),
2000            events,
2001        };
2002        let mut frame_size = BoundedFrameCounter::default();
2003        if let Err(error) = serde_json::to_writer(&mut frame_size, &request) {
2004            self.count_append_drop(
2005                count,
2006                if frame_size.exceeded {
2007                    "frame cap"
2008                } else {
2009                    "serialization"
2010                },
2011            );
2012            tracing::warn!(error = %error, "events append batch could not fit a wire frame");
2013            return;
2014        }
2015        let Some(reservation) = reserve_append_bytes(
2016            &self.counters,
2017            frame_size.bytes,
2018            self.append_queue_byte_budget,
2019        ) else {
2020            self.count_append_drop(count, "queue byte budget");
2021            return;
2022        };
2023        let batch = QueuedAppend {
2024            request,
2025            count,
2026            _reservation: reservation,
2027        };
2028        match self.append_tx.try_send(batch) {
2029            Ok(()) => {}
2030            Err(_) => {
2031                // The rejected batch drops here and releases its reservation.
2032                self.count_append_drop(count, "queue batch limit");
2033            }
2034        }
2035    }
2036
2037    fn count_append_drop(&self, count: u64, reason: &'static str) {
2038        self.counters
2039            .dropped_batches
2040            .fetch_add(1, Ordering::Relaxed);
2041        self.counters
2042            .dropped_events
2043            .fetch_add(count, Ordering::Relaxed);
2044        if !self.outage_logged.swap(true, Ordering::Relaxed) {
2045            tracing::warn!(
2046                dropped_events = count,
2047                reason,
2048                "events append queue rejected loss-tolerant batch"
2049            );
2050        }
2051    }
2052
2053    /// One synchronous framed round-trip with a bounded timeout, on its own
2054    /// fresh connection. No connection is shared between concurrent callers,
2055    /// so no caller ever waits on another's stall, and `REQUEST_TIMEOUT`
2056    /// bounds the whole attempt — connect, peer verification, write, read.
2057    async fn round_trip(&self, request: &EventsRequest) -> StorageResult<EventsResponse> {
2058        let op = "events daemon round-trip";
2059        let payload = serde_json::to_vec(request).map_err(|error| StorageError::Serialization {
2060            capability: khive_storage::StorageCapability::Events,
2061            message: format!("events request serialization failed: {error}"),
2062        })?;
2063
2064        let attempt = async {
2065            let mut stream = connect_verified(&self.socket_path)
2066                .await
2067                .map_err(|error| ("connect", error))?;
2068            write_frame(&mut stream, &payload)
2069                .await
2070                .map_err(|error| ("write", error))?;
2071            read_frame(&mut stream)
2072                .await
2073                .map_err(|error| ("read", error))
2074        };
2075        let bytes = match tokio::time::timeout(REQUEST_TIMEOUT, attempt).await {
2076            Ok(Ok(bytes)) => bytes,
2077            Ok(Err(("read", error))) => {
2078                return Err(StorageError::Serialization {
2079                    capability: khive_storage::StorageCapability::Events,
2080                    message: format!(
2081                        "events daemon closed or broke the response after connection at {}: {error}",
2082                        self.socket_path.display()
2083                    ),
2084                });
2085            }
2086            Ok(Err((stage, error))) => {
2087                return Err(StorageError::Pool {
2088                    operation: op.into(),
2089                    message: format!(
2090                        "events daemon {stage} failed at {}: {error}",
2091                        self.socket_path.display()
2092                    ),
2093                });
2094            }
2095            Err(_elapsed) => {
2096                return Err(StorageError::Timeout {
2097                    operation: op.into(),
2098                });
2099            }
2100        };
2101
2102        serde_json::from_slice::<EventsResponse>(&bytes).map_err(|error| {
2103            StorageError::Serialization {
2104                capability: khive_storage::StorageCapability::Events,
2105                message: format!("events response deserialization failed: {error}"),
2106            }
2107        })
2108    }
2109}
2110
2111/// Drain every batch currently sitting in the fire-and-forget queue,
2112/// counting each as a dropped batch/events and logging one summary line
2113/// (never one line per batch — a full queue at shutdown is exactly the
2114/// bursty case a per-batch log would flood). Non-blocking: `try_recv` only
2115/// consumes what is already queued, so it terminates as soon as the queue
2116/// (temporarily) empties even though the sender half is still live.
2117#[cfg(unix)]
2118fn drain_dropped_queue(
2119    rx: &mut tokio::sync::mpsc::Receiver<QueuedAppend>,
2120    counters: &ForwardingCounters,
2121) {
2122    let mut dropped_batches = 0u64;
2123    let mut dropped_events = 0u64;
2124    while let Ok(batch) = rx.try_recv() {
2125        dropped_batches += 1;
2126        dropped_events += batch.count;
2127        // `batch` drops its byte reservation at the end of this iteration.
2128    }
2129    if dropped_batches > 0 {
2130        counters
2131            .dropped_batches
2132            .fetch_add(dropped_batches, Ordering::Relaxed);
2133        counters
2134            .dropped_events
2135            .fetch_add(dropped_events, Ordering::Relaxed);
2136        tracing::warn!(
2137            dropped_batches,
2138            dropped_events,
2139            "events forwarder shutting down; dropping queued loss-tolerant batches"
2140        );
2141    }
2142}
2143
2144/// The background fire-and-forget forwarder: drains the bounded queue into
2145/// framed appends on its own connection, reconnecting with backoff. A batch
2146/// that cannot be delivered is dropped and counted — never retried, never
2147/// blocking the queue behind it.
2148#[cfg(unix)]
2149async fn run_forwarder(
2150    socket_path: PathBuf,
2151    mut rx: tokio::sync::mpsc::Receiver<QueuedAppend>,
2152    counters: Arc<ForwardingCounters>,
2153    outage_logged: Arc<AtomicBool>,
2154    delivery_timeout: Duration,
2155    shutdown: tokio_util::sync::CancellationToken,
2156) {
2157    // The client lives in a process-global registry, so its sender is never
2158    // dropped and `rx.recv()` alone would keep this tracked task alive
2159    // through daemon shutdown, forcing `drain()` to its full timeout.
2160    // Observe the shutdown token directly: queued batches at shutdown are a
2161    // counted drop, exactly the lane's loss-tolerant contract — enforced by
2162    // draining and counting `rx` below on every shutdown exit, not just
2163    // implied by the comment.
2164    let mut conn: Option<UnixStream> = None;
2165    loop {
2166        let batch = tokio::select! {
2167            _ = shutdown.cancelled() => {
2168                drain_dropped_queue(&mut rx, &counters);
2169                break;
2170            },
2171            received = rx.recv() => match received {
2172                Some(batch) => batch,
2173                None => break,
2174            },
2175        };
2176        let count = batch.count;
2177        let payload = match serde_json::to_vec(&batch.request) {
2178            Ok(payload) => payload,
2179            Err(error) => {
2180                counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2181                counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2182                tracing::error!(error = %error, "events forwarder serialization failed; batch dropped");
2183                continue;
2184            }
2185        };
2186
2187        // The delivery itself must stay both cancellable and bounded: a
2188        // connected but non-responding daemon would otherwise park this
2189        // tracked task inside write_frame/read_frame, beyond the reach of the
2190        // recv-side shutdown select — and daemon drain would wait its full
2191        // timeout on it. Shutdown mid-delivery abandons the batch (the
2192        // lane's loss-tolerant contract); a timeout poisons the connection,
2193        // since the peer may answer the abandoned frame later.
2194        let delivered = tokio::select! {
2195            _ = shutdown.cancelled() => {
2196                counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2197                counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2198                tracing::warn!(
2199                    dropped_events = count,
2200                    "events forwarder shutting down mid-delivery; in-flight batch dropped"
2201                );
2202                drain_dropped_queue(&mut rx, &counters);
2203                break;
2204            },
2205            outcome = tokio::time::timeout(
2206                delivery_timeout,
2207                deliver_batch(&socket_path, &mut conn, &payload),
2208            ) => match outcome {
2209                Ok(delivered) => delivered,
2210                Err(_elapsed) => {
2211                    conn = None;
2212                    false
2213                }
2214            },
2215        };
2216        if delivered {
2217            counters.forwarded_batches.fetch_add(1, Ordering::Relaxed);
2218            counters
2219                .forwarded_events
2220                .fetch_add(count, Ordering::Relaxed);
2221            if outage_logged.swap(false, Ordering::Relaxed) {
2222                tracing::info!("events daemon reachable again; forwarding resumed");
2223            }
2224        } else {
2225            counters.dropped_batches.fetch_add(1, Ordering::Relaxed);
2226            counters.dropped_events.fetch_add(count, Ordering::Relaxed);
2227            if !outage_logged.swap(true, Ordering::Relaxed) {
2228                tracing::warn!(
2229                    socket = %socket_path.display(),
2230                    "events daemon unreachable; dropping loss-tolerant events until it returns"
2231                );
2232            }
2233            tokio::select! {
2234                _ = shutdown.cancelled() => {
2235                    drain_dropped_queue(&mut rx, &counters);
2236                    break;
2237                },
2238                _ = tokio::time::sleep(FORWARDER_BACKOFF) => {}
2239            }
2240        }
2241    }
2242}
2243
2244/// Try to deliver one framed append over the forwarder connection,
2245/// (re)connecting at most once. Returns whether the daemon acknowledged.
2246/// Connections are peer-verified: a socket answered by a foreign uid is a
2247/// failed connect, never a delivery target.
2248#[cfg(unix)]
2249async fn deliver_batch(socket_path: &Path, conn: &mut Option<UnixStream>, payload: &[u8]) -> bool {
2250    for _attempt in 0..2u8 {
2251        if conn.is_none() {
2252            match connect_verified(socket_path).await {
2253                Ok(stream) => *conn = Some(stream),
2254                Err(_) => return false,
2255            }
2256        }
2257        let stream = conn.as_mut().expect("connection populated above");
2258        let ok = async {
2259            write_frame(stream, payload).await?;
2260            let bytes = read_frame(stream).await?;
2261            std::io::Result::Ok(bytes)
2262        }
2263        .await;
2264        match ok {
2265            Ok(bytes) => {
2266                return !matches!(
2267                    serde_json::from_slice::<EventsResponse>(&bytes),
2268                    Ok(EventsResponse::Error { .. }) | Err(_)
2269                );
2270            }
2271            Err(_) => {
2272                // Stale connection (daemon restarted): drop it and retry once
2273                // with a fresh connect; a second failure is a real outage.
2274                *conn = None;
2275            }
2276        }
2277    }
2278    false
2279}
2280
2281// ---------------------------------------------------------------------------
2282// The per-namespace EventStore handle
2283// ---------------------------------------------------------------------------
2284
2285/// [`EventStore`] implementation the runtime hands out when the events split
2286/// runs in daemon mode. Appends are fire-and-forget through the client's
2287/// bounded queue; the ADR-133 idempotent lane and all reads are synchronous
2288/// round-trips to the events daemon.
2289#[cfg(unix)]
2290#[derive(Debug)]
2291pub struct ForwardingEventStore {
2292    namespace: String,
2293    client: Arc<EventsSplitClient>,
2294}
2295
2296#[cfg(unix)]
2297impl ForwardingEventStore {
2298    pub fn new(namespace: impl Into<String>, client: Arc<EventsSplitClient>) -> Self {
2299        Self {
2300            namespace: namespace.into(),
2301            client,
2302        }
2303    }
2304
2305    fn unexpected(&self, op: &'static str, response: EventsResponse) -> StorageError {
2306        match response {
2307            EventsResponse::Error {
2308                message,
2309                retryable,
2310                writer_task_failure,
2311            } => {
2312                if let Some(failure) = writer_task_failure {
2313                    match failure {
2314                        WireWriterTaskFailure::RequestFailed { request_state } => {
2315                            let source = if retryable {
2316                                StorageError::Pool {
2317                                    operation: op.into(),
2318                                    message,
2319                                }
2320                            } else {
2321                                StorageError::InvalidInput {
2322                                    capability: khive_storage::StorageCapability::Events,
2323                                    operation: op.into(),
2324                                    message,
2325                                }
2326                            };
2327                            StorageError::WriterTaskRequestFailed {
2328                                request_state: request_state.into(),
2329                                source: Box::new(source),
2330                            }
2331                        }
2332                        WireWriterTaskFailure::TaskTerminated { request_state } => {
2333                            StorageError::WriterTaskTerminated {
2334                                request_state: request_state.into(),
2335                            }
2336                        }
2337                    }
2338                } else if retryable {
2339                    StorageError::Pool {
2340                        operation: op.into(),
2341                        message,
2342                    }
2343                } else {
2344                    StorageError::InvalidInput {
2345                        capability: khive_storage::StorageCapability::Events,
2346                        operation: op.into(),
2347                        message,
2348                    }
2349                }
2350            }
2351            other => StorageError::Serialization {
2352                capability: khive_storage::StorageCapability::Events,
2353                message: format!("events daemon returned mismatched response for {op}: {other:?}"),
2354            },
2355        }
2356    }
2357}
2358
2359#[cfg(unix)]
2360#[async_trait]
2361impl EventStore for ForwardingEventStore {
2362    async fn append_event(&self, event: Event) -> StorageResult<()> {
2363        self.client.enqueue(&self.namespace, vec![event]);
2364        Ok(())
2365    }
2366
2367    async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2368        let attempted = events.len() as u64;
2369        self.client.enqueue(&self.namespace, events);
2370        // Fire-and-forget: the hand-off succeeded or was counted as a drop;
2371        // either way the caller's contract is "accepted for forwarding".
2372        Ok(BatchWriteSummary {
2373            attempted,
2374            affected: attempted,
2375            ..BatchWriteSummary::default()
2376        })
2377    }
2378
2379    async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2380        let request = EventsRequest::GetEvent {
2381            protocol_version: EVENTS_PROTOCOL_VERSION,
2382            namespace: self.namespace.clone(),
2383            id,
2384        };
2385        match self.client.round_trip(&request).await? {
2386            EventsResponse::Event { event } => Ok(event),
2387            other => Err(self.unexpected("get_event", other)),
2388        }
2389    }
2390
2391    async fn query_events(
2392        &self,
2393        filter: EventFilter,
2394        page: PageRequest,
2395    ) -> StorageResult<Page<Event>> {
2396        let request = EventsRequest::QueryEvents {
2397            protocol_version: EVENTS_PROTOCOL_VERSION,
2398            namespace: self.namespace.clone(),
2399            filter,
2400            page,
2401        };
2402        match self.client.round_trip(&request).await? {
2403            EventsResponse::Pageful { page } => Ok(page),
2404            other => Err(self.unexpected("query_events", other)),
2405        }
2406    }
2407
2408    async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2409        let request = EventsRequest::CountEvents {
2410            protocol_version: EVENTS_PROTOCOL_VERSION,
2411            namespace: self.namespace.clone(),
2412            filter,
2413        };
2414        match self.client.round_trip(&request).await? {
2415            EventsResponse::Count { count } => Ok(count),
2416            other => Err(self.unexpected("count_events", other)),
2417        }
2418    }
2419
2420    async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
2421        crate::event_page::validate_page_query(&query)?;
2422        let request = EventsRequest::QueryEventPage {
2423            protocol_version: EVENTS_PROTOCOL_VERSION,
2424            namespace: self.namespace.clone(),
2425            query: query.clone(),
2426        };
2427        match self.client.round_trip(&request).await? {
2428            EventsResponse::EventPageWindow { window } => {
2429                crate::event_page::validate_window(&query, Some(&self.namespace), &window)?;
2430                Ok(window)
2431            }
2432            error @ EventsResponse::Error { .. } => Err(self.unexpected("query_event_page", error)),
2433            _ => Err(crate::event_page::page_error(
2434                "events page response kind mismatch",
2435            )),
2436        }
2437    }
2438
2439    fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2440        // Same validation code path as the daemon-side store, zero I/O.
2441        self.client.preflight_store.preflight_event(event)
2442    }
2443
2444    async fn append_events_idempotent(
2445        &self,
2446        events: Vec<Event>,
2447    ) -> StorageResult<IdempotentEventBatchResult> {
2448        let request = EventsRequest::AppendEventsIdempotent {
2449            protocol_version: EVENTS_PROTOCOL_VERSION,
2450            namespace: self.namespace.clone(),
2451            events,
2452        };
2453        match self.client.round_trip(&request).await? {
2454            EventsResponse::Idempotent { result } => Ok(result),
2455            other => Err(self.unexpected("append_events_idempotent", other)),
2456        }
2457    }
2458
2459    fn supports_idempotent_audit_batch(&self) -> bool {
2460        true
2461    }
2462}
2463
2464/// The store the runtime hands out when the events split is configured:
2465/// routes by APPEND CLASS rather than moving the whole event plane.
2466///
2467/// - The ADR-133 idempotent audit-batch lane — the measured bulk of event
2468///   write volume (verb-dispatch audit plus the config-lock rows that ride
2469///   the same flusher) — goes to the events lane (the events database, forwarded or
2470///   direct).
2471/// - Plain appends stay on the legacy store, because the legacy `events`
2472///   table has raw-SQL consumers whose correctness depends on finding those
2473///   rows there: the schedule drain's creator-provenance fence, the kg
2474///   projection worker's guarded event INSERT (transactional with main-db
2475///   state), and GraphQuery's cross-substrate UNION. Those events are
2476///   low-volume domain facts; the split's contention relief does not need
2477///   them moved, and moving them breaks the consumers by construction.
2478/// - Reads merge both stores so trait-level consumers
2479///   (`brain.event_counts`, event getters) observe one event plane.
2480pub struct SplitEventStore {
2481    legacy: Arc<dyn EventStore>,
2482    lane: Arc<dyn EventStore>,
2483}
2484
2485impl std::fmt::Debug for SplitEventStore {
2486    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2487        f.debug_struct("SplitEventStore").finish_non_exhaustive()
2488    }
2489}
2490
2491impl SplitEventStore {
2492    /// Bound on `offset + limit` for a merged `query_events` window. Offset
2493    /// pagination over two stores materializes the whole prefix in memory
2494    /// (see `query_events`), so an unbounded offset would let a single
2495    /// request buffer both stores wholesale. The bound comfortably admits
2496    /// the largest legitimate bounded window in the tree
2497    /// (`brain.event_counts`' 50k page); deep walks page with a `before`
2498    /// cursor at `offset: 0`, which never grows the materialized prefix.
2499    pub const MAX_MERGED_WINDOW_ROWS: u64 = 100_000;
2500
2501    pub fn new(legacy: Arc<dyn EventStore>, lane: Arc<dyn EventStore>) -> Self {
2502        Self { legacy, lane }
2503    }
2504}
2505
2506#[async_trait]
2507impl EventStore for SplitEventStore {
2508    async fn append_event(&self, event: Event) -> StorageResult<()> {
2509        self.legacy.append_event(event).await
2510    }
2511
2512    async fn append_events(&self, events: Vec<Event>) -> StorageResult<BatchWriteSummary> {
2513        self.legacy.append_events(events).await
2514    }
2515
2516    async fn get_event(&self, id: Uuid) -> StorageResult<Option<Event>> {
2517        // Domain events (legacy) are the likelier and cheaper hit; fall back
2518        // to the lane for audit-batch rows.
2519        if let Some(event) = self.legacy.get_event(id).await? {
2520            return Ok(Some(event));
2521        }
2522        self.lane.get_event(id).await
2523    }
2524
2525    async fn query_events(
2526        &self,
2527        filter: EventFilter,
2528        page: PageRequest,
2529    ) -> StorageResult<Page<Event>> {
2530        // Offset pagination cannot be split across two stores: fetch each
2531        // store's prefix covering the requested window, merge in the stores'
2532        // shared order (created_at DESC, id DESC), then window in memory.
2533        // Cost is O(offset + limit) rows materialized from each store — the
2534        // floor for offset semantics over two sources, since the split point
2535        // is unknowable without both prefixes. The sort below is a single
2536        // merge pass in practice (std's stable sort is adaptive on the
2537        // concatenation of two sorted runs). That materialization is why the
2538        // window is bounded below: without a bound, one authenticated
2539        // request with a pathological offset would make both stores buffer
2540        // and sort every matching row. Deep walks page with a `before`
2541        // cursor (strict `created_at <` bound in `EventFilter`) at
2542        // `offset: 0`, which stays inside the bound at any depth.
2543        let window = page.offset.saturating_add(u64::from(page.limit));
2544        if window > Self::MAX_MERGED_WINDOW_ROWS {
2545            return Err(StorageError::InvalidInput {
2546                capability: khive_storage::StorageCapability::Events,
2547                operation: "query_events".into(),
2548                message: format!(
2549                    "offset+limit ({window}) exceeds the merged event plane's window bound \
2550                     of {}; page deep windows with a `before` cursor at offset 0, or narrow \
2551                     the filter",
2552                    Self::MAX_MERGED_WINDOW_ROWS
2553                ),
2554            });
2555        }
2556        let prefix = PageRequest {
2557            offset: 0,
2558            limit: window.min(u64::from(u32::MAX)) as u32,
2559        };
2560        let legacy = self
2561            .legacy
2562            .query_events(filter.clone(), prefix.clone())
2563            .await?;
2564        let lane = self.lane.query_events(filter, prefix).await?;
2565        let total = match (legacy.total, lane.total) {
2566            (Some(a), Some(b)) => Some(a + b),
2567            _ => None,
2568        };
2569        let mut items = legacy.items;
2570        items.extend(lane.items);
2571        items.sort_by(|a, b| {
2572            b.created_at
2573                .cmp(&a.created_at)
2574                .then_with(|| b.id.cmp(&a.id))
2575        });
2576        let items = items
2577            .into_iter()
2578            .skip(page.offset as usize)
2579            .take(page.limit as usize)
2580            .collect();
2581        Ok(Page { items, total })
2582    }
2583
2584    async fn count_events(&self, filter: EventFilter) -> StorageResult<u64> {
2585        let legacy = self.legacy.count_events(filter.clone()).await?;
2586        let lane = self.lane.count_events(filter).await?;
2587        Ok(legacy + lane)
2588    }
2589
2590    async fn query_event_page(&self, query: EventPageQuery) -> StorageResult<EventPageWindow> {
2591        crate::event_page::split_page(self.legacy.as_ref(), self.lane.as_ref(), query).await
2592    }
2593
2594    fn preflight_event(&self, event: &Event) -> StorageResult<()> {
2595        // Validation is store-independent; the lane's implementation is the
2596        // local zero-I/O validator in forwarding mode.
2597        self.lane.preflight_event(event)
2598    }
2599
2600    async fn append_events_idempotent(
2601        &self,
2602        events: Vec<Event>,
2603    ) -> StorageResult<IdempotentEventBatchResult> {
2604        // A retried batch may have landed — fully or partially — on the
2605        // legacy store before the split cutover, and merged reads stay
2606        // duplicate-free only while an id lives in exactly one store. Probe
2607        // the legacy store for the batch's ids and route resident rows back
2608        // through the legacy store's own idempotent machinery, which compares
2609        // every persisted column without re-inserting; only genuinely new
2610        // rows reach the lane. Without the probe, a legacy-resident id would
2611        // be inserted a second time into the lane and every merged query and
2612        // count would double-count it.
2613        if events.is_empty() {
2614            return self.lane.append_events_idempotent(events).await;
2615        }
2616        let ids: Vec<Uuid> = events.iter().map(|event| event.id).collect();
2617        let probe_limit = u32::try_from(ids.len()).map_err(|_| StorageError::InvalidInput {
2618            capability: khive_storage::StorageCapability::Events,
2619            operation: "append_events_idempotent".into(),
2620            message: format!(
2621                "batch of {} rows exceeds the legacy probe window",
2622                ids.len()
2623            ),
2624        })?;
2625        let existing = self
2626            .legacy
2627            .query_events(
2628                EventFilter {
2629                    ids,
2630                    ..EventFilter::default()
2631                },
2632                PageRequest {
2633                    offset: 0,
2634                    limit: probe_limit,
2635                },
2636            )
2637            .await?;
2638        let legacy_resident: std::collections::HashSet<Uuid> =
2639            existing.items.iter().map(|event| event.id).collect();
2640        if legacy_resident.is_empty() {
2641            return self.lane.append_events_idempotent(events).await;
2642        }
2643        let mut legacy_rows = Vec::new();
2644        let mut lane_rows = Vec::new();
2645        let mut routed_to_legacy = Vec::with_capacity(events.len());
2646        for event in events {
2647            if legacy_resident.contains(&event.id) {
2648                routed_to_legacy.push(true);
2649                legacy_rows.push(event);
2650            } else {
2651                routed_to_legacy.push(false);
2652                lane_rows.push(event);
2653            }
2654        }
2655        let legacy_expected = legacy_rows.len();
2656        let lane_expected = lane_rows.len();
2657        let legacy_result = self.legacy.append_events_idempotent(legacy_rows).await?;
2658        let lane_result = if lane_expected == 0 {
2659            IdempotentEventBatchResult { rows: Vec::new() }
2660        } else {
2661            self.lane.append_events_idempotent(lane_rows).await?
2662        };
2663        if legacy_result.rows.len() != legacy_expected || lane_result.rows.len() != lane_expected {
2664            return Err(StorageError::Driver {
2665                capability: khive_storage::StorageCapability::Events,
2666                operation: "append_events_idempotent".into(),
2667                source: format!(
2668                    "idempotent sub-batch result length mismatch: legacy {}/{}, lane {}/{}",
2669                    legacy_result.rows.len(),
2670                    legacy_expected,
2671                    lane_result.rows.len(),
2672                    lane_expected
2673                )
2674                .into(),
2675            });
2676        }
2677        let mut legacy_iter = legacy_result.rows.into_iter();
2678        let mut lane_iter = lane_result.rows.into_iter();
2679        let rows = routed_to_legacy
2680            .into_iter()
2681            .map(|to_legacy| {
2682                if to_legacy {
2683                    legacy_iter.next().expect("length checked above")
2684                } else {
2685                    lane_iter.next().expect("length checked above")
2686                }
2687            })
2688            .collect();
2689        Ok(IdempotentEventBatchResult { rows })
2690    }
2691
2692    fn supports_idempotent_audit_batch(&self) -> bool {
2693        // The legacy-resident probe above routes retried rows through the
2694        // legacy store's idempotent path, so the split supports the audit
2695        // batch only when BOTH stores do.
2696        self.lane.supports_idempotent_audit_batch() && self.legacy.supports_idempotent_audit_batch()
2697    }
2698}
2699
2700#[cfg(test)]
2701mod tests {
2702    use super::*;
2703
2704    #[test]
2705    fn fixture_registry_cleanup_preserves_other_roots() {
2706        let live_dir = tempfile::tempdir().unwrap();
2707        let _live_guard = TestRegistryGuard::new(live_dir.path());
2708        let live_path = live_dir.path().join("events.db");
2709        let live_backend = direct_backend_for(&live_path).unwrap();
2710
2711        let finished_dir = tempfile::tempdir().unwrap();
2712        let finished_backend = {
2713            let _finished_guard = TestRegistryGuard::new(finished_dir.path());
2714            let backend = direct_backend_for(&finished_dir.path().join("events.db")).unwrap();
2715            let weak = Arc::downgrade(&backend);
2716            drop(backend);
2717            assert!(
2718                weak.upgrade().is_some(),
2719                "registry retains the live fixture"
2720            );
2721            weak
2722        };
2723
2724        assert!(
2725            finished_backend.upgrade().is_none(),
2726            "FD_FIXTURE_REGISTRY_RELEASED"
2727        );
2728        let same_live_backend = direct_backend_for(&live_path).unwrap();
2729        assert!(
2730            Arc::ptr_eq(&live_backend, &same_live_backend),
2731            "FD_FIXTURE_OTHER_ROOT_RETAINED"
2732        );
2733    }
2734
2735    #[cfg(unix)]
2736    mod target_filter_tests {
2737        include!("events_split_target_tests.rs");
2738    }
2739
2740    // The full daemon/forwarding test suite for this module lands with the
2741    // final slice of this series; these tests cover the naming and
2742    // classification contracts the module itself defines.
2743
2744    include!("events_sidecar_path_tests.rs");
2745
2746    #[cfg(unix)]
2747    #[test]
2748    fn deep_symlink_chains_within_the_kernel_bound_derive_the_target_sidecar() {
2749        // Linux resolves up to 40 links per lookup; a 35-hop dangling chain
2750        // therefore opens fine, so it must ALSO derive the target's sidecar
2751        // — a resolver stopping short would split audit records between the
2752        // alias and target spellings of one database.
2753        let dir = tempfile::tempdir().unwrap();
2754        let real = dir.path().join("real.db");
2755        let mut prev = real.clone();
2756        for i in 0..35 {
2757            let link = dir.path().join(format!("hop{i}.db"));
2758            std::os::unix::fs::symlink(&prev, &link).unwrap();
2759            prev = link;
2760        }
2761        assert_eq!(
2762            events_db_path_beside(&prev),
2763            events_db_path_beside(&real),
2764            "a 35-hop dangling chain must derive the target's sidecar"
2765        );
2766        // Control: an unrelated dangling path still derives its own name.
2767        assert_ne!(
2768            events_db_path_beside(&dir.path().join("unrelated.db")),
2769            events_db_path_beside(&real)
2770        );
2771    }
2772
2773    include!("events_sidecar_hardening_tests.rs");
2774
2775    #[cfg(unix)]
2776    fn write_owner_only_test_file(path: &Path) {
2777        use std::os::unix::fs::PermissionsExt;
2778        std::fs::write(path, b"").unwrap();
2779        std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).unwrap();
2780    }
2781
2782    #[cfg(unix)]
2783    #[test]
2784    fn unopened_check_refuses_same_mode_regular_file_replacement() {
2785        // Must fail with the old mode-only verifier: every replacement remains
2786        // a regular 0600 file. Creating it before rename guarantees a different
2787        // live inode, without relying on inode reuse after unlink or on timing.
2788        for target_index in 0..3 {
2789            let dir = tempfile::tempdir().unwrap();
2790            let db = dir.path().join("events.db");
2791            let targets = events_db_targets(&db);
2792            for path in &targets {
2793                write_owner_only_test_file(path);
2794            }
2795            let before_open = harden_events_db_sidecars(&db).unwrap();
2796            let replacement = dir.path().join("replacement");
2797            write_owner_only_test_file(&replacement);
2798            let replacement_id = EventsFileIdentity::from_metadata(
2799                &std::fs::symlink_metadata(&replacement).unwrap(),
2800            );
2801            assert_ne!(before_open[target_index], Some(replacement_id));
2802            // Keep a real SQLite connection open across the pathname change.
2803            // Do not run SQL against deliberately replaced files afterward.
2804            let _connection = rusqlite::Connection::open(&db).unwrap();
2805            verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2806            std::fs::rename(&replacement, &targets[target_index]).unwrap();
2807            let error = verify_events_db_owner_only_unopened(&db, &before_open)
2808                .unwrap_err()
2809                .to_string();
2810            assert!(error.contains("changed identity"), "{error}");
2811            assert!(
2812                error.contains(&targets[target_index].display().to_string()),
2813                "{error}"
2814            );
2815        }
2816    }
2817
2818    #[cfg(unix)]
2819    #[test]
2820    fn unopened_check_accepts_sidecars_missing_at_check_time() {
2821        let dir = tempfile::tempdir().unwrap();
2822        let db = dir.path().join("events.db");
2823        let targets = events_db_targets(&db);
2824        for path in &targets {
2825            write_owner_only_test_file(path);
2826        }
2827        let before_open = harden_events_db_sidecars(&db).unwrap();
2828        let _connection = rusqlite::Connection::open(&db).unwrap();
2829        for path in targets.iter().skip(1) {
2830            std::fs::remove_file(path).unwrap();
2831            verify_events_db_owner_only_unopened(&db, &before_open)
2832                .expect("SQLite sidecars may disappear without losing the main database");
2833        }
2834    }
2835
2836    #[cfg(unix)]
2837    #[test]
2838    fn unopened_check_refuses_main_database_removed_after_sqlite_open() {
2839        let dir = tempfile::tempdir().unwrap();
2840        let db = dir.path().join("events.db");
2841        ensure_events_db_owner_only(&db).unwrap();
2842        let before_open = harden_events_db_sidecars(&db).unwrap();
2843        let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2844        backend.events().unwrap();
2845        verify_events_db_owner_only_unopened(&db, &before_open).unwrap();
2846        std::fs::remove_file(&db).unwrap();
2847        // Must fail with the old NotFound-for-every-target behavior. The live
2848        // connection can still name an unlinked inode; it must not be served.
2849        let error = verify_events_db_owner_only_unopened(&db, &before_open)
2850            .unwrap_err()
2851            .to_string();
2852        assert!(
2853            error.contains("main database") && error.contains("disappeared"),
2854            "{error}"
2855        );
2856    }
2857
2858    #[cfg(unix)]
2859    #[test]
2860    fn unopened_check_accepts_fresh_sqlite_sidecars_and_unchanged_main() {
2861        let dir = tempfile::tempdir().unwrap();
2862        let db = dir.path().join("new.events.db");
2863        assert!(!db.exists());
2864        ensure_events_db_owner_only(&db).unwrap();
2865        let before_open = harden_events_db_sidecars(&db).unwrap();
2866        assert!(before_open[0].is_some());
2867        assert_eq!(&before_open[1..], &[None, None]);
2868        let backend = StorageBackend::sqlite_for_test(&db).unwrap();
2869        backend.events().unwrap();
2870        for path in events_db_targets(&db).iter().skip(1) {
2871            assert!(path.exists(), "SQLite creates {}", path.display());
2872        }
2873        verify_events_db_owner_only_unopened(&db, &before_open)
2874            .expect("new SQLite sidecars have no pre-open identity to contradict");
2875    }
2876
2877    #[test]
2878    fn socket_derives_beside_the_sidecar() {
2879        let db = events_db_path_beside(Path::new("/data/khive.db"));
2880        let socket = events_socket_path_beside(&db);
2881        assert!(socket.ends_with("khive.db.events.sock"), "got {socket:?}");
2882    }
2883
2884    #[test]
2885    fn bare_relative_main_db_yields_an_absolute_sidecar() {
2886        // A bare file name has `Some("")` for a parent; every downstream
2887        // consumer (lock parenting, socket-dir trust validation) needs a
2888        // real directory, so the sidecar path must come back anchored.
2889        let path = events_db_path_beside(Path::new("khive.db"));
2890        assert!(path.is_absolute(), "got relative {path:?}");
2891        assert!(path.ends_with("khive.db.events.db"), "got {path:?}");
2892    }
2893
2894    #[cfg(unix)]
2895    #[tokio::test]
2896    async fn over_cap_query_page_limit_is_refused_before_materialization() {
2897        // `PageRequest.limit` arrives client-controlled off the wire; the
2898        // daemon must refuse an over-cap page with a typed, actionable error
2899        // instead of materializing and serializing it — and must never
2900        // silently clamp, or the short page reads as end-of-data.
2901        let backend = Arc::new(StorageBackend::memory().unwrap());
2902        let stores: NamespaceStores =
2903            Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
2904        let request = |limit: u32| EventsRequest::QueryEvents {
2905            protocol_version: EVENTS_PROTOCOL_VERSION,
2906            namespace: "local".to_string(),
2907            filter: EventFilter::default(),
2908            page: PageRequest { offset: 0, limit },
2909        };
2910        let refused =
2911            dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS + 1), &backend, &stores)
2912                .await;
2913        match refused {
2914            EventsResponse::Error {
2915                message,
2916                retryable,
2917                writer_task_failure,
2918            } => {
2919                assert!(
2920                    message.contains(&MAX_QUERY_EVENTS_PAGE_ROWS.to_string()),
2921                    "refusal must name the cap: {message}"
2922                );
2923                assert!(!retryable, "an over-cap page is not transient");
2924                assert!(writer_task_failure.is_none());
2925            }
2926            other => panic!("over-cap query must be refused, got {other:?}"),
2927        }
2928        // Control: a request AT the cap passes admission and reaches the
2929        // store — an empty page, not a refusal.
2930        let at_cap =
2931            dispatch_events_request(request(MAX_QUERY_EVENTS_PAGE_ROWS), &backend, &stores).await;
2932        assert!(
2933            matches!(at_cap, EventsResponse::Pageful { .. }),
2934            "at-cap query must reach the store, got {at_cap:?}"
2935        );
2936    }
2937
2938    #[cfg(unix)]
2939    #[test]
2940    fn side_effects_unknown_state_crosses_the_wire_as_terminal_failure() {
2941        let response = storage_error_response(&StorageError::WriterTaskTerminated {
2942            request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
2943        });
2944        assert!(
2945            matches!(
2946                response,
2947                EventsResponse::Error {
2948                    writer_task_failure: Some(WireWriterTaskFailure::TaskTerminated {
2949                        request_state: WireWriterTaskState::SideEffectsUnknown,
2950                    }),
2951                    ..
2952                }
2953            ),
2954            "an unknown-commit-state termination must carry its state on the wire"
2955        );
2956        // Any other refusal must NOT carry a writer state.
2957        let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
2958        assert!(matches!(
2959            busy,
2960            EventsResponse::Error {
2961                writer_task_failure: None,
2962                ..
2963            }
2964        ));
2965    }
2966
2967    #[cfg(unix)]
2968    #[test]
2969    fn proven_rollback_crosses_the_wire_without_claiming_task_termination() {
2970        let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
2971            request_state: khive_storage::WriterTaskRequestState::TransactionRolledBack,
2972            source: Box::new(StorageError::Pool {
2973                operation: "writer_task_commit".into(),
2974                message: "commit refused".into(),
2975            }),
2976        });
2977        assert!(matches!(
2978            response,
2979            EventsResponse::Error {
2980                writer_task_failure: Some(WireWriterTaskFailure::RequestFailed {
2981                    request_state: WireWriterTaskState::TransactionRolledBack,
2982                }),
2983                ..
2984            }
2985        ));
2986    }
2987
2988    #[test]
2989    fn error_frames_without_writer_failure_still_parse() {
2990        // `writer_task_failure` is `#[serde(default)]`: an Error frame for a
2991        // non-writer failure (parse error, version refusal) omits it and the
2992        // client must read `None`, not fail the parse.
2993        let bytes = br#"{"kind":"error","message":"boom","retryable":true}"#;
2994        let parsed: EventsResponse = serde_json::from_slice(bytes).expect("stateless frame parses");
2995        assert!(matches!(
2996            parsed,
2997            EventsResponse::Error {
2998                retryable: true,
2999                writer_task_failure: None,
3000                ..
3001            }
3002        ));
3003    }
3004
3005    fn split_retry_event(verb: &str) -> Event {
3006        Event::new(
3007            "test",
3008            verb,
3009            khive_types::EventKind::RecallExecuted,
3010            khive_types::SubstrateKind::Note,
3011            "agent:test",
3012        )
3013    }
3014
3015    fn store_pair(dir: &Path) -> (Arc<dyn EventStore>, Arc<dyn EventStore>) {
3016        let legacy = direct_backend_for(&dir.join("legacy.db"))
3017            .expect("legacy backend")
3018            .events_for_namespace("test")
3019            .expect("legacy store");
3020        let lane = direct_backend_for(&dir.join("lane.db"))
3021            .expect("lane backend")
3022            .events_for_namespace("test")
3023            .expect("lane store");
3024        (legacy, lane)
3025    }
3026
3027    /// A retried audit batch whose first attempt landed on the legacy store
3028    /// before the cutover must be answered by the legacy store's comparison,
3029    /// never re-inserted into the lane — otherwise merged reads and counts
3030    /// double-count the id.
3031    #[tokio::test]
3032    async fn idempotent_retry_of_legacy_resident_rows_does_not_duplicate() {
3033        use khive_storage::event::EventAppendDisposition;
3034
3035        let dir = tempfile::tempdir().unwrap();
3036        let _registry_guard = TestRegistryGuard::new(dir.path());
3037        let (legacy, lane) = store_pair(dir.path());
3038
3039        let resident = split_retry_event("recall");
3040        legacy
3041            .append_events_idempotent(vec![resident.clone()])
3042            .await
3043            .expect("pre-cutover landing");
3044
3045        let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3046        let fresh = split_retry_event("search");
3047        let result = split
3048            .append_events_idempotent(vec![resident.clone(), fresh.clone()])
3049            .await
3050            .expect("mixed retry batch");
3051        assert_eq!(
3052            result.rows,
3053            vec![
3054                EventAppendDisposition::AlreadyPresentIdentical,
3055                EventAppendDisposition::Inserted,
3056            ],
3057            "input order must be preserved across the two sub-batches"
3058        );
3059        assert_eq!(
3060            lane.count_events(EventFilter::default()).await.unwrap(),
3061            1,
3062            "the legacy-resident row must not reach the lane"
3063        );
3064        assert_eq!(
3065            split.count_events(EventFilter::default()).await.unwrap(),
3066            2,
3067            "merged count must not double-count the retried id"
3068        );
3069
3070        // A conflicting retry of the resident row reports the conflict from
3071        // the legacy comparison and inserts nowhere.
3072        let mut mutated = resident.clone();
3073        mutated.verb = "other".to_string();
3074        let conflict = split
3075            .append_events_idempotent(vec![mutated])
3076            .await
3077            .expect("conflicting retry");
3078        assert_eq!(
3079            conflict.rows,
3080            vec![EventAppendDisposition::IdentityConflict]
3081        );
3082        assert_eq!(split.count_events(EventFilter::default()).await.unwrap(), 2);
3083    }
3084
3085    /// The merged plane refuses to materialize an unbounded offset prefix;
3086    /// the refusal names the cursor remedy, and the bound itself is usable.
3087    #[tokio::test]
3088    async fn merged_offset_window_is_bounded() {
3089        let dir = tempfile::tempdir().unwrap();
3090        let _registry_guard = TestRegistryGuard::new(dir.path());
3091        let (legacy, lane) = store_pair(dir.path());
3092        let split = SplitEventStore::new(legacy, lane);
3093
3094        let err = split
3095            .query_events(
3096                EventFilter::default(),
3097                PageRequest {
3098                    offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS,
3099                    limit: 1,
3100                },
3101            )
3102            .await
3103            .expect_err("a window past the bound must be refused");
3104        assert!(
3105            matches!(err, StorageError::InvalidInput { .. }),
3106            "got {err:?}"
3107        );
3108        assert!(
3109            err.to_string().contains("before"),
3110            "the refusal must name the cursor remedy: {err}"
3111        );
3112
3113        let at_bound = split
3114            .query_events(
3115                EventFilter::default(),
3116                PageRequest {
3117                    offset: SplitEventStore::MAX_MERGED_WINDOW_ROWS - 1,
3118                    limit: 1,
3119                },
3120            )
3121            .await
3122            .expect("a window at the bound is admitted");
3123        assert!(at_bound.items.is_empty());
3124    }
3125
3126    #[cfg(unix)]
3127    #[tokio::test]
3128    async fn stateless_non_retryable_error_maps_terminal_not_writer_terminated() {
3129        // The terminal arm of the client mapping, pinned deliberately: an
3130        // Error frame with `retryable: false` and NO `writer_task_failure` is
3131        // a non-writer failure (parse error, version refusal) and must map
3132        // to terminal `InvalidInput` — never to `WriterTaskTerminated`
3133        // (there is no state to reconstruct) and never to the retryable
3134        // `Pool` arm (retrying a version refusal cannot succeed).
3135        let client = EventsSplitClient::new(std::path::PathBuf::from("/tmp/never-bound.sock"))
3136            .expect("client builds");
3137        let store = ForwardingEventStore::new("test", client);
3138        let bytes = br#"{"kind":"error","message":"refused","retryable":false}"#;
3139        let parsed: EventsResponse = serde_json::from_slice(bytes).expect("frame parses");
3140        assert!(matches!(
3141            store.unexpected("append", parsed),
3142            StorageError::InvalidInput { .. }
3143        ));
3144        // Control: the same non-retryable frame WITH a writer state must
3145        // take the reconstruction arm instead — proving the terminal arm
3146        // above is selected by the field's absence, not by `retryable`.
3147        let with_state = br#"{"kind":"error","message":"died","retryable":false,"writer_task_failure":{"kind":"task_terminated","request_state":"side_effects_unknown"}}"#;
3148        let parsed: EventsResponse =
3149            serde_json::from_slice(with_state).expect("stateful frame parses");
3150        assert!(matches!(
3151            store.unexpected("append", parsed),
3152            StorageError::WriterTaskTerminated { .. }
3153        ));
3154    }
3155
3156    #[cfg(unix)]
3157    #[tokio::test]
3158    async fn writer_task_states_cross_the_wire_verbatim() {
3159        // The ADR-133 retry classifier decides per-state, so every
3160        // WriterTaskTerminated state must survive server → wire → client
3161        // reconstruction exactly. Before this field existed, NotStarted and
3162        // TransactionRolledBack crossed as a generic non-retryable refusal
3163        // and came back as terminal InvalidInput — a broken retry contract.
3164        use khive_storage::WriterTaskRequestState as S;
3165        let dir = tempfile::tempdir().unwrap();
3166        let client =
3167            EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3168        let store = ForwardingEventStore::new("test", client);
3169        for state in [
3170            S::NotStarted,
3171            S::TransactionRolledBack,
3172            S::SideEffectsUnknown,
3173        ] {
3174            let response = storage_error_response(&StorageError::WriterTaskTerminated {
3175                request_state: state,
3176            });
3177            // Round-trip through serde like the socket does.
3178            let bytes = serde_json::to_vec(&response).unwrap();
3179            let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3180            let err = store.unexpected("append_events_idempotent", parsed);
3181            assert!(
3182                matches!(
3183                    err,
3184                    StorageError::WriterTaskTerminated { request_state } if request_state == state
3185                ),
3186                "state {state:?} did not survive the socket: got {err:?}"
3187            );
3188        }
3189        // Control: a non-writer-task refusal must NOT carry the field.
3190        let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3191        assert!(matches!(
3192            busy,
3193            EventsResponse::Error {
3194                writer_task_failure: None,
3195                ..
3196            }
3197        ));
3198    }
3199
3200    #[cfg(unix)]
3201    #[tokio::test]
3202    async fn proven_rollback_round_trips_as_non_terminal_request_failure() {
3203        use khive_storage::WriterTaskRequestState as S;
3204
3205        let dir = tempfile::tempdir().unwrap();
3206        let client =
3207            EventsSplitClient::new(dir.path().join("never-bound.sock")).expect("client builds");
3208        let store = ForwardingEventStore::new("test", client);
3209        let response = storage_error_response(&StorageError::WriterTaskRequestFailed {
3210            request_state: S::TransactionRolledBack,
3211            source: Box::new(StorageError::Pool {
3212                operation: "writer_task_commit".into(),
3213                message: "commit refused".into(),
3214            }),
3215        });
3216        let bytes = serde_json::to_vec(&response).unwrap();
3217        let parsed: EventsResponse = serde_json::from_slice(&bytes).unwrap();
3218        let err = store.unexpected("append_events_idempotent", parsed);
3219        assert!(
3220            matches!(
3221                err,
3222                StorageError::WriterTaskRequestFailed {
3223                    request_state: S::TransactionRolledBack,
3224                    ..
3225                }
3226            ),
3227            "proven rollback must not reconstruct as a terminal writer: {err:?}"
3228        );
3229    }
3230
3231    #[cfg(unix)]
3232    #[test]
3233    fn direct_backend_hardens_preexisting_db_and_sidecars() {
3234        use std::os::unix::fs::PermissionsExt;
3235        let dir = tempfile::tempdir().unwrap();
3236        let _registry_guard = TestRegistryGuard::new(dir.path());
3237        let db = dir.path().join("pre-existing.events.db");
3238        let wal = dir.path().join("pre-existing.events.db-wal");
3239        std::fs::write(&db, b"").unwrap();
3240        std::fs::write(&wal, b"").unwrap();
3241        for path in [&db, &wal] {
3242            std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o644)).unwrap();
3243        }
3244        // Control: the loose mode is really in place before the open.
3245        assert_eq!(
3246            std::fs::metadata(&db).unwrap().permissions().mode() & 0o777,
3247            0o644
3248        );
3249        direct_backend_for(&db).expect("writable open succeeds");
3250        for path in [&db, &wal] {
3251            assert_eq!(
3252                std::fs::metadata(path).unwrap().permissions().mode() & 0o777,
3253                0o600,
3254                "pre-existing {} must be tightened to owner-only",
3255                path.display()
3256            );
3257        }
3258    }
3259
3260    #[cfg(unix)]
3261    #[test]
3262    fn namespace_store_cache_is_bounded_and_trim_normalized() {
3263        let backend = StorageBackend::memory().unwrap();
3264        let stores: NamespaceStores =
3265            Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3266        // cap=1: the first namespace caches; the second evicts it and takes
3267        // the slot — bounded, and never condemned to per-request rebuilds.
3268        namespace_store_with_cap(&backend, &stores, "alpha", 1).expect("first store");
3269        namespace_store_with_cap(&backend, &stores, "beta", 1).expect("second store evicts");
3270        {
3271            let map = stores.lock().unwrap();
3272            assert_eq!(map.len(), 1, "cache must not grow past its cap");
3273            assert!(
3274                map.contains_key("beta"),
3275                "the newest namespace must be admitted at the cap"
3276            );
3277        }
3278        // Trimmed spelling of a cached namespace hits the same entry rather
3279        // than minting a duplicate handle (the backend trims namespaces).
3280        namespace_store_with_cap(&backend, &stores, "beta", 1).expect("cached beta");
3281        namespace_store_with_cap(&backend, &stores, "  beta  ", 1).expect("trimmed spelling");
3282        {
3283            let map = stores.lock().unwrap();
3284            assert_eq!(
3285                map.len(),
3286                1,
3287                "spellings of one namespace must share one entry"
3288            );
3289            assert!(map.contains_key("beta"), "trimmed key is the cache key");
3290        }
3291    }
3292
3293    #[cfg(unix)]
3294    #[tokio::test]
3295    async fn frame_budget_admits_before_allocating_and_releases_after() {
3296        let (mut client, mut server) = UnixStream::pair().expect("socketpair");
3297        let budget = Arc::new(tokio::sync::Semaphore::new(1024));
3298        let payload = vec![7u8; 100];
3299        write_frame(&mut client, &payload).await.expect("write");
3300        let (bytes, permit) = read_frame_budgeted(&mut server, &budget)
3301            .await
3302            .expect("budgeted read");
3303        assert_eq!(bytes, payload);
3304        // The permit holds exactly the frame's bytes against the budget…
3305        assert_eq!(budget.available_permits(), 1024 - 100);
3306        // …and releases them when the request lifecycle ends.
3307        drop(permit);
3308        assert_eq!(budget.available_permits(), 1024);
3309        // A frame past the per-frame cap is refused before admission: the
3310        // budget is untouched, so an oversized declaration cannot starve it.
3311        let mut oversized = Vec::from(((crate::daemon::MAX_FRAME_BYTES + 1) as u32).to_be_bytes());
3312        oversized.extend_from_slice(&[0u8; 8]);
3313        use tokio::io::AsyncWriteExt;
3314        client.write_all(&oversized).await.expect("raw prefix");
3315        let err = read_frame_budgeted(&mut server, &budget)
3316            .await
3317            .expect_err("oversized declaration must refuse");
3318        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
3319        assert_eq!(
3320            budget.available_permits(),
3321            1024,
3322            "refusal must not consume budget"
3323        );
3324    }
3325
3326    #[cfg(unix)]
3327    #[tokio::test]
3328    async fn dispatch_rejects_invalid_wire_namespace() {
3329        let backend = Arc::new(StorageBackend::memory().unwrap());
3330        let stores: NamespaceStores =
3331            Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3332        // An attacker-shaped namespace: far past the 256-byte Namespace
3333        // bound. Must be refused as a typed non-retryable error before it
3334        // can become a cache key or a store.
3335        let huge = "n".repeat(64 * 1024);
3336        let response = dispatch_events_request(
3337            EventsRequest::CountEvents {
3338                protocol_version: EVENTS_PROTOCOL_VERSION,
3339                namespace: huge,
3340                filter: Default::default(),
3341            },
3342            &backend,
3343            &stores,
3344        )
3345        .await;
3346        assert!(
3347            matches!(
3348                &response,
3349                EventsResponse::Error {
3350                    retryable: false,
3351                    writer_task_failure: None,
3352                    ..
3353                }
3354            ),
3355            "oversized namespace must be a typed refusal, got {response:?}"
3356        );
3357        assert_eq!(
3358            stores.lock().unwrap().len(),
3359            0,
3360            "a rejected namespace must never enter the cache"
3361        );
3362        // Control: a valid namespace on the same dispatch path succeeds.
3363        let ok = dispatch_events_request(
3364            EventsRequest::CountEvents {
3365                protocol_version: EVENTS_PROTOCOL_VERSION,
3366                namespace: "local".to_string(),
3367                filter: Default::default(),
3368            },
3369            &backend,
3370            &stores,
3371        )
3372        .await;
3373        assert!(
3374            matches!(ok, EventsResponse::Count { .. }),
3375            "valid namespace must dispatch, got {ok:?}"
3376        );
3377    }
3378
3379    include!("events_split_shutdown_tests.rs");
3380
3381    #[cfg(unix)]
3382    #[tokio::test]
3383    async fn forwarder_abandons_hung_delivery() {
3384        // A daemon that accepts the connection and then never responds must
3385        // not park the forwarder forever: the delivery deadline fires, the
3386        // batch is dropped and counted, and the task stays live.
3387        let dir = tempfile::tempdir().unwrap();
3388        let socket = dir.path().join("hung.sock");
3389        let listener = tokio::net::UnixListener::bind(&socket).unwrap();
3390        // Accept and hold connections open without ever reading or replying.
3391        let _server = tokio::spawn(async move {
3392            let mut held = Vec::new();
3393            loop {
3394                if let Ok((stream, _)) = listener.accept().await {
3395                    held.push(stream);
3396                }
3397            }
3398        });
3399        let client = EventsSplitClient::new_with_queue_depth_and_delivery_timeout(
3400            socket.clone(),
3401            4,
3402            Duration::from_millis(100),
3403        )
3404        .expect("client builds");
3405        let event = Event::new(
3406            "test",
3407            "noop",
3408            khive_types::EventKind::Audit,
3409            khive_types::SubstrateKind::Event,
3410            "tester",
3411        );
3412        client.enqueue("test", vec![event]);
3413        let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
3414        loop {
3415            if client.metrics().dropped_batches >= 1 {
3416                break;
3417            }
3418            assert!(
3419                tokio::time::Instant::now() < deadline,
3420                "forwarder never abandoned the hung delivery: {:?}",
3421                client.metrics()
3422            );
3423            tokio::time::sleep(Duration::from_millis(20)).await;
3424        }
3425    }
3426
3427    #[cfg(unix)]
3428    #[test]
3429    fn wire_retryability_follows_the_storage_classifier() {
3430        let busy = storage_error_response(&StorageError::WriterTaskBusy { timeout_ms: 5 });
3431        assert!(
3432            matches!(
3433                busy,
3434                EventsResponse::Error {
3435                    retryable: true,
3436                    ..
3437                }
3438            ),
3439            "transient writer contention must stay retryable across the socket"
3440        );
3441        // A terminated writer task is not retryable per the storage layer's
3442        // own classifier; the wire response must agree rather than widen it.
3443        let terminated = storage_error_response(&StorageError::WriterTaskTerminated {
3444            request_state: khive_storage::WriterTaskRequestState::SideEffectsUnknown,
3445        });
3446        assert!(
3447            matches!(
3448                terminated,
3449                EventsResponse::Error {
3450                    retryable: false,
3451                    ..
3452                }
3453            ),
3454            "a terminated writer must not be reported transient"
3455        );
3456    }
3457
3458    use khive_storage::event::EventAppendDisposition;
3459    use khive_types::{EventKind, SubstrateKind};
3460
3461    fn test_event(namespace: &str) -> Event {
3462        Event::new(
3463            namespace,
3464            "test.verb",
3465            EventKind::Audit,
3466            SubstrateKind::Event,
3467            "actor:test",
3468        )
3469    }
3470
3471    /// Boot a real events daemon on temp paths, poll until its socket accepts.
3472    #[cfg(unix)]
3473    async fn boot_daemon(dir: &tempfile::TempDir) -> (PathBuf, PathBuf) {
3474        let db = dir.path().join("events.db");
3475        let socket = dir.path().join("events.sock");
3476        let (db_clone, socket_clone) = (db.clone(), socket.clone());
3477        tokio::spawn(async move {
3478            let _ = run_events_daemon(&db_clone, &socket_clone).await;
3479        });
3480        for _ in 0..100 {
3481            if UnixStream::connect(&socket).await.is_ok() {
3482                return (db, socket);
3483            }
3484            tokio::time::sleep(Duration::from_millis(20)).await;
3485        }
3486        panic!("events daemon did not come up on {}", socket.display());
3487    }
3488
3489    /// The split store's routing contract: plain appends land ONLY in the
3490    /// legacy store (whose raw-SQL consumers depend on finding them there),
3491    /// the idempotent audit lane lands ONLY in the events lane, and reads
3492    /// merge both sides into one event plane.
3493    #[tokio::test]
3494    async fn split_store_routes_plain_to_legacy_idempotent_to_lane_and_merges_reads() {
3495        let dir = tempfile::tempdir().expect("tempdir");
3496        let _registry_guard = TestRegistryGuard::new(dir.path());
3497        let legacy_backend =
3498            direct_backend_for(&dir.path().join("legacy.db")).expect("legacy backend");
3499        let lane_backend = direct_backend_for(&dir.path().join("lane.db")).expect("lane backend");
3500        let legacy = legacy_backend
3501            .events_for_namespace("local")
3502            .expect("legacy store");
3503        let lane = lane_backend
3504            .events_for_namespace("local")
3505            .expect("lane store");
3506        let split = SplitEventStore::new(Arc::clone(&legacy), Arc::clone(&lane));
3507
3508        let plain = test_event("local");
3509        let plain_id = plain.id;
3510        split.append_event(plain).await.expect("plain append");
3511
3512        let audit = test_event("local");
3513        let audit_id = audit.id;
3514        let result = split
3515            .append_events_idempotent(vec![audit])
3516            .await
3517            .expect("idempotent append");
3518        assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3519
3520        // Routing: each row is in exactly its own side.
3521        assert!(legacy
3522            .get_event(plain_id)
3523            .await
3524            .expect("legacy get")
3525            .is_some());
3526        assert!(lane.get_event(plain_id).await.expect("lane get").is_none());
3527        assert!(legacy
3528            .get_event(audit_id)
3529            .await
3530            .expect("legacy get")
3531            .is_none());
3532        assert!(lane.get_event(audit_id).await.expect("lane get").is_some());
3533
3534        // Merged reads observe one event plane whichever side holds the row.
3535        assert!(split
3536            .get_event(plain_id)
3537            .await
3538            .expect("split get")
3539            .is_some());
3540        assert!(split
3541            .get_event(audit_id)
3542            .await
3543            .expect("split get")
3544            .is_some());
3545        assert_eq!(
3546            split
3547                .count_events(EventFilter::default())
3548                .await
3549                .expect("split count"),
3550            2
3551        );
3552        let page = split
3553            .query_events(
3554                EventFilter::default(),
3555                PageRequest {
3556                    offset: 0,
3557                    limit: 10,
3558                },
3559            )
3560            .await
3561            .expect("split query");
3562        let ids: Vec<Uuid> = page.items.iter().map(|e| e.id).collect();
3563        assert!(ids.contains(&plain_id) && ids.contains(&audit_id));
3564
3565        // Windowing across the merge: page size 1 at offsets 0 and 1 yields
3566        // the two rows exactly once each.
3567        let mut seen = Vec::new();
3568        for offset in 0..2 {
3569            let page = split
3570                .query_events(EventFilter::default(), PageRequest { offset, limit: 1 })
3571                .await
3572                .expect("windowed query");
3573            assert_eq!(page.items.len(), 1);
3574            seen.push(page.items[0].id);
3575        }
3576        seen.sort();
3577        let mut expected = vec![plain_id, audit_id];
3578        expected.sort();
3579        assert_eq!(seen, expected);
3580    }
3581
3582    /// Mechanized co-residency guard: the amendment's load-bearing claim is
3583    /// that every class a raw-SQL consumer of the legacy `events` table
3584    /// depends on still LANDS in that table. Two of the three enumerated
3585    /// consumers (kg projection guard, graph-query union) fail silently if
3586    /// the classification drifts, so the classification is guarded here by
3587    /// the consumers' own access pattern — a raw SQL read against the legacy
3588    /// backend — rather than by enumeration. Arm 1 reddens if plain appends
3589    /// (the provenance/domain class) ever stop reaching the legacy table.
3590    /// Arm 2 is the boundary's other face and doubles as the positive
3591    /// control for the instrument: the same query DOES find lane-routed rows
3592    /// on the lane's own backend, so an empty arm-1 result could not be a
3593    /// broken query.
3594    #[tokio::test]
3595    async fn raw_sql_consumers_of_the_legacy_events_table_still_see_plain_appends() {
3596        use khive_storage::types::{SqlStatement, SqlValue};
3597
3598        let dir = tempfile::tempdir().expect("tempdir");
3599        let _registry_guard = TestRegistryGuard::new(dir.path());
3600        let legacy_backend =
3601            direct_backend_for(&dir.path().join("legacy-guard.db")).expect("legacy backend");
3602        let lane_backend =
3603            direct_backend_for(&dir.path().join("lane-guard.db")).expect("lane backend");
3604        let split = SplitEventStore::new(
3605            legacy_backend
3606                .events_for_namespace("local")
3607                .expect("legacy store"),
3608            lane_backend
3609                .events_for_namespace("local")
3610                .expect("lane store"),
3611        );
3612
3613        // The schedule pack's creator-provenance write is a plain append; the
3614        // audit flusher's write is an idempotent batch. Route one of each.
3615        let provenance = test_event("local");
3616        let provenance_id = provenance.id;
3617        split.append_event(provenance).await.expect("plain append");
3618        let audit = test_event("local");
3619        let audit_id = audit.id;
3620        split
3621            .append_events_idempotent(vec![audit])
3622            .await
3623            .expect("idempotent append");
3624
3625        let count_by_raw_sql = |backend: Arc<StorageBackend>, id: Uuid| async move {
3626            let mut reader = backend.sql().reader().await.expect("sql reader");
3627            let rows = reader
3628                .query_all(SqlStatement {
3629                    sql: "SELECT actor FROM events WHERE id = ?1".to_string(),
3630                    params: vec![SqlValue::Text(id.to_string())],
3631                    label: None,
3632                })
3633                .await
3634                .expect("raw events query");
3635            rows.len()
3636        };
3637
3638        // Arm 1: the co-residency contract. A raw-SQL consumer of the legacy
3639        // table finds the plain-append row there.
3640        assert_eq!(
3641            count_by_raw_sql(Arc::clone(&legacy_backend), provenance_id).await,
3642            1,
3643            "plain appends must stay visible to raw-SQL consumers of the legacy events table"
3644        );
3645        // Arm 2a: the moved class is genuinely gone from the legacy table —
3646        // this is the silent-failure face of the boundary, held visible.
3647        assert_eq!(
3648            count_by_raw_sql(Arc::clone(&legacy_backend), audit_id).await,
3649            0,
3650            "audit-lane rows must not land in the legacy events table"
3651        );
3652        // Arm 2b: positive control for the instrument — the identical query
3653        // finds the moved row on the lane backend, so arm 2a's zero (and any
3654        // future arm-1 zero) is a routing fact, not a dead query.
3655        assert_eq!(
3656            count_by_raw_sql(Arc::clone(&lane_backend), audit_id).await,
3657            1,
3658            "the raw query must prove it can find rows where they actually live"
3659        );
3660    }
3661
3662    /// End to end over the real socket: the idempotent lane returns true
3663    /// dispositions, and the row is readable back through the same store.
3664    #[cfg(unix)]
3665    #[tokio::test]
3666    async fn idempotent_append_round_trips_through_the_daemon() {
3667        let dir = tempfile::tempdir().expect("tempdir");
3668        let (_db, socket) = boot_daemon(&dir).await;
3669        let client = EventsSplitClient::new(socket).expect("client");
3670        let store = ForwardingEventStore::new("local", client);
3671
3672        let event = test_event("local");
3673        let id = event.id;
3674        let result = store
3675            .append_events_idempotent(vec![event.clone()])
3676            .await
3677            .expect("idempotent append over socket");
3678        assert_eq!(result.rows, vec![EventAppendDisposition::Inserted]);
3679
3680        // Retry of the identical row reports AlreadyPresentIdentical — the
3681        // daemon ran the real idempotent path, not a blind re-insert.
3682        let retry = store
3683            .append_events_idempotent(vec![event])
3684            .await
3685            .expect("idempotent retry");
3686        assert_eq!(
3687            retry.rows,
3688            vec![EventAppendDisposition::AlreadyPresentIdentical]
3689        );
3690
3691        let fetched = store.get_event(id).await.expect("get over socket");
3692        assert_eq!(fetched.map(|e| e.id), Some(id));
3693        let count = store
3694            .count_events(EventFilter::default())
3695            .await
3696            .expect("count over socket");
3697        assert_eq!(count, 1);
3698    }
3699
3700    /// The fire-and-forget lane delivers to the daemon store: enqueue returns
3701    /// immediately and the row becomes visible to a subsequent count.
3702    #[cfg(unix)]
3703    #[tokio::test]
3704    async fn fire_and_forget_append_lands_in_the_daemon_store() {
3705        let dir = tempfile::tempdir().expect("tempdir");
3706        let (_db, socket) = boot_daemon(&dir).await;
3707        let client = EventsSplitClient::new(socket).expect("client");
3708        let store = ForwardingEventStore::new("local", Arc::clone(&client));
3709
3710        store
3711            .append_event(test_event("local"))
3712            .await
3713            .expect("append_event is fire-and-forget");
3714
3715        // Poll for BOTH effects: the row landing daemon-side and the
3716        // forwarder's own counter — the row can be visible a beat before the
3717        // forwarder task is rescheduled to record the delivery.
3718        let (count, metrics) = tokio::time::timeout(Duration::from_secs(30), async {
3719            loop {
3720                let count = store
3721                    .count_events(EventFilter::default())
3722                    .await
3723                    .expect("count over socket");
3724                let metrics = client.metrics();
3725                if (count == 1 && metrics.forwarded_events >= 1) || metrics.dropped_events > 0 {
3726                    break (count, metrics);
3727                }
3728                tokio::time::sleep(Duration::from_millis(20)).await;
3729            }
3730        })
3731        .await
3732        .expect("forwarder must deliver or report a drop within 30 seconds");
3733        assert_eq!(
3734            metrics.dropped_events, 0,
3735            "forwarder dropped an event before it reached the daemon store"
3736        );
3737        assert_eq!(count, 1, "forwarded event must land in the daemon store");
3738        assert!(metrics.forwarded_events >= 1, "delivery must be counted");
3739    }
3740
3741    /// R4 overflow arm: with the daemon dead, appends beyond the queue bound
3742    /// are DROPPED and counted while every append still returns Ok — the
3743    /// domain path proceeds. The drop counter is the loss-tolerance contract.
3744    #[cfg(unix)]
3745    #[tokio::test]
3746    async fn queue_overflow_drops_and_counts_but_never_errors() {
3747        let dir = tempfile::tempdir().expect("tempdir");
3748        let dead_socket = dir.path().join("nobody-home.sock");
3749        let client = EventsSplitClient::new_with_queue_depth(dead_socket, 2).expect("client");
3750        let store = ForwardingEventStore::new("local", Arc::clone(&client));
3751
3752        for _ in 0..20 {
3753            store
3754                .append_event(test_event("local"))
3755                .await
3756                .expect("append_event must not error under overflow");
3757        }
3758        let metrics = client.metrics();
3759        assert!(
3760            metrics.dropped_events > 0,
3761            "overflow must register in the drop counter, got {metrics:?}"
3762        );
3763    }
3764
3765    /// R4 dead-socket arm: synchronous lanes fail with a typed retryable
3766    /// storage error (never hang, never a generic panic), and the offline
3767    /// preflight validator keeps working with zero I/O.
3768    #[cfg(unix)]
3769    #[tokio::test]
3770    async fn dead_socket_reads_fail_typed_and_preflight_stays_local() {
3771        let dir = tempfile::tempdir().expect("tempdir");
3772        let dead_socket = dir.path().join("nobody-home.sock");
3773        let client = EventsSplitClient::new(dead_socket).expect("client");
3774        let store = ForwardingEventStore::new("local", client);
3775
3776        let error = store
3777            .get_event(Uuid::new_v4())
3778            .await
3779            .expect_err("read against a dead socket must fail");
3780        assert!(
3781            matches!(error, StorageError::Pool { .. }),
3782            "expected the typed unreachable error, got {error:?}"
3783        );
3784
3785        // The audit-batch seam contract: preflight is local and functional
3786        // with the daemon down.
3787        assert!(store.supports_idempotent_audit_batch());
3788        store
3789            .preflight_event(&test_event("local"))
3790            .expect("offline preflight validates a well-formed event");
3791    }
3792
3793    /// A connected peer that consumed the request but sent no response is a
3794    /// broken response, not evidence that no daemon answered the socket.
3795    #[cfg(unix)]
3796    #[tokio::test]
3797    async fn connected_events_peer_closing_without_response_is_not_unreachable() {
3798        let dir = tempfile::tempdir().unwrap();
3799        let socket = dir.path().join("closes-after-request.sock");
3800        let listener = UnixListener::bind(&socket).unwrap();
3801        let server = tokio::spawn(async move {
3802            let (mut stream, _) = listener.accept().await.unwrap();
3803            read_frame(&mut stream)
3804                .await
3805                .expect("complete request arrives");
3806            // Drop the connection without writing a response frame.
3807        });
3808        let client = EventsSplitClient::new(socket).unwrap();
3809        let store = ForwardingEventStore::new("local", client);
3810        let error = store.get_event(Uuid::new_v4()).await.unwrap_err();
3811        assert!(
3812            matches!(error, StorageError::Serialization { .. }),
3813            "a post-connect response failure must not be Pool/unreachable: {error:?}"
3814        );
3815        server.await.unwrap();
3816    }
3817
3818    /// 4096 ordinary 3-KiB events obey the row limit yet exceed one IPC
3819    /// frame. The daemon must answer with a non-retryable size refusal so a
3820    /// caller can use a narrower page on the same healthy socket.
3821    #[cfg(unix)]
3822    #[tokio::test]
3823    async fn query_page_over_frame_cap_returns_typed_size_refusal() {
3824        let dir = tempfile::tempdir().unwrap();
3825        let socket = dir.path().join("large-query.sock");
3826        let backend = Arc::new(StorageBackend::memory().unwrap());
3827        let event_store = backend.events_for_namespace("local").unwrap();
3828        let events: Vec<_> = (0..MAX_QUERY_EVENTS_PAGE_ROWS)
3829            .map(|_| {
3830                test_event("local").with_payload(serde_json::json!({"data": "x".repeat(3 * 1024)}))
3831            })
3832            .collect();
3833        event_store
3834            .append_events(events)
3835            .await
3836            .expect("seed large page");
3837
3838        let listener = UnixListener::bind(&socket).unwrap();
3839        let stores: NamespaceStores =
3840            Arc::new(std::sync::Mutex::new(std::collections::HashMap::new()));
3841        let server = tokio::spawn(async move {
3842            let (stream, _) = listener.accept().await.unwrap();
3843            serve_events_conn(
3844                stream,
3845                backend,
3846                stores,
3847                Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_REQUEST_BYTES)),
3848            )
3849            .await;
3850        });
3851        let client = EventsSplitClient::new(socket).unwrap();
3852        let store = ForwardingEventStore::new("local", client);
3853        let error = store
3854            .query_events(
3855                EventFilter::default(),
3856                PageRequest {
3857                    offset: 0,
3858                    limit: MAX_QUERY_EVENTS_PAGE_ROWS,
3859                },
3860            )
3861            .await
3862            .expect_err("the full page cannot fit one frame");
3863        assert!(
3864            matches!(error, StorageError::InvalidInput { .. }),
3865            "expected non-retryable frame-size refusal, got {error:?}"
3866        );
3867        assert!(
3868            error.to_string().contains("response_frame_size_limit"),
3869            "{error}"
3870        );
3871        assert!(
3872            error.to_string().contains("request a narrower page"),
3873            "{error}"
3874        );
3875        server.abort();
3876    }
3877
3878    /// The byte budget counts the in-flight batch as well as queued batches.
3879    /// A stalled daemon therefore cannot retain another large batch just
3880    /// because the mpsc channel still has spare batch slots.
3881    #[cfg(unix)]
3882    #[tokio::test]
3883    async fn forwarding_queue_enforces_serialized_byte_and_frame_limits_with_drop_metrics() {
3884        let dir = tempfile::tempdir().unwrap();
3885        let socket = dir.path().join("stalled-forwarder.sock");
3886        let listener = UnixListener::bind(&socket).unwrap();
3887        let server = tokio::spawn(async move {
3888            let (stream, _) = listener.accept().await.unwrap();
3889            let _hold_open = stream;
3890            std::future::pending::<()>().await;
3891        });
3892        let event = test_event("local").with_payload(serde_json::json!({"data": "x".repeat(4096)}));
3893        let first = vec![event.clone(), event.clone()];
3894        let request_bytes = serde_json::to_vec(&EventsRequest::AppendEvents {
3895            protocol_version: EVENTS_PROTOCOL_VERSION,
3896            namespace: "local".into(),
3897            events: first.clone(),
3898        })
3899        .unwrap()
3900        .len();
3901        let byte_budget = request_bytes + 64;
3902        let client = EventsSplitClient::new_with_limits_and_delivery_timeout(
3903            socket,
3904            8,
3905            byte_budget,
3906            Duration::from_secs(30),
3907        )
3908        .unwrap();
3909        let store = ForwardingEventStore::new("local", Arc::clone(&client));
3910        store.append_events(first.clone()).await.unwrap();
3911        assert_eq!(client.metrics().queued_bytes, request_bytes);
3912
3913        store.append_events(first).await.unwrap();
3914        let metrics = client.metrics();
3915        assert_eq!(metrics.dropped_batches, 1);
3916        assert_eq!(metrics.dropped_events, 2);
3917        assert_eq!(metrics.queued_bytes, request_bytes);
3918        assert!(metrics.queued_bytes <= byte_budget);
3919
3920        // One append request larger than the daemon frame cannot ever be
3921        // delivered; it must be refused before occupying a queue slot.
3922        let oversized = test_event("local")
3923            .with_payload(serde_json::json!({"data": "x".repeat(crate::daemon::MAX_FRAME_BYTES)}));
3924        store.append_event(oversized).await.unwrap();
3925        let metrics = client.metrics();
3926        assert_eq!(metrics.dropped_batches, 2);
3927        assert_eq!(metrics.dropped_events, 3);
3928        assert_eq!(metrics.queued_bytes, request_bytes);
3929        server.abort();
3930    }
3931
3932    /// Version-skew arm: a frame carrying an unknown protocol version gets a
3933    /// typed non-retryable refusal, not a deserialization failure.
3934    #[cfg(unix)]
3935    #[tokio::test]
3936    async fn protocol_version_skew_is_a_typed_refusal() {
3937        let dir = tempfile::tempdir().expect("tempdir");
3938        let (_db, socket) = boot_daemon(&dir).await;
3939
3940        let mut stream = UnixStream::connect(&socket).await.expect("connect");
3941        let request = EventsRequest::CountEvents {
3942            protocol_version: EVENTS_PROTOCOL_VERSION + 1,
3943            namespace: "local".into(),
3944            filter: EventFilter::default(),
3945        };
3946        let payload = serde_json::to_vec(&request).expect("serialize");
3947        write_frame(&mut stream, &payload).await.expect("write");
3948        let bytes = read_frame(&mut stream).await.expect("read");
3949        let response: EventsResponse = serde_json::from_slice(&bytes).expect("parse");
3950        match response {
3951            EventsResponse::Error {
3952                message, retryable, ..
3953            } => {
3954                assert!(!retryable, "version skew is not retryable");
3955                assert!(message.contains("protocol version"), "message: {message}");
3956            }
3957            other => panic!("expected a typed refusal, got {other:?}"),
3958        }
3959    }
3960
3961    /// The per-socket daemon guard is exclusive while held and reusable after
3962    /// release — the mechanism that keeps a supervisor respawn race down to
3963    /// one surviving daemon.
3964    #[cfg(unix)]
3965    #[test]
3966    fn events_daemon_guard_is_exclusive_then_reusable() {
3967        let dir = tempfile::tempdir().expect("tempdir");
3968        let socket = dir.path().join("events.sock");
3969        let first = match acquire_events_daemon_guard_outcome(&socket) {
3970            EventsDaemonGuardAcquisition::Held(guard) => guard,
3971            other => panic!("first acquire must succeed: {other:?}"),
3972        };
3973        let second = acquire_events_daemon_guard_outcome(&socket);
3974        assert!(
3975            matches!(second, EventsDaemonGuardAcquisition::Contended),
3976            "second acquire must report contention while the first guard is held: {second:?}"
3977        );
3978        drop(first);
3979        // Another test may fork while the first description is open. A child
3980        // can retain that flock until exec or exit; retry only that diagnosed
3981        // contention, never an open/hardening error, and never without a bound.
3982        const MAX_ATTEMPTS: usize = 100;
3983        for attempt in 1..=MAX_ATTEMPTS {
3984            match acquire_events_daemon_guard_outcome(&socket) {
3985                EventsDaemonGuardAcquisition::Held(_) => return,
3986                EventsDaemonGuardAcquisition::Contended if attempt < MAX_ATTEMPTS => {
3987                    std::thread::sleep(Duration::from_millis(10));
3988                }
3989                other => panic!(
3990                    "acquire must succeed after release within {MAX_ATTEMPTS} attempts; \
3991                     attempt {attempt}: {other:?}"
3992                ),
3993            }
3994        }
3995        unreachable!("the final acquisition attempt returns or reports its refusal");
3996    }
3997
3998    #[cfg(unix)]
3999    #[test]
4000    fn events_daemon_guard_open_failure_is_not_contention() {
4001        let dir = tempfile::tempdir().expect("tempdir");
4002        let clean = dir.path().join("clean.sock");
4003        assert!(matches!(
4004            acquire_events_daemon_guard_outcome(&clean),
4005            EventsDaemonGuardAcquisition::Held(_)
4006        ));
4007
4008        assert!(
4009            try_acquire_events_daemon_guard(&clean).is_some(),
4010            "the public Option entrance must preserve a successful acquisition"
4011        );
4012
4013        // Opening a directory for writing must fail even when no process
4014        // holds a daemon lock. Exercise the real open, not an injected result.
4015        let socket = dir.path().join("blocked.sock");
4016        let lock_path = socket.with_extension("lock");
4017        std::fs::create_dir(&lock_path).expect("create directory at lock entry");
4018        let outcome = acquire_events_daemon_guard_outcome(&socket);
4019        assert!(
4020            matches!(outcome, EventsDaemonGuardAcquisition::OpenFailed(_)),
4021            "an actual lock-file open failure must not report contention: {outcome:?}"
4022        );
4023        assert!(
4024            try_acquire_events_daemon_guard(&socket).is_none(),
4025            "the public Option entrance must refuse an actual open failure"
4026        );
4027        assert!(
4028            lock_path.is_dir(),
4029            "the refused entry must remain a directory"
4030        );
4031    }
4032
4033    /// A read-only file-backed runtime with the split configured must never
4034    /// create or schema-initialize an events database. Missing events db →
4035    /// `events()` serves the legacy store alone and leaves the filesystem
4036    /// untouched; a pre-existing events db (minted by an earlier writable
4037    /// host) is opened read-only and its rows merge into reads.
4038    /// (unix-gated for the sidecar-freeze harness step only.)
4039    #[cfg(unix)]
4040    #[tokio::test]
4041    async fn read_only_runtime_never_creates_an_events_db() {
4042        use crate::{KhiveRuntime, Namespace, RuntimeConfig};
4043
4044        let dir = tempfile::tempdir().expect("tempdir");
4045        let _registry_guard = TestRegistryGuard::new(dir.path());
4046        let main_db = dir.path().join("main.db");
4047        // Materialize + migrate the main database with a writable runtime.
4048        drop(
4049            KhiveRuntime::new_for_test(RuntimeConfig {
4050                db_path: Some(main_db.clone()),
4051                ..RuntimeConfig::no_embeddings()
4052            })
4053            .expect("create main db"),
4054        );
4055        let events_db = events_db_path_beside(&main_db);
4056        assert!(!events_db.exists(), "precondition: no events db yet");
4057
4058        // Freeze the WAL sidecars the writable runtime left behind — the
4059        // read-only opener refuses a snapshot with a writable -shm.
4060        khive_storage::test_support::freeze_snapshot_sidecars(&main_db);
4061
4062        let split_config = |db: PathBuf| RuntimeConfig {
4063            db_path: Some(main_db.clone()),
4064            events_split: Some(EventsSplitConfig {
4065                db_path: db,
4066                socket_path: None,
4067            }),
4068            ..RuntimeConfig::no_embeddings()
4069        };
4070
4071        // Arm 1: read-only runtime, no events db on disk. Reads work through
4072        // the legacy store and nothing is minted.
4073        let ro = KhiveRuntime::new_readonly_for_test(split_config(events_db.clone()))
4074            .expect("read-only runtime");
4075        let token = ro.authorize(Namespace::local()).expect("token");
4076        let store = ro.events(&token).expect("events store");
4077        let count = store
4078            .count_events(EventFilter::default())
4079            .await
4080            .expect("count through legacy-only plane");
4081        assert_eq!(count, 0);
4082        assert!(
4083            !events_db.exists(),
4084            "a read-only runtime must not mint the events database"
4085        );
4086
4087        // Arm 2: a writable host mints the events db and lands a lane row;
4088        // the read-only runtime then opens it read-only and merges the row.
4089        // The writer is scoped (not the process-global registry) and its WAL
4090        // sidecars frozen after drop — the read-only opener rightly refuses a
4091        // snapshot with a live writer, and that refusal is not this test's
4092        // subject.
4093        {
4094            let lane_backend = StorageBackend::sqlite_for_test(&events_db).expect("writable lane");
4095            lane_backend
4096                .events_for_namespace("local")
4097                .expect("lane store")
4098                .append_event(test_event("local"))
4099                .await
4100                .expect("seed lane row");
4101        }
4102        khive_storage::test_support::freeze_snapshot_sidecars(&events_db);
4103        let store = ro.events(&token).expect("events store with lane present");
4104        let count = store
4105            .count_events(EventFilter::default())
4106            .await
4107            .expect("merged count");
4108        assert_eq!(count, 1, "the pre-existing lane row must merge into reads");
4109    }
4110
4111    /// Embedded (direct) mode: the process-global backend writes and reads
4112    /// `events.db` without any daemon.
4113    #[tokio::test]
4114    async fn direct_mode_appends_and_reads_without_a_daemon() {
4115        let dir = tempfile::tempdir().expect("tempdir");
4116        let _registry_guard = TestRegistryGuard::new(dir.path());
4117        let db = dir.path().join("events.db");
4118        let backend = direct_backend_for(&db).expect("direct backend");
4119        let store = backend.events_for_namespace("local").expect("store");
4120        store
4121            .append_event(test_event("local"))
4122            .await
4123            .expect("direct append");
4124        let count = store
4125            .count_events(EventFilter::default())
4126            .await
4127            .expect("count");
4128        assert_eq!(count, 1);
4129    }
4130}