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