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