Skip to main content

khive_runtime/
daemon.rs

1//! khived daemon server — persistent warm runtime over a Unix socket.
2//!
3//! The daemon binds `~/.khive/khived.sock`, accepts length-prefixed request
4//! frames, dispatches them through a [`DaemonDispatch`] implementor, and serves
5//! results back. It is transport-agnostic: the MCP crate provides the dispatch
6//! impl, but any future client (CLI, HTTP gateway) can reuse this server.
7//!
8//! The client side (forwarding, auto-spawn) lives in the transport crate
9//! (e.g. `khive-mcp`), not here.
10
11use std::sync::Arc;
12
13mod store_guard;
14mod store_identity;
15#[cfg(unix)]
16use store_guard::ensure_claimed_parent_identity;
17#[cfg(unix)]
18pub use store_guard::{acquire_daemon_store_guards, bind_daemon_store_files, claim_stores};
19pub use store_guard::{assert_daemon_store_identities, DaemonStoreGuard};
20#[cfg(unix)]
21pub use store_identity::claimed_daemon_store_identity;
22mod supervisor_marker;
23#[cfg(unix)]
24pub use supervisor_marker::supervisor_marker_path;
25
26#[cfg(unix)]
27use std::io::Write as _;
28#[cfg(unix)]
29use std::os::unix::fs::{MetadataExt, PermissionsExt};
30#[cfg(unix)]
31use std::os::unix::io::AsRawFd;
32use std::path::PathBuf;
33
34#[cfg(unix)]
35use async_trait::async_trait;
36#[cfg(unix)]
37use libc;
38use serde::{Deserialize, Serialize};
39#[cfg(unix)]
40use tokio::io::{AsyncReadExt, AsyncWriteExt};
41#[cfg(unix)]
42use tokio::net::{UnixListener, UnixStream};
43
44#[cfg(unix)]
45use crate::pack::RequestIdentity;
46#[cfg(unix)]
47use khive_db::{run_checkpoint_task, CheckpointConfig, CheckpointLifecycleOwner, ConnectionPool};
48
49mod load_limits;
50#[cfg(unix)]
51use load_limits::{admit_or_refuse_busy, ConnectionAdmission};
52pub use load_limits::{
53    recall_ledger_snapshot, track_recall_ledger_task, ConnectionCapSnapshot, RecallLedgerSnapshot,
54};
55
56/// Maximum frame size accepted in either direction.
57pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
58
59/// Wire protocol version for the daemon IPC framing.
60///
61/// Increment this constant whenever the request or response frame shape
62/// changes in a backward-incompatible way. The client sends its version
63/// in every request; the daemon rejects mismatches with an explicit error
64/// that names both sides so the operator knows exactly what to do
65/// (`make local` rebuilds the client binary).
66/// See `docs/api/daemon.md#protocol_version` for the version-by-version history.
67pub const PROTOCOL_VERSION: u32 = 8;
68
69/// ADR-049 Amendment 11's disclosed initial demand idle interval.
70pub const DEFAULT_DEMAND_IDLE_SECS: u64 = 1_800;
71
72/// A launch-time choice, never inferred from process ancestry or environment.
73#[derive(Serialize, Deserialize, Debug, Clone, Copy, Default, PartialEq, Eq)]
74#[serde(rename_all = "snake_case")]
75pub enum DaemonLifetime {
76    Demand,
77    #[default]
78    Persistent,
79}
80
81/// Immutable daemon options. Existing entry points use persistent mode.
82#[derive(Debug, Clone, Copy)]
83pub struct DaemonOptions {
84    pub lifetime: DaemonLifetime,
85    pub idle_interval: std::time::Duration,
86}
87
88impl Default for DaemonOptions {
89    fn default() -> Self {
90        Self {
91            lifetime: DaemonLifetime::Persistent,
92            idle_interval: std::time::Duration::from_secs(DEFAULT_DEMAND_IDLE_SECS),
93        }
94    }
95}
96
97/// Host-owned startup decisions disclosed by lifecycle diagnostics.
98#[derive(Debug, Clone, Default)]
99pub struct DaemonStartupReport {
100    pub skipped_components: Vec<String>,
101    /// A named unknown inventory or service obligation prevents retirement.
102    pub idle_ineligible_reasons: Vec<String>,
103}
104
105#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq)]
106#[serde(rename_all = "snake_case")]
107pub enum DaemonLifecyclePhase {
108    Serving,
109    Draining,
110    Stopped,
111}
112
113#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq)]
114#[serde(rename_all = "snake_case")]
115pub enum DaemonShutdownReason {
116    Idle,
117    Signal,
118}
119
120/// Additive diagnostics for one daemon incarnation.
121#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)]
122pub struct DaemonLifecycleSnapshot {
123    pub lifetime: DaemonLifetime,
124    pub instance_generation: String,
125    pub effective_idle_interval_ms: u64,
126    pub phase: DaemonLifecyclePhase,
127    pub shutdown_reason: Option<DaemonShutdownReason>,
128    pub skipped_components: Vec<String>,
129    pub idle_ineligible_reasons: Vec<String>,
130    pub ordinary_requests: usize,
131    pub idle_blockers: Vec<String>,
132}
133
134#[cfg(unix)]
135struct DaemonLifecycle {
136    options: DaemonOptions,
137    state: std::sync::Mutex<DaemonLifecycleState>,
138    /// Admission for new connections on the daemon socket.
139    connections: ConnectionAdmission,
140}
141
142#[cfg(unix)]
143struct DaemonLifecycleState {
144    snapshot: DaemonLifecycleSnapshot,
145    last_request_completion: Option<tokio::time::Instant>,
146}
147
148#[cfg(unix)]
149impl DaemonLifecycle {
150    fn new(options: DaemonOptions, report: DaemonStartupReport) -> Self {
151        Self {
152            options,
153            state: std::sync::Mutex::new(DaemonLifecycleState {
154                snapshot: DaemonLifecycleSnapshot {
155                    lifetime: options.lifetime,
156                    instance_generation: uuid::Uuid::new_v4().to_string(),
157                    effective_idle_interval_ms: options
158                        .idle_interval
159                        .as_millis()
160                        .min(u128::from(u64::MAX))
161                        as u64,
162                    phase: DaemonLifecyclePhase::Serving,
163                    shutdown_reason: None,
164                    skipped_components: report.skipped_components,
165                    idle_ineligible_reasons: report.idle_ineligible_reasons,
166                    ordinary_requests: 0,
167                    idle_blockers: Vec::new(),
168                },
169                last_request_completion: None,
170            }),
171            connections: ConnectionAdmission::from_env(),
172        }
173    }
174
175    fn snapshot(&self) -> DaemonLifecycleSnapshot {
176        self.state
177            .lock()
178            .unwrap_or_else(std::sync::PoisonError::into_inner)
179            .snapshot
180            .clone()
181    }
182
183    fn ready(&self) {
184        self.state
185            .lock()
186            .unwrap_or_else(std::sync::PoisonError::into_inner)
187            .last_request_completion = Some(tokio::time::Instant::now());
188    }
189
190    fn admit(self: &Arc<Self>) -> Option<OrdinaryRequestGuard> {
191        let mut state = self
192            .state
193            .lock()
194            .unwrap_or_else(std::sync::PoisonError::into_inner);
195        if state.snapshot.phase != DaemonLifecyclePhase::Serving {
196            return None;
197        }
198        state.snapshot.ordinary_requests += 1;
199        Some(OrdinaryRequestGuard(Arc::clone(self)))
200    }
201
202    /// The same mutex orders ordinary admission and the irreversible idle decision.
203    fn try_idle(&self, blockers: impl FnOnce() -> Vec<String>) -> bool {
204        let mut state = self
205            .state
206            .lock()
207            .unwrap_or_else(std::sync::PoisonError::into_inner);
208        if self.options.lifetime != DaemonLifetime::Demand
209            || state.snapshot.phase != DaemonLifecyclePhase::Serving
210            || state.snapshot.ordinary_requests != 0
211            || !state.snapshot.idle_ineligible_reasons.is_empty()
212            || state
213                .last_request_completion
214                .is_none_or(|last| last.elapsed() < self.options.idle_interval)
215        {
216            return false;
217        }
218        state.snapshot.idle_blockers = blockers();
219        if !state.snapshot.idle_blockers.is_empty() {
220            return false;
221        }
222        state.snapshot.phase = DaemonLifecyclePhase::Draining;
223        state.snapshot.shutdown_reason = Some(DaemonShutdownReason::Idle);
224        true
225    }
226
227    fn draining(&self, reason: DaemonShutdownReason) {
228        let mut state = self
229            .state
230            .lock()
231            .unwrap_or_else(std::sync::PoisonError::into_inner);
232        if state.snapshot.phase != DaemonLifecyclePhase::Stopped {
233            state.snapshot.phase = DaemonLifecyclePhase::Draining;
234            state.snapshot.shutdown_reason = Some(reason);
235        }
236    }
237
238    fn stopped(&self) {
239        self.state
240            .lock()
241            .unwrap_or_else(std::sync::PoisonError::into_inner)
242            .snapshot
243            .phase = DaemonLifecyclePhase::Stopped;
244    }
245}
246
247#[cfg(unix)]
248struct OrdinaryRequestGuard(Arc<DaemonLifecycle>);
249
250#[cfg(unix)]
251impl Drop for OrdinaryRequestGuard {
252    fn drop(&mut self) {
253        let mut state = self
254            .0
255            .state
256            .lock()
257            .unwrap_or_else(std::sync::PoisonError::into_inner);
258        state.snapshot.ordinary_requests -= 1;
259        // The guard covers response transport and cleanup, not just dispatch.
260        state.last_request_completion = Some(tokio::time::Instant::now());
261    }
262}
263
264/// Internal signal carried in a dispatch result until the daemon moves it to
265/// response-frame metadata. It must never be sent in `result`: older v8
266/// clients publish that string without inspecting its contents.
267#[doc(hidden)]
268pub const DAEMON_LEXICAL_TIMEOUT_MARKER: &str = "__khive_daemon_lexical_timeout";
269
270const DEFAULT_DRAIN_TIMEOUT_SECS: u64 = 10;
271/// An accepted local socket must finish its first frame within this window.
272/// Dispatch deadlines start only after decoding, so they cannot reap peers
273/// that connect and then stop sending request bytes.
274#[cfg(unix)]
275const INITIAL_FRAME_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
276
277#[cfg(unix)]
278fn next_accept_error_backoff(previous: Option<std::time::Duration>) -> std::time::Duration {
279    previous
280        .map(|delay| delay.saturating_mul(2))
281        .unwrap_or_else(|| std::time::Duration::from_millis(10))
282        .min(std::time::Duration::from_secs(1))
283}
284
285// ── paths ─────────────────────────────────────────────────────────────────────
286
287/// Base `.khive` directory used to anchor every advisory lock/socket/pid
288/// path below. Pure path computation — portable on every target, even
289/// though most of its callers (socket/pid paths) are unix-only.
290///
291/// Resolution order: `HOME`, then `USERPROFILE` (the conventional Windows
292/// home variable), then a platform-specific last resort. On non-unix targets
293/// the last resort is the OS temp directory (per-user on Windows), so the
294/// lock path can never be working-directory-relative there: two processes
295/// opening the same database from different working directories must resolve
296/// the same lock file. On unix the last resort stays the historical `"."` —
297/// a shared world-writable anchor such as `/tmp` would be worse, since a
298/// local attacker could pre-claim the directory and the socket/lock files
299/// under it before the daemon's first run.
300fn khive_dir() -> PathBuf {
301    khive_root_from(
302        std::env::var("HOME").ok(),
303        std::env::var("USERPROFILE").ok(),
304    )
305}
306
307/// Env-free core of [`khive_dir`], split out so the fallback chain is
308/// testable without mutating process-global environment variables.
309fn khive_root_from(home: Option<String>, userprofile: Option<String>) -> PathBuf {
310    home.filter(|v| !v.trim().is_empty())
311        .or_else(|| userprofile.filter(|v| !v.trim().is_empty()))
312        .map(PathBuf::from)
313        .unwrap_or_else(last_resort_root)
314        .join(".khive")
315}
316
317/// The directory for the SQLite volume lock files, from the single rule in
318/// [`khive_db::default_volume_lock_dir`]: `KHIVE_VOLUME_LOCK_DIR` when set,
319/// else `<home>/.khive/sqlite-volume-locks`. Unlike `khive_dir` it has no
320/// last-resort root: without a home directory the result is a configuration
321/// error, because a working-directory-relative lock directory would give two
322/// processes two different lock files.
323pub fn volume_lock_dir() -> Result<PathBuf, khive_db::SqliteError> {
324    khive_db::default_volume_lock_dir()
325}
326
327/// See [`khive_dir`] for why the two arms differ.
328#[cfg(unix)]
329fn last_resort_root() -> PathBuf {
330    PathBuf::from(".")
331}
332
333#[cfg(not(unix))]
334fn last_resort_root() -> PathBuf {
335    std::env::temp_dir()
336}
337
338/// Env var overriding the socket half of the daemon rendezvous.
339#[cfg(unix)]
340const SOCKET_PATH_ENV: &str = "KHIVE_SOCKET";
341
342/// Env var overriding the PID-file half of the daemon rendezvous.
343#[cfg(unix)]
344const PID_PATH_ENV: &str = "KHIVE_PID";
345
346/// Read a path override, treating an empty value as unset.
347///
348/// One predicate for "the operator set this variable", shared by the path
349/// resolvers and [`ensure_rendezvous_overrides_paired`]. A pairing check that
350/// disagreed with the resolvers about what counts as set would either refuse
351/// boots that resolve consistently, or admit the split rendezvous it exists
352/// to stop.
353#[cfg(unix)]
354fn path_override(key: &str) -> Option<PathBuf> {
355    match std::env::var(key) {
356        Ok(p) if !p.is_empty() => Some(PathBuf::from(p)),
357        _ => None,
358    }
359}
360
361#[cfg(unix)]
362fn default_socket_path() -> PathBuf {
363    khive_dir().join("khived.sock")
364}
365
366#[cfg(unix)]
367fn default_pid_path() -> PathBuf {
368    khive_dir().join("khived.pid")
369}
370
371/// Unix socket path the daemon binds and clients connect to.
372///
373/// Overridable via the `KHIVE_SOCKET` env var (for tests and ops), which must
374/// be set together with `KHIVE_PID`: the daemon refuses to boot when exactly
375/// one of the two is set.
376#[cfg(unix)]
377pub fn socket_path() -> PathBuf {
378    path_override(SOCKET_PATH_ENV).unwrap_or_else(default_socket_path)
379}
380
381/// PID file path written by the daemon.
382///
383/// Overridable via the `KHIVE_PID` env var, which must be set together with
384/// `KHIVE_SOCKET`: the daemon refuses to boot when exactly one of the two is
385/// set.
386#[cfg(unix)]
387pub fn pid_path() -> PathBuf {
388    path_override(PID_PATH_ENV).unwrap_or_else(default_pid_path)
389}
390
391/// Refuse to boot when exactly one of `KHIVE_SOCKET` / `KHIVE_PID` is set
392/// (#2656).
393///
394/// The socket and the PID file are the two halves of one rendezvous, but they
395/// resolve independently: with only `KHIVE_SOCKET` set a daemon binds a
396/// private socket while still claiming the shared PID file, and with only
397/// `KHIVE_PID` set it writes a private PID file while binding the shared
398/// socket. Either way [`cleanup_stale_daemon`] reads an incumbent's pid out of
399/// one instance's file and judges it by probing the other instance's socket,
400/// so both of its branches are wrong: a live incumbent produces a refusal
401/// naming a pid that has nothing to do with the socket being started, and a
402/// pid that is no longer running makes this process delete a rendezvous file
403/// another daemon's `shutdown_cleanup_if_owned` still expects to own.
404///
405/// Setting both variables (a fully private rendezvous) and setting neither
406/// (the default rendezvous) are both unchanged.
407#[cfg(unix)]
408fn ensure_rendezvous_overrides_paired() -> anyhow::Result<()> {
409    match (path_override(SOCKET_PATH_ENV), path_override(PID_PATH_ENV)) {
410        (Some(socket), None) => anyhow::bail!(
411            "refusing to start: {SOCKET_PATH_ENV} is set to {} but {PID_PATH_ENV} is not set. \
412             The socket and the PID file are two halves of one daemon rendezvous and must move \
413             together: with only {SOCKET_PATH_ENV} set, this daemon would bind a private socket \
414             while claiming the shared PID file at {}, which belongs to the default rendezvous \
415             served on {}. Set {PID_PATH_ENV} to a private path beside the socket, or unset \
416             {SOCKET_PATH_ENV} to share the default rendezvous.",
417            socket.display(),
418            default_pid_path().display(),
419            default_socket_path().display(),
420        ),
421        (None, Some(pid)) => anyhow::bail!(
422            "refusing to start: {PID_PATH_ENV} is set to {} but {SOCKET_PATH_ENV} is not set. \
423             The socket and the PID file are two halves of one daemon rendezvous and must move \
424             together: with only {PID_PATH_ENV} set, this daemon would write a private PID file \
425             while binding the shared socket at {}, the default rendezvous whose owner is \
426             recorded in {}. Set {SOCKET_PATH_ENV} to a private path beside the PID file, or \
427             unset {PID_PATH_ENV} to share the default rendezvous.",
428            pid.display(),
429            default_socket_path().display(),
430            default_pid_path().display(),
431        ),
432        _ => Ok(()),
433    }
434}
435
436/// Advisory lock file used to serialize stale-daemon recovery across concurrent
437/// clients (flock/`File::lock` on the file; released when the lock file
438/// handle is dropped). Path computation is portable; `kkernel exec`'s
439/// non-unix local-construction guard (`kkernel::exec::acquire_local_construction_guard`)
440/// shares this exact path with the unix daemon-boot guard so the two stay
441/// mutually exclusive.
442///
443/// Overridable via the `KHIVE_LOCK` env var (for tests).
444pub fn lock_path() -> PathBuf {
445    if let Ok(p) = std::env::var("KHIVE_LOCK") {
446        if !p.is_empty() {
447            return PathBuf::from(p);
448        }
449    }
450    khive_dir().join("khived.recovery.lock")
451}
452
453/// Advisory lock file used to serialize RECOVERY (kill+respawn) attempts
454/// across concurrent clients only — the daemon's own boot sequence never
455/// acquires this file ([`lock_path`] / [`acquire_daemon_boot_guard`] is the
456/// boot-side lock). A recoverer holding this lock across dead-confirmation
457/// → kill → spawn (khive-mcp's `kill_and_respawn`) therefore can never
458/// deadlock against a peer daemon's boot, unlike holding the shared boot
459/// lock for that whole span would.
460///
461/// Overridable via the `KHIVE_RECOVERER_LOCK` env var (for tests).
462#[cfg(unix)]
463pub fn recoverer_lock_path() -> PathBuf {
464    if let Ok(p) = std::env::var("KHIVE_RECOVERER_LOCK") {
465        if !p.is_empty() {
466            return PathBuf::from(p);
467        }
468    }
469    khive_dir().join("khived.recoverer.lock")
470}
471
472#[cfg(unix)]
473pub const SUPERVISOR_CLAIM_ENV: &str = "KHIVE_SUPERVISOR_CLAIM";
474
475#[cfg(unix)]
476fn read_supervisor_marker_claim() -> Option<(u32, String)> {
477    use std::io::Read;
478    use std::os::unix::fs::OpenOptionsExt;
479
480    let file = std::fs::OpenOptions::new()
481        .read(true)
482        .custom_flags(libc::O_NONBLOCK | libc::O_NOFOLLOW)
483        .open(supervisor_marker_path())
484        .ok()?;
485    if !file.metadata().ok()?.is_file() {
486        return None;
487    }
488    let mut marker = String::new();
489    file.take(4097).read_to_string(&mut marker).ok()?;
490    if marker.len() > 4096 {
491        return None;
492    }
493    let mut lines = marker.lines();
494    if lines.next()?.is_empty() {
495        return None;
496    }
497    let pid = lines.next()?.parse::<u32>().ok().filter(|pid| *pid > 0)?;
498    lines
499        .next()?
500        .parse::<u64>()
501        .ok()
502        .filter(|seconds| *seconds > 0)?;
503    let claim = lines.next()?.to_string();
504    if lines.next().is_some() {
505        return None;
506    }
507    let parsed = uuid::Uuid::parse_str(&claim).ok()?;
508    if parsed.get_version() != Some(uuid::Version::Random) || parsed.to_string() != claim {
509        return None;
510    }
511    Some((pid, claim))
512}
513
514#[cfg(unix)]
515fn current_supervisor_claim() -> Option<String> {
516    let claim = std::env::var(SUPERVISOR_CLAIM_ENV).ok()?;
517    let (pid, published_claim) = read_supervisor_marker_claim()?;
518    (pid == std::process::id() && claim == published_claim).then_some(claim)
519}
520
521#[cfg(unix)]
522fn open_lock_file(path: &std::path::Path) -> std::io::Result<std::fs::File> {
523    if let Some(parent) = path.parent() {
524        let _ = std::fs::create_dir_all(parent);
525    }
526    std::fs::OpenOptions::new()
527        .create(true)
528        .truncate(false)
529        .write(true)
530        .open(path)
531}
532
533#[cfg(unix)]
534fn acquire_flock_blocking(path: &std::path::Path, label: &str) -> Option<std::fs::File> {
535    let file = match open_lock_file(path) {
536        Ok(f) => f,
537        Err(e) => {
538            tracing::warn!(error = %e, path = ?path, "cannot open {label} lock file");
539            return None;
540        }
541    };
542    // SAFETY: flock is a POSIX advisory lock with no memory side-effects.
543    let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) };
544    if rc != 0 {
545        tracing::warn!("flock LOCK_EX failed on {label} lock");
546        return None;
547    }
548    Some(file)
549}
550
551/// Acquire an exclusive advisory flock on the recovery/startup lock file.
552///
553/// The returned `File` holds the lock for its lifetime; dropping it releases
554/// it.  Used by both the client (serializing kill+spawn) and the daemon server
555/// (serializing cleanup+bind+pid-write) so the two critical sections are
556/// mutually exclusive across processes.
557#[cfg(unix)]
558pub fn acquire_recovery_lock() -> Option<std::fs::File> {
559    acquire_flock_blocking(&lock_path(), "recovery")
560}
561
562/// Attempt to acquire an exclusive advisory flock on `path`, retrying with a
563/// non-blocking `flock(LOCK_NB)` until `deadline` elapses. Bounded alternative
564/// to `acquire_recovery_lock`/`acquire_daemon_boot_guard`'s unbounded blocking
565/// flock — see `docs/api/daemon.md#try_acquire_flock_until` for why a caller
566/// merely detecting lock freedom needs a deadline instead.
567///
568/// - `Ok(Some(file))` — the lock was free within the deadline.
569/// - `Ok(None)` — `deadline` elapsed while the lock stayed held; an explicit
570///   "could not confirm" outcome, distinct from a hard I/O error.
571/// - `Err(_)` — the lock file could not be opened, or `flock` failed for a
572///   reason other than contention.
573///
574/// Blocking (paces retries with `std::thread::sleep`) — async callers must
575/// run this via `spawn_blocking`.
576#[cfg(unix)]
577fn try_acquire_flock_until(
578    path: &std::path::Path,
579    deadline: std::time::Instant,
580) -> std::io::Result<Option<std::fs::File>> {
581    let file = open_lock_file(path)?;
582    let poll_interval = std::time::Duration::from_millis(10);
583    loop {
584        // SAFETY: flock is a POSIX advisory lock with no memory side-effects.
585        let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
586        if rc == 0 {
587            return Ok(Some(file));
588        }
589        let err = std::io::Error::last_os_error();
590        if err.raw_os_error() != Some(libc::EWOULDBLOCK) {
591            return Err(err);
592        }
593        let now = std::time::Instant::now();
594        if now >= deadline {
595            return Ok(None);
596        }
597        std::thread::sleep(poll_interval.min(deadline - now));
598    }
599}
600
601/// Bounded, deadline-aware variant of [`acquire_daemon_boot_guard`]: attempts
602/// the SAME boot/recovery lock ([`lock_path`]) but gives up at `deadline`
603/// instead of blocking forever. For callers that need to detect "is a boot in
604/// progress right now" without risking an unbounded wait behind a wedged
605/// holder: e.g. khive-mcp's `confirm_genuinely_dead` re-probing rounds,
606/// where `DEAD_CONFIRM_ROUNDS` must bound elapsed time, not just probe count.
607#[cfg(unix)]
608pub fn try_acquire_daemon_boot_guard_until(
609    deadline: std::time::Instant,
610) -> std::io::Result<Option<DaemonBootGuard>> {
611    try_acquire_flock_until(&lock_path(), deadline)
612}
613
614/// Bounded, deadline-aware acquisition of the recoverer-only lock
615/// ([`recoverer_lock_path`]). See [`try_acquire_daemon_boot_guard_until`] for
616/// the shared rationale — a second recoverer waiting for a peer's dead
617/// confirmation/kill/spawn critical section must give up and report
618/// "uncertain" rather than block forever if that peer is itself wedged.
619#[cfg(unix)]
620pub fn try_acquire_recoverer_lock_until(
621    deadline: std::time::Instant,
622) -> std::io::Result<Option<std::fs::File>> {
623    try_acquire_flock_until(&recoverer_lock_path(), deadline)
624}
625
626/// Guard returned by [`acquire_daemon_boot_guard`], held across cold-boot
627/// schema initialization (migrations + pack schema plans / FTS DDL) through
628/// daemon bind + pid-write.
629#[cfg(unix)]
630pub type DaemonBootGuard = std::fs::File;
631
632/// Acquire the recovery/boot lock, treating failure as fatal.
633///
634/// Unlike [`acquire_recovery_lock`] (best-effort, `None` on failure: used by
635/// shutdown cleanup, where skipping unlink is safer than blocking forever),
636/// daemon-mode boot must hold this lock across migrations/FTS DDL through
637/// bind+pid-write. Silently continuing with no lock reopens the cold-boot FTS
638/// race this guard exists to close, so callers that are about to run
639/// daemon-mode boot (or wait for one to quiesce) must fail loudly instead of
640/// proceeding unguarded.
641#[cfg(unix)]
642pub fn acquire_daemon_boot_guard() -> anyhow::Result<DaemonBootGuard> {
643    acquire_recovery_lock()
644        .ok_or_else(|| anyhow::anyhow!("failed to acquire daemon boot/recovery lock"))
645}
646
647/// Identity of a bound Unix socket path, used to tell "the socket I bound" apart
648/// from "a same-path socket some other daemon bound after mine was removed".
649///
650/// A socket path can be recreated by a different process between the time
651/// this daemon captures its identity and the time it later checks it, so `dev`
652/// and `ino` (not the path) are what must match for cleanup to be safe.
653#[cfg(unix)]
654#[derive(Clone, Copy, PartialEq, Eq)]
655struct SocketIdentity {
656    dev: u64,
657    ino: u64,
658}
659
660#[cfg(unix)]
661fn socket_identity(path: &std::path::Path) -> Option<SocketIdentity> {
662    use std::os::unix::fs::MetadataExt;
663    let meta = std::fs::metadata(path).ok()?;
664    Some(SocketIdentity {
665        dev: meta.dev(),
666        ino: meta.ino(),
667    })
668}
669
670// ── connection principal ──────────────────────────────────────────────────────
671
672/// The uid on the other end of an accepted connection, read from the kernel.
673///
674/// This is the only identity on this socket the caller cannot choose. Every
675/// identity field on the request frame — `namespace`, `actor_id`,
676/// `visible_namespaces`, `config_id` — is supplied by the connecting process,
677/// so none of them can answer "who is this". A check reading self-asserted
678/// fields is not a weak gate, it is not a gate: anyone who wants to pass it
679/// asserts the passing values.
680///
681/// `getpeereid(2)` on macOS/BSD, `SO_PEERCRED` on Linux. Both report the peer's
682/// credentials as recorded by the kernel at connect time.
683#[cfg(unix)]
684pub(crate) fn peer_uid(stream: &UnixStream) -> std::io::Result<u32> {
685    use std::os::fd::AsRawFd;
686    let fd = stream.as_raw_fd();
687
688    #[cfg(any(target_os = "macos", target_os = "ios", target_vendor = "apple"))]
689    {
690        let mut uid: libc::uid_t = 0;
691        let mut gid: libc::gid_t = 0;
692        // SAFETY: `fd` is a live connected socket owned by `stream` for the
693        // duration of this call; both out-params are valid initialized locals.
694        let rc = unsafe { libc::getpeereid(fd, &mut uid, &mut gid) };
695        if rc != 0 {
696            return Err(std::io::Error::last_os_error());
697        }
698        Ok(uid as u32)
699    }
700
701    #[cfg(target_os = "linux")]
702    {
703        let mut cred = libc::ucred {
704            pid: 0,
705            uid: 0,
706            gid: 0,
707        };
708        let mut len = std::mem::size_of::<libc::ucred>() as libc::socklen_t;
709        // SAFETY: `fd` is a live connected socket owned by `stream`; `cred` is
710        // an initialized local of exactly `len` bytes, which is what
711        // SO_PEERCRED writes.
712        let rc = unsafe {
713            libc::getsockopt(
714                fd,
715                libc::SOL_SOCKET,
716                libc::SO_PEERCRED,
717                (&mut cred as *mut libc::ucred).cast::<libc::c_void>(),
718                &mut len,
719            )
720        };
721        if rc != 0 {
722            return Err(std::io::Error::last_os_error());
723        }
724        Ok(cred.uid)
725    }
726
727    #[cfg(not(any(
728        target_os = "linux",
729        target_os = "macos",
730        target_os = "ios",
731        target_vendor = "apple"
732    )))]
733    {
734        let _ = fd;
735        Err(std::io::Error::new(
736            std::io::ErrorKind::Unsupported,
737            "peer-credential capture is not implemented for this platform",
738        ))
739    }
740}
741
742/// Whether a connection from `uid` may be served by this daemon.
743///
744/// ADR-096 accepted per-request identity threading **for the single-principal
745/// owner-only socket only**, and rested that on three things: the `0600`
746/// socket, all connections being the same uid, and the database being already
747/// same-uid-accessible. The socket mode is asserted at bind. This asserts the
748/// second, which previously had no representation in the code at all — nothing
749/// read peer identity, so nothing could notice when it stopped being true.
750///
751/// **Principal is not attribution.** Many `actor_id`s over one socket is
752/// exactly what ADR-096 shipped and what every seat on a normal host does;
753/// refusing a second distinct actor would break the accepted design. The
754/// principal is the uid, and this refuses only a genuinely foreign one.
755///
756/// **There is deliberately no configuration escape hatch.** A flag permitting
757/// other uids would not weaken this assertion, it would delete it, in the way
758/// hardest to notice later: the check still exists, its tests still pass, and
759/// the deployment that matters has it off. A deployment that genuinely needs
760/// multiple uids needs a code change and a gated ADR — which is precisely the
761/// decision that should be impossible to make by accident.
762#[cfg(unix)]
763pub(crate) fn uid_is_permitted(peer: u32, daemon_euid: u32) -> bool {
764    peer == daemon_euid
765}
766
767// ── wire types ────────────────────────────────────────────────────────────────
768
769mod config_id;
770#[cfg(test)]
771use config_id::parse_config_id;
772pub use config_id::{
773    config_id_extra_embedder_exclusions, config_ids_compatible, first_config_mismatch_field,
774};
775
776/// Request frame sent from a client to the daemon.
777#[derive(Serialize, Deserialize, Default)]
778pub struct DaemonRequestFrame {
779    pub ops: String,
780    /// Parse and inspect the catalog without dispatch, identity, or storage access.
781    #[serde(default)]
782    pub plan: bool,
783    #[serde(skip_serializing_if = "Option::is_none")]
784    pub presentation: Option<String>,
785    #[serde(skip_serializing_if = "Option::is_none")]
786    pub presentation_per_op: Option<Vec<Option<String>>>,
787    /// The client's resolved storage/gate default namespace for this request.
788    ///
789    /// As of protocol version 3 (ADR-096) the daemon serves the request under
790    /// this namespace instead of rejecting on mismatch: a per-request
791    /// identity input, not a same-process-identity assertion.
792    pub namespace: String,
793    /// The client's resolved write-stamp / gate actor identity (ADR-057),
794    /// carried on the frame so the warm daemon stamps writes with the
795    /// *caller's* actor instead of its own baked `actor_id` (ADR-096). `None`
796    /// mints `ActorRef::anonymous()`, matching an unconfigured actor.
797    #[serde(default)]
798    pub actor_id: Option<String>,
799    /// Opaque process provenance resolved in the originating client process.
800    /// It is carried per request because a shared warm daemon's environment
801    /// does not identify the worker that submitted the operation. Protocol v4
802    /// makes this field part of dispatch semantics: a v3 daemon must reject the
803    /// request rather than execute it while silently discarding provenance.
804    #[serde(default, skip_serializing_if = "Option::is_none")]
805    pub process_ref: Option<String>,
806    /// The client's resolved extra read-visibility namespaces (ADR-007 Rule
807    /// 3b), carried on the frame so the warm daemon widens read scope to
808    /// match the caller's own configuration rather than its own baked
809    /// `visible_namespaces` (ADR-096). A non-`local` `actor_id` joins default
810    /// reads where the registry mints the token (ADR-007 Rev 4 Rule 3b), so
811    /// an empty list still includes that actor in default reads. Explicit
812    /// `namespace=` operations remain scoped to exactly that namespace.
813    #[serde(default)]
814    pub visible_namespaces: Vec<String>,
815    /// Fingerprint of the client's engine-coherence config: packs, db target,
816    /// embedders, backend routing, and construction-baked outbound policy.
817    /// Identity fields are carried separately in this frame. The daemon rejects
818    /// requests whose configuration differs, except when its extra-embedder set
819    /// is a superset of the client's and every other field matches. See
820    /// ADR-027 / ADR-049 / ADR-096.
821    #[serde(default)]
822    pub config_id: String,
823    /// IPC protocol version sent by the client. Pre-versioning clients omit
824    /// this field (deserializes to 0). The daemon compares against
825    /// [`PROTOCOL_VERSION`] and rejects mismatches with an explicit error.
826    #[serde(default)]
827    pub protocol_version: u32,
828    /// When `true`, the daemon returns an identity frame (ok=true, result=None)
829    /// immediately after identity validation — without calling the dispatcher.
830    /// Used by the client's under-lock recovery probe to confirm a daemon is
831    /// alive and identity-matching without dispatching any mutating verb.
832    /// Pre-probe clients omit this field (deserializes to false → normal dispatch).
833    #[serde(default)]
834    pub probe_only: bool,
835    /// When `true`, the daemon returns a point-in-time [`MetricsSnapshot`] of
836    /// its server-side gauges (a read-only measurement surface for the
837    /// load/perf harness) instead of dispatching any op. Handled before the
838    /// `config_id` equality reject: a gauge read is process-global and
839    /// namespace/config-agnostic, not a namespaced record operation.
840    /// READ-ONLY — this field is the only input the frame accepts for a
841    /// metrics request; there is no reset or mutation reachable over the
842    /// wire. Pre-metrics clients omit this field (deserializes to `false` →
843    /// normal dispatch, unaffected).
844    #[serde(default)]
845    pub metrics_only: bool,
846    /// Output format for this request (ADR-078). Forwarded to the daemon's
847    /// serialization seam. `None` means use the daemon's resolved default.
848    #[serde(default)]
849    #[serde(skip_serializing_if = "Option::is_none")]
850    pub format: Option<String>,
851    /// Per-operation output format overrides (ADR-078).
852    #[serde(default)]
853    #[serde(skip_serializing_if = "Option::is_none")]
854    pub format_per_op: Option<Vec<Option<String>>>,
855    /// Whether this request originated from the agent-facing MCP `request`
856    /// tool (the wire surface). When `true`, the daemon rejects
857    /// `Visibility::Subhandler` verbs: agents must not invoke internal
858    /// subhandlers. When `false` (the default, and the only value any
859    /// operator path sends), subhandlers are allowed: `kkernel exec` and
860    /// other in-process callers are trusted operator surfaces.
861    ///
862    /// This is the origin discriminator, not a daemon-vs-local one: operator
863    /// requests flow through the daemon by default too, so the gate cannot
864    /// key on transport.
865    #[serde(default)]
866    pub from_wire: bool,
867    /// Request-group correlation id (khive#948), echoed back unchanged on
868    /// [`DaemonResponseFrame::request_id`] and stamped into the dispatch's
869    /// audit event (`resource.request_id`) so a benchmark harness can join
870    /// its own pre-send sample to the server-side audit row for the same
871    /// request. Agent-facing MCP requests always carry one: the bridge keeps a
872    /// caller-supplied value or mints an opaque nonzero value when absent.
873    /// Operator-built/probe frames may still use `None`. Purely additive —
874    /// `#[serde(default)]` matches `metrics_only`/`format`/`format_per_op`
875    /// precedent, with no `PROTOCOL_VERSION` bump.
876    #[serde(default)]
877    #[serde(skip_serializing_if = "Option::is_none")]
878    pub request_id: Option<u64>,
879}
880
881/// A dispatch failure whose domain outcome remains available to the transport.
882#[derive(Debug, Clone)]
883pub struct DaemonDispatchError {
884    pub message: String,
885    pub error_detail: serde_json::Value,
886}
887
888/// Per-field container limit, asserted equal to the request parser's bound by MCP.
889pub const ERROR_DETAIL_NESTING_DEPTH_LIMIT: usize = 64;
890
891fn error_detail_value_within_limit(value: &serde_json::Value) -> bool {
892    let mut pending = vec![(value, 0_usize)];
893    while let Some((value, depth)) = pending.pop() {
894        match value {
895            serde_json::Value::Array(items) if depth < ERROR_DETAIL_NESTING_DEPTH_LIMIT => {
896                pending.extend(items.iter().map(|child| (child, depth + 1)));
897            }
898            serde_json::Value::Object(fields) if depth < ERROR_DETAIL_NESTING_DEPTH_LIMIT => {
899                pending.extend(fields.values().map(|child| (child, depth + 1)));
900            }
901            serde_json::Value::Array(_) | serde_json::Value::Object(_) => return false,
902            _ => {}
903        }
904    }
905    true
906}
907
908fn drop_error_detail_iteratively(value: serde_json::Value) {
909    let mut pending = vec![value];
910    while let Some(value) = pending.pop() {
911        match value {
912            serde_json::Value::Array(items) => pending.extend(items),
913            serde_json::Value::Object(fields) => pending.extend(fields.into_values()),
914            _ => {}
915        }
916    }
917}
918
919impl DaemonDispatchError {
920    /// Missing or unrecognized disposition from a legacy implementation is unknown.
921    pub fn new(message: impl Into<String>, error_detail: Option<serde_json::Value>) -> Self {
922        let message = message.into();
923        let mut fields = match error_detail {
924            Some(serde_json::Value::Object(fields)) => fields,
925            Some(data) => serde_json::Map::from_iter([("data".to_string(), data)]),
926            None => serde_json::Map::new(),
927        };
928        let disposition = match fields
929            .get("domain_disposition")
930            .and_then(serde_json::Value::as_str)
931        {
932            Some("committed") => crate::DomainDisposition::Committed,
933            Some("not_committed") => crate::DomainDisposition::NotCommitted,
934            _ => crate::DomainDisposition::Unknown,
935        };
936        if disposition != crate::DomainDisposition::Committed {
937            if let Some(result) = fields.remove("domain_result") {
938                drop_error_detail_iteratively(result);
939            }
940        }
941        let rejected: Vec<String> = fields
942            .iter()
943            .filter(|(_, value)| !error_detail_value_within_limit(value))
944            .map(|(name, _)| name.clone())
945            .collect();
946        let omitted_result = rejected.iter().any(|name| name == "domain_result");
947        let omitted_detail = !rejected.is_empty();
948        for name in rejected {
949            if let Some(value) = fields.remove(&name) {
950                drop_error_detail_iteratively(value);
951            }
952        }
953        let mut error_detail = serde_json::Value::Object(fields);
954        if error_detail["kind"].as_str().is_none() {
955            error_detail["kind"] = serde_json::json!("internal");
956        }
957        if error_detail["message"].as_str().is_none() {
958            error_detail["message"] = serde_json::json!(message);
959        }
960        error_detail["domain_disposition"] = serde_json::json!(disposition.as_str());
961        if omitted_detail {
962            error_detail["code"] = serde_json::json!(if omitted_result {
963                "result_too_deep"
964            } else {
965                "error_detail_too_deep"
966            });
967        }
968        Self {
969            message,
970            error_detail,
971        }
972    }
973}
974
975/// Response frame sent from the daemon back to a client.
976#[derive(Serialize, Deserialize, Debug)]
977pub struct DaemonResponseFrame {
978    pub ok: bool,
979    pub result: Option<String>,
980    pub error: Option<String>,
981    /// Additive failure metadata; legacy protocol-v4 peers still read `error` as text.
982    /// On a successful response, `{"lexical_timeout":true}` is a daemon-only
983    /// diagnostic that old clients ignore and new clients log locally.
984    #[serde(default, skip_serializing_if = "Option::is_none")]
985    pub error_detail: Option<serde_json::Value>,
986    pub namespace_mismatch: bool,
987    /// Set when the request's `config_id` does not match the daemon's. Like
988    /// `namespace_mismatch`, this signals the client to fall back to local
989    /// dispatch rather than execute under a different runtime/config.
990    #[serde(default)]
991    pub config_mismatch: bool,
992    /// The `config_id` the daemon dispatched under, echoed back so the client
993    /// can positively confirm the result came from a matching runtime. A
994    /// pre-`config_id` daemon omits this field (deserializes to `None`), which
995    /// the client treats as a mismatch and falls back to local dispatch — this
996    /// closes the upgrade window where a new restricted client could otherwise
997    /// trust a still-warm legacy daemon's broader registry.
998    #[serde(default)]
999    pub served_config_id: Option<String>,
1000    /// Set when the client's `protocol_version` does not match the daemon's
1001    /// [`PROTOCOL_VERSION`]. The client must treat this as a hard error and
1002    /// surface the human-readable `error` field rather than falling back to
1003    /// local dispatch (which would hide the version skew).
1004    #[serde(default)]
1005    pub version_mismatch: bool,
1006    /// The daemon's [`PROTOCOL_VERSION`], echoed in error responses so the
1007    /// client can include both sides in the diagnostic message. Pre-versioning
1008    /// daemons omit this field (deserializes to 0).
1009    #[serde(default)]
1010    pub daemon_protocol_version: u32,
1011    /// Populated when the request set `metrics_only: true`: a point-in-time
1012    /// snapshot of the daemon's server-side gauges. `None` on every other
1013    /// response, and on any response from a daemon that predates this field
1014    /// (client-side back-compat via `#[serde(default)]`, matching
1015    /// `served_config_id`'s upgrade-window handling above).
1016    #[serde(default, skip_serializing_if = "Option::is_none")]
1017    pub metrics: Option<MetricsSnapshot>,
1018    /// Echo of the request's `request_id` (khive#948), present whenever the
1019    /// frame that produced this response carried one — including on every
1020    /// error/denied arm, not only success, so a client can join a failure
1021    /// the same way it joins a success. `#[serde(default)]` so an older
1022    /// daemon's response (predating this field) deserializes to `None`
1023    /// rather than a parse error.
1024    #[serde(default)]
1025    pub request_id: Option<u64>,
1026}
1027
1028/// Move the private dispatch signal out of the result before any client can
1029/// observe it. Unmarked results retain their exact bytes. The marked result
1030/// was serialized from a JSON Value by the MCP server, so reserializing after
1031/// removal reproduces its public envelope. The marker's escaped frame cost
1032/// equals `error_detail:{"lexical_timeout":true}`, keeping the server's exact
1033/// frame-fit calculation valid after this move.
1034#[cfg(unix)]
1035fn take_daemon_lexical_timeout_marker(raw: String) -> (String, Option<serde_json::Value>) {
1036    if !raw.contains(DAEMON_LEXICAL_TIMEOUT_MARKER) {
1037        return (raw, None);
1038    }
1039    let Ok(mut value) = serde_json::from_str::<serde_json::Value>(&raw) else {
1040        return (raw, None);
1041    };
1042    let Some(fields) = value.as_object_mut() else {
1043        return (raw, None);
1044    };
1045    if !fields
1046        .get("results")
1047        .is_some_and(serde_json::Value::is_array)
1048    {
1049        return (raw, None);
1050    }
1051    let Some(marker) = fields.remove(DAEMON_LEXICAL_TIMEOUT_MARKER) else {
1052        return (raw, None);
1053    };
1054    let detail =
1055        (marker.as_bool() == Some(true)).then(|| serde_json::json!({"lexical_timeout": true}));
1056    (
1057        serde_json::to_string(&value).expect("serde_json::Value is serializable"),
1058        detail,
1059    )
1060}
1061
1062/// One checkpoint store in this daemon's fixed topology. IDs are process-local:
1063/// `main` or `secondary:<index>` in dispatcher order. The basename is display-only,
1064/// not an identity; no directory path is exposed. Restart/topology changes reset
1065/// the interpretation of interval deltas.
1066#[derive(Serialize, Deserialize, Debug, Clone, Default, PartialEq)]
1067#[serde(default)]
1068pub struct CheckpointStoreMetrics {
1069    pub store_id: String,
1070    pub role: String,
1071    pub database: Option<String>,
1072    #[serde(flatten)]
1073    pub timing: khive_db::checkpoint::CheckpointTiming,
1074}
1075
1076/// Point-in-time snapshot of the daemon's server-side gauges — the
1077/// load/perf harness read-surface (measurement substrate, not a product feature).
1078///
1079/// Every field here is a **server-side** gauge reachable from `handle_conn`
1080/// without any mutation: [`khive_storage::tx_registry`] (ADR-091 Plank 0,
1081/// process-global singleton), the main pool's backend-keyed routine WAL
1082/// sample and process TRUNCATE counters (`khive_db::checkpoint`), and the
1083/// ADR-067 Component A write queue depth/latest writer-stage sample. There is
1084/// no reset reachable through this type or through [`DaemonRequestFrame`] —
1085/// gauges out, nothing in.
1086#[derive(Serialize, Deserialize, Debug, Clone, Default, PartialEq)]
1087pub struct MetricsSnapshot {
1088    /// Launch mode and the current lifecycle of this daemon incarnation.
1089    #[serde(default, skip_serializing_if = "Option::is_none")]
1090    pub lifecycle: Option<DaemonLifecycleSnapshot>,
1091    /// Last-observed WAL page count from the periodic checkpoint tick.
1092    /// `None` when the checkpoint task has never ticked in this process
1093    /// (for example, an in-memory dispatcher with no pool, or a daemon that
1094    /// just started and hasn't hit its first tick yet).
1095    pub wal_pages: Option<u64>,
1096    /// Logical frames present in the WAL at the main backend's most recent
1097    /// periodic PASSIVE checkpoint. Kept separate from physical allocation.
1098    #[serde(default)]
1099    pub wal_log_frames: Option<u64>,
1100    /// Frames backfilled by that same periodic PASSIVE pass.
1101    #[serde(default)]
1102    pub wal_checkpointed_frames: Option<u64>,
1103    /// Logical frames still pending after that pass (`log - checkpointed`).
1104    #[serde(default)]
1105    pub wal_pending_frames: Option<u64>,
1106    /// Physical `-wal` sidecar bytes captured at the same routine tick.
1107    #[serde(default)]
1108    pub wal_physical_bytes: Option<u64>,
1109    /// Wall-clock timestamp of the routine WAL sample, for staleness checks.
1110    #[serde(default)]
1111    pub wal_observed_at_unix_ms: Option<u64>,
1112    /// Cumulative actual routine PASSIVE calls for each store this daemon
1113    /// checkpoints. Microseconds preserve sub-millisecond costs; max is since
1114    /// process start, not an interval maximum. Counters saturate rather than wrap.
1115    #[serde(default)]
1116    pub wal_checkpoint_stores: Vec<CheckpointStoreMetrics>,
1117    /// Total WAL TRUNCATE escalation attempts (ADR-091 Plank 2) made in this
1118    /// process's lifetime, regardless of whether they succeeded in reclaiming
1119    /// pages.
1120    pub wal_truncate_attempts: u64,
1121    /// Current consecutive-failure count for TRUNCATE attempts that failed to
1122    /// bring the WAL back below `warn_pages`; resets to 0 the next time an
1123    /// attempt clears it.
1124    pub wal_truncate_consecutive_failures: u64,
1125    /// Total checkpoint ticks skipped because the dedicated checkpoint
1126    /// connection was unavailable (ADR-091 checkpoint-pressure telemetry),
1127    /// across this process's lifetime. `#[serde(default)]` so an older client
1128    /// decoding a newer daemon's snapshot (or vice versa) does not fail.
1129    #[serde(default)]
1130    pub wal_checkpoint_skipped_ticks: u64,
1131    /// Current consecutive-skip run length; 0 once the next tick is observed.
1132    #[serde(default)]
1133    pub wal_checkpoint_consecutive_skips: u64,
1134    /// WAL page count last known at the time of the most recent skip, if any
1135    /// skip has occurred yet in this process.
1136    #[serde(default)]
1137    pub wal_checkpoint_last_skip_wal_pages: Option<u64>,
1138    /// Age, in microseconds, of the oldest currently-open transaction
1139    /// registry entry (ADR-091 Plank 0). `None` when no transaction is
1140    /// currently open.
1141    pub oldest_pinned_tx_micros: Option<u64>,
1142    /// Diagnostic label of the oldest currently-open transaction registry
1143    /// entry, if any and if it was registered with one.
1144    pub oldest_pinned_tx_label: Option<String>,
1145    /// Number of currently open transaction registry entries.
1146    pub open_tx_count: usize,
1147    /// Current write-queue backlog depth (ADR-067 Component A): requests
1148    /// enqueued but not yet accepted by the `WriterTask` drain loop. `None`
1149    /// unless the write queue is enabled (`KHIVE_WRITE_QUEUE=1`) and a
1150    /// file-backed pool is available.
1151    pub write_queue_depth: Option<usize>,
1152    /// The write queue's configured bounded capacity
1153    /// (`PoolConfig::write_queue_capacity`), gated the same as
1154    /// `write_queue_depth`.
1155    pub write_queue_capacity: Option<usize>,
1156    /// Latest completed writer-task span: bounded-channel admission/backlog.
1157    #[serde(default)]
1158    pub write_last_queue_wait_micros: Option<u64>,
1159    /// Latest completed writer-task span: `BEGIN IMMEDIATE` acquisition.
1160    #[serde(default)]
1161    pub write_last_transaction_acquire_micros: Option<u64>,
1162    /// Latest completed writer-task span: application transaction body.
1163    #[serde(default)]
1164    pub write_last_body_micros: Option<u64>,
1165    /// Latest completed writer-task span: SQLite COMMIT/fsync phase.
1166    #[serde(default)]
1167    pub write_last_commit_micros: Option<u64>,
1168    /// Whole latest writer-task request span, retained for compatibility and
1169    /// comparison with the decomposed stages.
1170    #[serde(default)]
1171    pub write_last_total_micros: Option<u64>,
1172    /// Wall-clock timestamp of the writer-stage sample.
1173    #[serde(default)]
1174    pub write_last_observed_at_unix_ms: Option<u64>,
1175    /// Connection cap of the daemon socket this snapshot was served from.
1176    /// `None` from a daemon that predates the cap and when no listener owns
1177    /// the connection.
1178    #[serde(default, skip_serializing_if = "Option::is_none")]
1179    pub connections: Option<ConnectionCapSnapshot>,
1180    /// Bounds and counts of the best-effort recall serve-ledger tasks.
1181    /// `None` from a daemon that predates the bound.
1182    #[serde(default, skip_serializing_if = "Option::is_none")]
1183    pub recall_ledger: Option<RecallLedgerSnapshot>,
1184}
1185
1186// ── framing ───────────────────────────────────────────────────────────────────
1187
1188/// Read one length-prefixed frame (4-byte BE u32 length + JSON bytes).
1189#[cfg(unix)]
1190pub async fn read_frame<R>(stream: &mut R) -> std::io::Result<Vec<u8>>
1191where
1192    R: tokio::io::AsyncRead + Unpin,
1193{
1194    let mut len_buf = [0u8; 4];
1195    stream.read_exact(&mut len_buf).await?;
1196    let len = u32::from_be_bytes(len_buf) as usize;
1197    if len > MAX_FRAME_BYTES {
1198        return Err(std::io::Error::new(
1199            std::io::ErrorKind::InvalidData,
1200            format!("daemon frame of {len} bytes exceeds {MAX_FRAME_BYTES} cap"),
1201        ));
1202    }
1203    let mut buf = vec![0u8; len];
1204    stream.read_exact(&mut buf).await?;
1205    Ok(buf)
1206}
1207
1208#[cfg(unix)]
1209fn initial_frame_timeout_error() -> std::io::Error {
1210    std::io::Error::new(
1211        std::io::ErrorKind::TimedOut,
1212        "daemon initial request frame read timed out",
1213    )
1214}
1215
1216#[cfg(unix)]
1217async fn read_initial_frame<R>(
1218    stream: &mut R,
1219    deadline: tokio::time::Instant,
1220) -> std::io::Result<Vec<u8>>
1221where
1222    R: tokio::io::AsyncRead + Unpin,
1223{
1224    // Tokio polls the inner future before checking its timer. A frame already
1225    // buffered when a delayed connection task first runs would otherwise pass
1226    // even though its acceptance-time deadline has expired.
1227    if tokio::time::Instant::now() >= deadline {
1228        return Err(initial_frame_timeout_error());
1229    }
1230    let raw = tokio::time::timeout_at(deadline, read_frame(stream))
1231        .await
1232        .map_err(|_| initial_frame_timeout_error())??;
1233    if tokio::time::Instant::now() >= deadline {
1234        return Err(initial_frame_timeout_error());
1235    }
1236    Ok(raw)
1237}
1238
1239/// Write one length-prefixed frame.
1240#[cfg(unix)]
1241pub async fn write_frame<W>(stream: &mut W, payload: &[u8]) -> std::io::Result<()>
1242where
1243    W: tokio::io::AsyncWrite + Unpin,
1244{
1245    if payload.len() > MAX_FRAME_BYTES {
1246        return Err(std::io::Error::new(
1247            std::io::ErrorKind::InvalidData,
1248            format!(
1249                "daemon frame of {} bytes exceeds {MAX_FRAME_BYTES} cap",
1250                payload.len()
1251            ),
1252        ));
1253    }
1254    let len = (payload.len() as u32).to_be_bytes();
1255    stream.write_all(&len).await?;
1256    stream.write_all(payload).await?;
1257    stream.flush().await?;
1258    Ok(())
1259}
1260
1261// ── dispatch trait ────────────────────────────────────────────────────────────
1262
1263/// Transport-agnostic dispatch interface for the daemon server.
1264///
1265/// The MCP crate implements this by dispatching through the shared request body
1266/// while honoring [`DaemonRequestFrame::from_wire`] (so subhandler visibility is
1267/// gated by request origin, not by transport); any future transport can do the
1268/// same.
1269#[cfg(unix)]
1270#[async_trait]
1271pub trait DaemonDispatch: Clone + Send + Sync + 'static {
1272    /// Named retained resources or unknown inventory that prevents idle exit.
1273    /// An implementor must explicitly account for its resources before retiring.
1274    fn idle_retirement_blockers(&self) -> Vec<String> {
1275        vec!["dispatcher_resource_inventory_unknown".to_owned()]
1276    }
1277
1278    /// Describe syntax and loaded catalog membership without dispatching.
1279    fn plan(&self, ops: &str) -> String;
1280
1281    /// Dispatch a verb-DSL request string and return the rendered result.
1282    ///
1283    /// `from_wire` carries the origin discriminator from
1284    /// [`DaemonRequestFrame::from_wire`]: when `true`, the implementor enforces
1285    /// verb visibility (rejects `Visibility::Subhandler` verbs); when `false`,
1286    /// the request is from a trusted operator surface and subhandlers pass.
1287    ///
1288    /// `identity` is the per-request identity context threaded from the frame
1289    /// (ADR-096): `Some(..)` when serving a request forwarded over the
1290    /// daemon socket (built from `frame.namespace` / `frame.actor_id` /
1291    /// `frame.visible_namespaces` by the connection handler), `None` for any
1292    /// other dispatch path. Implementors should mint the storage/gate token from
1293    /// `identity` when present and fall back to their own construction-baked
1294    /// identity when absent, so pure local (non-daemon) dispatch is unchanged.
1295    #[allow(clippy::too_many_arguments)]
1296    async fn dispatch(
1297        &self,
1298        ops: String,
1299        presentation: Option<String>,
1300        presentation_per_op: Option<Vec<Option<String>>>,
1301        format: Option<String>,
1302        format_per_op: Option<Vec<Option<String>>>,
1303        from_wire: bool,
1304        identity: Option<RequestIdentity>,
1305    ) -> Result<String, String>;
1306
1307    /// Read-deadline ceiling for one request. The default is the operator
1308    /// ceiling; an implementor that understands the request may grant a
1309    /// longer bounded allowance (a long poll's declared wait plus a margin).
1310    fn request_read_timeout(&self, _ops: &str) -> std::time::Duration {
1311        khive_storage::request_read_timeout_from_env()
1312    }
1313
1314    /// Preserve structured dispatch errors without breaking string-only implementors.
1315    #[allow(clippy::too_many_arguments)]
1316    async fn dispatch_with_error_detail(
1317        &self,
1318        ops: String,
1319        presentation: Option<String>,
1320        presentation_per_op: Option<Vec<Option<String>>>,
1321        format: Option<String>,
1322        format_per_op: Option<Vec<Option<String>>>,
1323        from_wire: bool,
1324        identity: Option<RequestIdentity>,
1325    ) -> Result<String, DaemonDispatchError> {
1326        self.dispatch(
1327            ops,
1328            presentation,
1329            presentation_per_op,
1330            format,
1331            format_per_op,
1332            from_wire,
1333            identity,
1334        )
1335        .await
1336        .map_err(|message| DaemonDispatchError::new(message, None))
1337    }
1338
1339    /// Warm every pack's in-memory state (ANN indexes, etc.).
1340    async fn warm_all(&self);
1341
1342    /// The namespace this dispatcher was configured for.
1343    fn namespace(&self) -> &str;
1344
1345    /// Fingerprint of this dispatcher's resolved runtime config (packs, db
1346    /// target, embedders). Used to reject forwarded requests from clients whose
1347    /// config differs, so a restricted client cannot dispatch through a broader
1348    /// daemon.
1349    fn config_id(&self) -> &str;
1350
1351    /// Return the pool to use for background WAL checkpointing, if available.
1352    ///
1353    /// Implementors backed by a file-based SQLite database should return
1354    /// `Some(pool_arc)`. In-memory or test dispatchers that have no pool
1355    /// return `None` and the checkpoint task is not spawned.
1356    ///
1357    /// The default implementation returns `None`.
1358    fn pool_for_checkpoint(&self) -> Option<Arc<ConnectionPool>> {
1359        None
1360    }
1361
1362    /// File-backed backend pools beyond [`Self::pool_for_checkpoint`]'s pool
1363    /// (ADR-091 Amendment 3): one checkpoint task is spawned per entry here,
1364    /// in addition to the one spawned for the primary pool, so a
1365    /// multi-backend deployment gets PASSIVE/TRUNCATE checkpointing and
1366    /// sidecar enumeration on every file-backed backend it wired, not only
1367    /// the main one.
1368    ///
1369    /// The default implementation returns an empty `Vec` — an implementor
1370    /// with only one backend (or none) needs no override.
1371    fn secondary_pools_for_checkpoint(&self) -> Vec<Arc<ConnectionPool>> {
1372        Vec::new()
1373    }
1374
1375    /// Return the audit `EventStore` the checkpoint task should append
1376    /// ADR-094 lifecycle events (`CheckpointOutcomeRecorded`) to, if any.
1377    ///
1378    /// Mirrors [`Self::pool_for_checkpoint`]'s default-`None` shape: an
1379    /// implementor with no configured event store (or no pool at all) simply
1380    /// gets a checkpoint task that never appends events — the checkpoint
1381    /// task itself remains fully functional either way.
1382    ///
1383    /// The default implementation returns `None`.
1384    fn event_store_for_checkpoint(&self) -> Option<Arc<dyn khive_storage::EventStore>> {
1385        None
1386    }
1387}
1388
1389#[cfg(unix)]
1390struct CheckpointTaskSpec {
1391    pool: Arc<ConnectionPool>,
1392    lifecycle_owner: Option<CheckpointLifecycleOwner>,
1393    is_main: bool,
1394}
1395
1396/// Build checkpoint-task fan-out and designate one lifecycle owner.
1397///
1398/// The main checkpoint task owns lifecycle emission when it exists. If the
1399/// main backend is in-memory and therefore has no checkpoint task, the first
1400/// file-backed secondary owns emission instead. All remaining tasks are
1401/// explicit non-owners.
1402#[cfg(unix)]
1403fn checkpoint_task_specs(
1404    main_pool: Option<Arc<ConnectionPool>>,
1405    secondary_pools: Vec<Arc<ConnectionPool>>,
1406    event_store: Option<Arc<dyn khive_storage::EventStore>>,
1407    namespace: String,
1408) -> Vec<CheckpointTaskSpec> {
1409    let mut tasks = Vec::with_capacity(usize::from(main_pool.is_some()) + secondary_pools.len());
1410    if let Some(pool) = main_pool {
1411        tasks.push(CheckpointTaskSpec {
1412            pool,
1413            lifecycle_owner: None,
1414            is_main: true,
1415        });
1416    }
1417    tasks.extend(secondary_pools.into_iter().map(|pool| CheckpointTaskSpec {
1418        pool,
1419        lifecycle_owner: None,
1420        is_main: false,
1421    }));
1422
1423    if let (Some(task), Some(event_store)) = (tasks.first_mut(), event_store) {
1424        task.lifecycle_owner = Some(CheckpointLifecycleOwner::new(event_store, namespace));
1425    }
1426    tasks
1427}
1428
1429// ── tracked background tasks ─────────────────────────────────────────────────
1430//
1431// Pack handlers (e.g. memory.recall's ADR-081 serve-ledger append) fire
1432// fire-and-forget `tokio::spawn`ed work off the response path so the caller
1433// never waits on a cross-pack dispatch or a SQL write. Left untracked, that
1434// work is invisible to `drain()`: a SIGTERM landing between the response
1435// returning and the spawned task completing can abort it mid-flight with no
1436// log and no row. `track_background_task` gives such spawns a process-wide
1437// presence that `drain()` waits on, exactly like the `active` counter does
1438// for in-flight connections: the caller still only pays for the spawn +
1439// counter increment, never the task's own work.
1440/// Set once by the boot path that takes the daemon role, and never cleared: a
1441/// process that is not the warm daemon has no path to becoming one except exec.
1442static WARM_INDEX_HOST: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false);
1443
1444/// Declare this process the warm index host. Called by the serve path as soon as
1445/// the daemon role is decided, before any runtime is built, so nothing warms
1446/// under the wrong answer.
1447pub fn mark_warm_index_host() {
1448    WARM_INDEX_HOST.store(true, std::sync::atomic::Ordering::Release);
1449}
1450
1451/// Whether this process is the warm index host.
1452///
1453/// Building an ANN index from the full corpus is minutes of CPU and hundreds of
1454/// megabytes of segment rewrite, and it pays for itself only across a process
1455/// that outlives the request. A short-lived client that does it pays the whole
1456/// cost, discards the result at exit, and publishes a checkpoint that every
1457/// other reader on the root must then re-read. Consumers use this to decide
1458/// whether to build or to serve degraded and let the daemon build.
1459pub fn is_warm_index_host() -> bool {
1460    WARM_INDEX_HOST.load(std::sync::atomic::Ordering::Acquire)
1461}
1462
1463static BACKGROUND_TASKS: std::sync::OnceLock<Arc<std::sync::atomic::AtomicUsize>> =
1464    std::sync::OnceLock::new();
1465
1466fn background_tasks() -> &'static Arc<std::sync::atomic::AtomicUsize> {
1467    BACKGROUND_TASKS.get_or_init(|| Arc::new(std::sync::atomic::AtomicUsize::new(0)))
1468}
1469
1470/// Decrements the shared background-task counter from `Drop`, so the count
1471/// comes back down whether the tracked future returns normally, panics, or
1472/// is cancelled — a plain post-`await` `fetch_sub` only covers the return
1473/// path and leaks the count forever on a panic, since unwinding skips every
1474/// statement after the panic point.
1475// ── outstanding background-task names ────────────────────────────────────────
1476//
1477// Names of the tasks the background counter is currently holding, so a drain
1478// timeout can say which ones held it open instead of printing a bare count.
1479// Registered and released at exactly the points the counter is incremented and
1480// decremented, and in the order that keeps the counter authoritative: the name
1481// goes in after the increment and comes out before the decrement, so a reported
1482// name always belongs to a task the counter already holds. The reverse ordering
1483// would let the warning name a task that had already finished, which is the one
1484// reading that would send someone looking in the wrong place.
1485static BACKGROUND_TASK_NAMES: std::sync::OnceLock<
1486    std::sync::Mutex<std::collections::HashMap<&'static str, usize>>,
1487> = std::sync::OnceLock::new();
1488
1489fn background_task_names_registry(
1490) -> &'static std::sync::Mutex<std::collections::HashMap<&'static str, usize>> {
1491    BACKGROUND_TASK_NAMES.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
1492}
1493
1494/// Name recorded for tasks spawned through the unnamed entry points, which
1495/// stay on the public API of a published crate.
1496pub const UNNAMED_BACKGROUND_TASK: &str = "unnamed";
1497
1498fn register_background_task_name(name: &'static str) {
1499    let mut names = background_task_names_registry()
1500        .lock()
1501        .unwrap_or_else(std::sync::PoisonError::into_inner);
1502    *names.entry(name).or_insert(0) += 1;
1503}
1504
1505fn release_background_task_name(name: &'static str) {
1506    let mut names = background_task_names_registry()
1507        .lock()
1508        .unwrap_or_else(std::sync::PoisonError::into_inner);
1509    if let Some(count) = names.get_mut(name) {
1510        *count -= 1;
1511        if *count == 0 {
1512            names.remove(name);
1513        }
1514    }
1515}
1516
1517/// Names of the in-flight tracked background tasks, sorted and deduplicated.
1518/// A diagnostic beside [`background_task_count`], never a substitute for it:
1519/// the count is what drain waits on.
1520pub fn background_task_names() -> Vec<String> {
1521    let names = background_task_names_registry()
1522        .lock()
1523        .unwrap_or_else(std::sync::PoisonError::into_inner);
1524    let mut out: Vec<String> = names.keys().map(|name| (*name).to_string()).collect();
1525    out.sort();
1526    out
1527}
1528
1529#[cfg(unix)]
1530fn idle_retirement_blockers<D: DaemonDispatch>(dispatcher: &D) -> Vec<String> {
1531    let mut blockers = dispatcher.idle_retirement_blockers();
1532    if !khive_storage::tx_registry::snapshot().is_empty() {
1533        blockers.push("open_sql_transaction".to_owned());
1534    }
1535    blockers.extend(
1536        active_phase_names()
1537            .into_iter()
1538            .map(|name| format!("active_phase:{name}")),
1539    );
1540    let count = background_task_count();
1541    let names = background_task_names_registry()
1542        .lock()
1543        .unwrap_or_else(std::sync::PoisonError::into_inner);
1544    if names.values().sum::<usize>() != count {
1545        blockers.push("tracked_worker_inventory_unsettled".to_owned());
1546    }
1547    // Only these inspected loops maintain replaceable caches/checkpoint state.
1548    // Their tracked lifetime still participates in the final drain.
1549    for name in names.keys() {
1550        if !matches!(
1551            *name,
1552            "wal_checkpoint" | "memory_ann_rotation_watch" | "knowledge_ann_rotation_watch"
1553        ) {
1554            blockers.push(format!("unsettled_worker:{name}"));
1555        }
1556    }
1557    blockers.sort();
1558    blockers.dedup();
1559    blockers
1560}
1561
1562#[cfg(unix)]
1563async fn wait_for_idle<D: DaemonDispatch>(dispatcher: &D, lifecycle: &DaemonLifecycle) {
1564    if lifecycle.options.lifetime == DaemonLifetime::Persistent {
1565        std::future::pending::<()>().await;
1566    }
1567    loop {
1568        if lifecycle.try_idle(|| idle_retirement_blockers(dispatcher)) {
1569            return;
1570        }
1571        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
1572    }
1573}
1574
1575struct BackgroundTaskGuard {
1576    counter: Arc<std::sync::atomic::AtomicUsize>,
1577    name: &'static str,
1578}
1579
1580impl Drop for BackgroundTaskGuard {
1581    fn drop(&mut self) {
1582        release_background_task_name(self.name);
1583        self.counter
1584            .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
1585    }
1586}
1587
1588/// Spawn a task that daemon shutdown's `drain()` waits for and return its join
1589/// handle. Retaining the handle lets boot coordinators form an explicit barrier;
1590/// dropping it deliberately detaches the task while the background counter still
1591/// keeps daemon drain aware of its lifetime. The decrement happens via
1592/// `BackgroundTaskGuard`'s `Drop`, including panic and cancellation paths.
1593pub fn spawn_tracked_task<F, T>(fut: F) -> tokio::task::JoinHandle<T>
1594where
1595    F: std::future::Future<Output = T> + Send + 'static,
1596    T: Send + 'static,
1597{
1598    spawn_named_tracked_task(UNNAMED_BACKGROUND_TASK, fut)
1599}
1600
1601/// [`spawn_tracked_task`] with a name that a drain timeout can print.
1602///
1603/// The name is a short static string describing the task, never a formatted or
1604/// caller-supplied value: it is read by an operator staring at a shutdown that
1605/// would not finish, so it is a label for a call site, not a record of one
1606/// occurrence.
1607pub fn spawn_named_tracked_task<F, T>(name: &'static str, fut: F) -> tokio::task::JoinHandle<T>
1608where
1609    F: std::future::Future<Output = T> + Send + 'static,
1610    T: Send + 'static,
1611{
1612    background_tasks().fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1613    register_background_task_name(name);
1614    let guard = BackgroundTaskGuard {
1615        counter: background_tasks().clone(),
1616        name,
1617    };
1618    tokio::spawn(async move {
1619        let _guard = guard;
1620        fut.await
1621    })
1622}
1623
1624/// Spawn a fire-and-forget task through [`spawn_tracked_task`].
1625///
1626/// Callers that need a boot or shutdown barrier should retain and await the
1627/// returned handle from [`spawn_tracked_task`] instead of detaching it here.
1628pub fn track_background_task<F>(fut: F)
1629where
1630    F: std::future::Future<Output = ()> + Send + 'static,
1631{
1632    track_named_background_task(UNNAMED_BACKGROUND_TASK, fut);
1633}
1634
1635/// [`track_background_task`] with a name that a drain timeout can print. See
1636/// [`spawn_named_tracked_task`] for what belongs in the name.
1637pub fn track_named_background_task<F>(name: &'static str, fut: F)
1638where
1639    F: std::future::Future<Output = ()> + Send + 'static,
1640{
1641    drop(spawn_named_tracked_task(name, fut));
1642}
1643
1644/// Current count of in-flight tasks started via [`track_background_task`].
1645/// Exposed for tests; `drain()` reads the shared counter directly.
1646pub fn background_task_count() -> usize {
1647    background_tasks().load(std::sync::atomic::Ordering::SeqCst)
1648}
1649
1650/// Process-wide daemon shutdown signal (ADR-119).
1651///
1652/// Cancelled exactly once, when the daemon's unified shutdown future resolves
1653/// — before `drain()` begins waiting on tracked tasks — so long-running
1654/// daemon components supervised outside this module observe shutdown through
1655/// the same path the daemon itself does, rather than inventing their own.
1656/// Clones share the underlying token; child tokens derived from it are
1657/// cancelled transitively.
1658///
1659/// In non-daemon processes the token simply never fires.
1660pub fn daemon_shutdown_token() -> tokio_util::sync::CancellationToken {
1661    static TOKEN: std::sync::OnceLock<tokio_util::sync::CancellationToken> =
1662        std::sync::OnceLock::new();
1663    TOKEN
1664        .get_or_init(tokio_util::sync::CancellationToken::new)
1665        .clone()
1666}
1667
1668// ── active background phase names (ADR-103) ──────────────────────────────────
1669//
1670// A lightweight, best-effort process-wide gauge of which named background
1671// phases (e.g. `ann_warm`) are in flight right now, read by `comm.health`'s
1672// resource self-report so a caller can see "what is the daemon doing" at a
1673// glance without correlating timestamps across the event log itself. Counted
1674// per name rather than boolean, since more than one occurrence of the same
1675// named phase can legitimately overlap (e.g. two embedding models warming
1676// concurrently) — the name only drops out of the reported set once every
1677// concurrent occurrence has ended.
1678static ACTIVE_PHASES: std::sync::OnceLock<
1679    std::sync::Mutex<std::collections::HashMap<String, usize>>,
1680> = std::sync::OnceLock::new();
1681
1682fn active_phases() -> &'static std::sync::Mutex<std::collections::HashMap<String, usize>> {
1683    ACTIVE_PHASES.get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new()))
1684}
1685
1686/// RAII guard for one occurrence of a named background phase. Increments the
1687/// phase's count on creation (see [`register_active_phase`]); decrements on
1688/// `Drop`, so the count comes back down whether the guarded work returns
1689/// normally, panics, or is cancelled — the same rationale as
1690/// `BackgroundTaskGuard` above.
1691pub struct PhaseGuard {
1692    name: String,
1693}
1694
1695impl Drop for PhaseGuard {
1696    fn drop(&mut self) {
1697        let mut map = active_phases()
1698            .lock()
1699            .unwrap_or_else(std::sync::PoisonError::into_inner);
1700        if let Some(count) = map.get_mut(&self.name) {
1701            *count -= 1;
1702            if *count == 0 {
1703                map.remove(&self.name);
1704            }
1705        }
1706    }
1707}
1708
1709/// Register one occurrence of a named background phase as currently active.
1710/// Returns a guard: drop it (or let it fall out of scope) when the phase
1711/// ends. Best-effort process-wide gauge only, read by `comm.health` — never
1712/// load-bearing for correctness.
1713pub fn register_active_phase(name: &str) -> PhaseGuard {
1714    let mut map = active_phases()
1715        .lock()
1716        .unwrap_or_else(std::sync::PoisonError::into_inner);
1717    *map.entry(name.to_string()).or_insert(0) += 1;
1718    PhaseGuard {
1719        name: name.to_string(),
1720    }
1721}
1722
1723/// Currently active background-phase names, sorted for deterministic output.
1724/// Empty when no tracked phase is in flight.
1725pub fn active_phase_names() -> Vec<String> {
1726    let map = active_phases()
1727        .lock()
1728        .unwrap_or_else(std::sync::PoisonError::into_inner);
1729    let mut names: Vec<String> = map.keys().cloned().collect();
1730    names.sort();
1731    names
1732}
1733
1734// ── server ────────────────────────────────────────────────────────────────────
1735
1736/// Build a point-in-time [`MetricsSnapshot`] of this process's server-side
1737/// gauges. Called only from `handle_conn`'s `metrics_only` arm — a
1738/// process-global, read-only assembly with no side effects of its own.
1739/// See `docs/api/daemon.md#build_metrics_snapshot` for where each gauge is sourced
1740/// from and why.
1741#[cfg(unix)]
1742fn build_metrics_snapshot<D: DaemonDispatch>(dispatcher: &D) -> MetricsSnapshot {
1743    let open_tx_count = khive_storage::tx_registry::snapshot().len();
1744    // ADR-091 Amendment 3: deliberately the process-wide aggregate, not a
1745    // backend-attributed view — this gauge reports "oldest pinned tx in this
1746    // process" across every wired backend, not an attribution claim about
1747    // which database it belongs to. Attribution consumers (the session sweep,
1748    // the per-backend checkpoint tasks) use the scoped `oldest_for` views in
1749    // khive-db instead.
1750    let (oldest_pinned_tx_micros, oldest_pinned_tx_label) =
1751        match khive_storage::tx_registry::oldest() {
1752            Some((_id, age, label)) => (Some(age.as_micros() as u64), label),
1753            None => (None, None),
1754        };
1755
1756    let checkpoint_pool = dispatcher.pool_for_checkpoint();
1757    let mut secondary_index = 0;
1758    let wal_checkpoint_stores = checkpoint_task_specs(
1759        checkpoint_pool.clone(),
1760        dispatcher.secondary_pools_for_checkpoint(),
1761        None,
1762        String::new(),
1763    )
1764    .into_iter()
1765    .map(|task| {
1766        let (store_id, role) = if task.is_main {
1767            ("main".to_string(), "main".to_string())
1768        } else {
1769            let store_id = format!("secondary:{secondary_index}");
1770            secondary_index += 1;
1771            (store_id, "secondary".to_string())
1772        };
1773        CheckpointStoreMetrics {
1774            store_id,
1775            role,
1776            database: task
1777                .pool
1778                .canonical_path()
1779                .and_then(std::path::Path::file_name)
1780                .map(|name| name.to_string_lossy().into_owned()),
1781            timing: khive_db::checkpoint::checkpoint_timing(&task.pool),
1782        }
1783    })
1784    .collect();
1785    let routine_wal = checkpoint_pool
1786        .as_deref()
1787        .and_then(khive_db::checkpoint::routine_wal_observation);
1788    let writer_stages = checkpoint_pool
1789        .as_deref()
1790        .and_then(khive_db::writer_task::last_writer_stage_observation);
1791    let (write_queue_depth, write_queue_capacity) = checkpoint_pool
1792        .as_ref()
1793        .and_then(|pool| pool.writer_task_handle().ok().flatten())
1794        .map(|handle| (Some(handle.queue_depth()), Some(handle.capacity())))
1795        .unwrap_or((None, None));
1796
1797    MetricsSnapshot {
1798        lifecycle: None,
1799        wal_pages: routine_wal.as_ref().map(|sample| sample.log_frames),
1800        wal_log_frames: routine_wal.as_ref().map(|sample| sample.log_frames),
1801        wal_checkpointed_frames: routine_wal
1802            .as_ref()
1803            .map(|sample| sample.checkpointed_frames),
1804        wal_pending_frames: routine_wal.as_ref().map(|sample| sample.pending_frames),
1805        wal_physical_bytes: routine_wal
1806            .as_ref()
1807            .and_then(|sample| sample.physical_wal_bytes),
1808        wal_observed_at_unix_ms: routine_wal
1809            .as_ref()
1810            .map(|sample| sample.observed_at_unix_ms),
1811        wal_checkpoint_stores,
1812        wal_truncate_attempts: khive_db::checkpoint::truncate_attempts(),
1813        wal_truncate_consecutive_failures: khive_db::checkpoint::truncate_consecutive_failures(),
1814        wal_checkpoint_skipped_ticks: khive_db::checkpoint::checkpoint_skipped_ticks(),
1815        wal_checkpoint_consecutive_skips: khive_db::checkpoint::checkpoint_consecutive_skips(),
1816        wal_checkpoint_last_skip_wal_pages: khive_db::checkpoint::checkpoint_last_skip_wal_pages(),
1817        oldest_pinned_tx_micros,
1818        oldest_pinned_tx_label,
1819        open_tx_count,
1820        write_queue_depth,
1821        write_queue_capacity,
1822        write_last_queue_wait_micros: writer_stages
1823            .as_ref()
1824            .map(|sample| sample.queue_wait_micros),
1825        write_last_transaction_acquire_micros: writer_stages
1826            .as_ref()
1827            .map(|sample| sample.transaction_acquire_micros),
1828        write_last_body_micros: writer_stages.as_ref().map(|sample| sample.body_micros),
1829        write_last_commit_micros: writer_stages.as_ref().map(|sample| sample.commit_micros),
1830        write_last_total_micros: writer_stages.as_ref().map(|sample| sample.total_micros),
1831        write_last_observed_at_unix_ms: writer_stages
1832            .as_ref()
1833            .map(|sample| sample.observed_at_unix_ms),
1834        connections: None,
1835        recall_ledger: Some(recall_ledger_snapshot()),
1836    }
1837}
1838
1839#[cfg(unix)]
1840async fn write_response_frame<W>(stream: &mut W, payload: &[u8]) -> std::io::Result<()>
1841where
1842    W: tokio::io::AsyncWrite + Unpin,
1843{
1844    tokio::time::timeout(INITIAL_FRAME_READ_TIMEOUT, write_frame(stream, payload))
1845        .await
1846        .map_err(|_| {
1847            std::io::Error::new(
1848                std::io::ErrorKind::TimedOut,
1849                "daemon response write timed out",
1850            )
1851        })?
1852}
1853
1854#[cfg(unix)]
1855async fn wait_for_peer_disconnect(read: &mut tokio::net::unix::OwnedReadHalf) {
1856    let mut byte = [0u8; 1];
1857    // One request is admitted per connection. EOF, a read error, or any
1858    // subsequent byte all make this connection no longer a valid response
1859    // peer, so the first completed read is the entire observation.
1860    let _ = read.read(&mut byte).await;
1861}
1862
1863#[cfg(all(unix, test))]
1864async fn handle_conn<D: DaemonDispatch>(stream: UnixStream, dispatcher: D) {
1865    handle_conn_with_shutdown(
1866        stream,
1867        dispatcher,
1868        None,
1869        tokio::time::Instant::now() + INITIAL_FRAME_READ_TIMEOUT,
1870    )
1871    .await;
1872}
1873
1874#[cfg(all(unix, feature = "fault-injection"))]
1875#[doc(hidden)]
1876pub async fn handle_conn_for_test<D: DaemonDispatch>(stream: UnixStream, dispatcher: D) {
1877    handle_conn_with_shutdown(
1878        stream,
1879        dispatcher,
1880        None,
1881        tokio::time::Instant::now() + INITIAL_FRAME_READ_TIMEOUT,
1882    )
1883    .await;
1884}
1885
1886#[cfg(unix)]
1887fn plan_frame_companion(raw: &[u8]) -> Option<&'static str> {
1888    let value: serde_json::Value = serde_json::from_slice(raw).ok()?;
1889    if value.get("plan").and_then(serde_json::Value::as_bool) != Some(true)
1890        || value
1891            .get("protocol_version")
1892            .and_then(serde_json::Value::as_u64)
1893            != Some(u64::from(PROTOCOL_VERSION))
1894    {
1895        return None;
1896    }
1897    [
1898        "presentation",
1899        "presentation_per_op",
1900        "format",
1901        "format_per_op",
1902        "request_id",
1903    ]
1904    .into_iter()
1905    .find(|field| value.get(*field).is_some())
1906}
1907
1908#[cfg(all(
1909    unix,
1910    any(test, feature = "fault-injection", feature = "test-internals")
1911))]
1912async fn handle_conn_with_shutdown<D: DaemonDispatch>(
1913    stream: UnixStream,
1914    dispatcher: D,
1915    shutdown: Option<tokio::sync::watch::Receiver<bool>>,
1916    initial_frame_deadline: tokio::time::Instant,
1917) {
1918    handle_conn_with_lifecycle(stream, dispatcher, shutdown, initial_frame_deadline, None).await;
1919}
1920
1921#[cfg(unix)]
1922async fn handle_conn_with_lifecycle<D: DaemonDispatch>(
1923    mut stream: UnixStream,
1924    dispatcher: D,
1925    shutdown: Option<tokio::sync::watch::Receiver<bool>>,
1926    initial_frame_deadline: tokio::time::Instant,
1927    lifecycle: Option<Arc<DaemonLifecycle>>,
1928) {
1929    let mut ordinary_admission = None;
1930    let production_shutdown = shutdown.is_some();
1931    // A handover is delivered over the probed connection. The production
1932    // listener enforces same-uid admission, and direct handler tests cannot
1933    // self-signal because they do not supply the daemon shutdown receiver.
1934    let handover_peer_allowed = peer_uid(&stream)
1935        .ok()
1936        .is_some_and(|uid| uid == unsafe { libc::geteuid() } as u32);
1937    let (local_shutdown_tx, local_shutdown_rx) = tokio::sync::watch::channel(false);
1938    let shutdown = shutdown.unwrap_or(local_shutdown_rx);
1939    // Keeps the fallback receiver open in direct/test calls. Production owns
1940    // a sender at the daemon-run scope and passes its receiver above.
1941    let _local_shutdown_tx = local_shutdown_tx;
1942    let raw = match read_initial_frame(&mut stream, initial_frame_deadline).await {
1943        Ok(r) => r,
1944        Err(e) => {
1945            tracing::debug!(error = %e, "failed to read daemon request frame");
1946            return;
1947        }
1948    };
1949    #[derive(Deserialize)]
1950    struct SupervisorRequestEnvelope {
1951        #[serde(flatten)]
1952        frame: DaemonRequestFrame,
1953        #[serde(default)]
1954        supervisor_handover: bool,
1955    }
1956    let decoded: Result<SupervisorRequestEnvelope, _> = serde_json::from_slice(&raw);
1957    if decoded.as_ref().ok().is_none_or(|item| item.frame.plan) {
1958        if let Some(field) = plan_frame_companion(&raw) {
1959            let response = DaemonResponseFrame {
1960                ok: false,
1961                result: None,
1962                error: Some(format!(
1963                    "invalid_params: plan=true cannot be combined with {field}"
1964                )),
1965                error_detail: Some(serde_json::json!({
1966                    "kind": "protocol",
1967                    "code": "invalid_params",
1968                    "message": format!("plan=true cannot be combined with {field}"),
1969                    "domain_disposition": crate::DomainDisposition::NotCommitted.as_str(),
1970                })),
1971                namespace_mismatch: false,
1972                config_mismatch: false,
1973                served_config_id: Some(dispatcher.config_id().to_string()),
1974                version_mismatch: false,
1975                daemon_protocol_version: PROTOCOL_VERSION,
1976                metrics: None,
1977                request_id: None,
1978            };
1979            if let Ok(payload) = serde_json::to_vec(&response) {
1980                if let Err(error) = write_response_frame(&mut stream, &payload).await {
1981                    tracing::debug!(%error, "failed to write plan envelope refusal");
1982                }
1983            }
1984            return;
1985        }
1986    }
1987    let (frame, handover_requested) = match decoded {
1988        Ok(item) => (item.frame, item.supervisor_handover),
1989        Err(e) => {
1990            tracing::debug!(error = %e, "failed to decode daemon request frame");
1991            return;
1992        }
1993    };
1994    let supervisor_probe = frame.probe_only;
1995    let handover_accepted = handover_requested
1996        && supervisor_probe
1997        && production_shutdown
1998        && handover_peer_allowed
1999        && frame.protocol_version == PROTOCOL_VERSION
2000        && !frame.plan
2001        && !frame.metrics_only
2002        && frame.ops.is_empty()
2003        && read_supervisor_marker_claim().is_some_and(|(pid, _)| pid != std::process::id());
2004    let (mut peer_read, mut peer_write) = stream.into_split();
2005
2006    let served_config_id = Some(dispatcher.config_id().to_string());
2007    let resp = if frame.protocol_version != PROTOCOL_VERSION {
2008        let msg = format!(
2009            "daemon protocol mismatch: client={} daemon={} — \
2010             rebuild/update the client binary (make local)",
2011            frame.protocol_version, PROTOCOL_VERSION,
2012        );
2013        tracing::warn!(
2014            client_version = frame.protocol_version,
2015            daemon_version = PROTOCOL_VERSION,
2016            "daemon protocol version mismatch"
2017        );
2018        DaemonResponseFrame {
2019            ok: false,
2020            result: None,
2021            error: Some(msg.clone()),
2022            error_detail: Some(serde_json::json!({
2023                "kind": "protocol",
2024                "code": "version_mismatch",
2025                "message": msg,
2026                "domain_disposition": crate::DomainDisposition::Unknown.as_str(),
2027            })),
2028            namespace_mismatch: false,
2029            config_mismatch: false,
2030            served_config_id,
2031            // A client below this protocol is a bridge that predates the binary this
2032            // daemon was spawned from. Through protocol 5 the bridge treats an
2033            // explicit `version_mismatch` from a higher-numbered daemon as a terminal
2034            // error it repeats on every request, and re-execs itself onto the on-disk
2035            // binary only for the implicit shape: an unequal `daemon_protocol_version`
2036            // with the flag clear. Answering older clients in that shape, still
2037            // refused and still carrying the code in `error_detail`, lets every
2038            // pre-swap bridge replace itself on its first request instead of staying
2039            // refused until a person reconnects the session. A client above this
2040            // protocol keeps the explicit flag. Remove once no live bridge predates
2041            // the two-direction re-exec in khive-mcp (an inode census of `kkernel mcp`
2042            // processes before the swap): bridges built with it no longer read the
2043            // flag, but a bump while older bridges still run must keep this shape so
2044            // they replace themselves too.
2045            version_mismatch: frame.protocol_version > PROTOCOL_VERSION,
2046            daemon_protocol_version: PROTOCOL_VERSION,
2047            metrics: None,
2048            request_id: frame.request_id,
2049        }
2050    } else if handover_requested {
2051        if handover_accepted {
2052            DaemonResponseFrame {
2053                ok: true,
2054                result: None,
2055                error: None,
2056                error_detail: None,
2057                namespace_mismatch: false,
2058                config_mismatch: false,
2059                served_config_id,
2060                version_mismatch: false,
2061                daemon_protocol_version: PROTOCOL_VERSION,
2062                metrics: None,
2063                request_id: frame.request_id,
2064            }
2065        } else {
2066            DaemonResponseFrame {
2067                ok: false,
2068                result: None,
2069                error: Some("supervisor handover refused".to_string()),
2070                error_detail: Some(serde_json::json!({
2071                    "kind": "protocol",
2072                    "code": "supervisor_handover_refused",
2073                    "message": "supervisor handover refused",
2074                    "domain_disposition": crate::DomainDisposition::NotCommitted.as_str(),
2075                })),
2076                namespace_mismatch: false,
2077                config_mismatch: false,
2078                served_config_id,
2079                version_mismatch: false,
2080                daemon_protocol_version: PROTOCOL_VERSION,
2081                metrics: None,
2082                request_id: frame.request_id,
2083            }
2084        }
2085    } else if frame.metrics_only && !frame.plan {
2086        // Process-global gauge read: namespace/config-agnostic, so this is
2087        // handled BEFORE the `config_id` equality reject below (unlike every
2088        // other arm) — a metrics probe must work regardless of which
2089        // client's config is asking, since it never touches the dispatcher's
2090        // packs/db/embed registry. READ-ONLY: builds a snapshot and returns
2091        // immediately, never reaching the ops-dispatch arm.
2092        DaemonResponseFrame {
2093            ok: true,
2094            result: None,
2095            error: None,
2096            error_detail: None,
2097            namespace_mismatch: false,
2098            config_mismatch: false,
2099            served_config_id,
2100            version_mismatch: false,
2101            daemon_protocol_version: PROTOCOL_VERSION,
2102            metrics: Some({
2103                let mut metrics = build_metrics_snapshot(&dispatcher);
2104                metrics.lifecycle = lifecycle.as_ref().map(|state| {
2105                    let mut snapshot = state.snapshot();
2106                    snapshot.idle_blockers = idle_retirement_blockers(&dispatcher);
2107                    snapshot
2108                });
2109                metrics.connections = lifecycle.as_ref().map(|state| state.connections.snapshot());
2110                metrics
2111            }),
2112            request_id: frame.request_id,
2113        }
2114    // There is no `frame.namespace != dispatcher.namespace()` reject here.
2115    // The daemon accepts and serves the request under the frame's own
2116    // identity (namespace / actor / visible_namespaces, built into a
2117    // `RequestIdentity` below) over its one shared warm registry, rather
2118    // than rejecting a differently-attributed same-uid connection to a cold
2119    // local-dispatch fallback. `config_id`: which governs packs/db/embed
2120    // coherence for the shared warm engine: remains a hard reject for every
2121    // field other than a daemon-side superset of the client's extra embedders.
2122    } else if !config_ids_compatible(&frame.config_id, dispatcher.config_id()) {
2123        DaemonResponseFrame {
2124            ok: false,
2125            result: None,
2126            error: None,
2127            error_detail: Some(serde_json::json!({
2128                "kind": "protocol",
2129                "code": "config_mismatch",
2130                "message": "daemon configuration does not match the request",
2131                "domain_disposition": crate::DomainDisposition::NotCommitted.as_str(),
2132            })),
2133            namespace_mismatch: false,
2134            config_mismatch: true,
2135            served_config_id,
2136            version_mismatch: false,
2137            daemon_protocol_version: PROTOCOL_VERSION,
2138            metrics: None,
2139            request_id: frame.request_id,
2140        }
2141    } else if frame.plan {
2142        DaemonResponseFrame {
2143            ok: true,
2144            result: Some(dispatcher.plan(&frame.ops)),
2145            error: None,
2146            error_detail: None,
2147            namespace_mismatch: false,
2148            config_mismatch: false,
2149            served_config_id,
2150            version_mismatch: false,
2151            daemon_protocol_version: PROTOCOL_VERSION,
2152            metrics: None,
2153            request_id: None,
2154        }
2155    } else if frame.probe_only {
2156        // Probe-only request: identity checks passed; return immediately without
2157        // dispatching any verb. The client uses this to confirm the daemon is
2158        // alive and identity-matching without triggering any mutation.
2159        DaemonResponseFrame {
2160            ok: true,
2161            result: None,
2162            error: None,
2163            error_detail: None,
2164            namespace_mismatch: false,
2165            config_mismatch: false,
2166            served_config_id,
2167            version_mismatch: false,
2168            daemon_protocol_version: PROTOCOL_VERSION,
2169            metrics: None,
2170            request_id: frame.request_id,
2171        }
2172    } else {
2173        if let Some(lifecycle) = &lifecycle {
2174            ordinary_admission = lifecycle.admit();
2175            if ordinary_admission.is_none() {
2176                let refusal = DaemonResponseFrame {
2177                    ok: false,
2178                    result: None,
2179                    error: Some("daemon is draining; request was not admitted".to_owned()),
2180                    error_detail: Some(serde_json::json!({
2181                        "kind": "runtime", "code": "daemon_draining",
2182                        "domain_disposition": crate::DomainDisposition::NotCommitted.as_str(),
2183                    })),
2184                    request_id: frame.request_id,
2185                    daemon_protocol_version: PROTOCOL_VERSION,
2186                    namespace_mismatch: false,
2187                    config_mismatch: false,
2188                    served_config_id: Some(dispatcher.config_id().to_owned()),
2189                    version_mismatch: false,
2190                    metrics: None,
2191                };
2192                if let Ok(payload) = serde_json::to_vec(&refusal) {
2193                    let _ = write_response_frame(&mut peer_write, &payload).await;
2194                }
2195                return;
2196            }
2197        }
2198        // Build the per-request identity context from the frame so the
2199        // implementor mints the storage/gate token from the CALLER's
2200        // identity, not the dispatcher's own construction-baked scalars.
2201        // This is always `Some` here: every frame that reaches this arm
2202        // carries a `namespace` (required on the wire) plus whatever
2203        // `actor_id`/`visible_namespaces` the client resolved (defaulting to
2204        // `None`/`vec![]` for an older, field-absent payload, which is
2205        // exactly the prior anonymous/no-extra-visibility behavior).
2206        // The caller's actor namespace joins default reads where the registry
2207        // mints the token (ADR-007 Rev 4 Rule 3b), the one seam every identity
2208        // path shares; the frame's list is forwarded as sent.
2209        let identity = RequestIdentity {
2210            namespace: frame.namespace.clone(),
2211            actor_id: frame.actor_id.clone(),
2212            visible_namespaces: frame.visible_namespaces.clone(),
2213            process_ref: frame.process_ref.clone(),
2214            request_id: frame.request_id,
2215        };
2216        tracing::debug!(
2217            request_id = frame.request_id,
2218            "daemon RequestIdentity constructed"
2219        );
2220        let (read_cancel_tx, read_cancel_rx) = tokio::sync::watch::channel(false);
2221        // The connection's own ceiling nests inside the dispatcher's, and a
2222        // nested scope keeps the earlier deadline, so the allowance must be
2223        // granted here or a long poll times out at the operator ceiling.
2224        let read_timeout = dispatcher.request_read_timeout(&frame.ops);
2225        let excluded_embedder_names =
2226            config_id_extra_embedder_exclusions(&frame.config_id, dispatcher.config_id())
2227                .expect("compatible configuration ids must expose their extra embedder sets");
2228        let dispatch = crate::runtime::scope_request_embedder_exclusions(
2229            excluded_embedder_names,
2230            khive_storage::scope_request_read_cancellation(
2231                shutdown,
2232                khive_storage::scope_request_read_cancellation(
2233                    read_cancel_rx,
2234                    khive_storage::scope_request_read_deadline(
2235                        read_timeout,
2236                        dispatcher.dispatch_with_error_detail(
2237                            frame.ops,
2238                            frame.presentation,
2239                            frame.presentation_per_op,
2240                            frame.format,
2241                            frame.format_per_op,
2242                            frame.from_wire,
2243                            Some(identity),
2244                        ),
2245                    ),
2246                ),
2247            ),
2248        );
2249        tokio::pin!(dispatch);
2250        let dispatch_result = tokio::select! {
2251            result = &mut dispatch => result,
2252            _ = wait_for_peer_disconnect(&mut peer_read) => {
2253                let _ = read_cancel_tx.send(true);
2254                dispatch.await
2255            }
2256        };
2257        match dispatch_result {
2258            Ok(result) => {
2259                let (result, detail) = take_daemon_lexical_timeout_marker(result);
2260                DaemonResponseFrame {
2261                    ok: true,
2262                    result: Some(result),
2263                    error: None,
2264                    error_detail: detail,
2265                    namespace_mismatch: false,
2266                    config_mismatch: false,
2267                    served_config_id,
2268                    version_mismatch: false,
2269                    daemon_protocol_version: PROTOCOL_VERSION,
2270                    metrics: None,
2271                    request_id: frame.request_id,
2272                }
2273            }
2274            Err(error) => {
2275                let error = DaemonDispatchError::new(error.message, Some(error.error_detail));
2276                DaemonResponseFrame {
2277                    ok: false,
2278                    result: None,
2279                    error: Some(error.message),
2280                    error_detail: Some(error.error_detail),
2281                    namespace_mismatch: false,
2282                    config_mismatch: false,
2283                    served_config_id,
2284                    version_mismatch: false,
2285                    daemon_protocol_version: PROTOCOL_VERSION,
2286                    metrics: None,
2287                    request_id: frame.request_id,
2288                }
2289            }
2290        }
2291    };
2292
2293    let payload = if supervisor_probe {
2294        serde_json::to_value(&resp).and_then(|mut value| {
2295            if let Some(claim) = current_supervisor_claim() {
2296                value["supervisor_claim"] = serde_json::Value::String(claim);
2297            }
2298            if handover_accepted {
2299                value["supervisor_handover_accepted"] = serde_json::Value::Bool(true);
2300            }
2301            serde_json::to_vec(&value)
2302        })
2303    } else {
2304        serde_json::to_vec(&resp)
2305    };
2306    let mut handover_ack_written = false;
2307    match payload {
2308        Ok(payload) => {
2309            if payload.len() > MAX_FRAME_BYTES {
2310                // The serialized response exceeds the IPC frame cap.  Send a
2311                // small explicit error frame so the client can distinguish a
2312                // per-request payload-size failure from a daemon crash.  A
2313                // client that receives this error frame will NOT trigger
2314                // stale-daemon kill/respawn (ParseFailure requires a read_frame
2315                // error, not an ok=false result).
2316                tracing::warn!(
2317                    bytes = payload.len(),
2318                    limit = MAX_FRAME_BYTES,
2319                    "daemon response exceeds MAX_FRAME_BYTES; sending explicit error frame"
2320                );
2321                let message = format!(
2322                    "response too large: {} bytes exceeds {} byte IPC cap",
2323                    payload.len(),
2324                    MAX_FRAME_BYTES,
2325                );
2326                // One frame may aggregate successful, failed, and aborted operations.
2327                let err_resp = DaemonResponseFrame {
2328                    ok: false,
2329                    result: None,
2330                    error: Some(message.clone()),
2331                    error_detail: Some(serde_json::json!({
2332                        "kind": "transport",
2333                        "code": "response_frame_size_limit",
2334                        "message": message,
2335                        "domain_disposition": crate::DomainDisposition::Unknown.as_str(),
2336                    })),
2337                    namespace_mismatch: false,
2338                    config_mismatch: false,
2339                    served_config_id: resp.served_config_id,
2340                    version_mismatch: false,
2341                    daemon_protocol_version: PROTOCOL_VERSION,
2342                    metrics: None,
2343                    request_id: resp.request_id,
2344                };
2345                if let Ok(err_payload) = serde_json::to_vec(&err_resp) {
2346                    if let Err(e) = write_response_frame(&mut peer_write, &err_payload).await {
2347                        tracing::debug!(error = %e, "failed to write oversized-response error frame");
2348                    }
2349                }
2350            } else {
2351                match write_response_frame(&mut peer_write, &payload).await {
2352                    Ok(()) => handover_ack_written = true,
2353                    Err(e) => tracing::debug!(error = %e, "failed to write daemon response frame"),
2354                }
2355            }
2356        }
2357        Err(e) => tracing::warn!(error = %e, "failed to serialize daemon response frame"),
2358    }
2359    if handover_accepted && handover_ack_written {
2360        // Signal this process, not a PID observed earlier over a socket.
2361        // Production installed its SIGTERM handler before binding the socket.
2362        if unsafe { libc::raise(libc::SIGTERM) } != 0 {
2363            tracing::error!(error = %std::io::Error::last_os_error(), "self-directed handover signal failed");
2364        }
2365    }
2366    drop(ordinary_admission);
2367}
2368
2369/// An accepted connection owns a drain slot from the accept loop until its
2370/// handler future is dropped. The claim happens before spawning so shutdown
2371/// cannot observe zero in the gap between `accept()` and the task's first poll.
2372#[cfg(unix)]
2373struct ActiveConnectionGuard {
2374    active: Arc<std::sync::atomic::AtomicUsize>,
2375}
2376
2377#[cfg(unix)]
2378impl ActiveConnectionGuard {
2379    fn claim(active: Arc<std::sync::atomic::AtomicUsize>) -> Self {
2380        active.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2381        Self { active }
2382    }
2383}
2384
2385#[cfg(unix)]
2386impl Drop for ActiveConnectionGuard {
2387    fn drop(&mut self) {
2388        self.active
2389            .fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
2390    }
2391}
2392
2393#[cfg(unix)]
2394fn spawn_connection_task<F>(
2395    active: Arc<std::sync::atomic::AtomicUsize>,
2396    future: F,
2397) -> tokio::task::JoinHandle<()>
2398where
2399    F: std::future::Future<Output = ()> + Send + 'static,
2400{
2401    let guard = ActiveConnectionGuard::claim(active);
2402    tokio::spawn(async move {
2403        let _guard = guard;
2404        future.await;
2405    })
2406}
2407
2408/// Run the daemon: bind the socket, warm in the background, serve request
2409/// frames until SIGTERM/SIGINT.
2410///
2411/// Fatally acquires its own startup lock, which only protects
2412/// cleanup→pid-claim→bind — `dispatcher` has already run migrations and
2413/// applied pack schema plans while constructing itself, unguarded. Production
2414/// boot must go through [`run_daemon_with_boot_guard`] instead, which extends
2415/// the same lock back over construction. This entry point is for callers
2416/// (and tests) that build the dispatcher and start serving as one atomic
2417/// step with no separate boot-guard window to protect.
2418#[cfg(unix)]
2419pub async fn run_daemon<D: DaemonDispatch>(dispatcher: D) -> anyhow::Result<()> {
2420    let boot_guard = Some(acquire_daemon_boot_guard()?);
2421    run_daemon_with_boot_guard_inner(
2422        dispatcher,
2423        boot_guard,
2424        false,
2425        DaemonOptions::default(),
2426        |_| DaemonStartupReport::default(),
2427    )
2428    .await
2429}
2430
2431/// Run a real daemon server for an in-process multi-launch test.
2432///
2433/// Separate production daemon candidates have distinct PIDs, so the boot fence
2434/// recognizes a live incumbent and makes later candidates exit. Parallel
2435/// test launchers share one OS process and therefore one PID; this explicit
2436/// fault-injection entry point preserves the production fence semantics by
2437/// allowing a live same-PID incumbent to win. Ordinary daemon startup
2438/// continues to treat a same-PID rendezvous as stale, protecting PID-reuse
2439/// cleanup behavior.
2440///
2441/// A losing candidate still follows the ordinary daemon-exit path and cancels
2442/// the process-wide component shutdown token. Callers must therefore use a
2443/// component-free dispatcher; this seam validates socket/PID ownership, not
2444/// multi-candidate component lifecycle.
2445#[cfg(all(unix, any(test, feature = "fault-injection")))]
2446#[doc(hidden)]
2447pub async fn run_daemon_in_process_test<D: DaemonDispatch>(dispatcher: D) -> anyhow::Result<()> {
2448    let boot_guard = Some(acquire_daemon_boot_guard()?);
2449    run_daemon_with_boot_guard_inner(
2450        dispatcher,
2451        boot_guard,
2452        true,
2453        DaemonOptions::default(),
2454        |_| DaemonStartupReport::default(),
2455    )
2456    .await
2457}
2458
2459#[cfg(unix)]
2460#[derive(Clone, Copy, PartialEq, Eq)]
2461enum RendezvousPathRole {
2462    Socket,
2463    PidFile,
2464}
2465
2466#[cfg(unix)]
2467impl RendezvousPathRole {
2468    fn env_name(self) -> &'static str {
2469        match self {
2470            Self::Socket => SOCKET_PATH_ENV,
2471            Self::PidFile => PID_PATH_ENV,
2472        }
2473    }
2474
2475    fn directory_name(self) -> &'static str {
2476        match self {
2477            Self::Socket => "socket directory",
2478            Self::PidFile => "PID-file directory",
2479        }
2480    }
2481
2482    fn path_name(self) -> &'static str {
2483        match self {
2484            Self::Socket => "socket path",
2485            Self::PidFile => "PID-file path",
2486        }
2487    }
2488
2489    fn path_component_name(self) -> &'static str {
2490        match self {
2491            Self::Socket => "socket-path",
2492            Self::PidFile => "PID-file-path",
2493        }
2494    }
2495}
2496
2497/// Vet the socket's parent directory, re-permissioning it only when it is the
2498/// directory khive owns by convention.
2499///
2500/// `KHIVE_SOCKET` takes an arbitrary path, and the previous unconditional
2501/// chmod-0700 of its parent was wrong in both directions: pointed at a shared
2502/// parent like `/tmp` it either failed outright for an ordinary user, or —
2503/// worse — succeeded when privileged and stripped access for every other
2504/// process on the machine. A directory we did not create is never modified.
2505///
2506/// What the directory must actually prevent is a *takeover of the socket
2507/// path*: the connection gate is the 0600 socket plus the accept-time
2508/// peer-uid check, but both defend the daemon's own socket — neither helps
2509/// once another local user can put *their* listener at the path clients
2510/// resolve. There are two ways a shared directory allows that, and the
2511/// sticky bit closes neither: a writer can *pre-bind* the predictable path
2512/// before this daemon starts (the sticky bit restricts unlinking, not
2513/// creating), and the directory's *owner* can unlink and rebind even in a
2514/// 1777 directory (sticky exempts the directory owner). So a caller-chosen
2515/// directory is served only when it is trusted end to end: owned by this
2516/// daemon's euid or root, and not writable by group or other at all. The
2517/// umask-default 0755 stays acceptable; shared sticky directories like
2518/// `/tmp` do not.
2519#[cfg(unix)]
2520pub(crate) fn ensure_socket_dir_is_trusted(parent: &std::path::Path) -> anyhow::Result<()> {
2521    // SAFETY: `geteuid` is always successful and takes no arguments.
2522    let daemon_euid = unsafe { libc::geteuid() } as u32;
2523    ensure_rendezvous_dir_is_trusted(parent, RendezvousPathRole::Socket, daemon_euid, true)
2524}
2525
2526/// Vet the parent directory of a PID file before reading, locking, or writing
2527/// it. The file's lock only protects the inode currently named by its path;
2528/// every directory component must therefore be as swap-resistant as the
2529/// socket rendezvous.
2530#[cfg(unix)]
2531pub fn ensure_pid_file_dir_is_trusted(pid_file: &std::path::Path) -> anyhow::Result<()> {
2532    let parent = pid_file
2533        .parent()
2534        .filter(|parent| !parent.as_os_str().is_empty())
2535        .unwrap_or_else(|| std::path::Path::new("."));
2536    // SAFETY: `geteuid` is always successful and takes no arguments.
2537    let daemon_euid = unsafe { libc::geteuid() } as u32;
2538    ensure_rendezvous_dir_is_trusted(parent, RendezvousPathRole::PidFile, daemon_euid, false)
2539}
2540
2541#[cfg(unix)]
2542fn ensure_rendezvous_dir_is_trusted(
2543    parent: &std::path::Path,
2544    role: RendezvousPathRole,
2545    daemon_euid: u32,
2546    repair_owned_default: bool,
2547) -> anyhow::Result<()> {
2548    let env_name = role.env_name();
2549    let directory_name = role.directory_name();
2550
2551    if repair_owned_default && parent == khive_dir() {
2552        std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700)).map_err(|e| {
2553            anyhow::anyhow!(
2554                "refusing to start: cannot chmod 0700 {}: {e}. The khive directory must be \
2555                 owner-only as the {directory_name} for {env_name}; it is part of the \
2556                 same-uid guarantee this daemon enforces.",
2557                parent.display()
2558            )
2559        })?;
2560        return ensure_rendezvous_path_is_swap_resistant(parent, daemon_euid, role);
2561    }
2562
2563    // Fail closed on the stat itself: not being able to read the metadata is
2564    // not the same as the directory passing.
2565    let meta = std::fs::metadata(parent).map_err(|e| {
2566        anyhow::anyhow!(
2567            "refusing to start: cannot stat {directory_name} {} for {env_name}: {e}. \
2568             It gates rendezvous-path safety, and unreadable metadata is not a passing state.",
2569            parent.display()
2570        )
2571    })?;
2572
2573    use std::os::unix::fs::MetadataExt;
2574    let owner = meta.uid();
2575    if owner != daemon_euid && owner != 0 {
2576        anyhow::bail!(
2577            "refusing to start: {directory_name} {} for {env_name} is owned by uid {owner}, \
2578             not this daemon's uid ({daemon_euid}) or root. A directory owner can replace \
2579             the rendezvous path regardless of mode bits. Point {env_name} at a directory \
2580             you own, or unset it for the default.",
2581            parent.display()
2582        );
2583    }
2584
2585    let mode = meta.permissions().mode();
2586    if mode & 0o022 != 0 {
2587        anyhow::bail!(
2588            "refusing to start: {directory_name} {} for {env_name} is mode {:04o} — writable \
2589             by group or other, so another local user could replace the rendezvous path. \
2590             Use a directory only you can write, or unset {env_name} for the default. \
2591             This daemon is not changing the permissions of a directory it does not own.",
2592            parent.display(),
2593            mode & 0o7777
2594        );
2595    }
2596
2597    ensure_rendezvous_path_is_swap_resistant(parent, daemon_euid, role)
2598}
2599
2600/// Socket-role form of [`ensure_rendezvous_path_is_swap_resistant`], used by the
2601/// socket-path tests.
2602#[cfg(all(unix, test))]
2603fn ensure_socket_path_is_swap_resistant(
2604    parent: &std::path::Path,
2605    daemon_euid: u32,
2606) -> anyhow::Result<()> {
2607    ensure_rendezvous_path_is_swap_resistant(parent, daemon_euid, RendezvousPathRole::Socket)
2608}
2609
2610/// Walk the socket directory path exactly as the kernel will traverse it at
2611/// bind time — component by component, following symlinks — and refuse any
2612/// node another local user could swap after this validation.
2613///
2614/// Checking only the canonicalized result is not enough: bind and every
2615/// client traverse the *original* path, so a symlink component owned by
2616/// another user can resolve somewhere trusted while it is being validated
2617/// and be retargeted before the socket is bound. Validating the nodes the
2618/// traversal actually visits closes that gap: a symlink component is
2619/// acceptable only when its owner — the only party besides root who can
2620/// retarget it — is this daemon's euid or root, and every directory
2621/// component is held to the ancestor rule below.
2622///
2623/// The directory rule is deliberately weaker than the immediate parent's.
2624/// The parent hosts the socket file, where the threat is *creation* of the
2625/// predictable path — the sticky bit does not restrict creating, so no
2626/// sticky exception is sound there. An ancestor only threatens via
2627/// *rename/unlink of an existing entry we own*, which the sticky bit does
2628/// restrict: in a sticky directory, only the entry's owner, the directory's
2629/// owner, or root may rename it. So a root-owned `/tmp` (1777) is an
2630/// acceptable ancestor of a user-owned 0700 socket directory, while a
2631/// non-sticky group/other-writable ancestor, or one owned by a third uid
2632/// (who could rename the entry, or chmod the directory first), is refused.
2633/// Every stat failure fails closed.
2634#[cfg(unix)]
2635fn ensure_rendezvous_path_is_swap_resistant(
2636    parent: &std::path::Path,
2637    daemon_euid: u32,
2638    role: RendezvousPathRole,
2639) -> anyhow::Result<()> {
2640    use std::os::unix::fs::MetadataExt;
2641
2642    let env_name = role.env_name();
2643    let directory_name = role.directory_name();
2644    let path_name = role.path_name();
2645    let component_name = role.path_component_name();
2646
2647    let absolute = if parent.is_absolute() {
2648        parent.to_path_buf()
2649    } else {
2650        std::env::current_dir()
2651            .map_err(|e| {
2652                anyhow::anyhow!(
2653                    "refusing to start: cannot resolve the working directory to absolutize \
2654                     {directory_name} {} for {env_name}: {e}.",
2655                    parent.display()
2656                )
2657            })?
2658            .join(parent)
2659    };
2660
2661    fn push_components(stack: &mut Vec<std::ffi::OsString>, path: &std::path::Path) {
2662        let components: Vec<_> = path
2663            .components()
2664            .map(|c| c.as_os_str().to_os_string())
2665            .collect();
2666        stack.extend(components.into_iter().rev());
2667    }
2668
2669    let mut stack: Vec<std::ffi::OsString> = Vec::new();
2670    push_components(&mut stack, &absolute);
2671    let mut resolved = std::path::PathBuf::new();
2672    let mut symlinks_followed = 0u32;
2673
2674    while let Some(component) = stack.pop() {
2675        if component == "/" {
2676            resolved = std::path::PathBuf::from("/");
2677            continue;
2678        }
2679        if component == "." {
2680            continue;
2681        }
2682        if component == ".." {
2683            resolved.pop();
2684            continue;
2685        }
2686        let candidate = resolved.join(&component);
2687        let meta = std::fs::symlink_metadata(&candidate).map_err(|e| {
2688            anyhow::anyhow!(
2689                "refusing to start: cannot stat {component_name} component {} for {env_name}: \
2690                 {e}. An unreadable component is not a passing one.",
2691                candidate.display()
2692            )
2693        })?;
2694        let owner = meta.uid();
2695
2696        if meta.file_type().is_symlink() {
2697            symlinks_followed += 1;
2698            if symlinks_followed > 40 {
2699                anyhow::bail!(
2700                    "refusing to start: {path_name} for {env_name} resolves through more than \
2701                     40 symlinks at {} — treating this as a loop.",
2702                    candidate.display()
2703                );
2704            }
2705            if owner != daemon_euid && owner != 0 {
2706                anyhow::bail!(
2707                    "refusing to start: {component_name} symlink component {} for {env_name} \
2708                     is owned by uid {owner}, not this daemon's uid ({daemon_euid}) or root — \
2709                     its owner could retarget it after this check and re-root the {path_name}. \
2710                     Point {env_name} somewhere trusted end to end, or unset it for the default.",
2711                    candidate.display()
2712                );
2713            }
2714            let target = std::fs::read_link(&candidate).map_err(|e| {
2715                anyhow::anyhow!(
2716                    "refusing to start: cannot read {component_name} symlink component {} \
2717                     for {env_name}: {e}.",
2718                    candidate.display()
2719                )
2720            })?;
2721            push_components(&mut stack, &target);
2722            continue;
2723        }
2724
2725        if meta.is_dir() {
2726            let mode = meta.permissions().mode();
2727            let sticky = mode & 0o1000 != 0;
2728            if owner != daemon_euid && owner != 0 {
2729                anyhow::bail!(
2730                    "refusing to start: {component_name} ancestor {} for {env_name} is owned by \
2731                     uid {owner}, not this daemon's uid ({daemon_euid}) or root — its owner \
2732                     could rename the next path component and re-root the {path_name}. Point \
2733                     {env_name} somewhere trusted end to end, or unset it for the default.",
2734                    candidate.display()
2735                );
2736            }
2737            if mode & 0o022 != 0 && !sticky {
2738                anyhow::bail!(
2739                    "refusing to start: {component_name} ancestor {} for {env_name} is mode \
2740                     {:04o} — writable by group or other without the sticky bit, so another \
2741                     local user could rename the next path component and re-root the {path_name}. \
2742                     Point {env_name} somewhere trusted end to end, or unset it for the default.",
2743                    candidate.display(),
2744                    mode & 0o7777
2745                );
2746            }
2747            resolved = candidate;
2748            continue;
2749        }
2750
2751        anyhow::bail!(
2752            "refusing to start: {component_name} component {} for {env_name} is neither a \
2753             directory nor a symlink — the {path_name} cannot traverse it.",
2754            candidate.display()
2755        );
2756    }
2757
2758    Ok(())
2759}
2760
2761/// Run the daemon using a startup lock acquired by the caller *before*
2762/// building `dispatcher`, so a second process racing to boot (e.g. two
2763/// `kkernel mcp --daemon` spawns before either has bound its socket) cannot
2764/// run migrations/FTS DDL concurrently against the same database file.
2765/// `boot_guard` is only `None` on non-unix targets, where there is no
2766/// advisory boot lock to hold in the first place; every unix daemon-mode
2767/// caller passes `Some`.
2768///
2769/// The guard is held across cleanup → pid-claim → bind, then dropped. The
2770/// caller must not still be holding a *different* handle to the same lock
2771/// file when this function is entered — see the "Deadlock note" on the
2772/// `_startup_lock` binding below for why that would self-deadlock on `flock`.
2773#[cfg(unix)]
2774pub async fn run_daemon_with_boot_guard<D: DaemonDispatch>(
2775    dispatcher: D,
2776    boot_guard: Option<std::fs::File>,
2777) -> anyhow::Result<()> {
2778    run_daemon_with_boot_guard_inner(
2779        dispatcher,
2780        boot_guard,
2781        false,
2782        DaemonOptions::default(),
2783        |_| DaemonStartupReport::default(),
2784    )
2785    .await
2786}
2787
2788/// Run the daemon and start host-owned background work only after the socket
2789/// is bound, permissions are restricted, and this process owns the PID file.
2790/// The callback runs once while the startup lock and teardown guard are held;
2791/// setup failures never invoke it. Work started by the callback must use the
2792/// daemon shutdown token and tracked-task drain contract.
2793#[cfg(unix)]
2794pub async fn run_daemon_with_boot_guard_and_start<D, F>(
2795    dispatcher: D,
2796    boot_guard: Option<std::fs::File>,
2797    start: F,
2798) -> anyhow::Result<()>
2799where
2800    D: DaemonDispatch,
2801    F: FnOnce(&D) + Send,
2802{
2803    run_daemon_with_options_and_boot_guard_and_start(
2804        dispatcher,
2805        boot_guard,
2806        DaemonOptions::default(),
2807        |dispatcher| {
2808            start(dispatcher);
2809            DaemonStartupReport::default()
2810        },
2811    )
2812    .await
2813}
2814
2815/// Start with an explicit launch mode and collect the host's startup inventory.
2816#[cfg(unix)]
2817pub async fn run_daemon_with_options_and_boot_guard_and_start<D, F>(
2818    dispatcher: D,
2819    boot_guard: Option<std::fs::File>,
2820    options: DaemonOptions,
2821    start: F,
2822) -> anyhow::Result<()>
2823where
2824    D: DaemonDispatch,
2825    F: FnOnce(&D) -> DaemonStartupReport + Send,
2826{
2827    anyhow::ensure!(
2828        !options.idle_interval.is_zero(),
2829        "daemon idle interval must be positive"
2830    );
2831    run_daemon_with_boot_guard_inner(dispatcher, boot_guard, false, options, start).await
2832}
2833
2834#[cfg(unix)]
2835async fn run_daemon_with_boot_guard_inner<D, F>(
2836    dispatcher: D,
2837    boot_guard: Option<std::fs::File>,
2838    allow_same_process_incumbent: bool,
2839    options: DaemonOptions,
2840    start: F,
2841) -> anyhow::Result<()>
2842where
2843    D: DaemonDispatch,
2844    F: FnOnce(&D) -> DaemonStartupReport + Send,
2845{
2846    // Cancel on every exit, including setup failure and unwinding from the
2847    // post-ownership startup callback. The guard precedes all fallible work
2848    // so even a process that never becomes the daemon relinquishes its
2849    // process-lifetime shutdown token; restarting requires exec.
2850    struct ComponentTeardown;
2851    impl Drop for ComponentTeardown {
2852        fn drop(&mut self) {
2853            daemon_shutdown_token().cancel();
2854        }
2855    }
2856    let _component_teardown = ComponentTeardown;
2857
2858    // Placed after the teardown guard so this refusal keeps the ADR-119
2859    // contract every other pre-bind error path has, and before the paths are
2860    // resolved so a split rendezvous never reaches cleanup/bind/pid-write.
2861    ensure_rendezvous_overrides_paired()?;
2862
2863    let sock = socket_path();
2864    let pid_file = pid_path();
2865    let socket_parent = sock.parent();
2866    let pid_parent = pid_file.parent();
2867
2868    if let Some(parent) = socket_parent {
2869        std::fs::create_dir_all(parent)?;
2870        ensure_socket_dir_is_trusted(parent)?;
2871    }
2872    // Identical parent paths traverse the same components, so the socket
2873    // check above also vets the PID-file parent. Aliased paths are checked
2874    // independently because each original path is traversed by file access.
2875    if pid_parent != socket_parent {
2876        ensure_pid_file_dir_is_trusted(&pid_file)?;
2877    }
2878
2879    // Hold the startup lock across cleanup → pid-claim → bind so a concurrent
2880    // client's kill_and_respawn (which also holds this lock) cannot remove the
2881    // rendezvous paths during setup. The PID file's own lock is retained after
2882    // this shared startup lock is released, including while the listener drains
2883    // during shutdown.
2884    //
2885    // Deadlock note: the client holds this lock only during kill+spawn and
2886    // releases it before the spawned daemon process starts (the lock guard is
2887    // dropped when kill_and_respawn returns, before the readiness probe loop).
2888    // The daemon holds exactly one handle to this lock for its whole boot
2889    // sequence (received as `boot_guard`, extended from before `dispatcher`
2890    // was constructed) — never a second, independently-acquired handle in the
2891    // same process, which would self-deadlock on `flock`.
2892    let _startup_lock = boot_guard;
2893
2894    // A second daemon must refuse loudly rather than exit successfully while
2895    // another daemon owns this rendezvous. Only `Stale` lets the caller proceed.
2896    match cleanup_stale_daemon(
2897        &sock,
2898        &pid_file,
2899        allow_same_process_incumbent,
2900        dispatcher.config_id(),
2901    )
2902    .await
2903    {
2904        Incumbent::Serving(incumbent_pid) => {
2905            tracing::error!(
2906                pid = incumbent_pid,
2907                socket = ?sock,
2908                "refusing to start: a khived instance is already serving this socket"
2909            );
2910            anyhow::bail!(
2911                "refusing to start: khived is already running as pid {incumbent_pid}, \
2912                 serving socket {}. Stop that instance first if you intend to replace it.",
2913                sock.display()
2914            );
2915        }
2916        Incumbent::Live(incumbent_pid) => {
2917            tracing::error!(
2918                pid = incumbent_pid,
2919                socket = ?sock,
2920                "refusing to start: a live process owns the PID file but no khived answered"
2921            );
2922            anyhow::bail!(
2923                "refusing to start: pid {incumbent_pid} owns the daemon PID file and is alive, \
2924                 but nothing answered the khived protocol on {}. It may be draining. Nothing \
2925                 was removed; stop that process first if you intend to replace it.",
2926                sock.display()
2927            );
2928        }
2929        Incumbent::Stale => {}
2930    }
2931
2932    // Install signal streams before publishing either rendezvous file. A
2933    // supervisor may stop us as soon as connect/pid checks succeed, before
2934    // the accept loop or shutdown future has been polled. Keep these streams
2935    // alive so a signal during the rest of startup reaches normal cleanup.
2936    let mut sigterm = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())?;
2937    let mut sigint = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt())?;
2938
2939    let pid_file_guard = match write_pid_file_exclusive(&pid_file) {
2940        Ok(guard) => guard,
2941        Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
2942            // A PID file appeared between cleanup and our claim. Never touch
2943            // the winner's files; defer only if it already answers as khived.
2944            if pid_file_names_a_reachable_daemon(
2945                &pid_file,
2946                &sock,
2947                allow_same_process_incumbent,
2948                dispatcher.config_id(),
2949            )
2950            .await
2951            {
2952                tracing::info!(
2953                    "a replacement khived already claimed the pid/socket rendezvous; exiting"
2954                );
2955                return Ok(());
2956            }
2957            anyhow::bail!(
2958                "failed to claim daemon pid file at {pid_file:?}: it already exists \
2959                 and does not name a reachable daemon"
2960            );
2961        }
2962        Err(e) => return Err(e.into()),
2963    };
2964
2965    let listener = match UnixListener::bind(&sock) {
2966        Ok(listener) => listener,
2967        Err(e) => {
2968            remove_pid_file_if_owned(&pid_file, &pid_file_guard);
2969            return Err(e.into());
2970        }
2971    };
2972    // Fail closed, same reason as the directory above. If this chmod fails the
2973    // socket is world-reachable in a way the accepted design never covered, so
2974    // the bound listener is dropped and the entry removed rather than served.
2975    if let Err(e) = std::fs::set_permissions(&sock, std::fs::Permissions::from_mode(0o600)) {
2976        drop(listener);
2977        let _ = std::fs::remove_file(&sock);
2978        remove_pid_file_if_owned(&pid_file, &pid_file_guard);
2979        return Err(anyhow::anyhow!(
2980            "refusing to start: cannot chmod 0600 {}: {e}. The daemon socket must be owner-only \
2981             — it is half of the single-principal guarantee this daemon enforces.",
2982            sock.display()
2983        ));
2984    }
2985    // Captured while still holding the startup lock, immediately after
2986    // bind, so shutdown cleanup can later prove "this is still the same socket
2987    // I bound" rather than trusting the path alone.
2988    let bound_identity = socket_identity(&sock);
2989
2990    let lifecycle = Arc::new(DaemonLifecycle::new(options, start(&dispatcher)));
2991
2992    // Release the shared startup lock now that the listener is bound. The
2993    // locked PID file continues to identify this daemon through shutdown.
2994    drop(_startup_lock);
2995    tracing::info!(
2996        socket = ?sock,
2997        pid = std::process::id(),
2998        source_revision = crate::BUILD_INFO.source_revision,
2999        build_time = crate::BUILD_INFO.build_time,
3000        "khived listening"
3001    );
3002
3003    {
3004        let warm = dispatcher.clone();
3005        track_named_background_task("daemon_warmup", async move {
3006            warm.warm_all().await;
3007        });
3008    }
3009
3010    // The checkpoint task's own strong-count-based exit is unreachable
3011    // whenever `event_store_for_checkpoint()` returns `Some` (the ordinary
3012    // production shape), because the `SqlEventStore` it wraps retains its
3013    // own clone of the same pool. An explicit watch channel replaces that
3014    // mechanism: the sender is held for the remainder of this function's
3015    // scope and signalled as the first action once shutdown is observed,
3016    // below.
3017    let (checkpoint_shutdown_tx, checkpoint_shutdown_rx) = tokio::sync::watch::channel(());
3018    // ADR-091 Amendment 3: one checkpoint task per file-backed backend the
3019    // dispatcher wired — the primary pool plus every entry
3020    // `secondary_pools_for_checkpoint` returns — sharing this one shutdown
3021    // channel (the sender broadcasts to every receiver clone), so the single
3022    // send below stops every spawned task before `drain()`.
3023    let checkpoint_tasks = checkpoint_task_specs(
3024        dispatcher.pool_for_checkpoint(),
3025        dispatcher.secondary_pools_for_checkpoint(),
3026        dispatcher.event_store_for_checkpoint(),
3027        dispatcher.namespace().to_string(),
3028    );
3029    if !checkpoint_tasks.is_empty() {
3030        let cfg = CheckpointConfig::from_env();
3031        let checkpoint_task_count = checkpoint_tasks.len();
3032        for task in checkpoint_tasks {
3033            track_named_background_task(
3034                "wal_checkpoint",
3035                run_checkpoint_task(
3036                    task.pool,
3037                    cfg.clone(),
3038                    task.lifecycle_owner,
3039                    checkpoint_shutdown_rx.clone(),
3040                    task.is_main,
3041                ),
3042            );
3043        }
3044        tracing::info!(checkpoint_task_count, "WAL checkpoint task(s) started");
3045    }
3046
3047    let active = Arc::new(std::sync::atomic::AtomicUsize::new(0));
3048    let connection_tasks = Arc::new(std::sync::Mutex::new(
3049        Vec::<tokio::task::JoinHandle<()>>::new(),
3050    ));
3051    let (request_shutdown_tx, request_shutdown_rx) = tokio::sync::watch::channel(false);
3052
3053    let shutdown = async {
3054        tokio::select! {
3055            _ = sigterm.recv() => tracing::info!("received SIGTERM"),
3056            _ = sigint.recv() => tracing::info!("received SIGINT"),
3057        }
3058        // Tokio retains its process-wide handlers after the streams are dropped.
3059        // This daemon cannot restart without exec; a repeat signal must terminate
3060        // even if shutdown is blocked in synchronous recovery-lock acquisition.
3061        for signal in [libc::SIGTERM, libc::SIGINT] {
3062            // SAFETY: setting SIG_DFL for these valid signals needs no handler
3063            // pointer or shared Rust state and applies to the whole process.
3064            if unsafe { libc::signal(signal, libc::SIG_DFL) } == libc::SIG_ERR {
3065                return Err(std::io::Error::last_os_error());
3066            }
3067        }
3068        Ok::<(), std::io::Error>(())
3069    };
3070    tokio::pin!(shutdown);
3071    lifecycle.ready();
3072
3073    // SAFETY: `geteuid` is always successful and takes no arguments.
3074    let daemon_euid = unsafe { libc::geteuid() } as u32;
3075
3076    let reason = tokio::select! {
3077        _ = async {
3078            let mut accept_error_backoff = None;
3079            let mut last_accept_error_log: Option<std::time::Instant> = None;
3080            loop {
3081                match listener.accept().await {
3082                    Ok((mut stream, _)) => {
3083                        let initial_frame_deadline =
3084                            tokio::time::Instant::now() + INITIAL_FRAME_READ_TIMEOUT;
3085                        accept_error_backoff = None;
3086                        last_accept_error_log = None;
3087                        // Refuse a foreign uid before any frame is read.
3088                        // Fails CLOSED: an error reading peer credentials is
3089                        // "cannot prove same-uid", which is the same answer as
3090                        // "is not same-uid" — never a pass.
3091                        //
3092                        // Scope: this is a same-UID check, which is strictly
3093                        // weaker than same-PRINCIPAL. Several distinct actors
3094                        // running under one uid all pass here, so this does not
3095                        // by itself establish that the daemon serves a single
3096                        // principal, and nothing downstream refuses or degrades
3097                        // on observing more than one. Do not describe it as a
3098                        // multi-principal guard — it bounds the process boundary,
3099                        // not the identity one.
3100                        match peer_uid(&stream) {
3101                            Ok(peer) if uid_is_permitted(peer, daemon_euid) => {}
3102                            Ok(peer) => {
3103                                tracing::error!(
3104                                    peer_uid = peer,
3105                                    daemon_euid,
3106                                    "refusing connection from a foreign uid: this daemon accepts \
3107                                     only peers running as its own uid"
3108                                );
3109                                drop(stream);
3110                                continue;
3111                            }
3112                            Err(e) => {
3113                                tracing::error!(
3114                                    error = %e,
3115                                    "refusing connection: cannot read peer credentials, so \
3116                                     same-uid cannot be proven"
3117                                );
3118                                drop(stream);
3119                                continue;
3120                            }
3121                        }
3122                        // One permit per connection, taken before the task is
3123                        // spawned. Past the cap the peer is answered with a busy
3124                        // error and the stream is dropped; connections already
3125                        // admitted keep their permits and are not affected.
3126                        let Some(permit) = admit_or_refuse_busy(
3127                            &lifecycle.connections,
3128                            &mut stream,
3129                            dispatcher.config_id(),
3130                        )
3131                        .await
3132                        else {
3133                            continue;
3134                        };
3135                        // Keep the acceptance-time deadline across the
3136                        // credential check and connection-task scheduling.
3137                        let d = dispatcher.clone();
3138                        let shutdown = request_shutdown_rx.clone();
3139                        let lifecycle = Arc::clone(&lifecycle);
3140                        let handle = spawn_connection_task(Arc::clone(&active), async move {
3141                            // Released when the handler ends, whether it returns,
3142                            // panics or is aborted at shutdown.
3143                            let _permit = permit;
3144                            handle_conn_with_lifecycle(
3145                                stream,
3146                                d,
3147                                Some(shutdown),
3148                                initial_frame_deadline,
3149                                Some(lifecycle),
3150                            )
3151                            .await;
3152                        });
3153                        let mut tasks = connection_tasks
3154                            .lock()
3155                            .unwrap_or_else(std::sync::PoisonError::into_inner);
3156                        tasks.retain(|task| !task.is_finished());
3157                        tasks.push(handle);
3158                    }
3159                    Err(e) => {
3160                        let delay = next_accept_error_backoff(accept_error_backoff);
3161                        accept_error_backoff = Some(delay);
3162                        let capacity_exhausted = matches!(
3163                            e.raw_os_error(),
3164                            Some(libc::EMFILE) | Some(libc::ENFILE)
3165                        );
3166                        if last_accept_error_log.is_none_or(|last| {
3167                            last.elapsed() >= std::time::Duration::from_secs(30)
3168                        }) {
3169                            tracing::error!(
3170                                error = %e,
3171                                capacity_exhausted,
3172                                retry_ms = delay.as_millis(),
3173                                "daemon accept failed; retrying with bounded backoff"
3174                            );
3175                            last_accept_error_log = Some(std::time::Instant::now());
3176                        }
3177                        tokio::time::sleep(delay).await;
3178                    }
3179                }
3180            }
3181        } => DaemonShutdownReason::Signal,
3182        result = &mut shutdown => { result?; DaemonShutdownReason::Signal },
3183        _ = wait_for_idle(&dispatcher, &lifecycle) => DaemonShutdownReason::Idle,
3184    };
3185
3186    lifecycle.draining(reason);
3187
3188    // A listening backlog is not admitted work. Close it before draining so
3189    // new clients cannot finish writing to a socket nobody will accept.
3190    drop(listener);
3191
3192    // Signal the checkpoint task to exit before draining, so `drain()`
3193    // actually waits on it via `track_background_task` rather than the
3194    // task outliving the drain window (or the process) unsignalled.
3195    let _ = checkpoint_shutdown_tx.send(());
3196
3197    // Per-run signal: read scopes stop promptly, admitted writes ignore it and
3198    // retain the rest of the configured drain window to commit or roll back.
3199    if reason == DaemonShutdownReason::Signal {
3200        let _ = request_shutdown_tx.send(true);
3201    }
3202
3203    // Same ordering contract for ADR-119 daemon components: cancel before
3204    // drain, so each component's supervisor (itself a tracked task) can run
3205    // its bounded shutdown inside the drain wait.
3206    daemon_shutdown_token().cancel();
3207
3208    let drained = if reason == DaemonShutdownReason::Idle {
3209        tokio::select! {
3210            _ = drain_for_idle(&active, drain_timeout()) => true,
3211            result = &mut shutdown => {
3212                result?;
3213                lifecycle.draining(DaemonShutdownReason::Signal);
3214                let _ = request_shutdown_tx.send(true);
3215                drain(&active).await
3216            }
3217        }
3218    } else {
3219        drain(&active).await
3220    };
3221    let tasks = {
3222        let mut retained = connection_tasks
3223            .lock()
3224            .unwrap_or_else(std::sync::PoisonError::into_inner);
3225        std::mem::take(&mut *retained)
3226    };
3227    finish_connection_tasks(tasks, drained).await;
3228
3229    // A concurrent client's `kill_and_respawn` may have already decided
3230    // this daemon looked stale, killed it, and spawned a replacement that
3231    // bound the same socket/PID paths while this daemon was draining above.
3232    // Reacquire the recovery lock (the same one that serializes startup) and
3233    // only unlink if the PID file still names this process AND the socket at
3234    // `sock` is still the exact one this daemon bound — otherwise a
3235    // replacement daemon owns those paths now and unlinking would delete its
3236    // live socket/PID out from under it.
3237    match acquire_recovery_lock() {
3238        Some(_shutdown_lock) => {
3239            shutdown_cleanup_if_owned(&sock, &pid_file, bound_identity);
3240        }
3241        None => {
3242            tracing::warn!(
3243                "could not acquire recovery lock for shutdown cleanup; \
3244                 skipping unlink to avoid deleting a replacement daemon's paths"
3245            );
3246        }
3247    }
3248    lifecycle.stopped();
3249    tracing::info!("khived stopped");
3250    Ok(())
3251}
3252
3253/// Remove `sock`/`pid_file` only if they still belong to this process: the PID
3254/// file must name `std::process::id()` AND the socket currently at `sock` must
3255/// still be the exact one identified by `bound_identity` (dev/ino, not path).
3256///
3257/// Returns `true` if cleanup ran, `false` if it was skipped because a
3258/// replacement daemon already owns those paths. The caller must hold
3259/// the recovery lock across this call — the same lock daemon startup holds
3260/// across cleanup+bind+pid-write — so no replacement can bind between this
3261/// function's checks and its unlinks.
3262#[cfg(unix)]
3263fn shutdown_cleanup_if_owned(
3264    sock: &std::path::Path,
3265    pid_file: &std::path::Path,
3266    bound_identity: Option<SocketIdentity>,
3267) -> bool {
3268    let pid_is_ours = std::fs::read_to_string(pid_file)
3269        .ok()
3270        .and_then(|s| s.trim().parse::<u32>().ok())
3271        == Some(std::process::id());
3272    let socket_is_ours = bound_identity.is_some() && socket_identity(sock) == bound_identity;
3273    if pid_is_ours && socket_is_ours {
3274        let _ = std::fs::remove_file(sock);
3275        let _ = std::fs::remove_file(pid_file);
3276        true
3277    } else {
3278        tracing::warn!(
3279            socket = ?sock,
3280            pid_file = ?pid_file,
3281            "skipping shutdown cleanup — a replacement daemon already owns this socket/PID"
3282        );
3283        false
3284    }
3285}
3286
3287// ── helpers ───────────────────────────────────────────────────────────────────
3288
3289/// Liveness verdict for a `kill(pid, 0)` probe.
3290#[cfg(unix)]
3291#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3292enum PidLiveness {
3293    /// errno 0 — signal delivery succeeded, the process exists and this
3294    /// caller may signal it.
3295    Alive,
3296    /// ESRCH (or any other non-EPERM errno) — no such process.
3297    Dead,
3298    /// EPERM — the process exists but this caller lacks permission to
3299    /// signal it. Unknown-safe: treated as running so stale-daemon cleanup
3300    /// never unlinks a live daemon's socket/PID file just because it is
3301    /// owned by a different user/uid.
3302    PermissionDenied,
3303}
3304
3305#[cfg(unix)]
3306impl PidLiveness {
3307    fn is_running(self) -> bool {
3308        !matches!(self, PidLiveness::Dead)
3309    }
3310}
3311
3312/// Maps a `kill(pid, 0)` outcome (return code + errno) to a [`PidLiveness`].
3313/// Pure and side-effect-free so the errno mapping can be unit tested without
3314/// a real process probe.
3315#[cfg(unix)]
3316fn classify_kill_result(rc: i32, errno: i32) -> PidLiveness {
3317    if rc == 0 {
3318        return PidLiveness::Alive;
3319    }
3320    match errno {
3321        libc::EPERM => PidLiveness::PermissionDenied,
3322        _ => PidLiveness::Dead,
3323    }
3324}
3325
3326#[cfg(unix)]
3327fn is_process_running(pid: u32) -> bool {
3328    let Ok(pid) = i32::try_from(pid) else {
3329        return false;
3330    };
3331    if pid <= 0 {
3332        return false;
3333    }
3334    // SAFETY: signal 0 is an existence/permission probe with no side effects.
3335    let rc = unsafe { libc::kill(pid, 0) };
3336    let errno = std::io::Error::last_os_error().raw_os_error().unwrap_or(0);
3337    classify_kill_result(rc, errno).is_running()
3338}
3339
3340/// Whether a PID may identify an incumbent from this candidate's point of view.
3341///
3342/// Production rejects the current PID so a stale rendezvous left by a prior
3343/// process whose PID was reused cannot protect an unrelated socket. The
3344/// in-process daemon harness opts in to the same-PID case because all of its
3345/// otherwise independent boot candidates necessarily share one OS process.
3346#[cfg(unix)]
3347fn pid_can_name_incumbent(pid: u32, current_pid: u32, allow_same_process_incumbent: bool) -> bool {
3348    allow_same_process_incumbent || pid != current_pid
3349}
3350
3351/// Bounded timeout for the protocol-identity probe used by duplicate-daemon
3352/// detection. Short enough that a hung or foreign listener does not stall
3353/// startup; long enough for a live khived under normal load to answer a
3354/// `probe_only` frame.
3355#[cfg(unix)]
3356const DUPLICATE_PROBE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(500);
3357
3358/// Whether the listener at `sock` actually speaks the khived wire protocol
3359/// **as a khived this process can defer to** — identified by a configuration
3360/// compatible with `expected_config_id`.
3361///
3362/// A live PID plus an accepting Unix socket is not proof of khived: any
3363/// unrelated process that happens to have bound the same path also answers
3364/// `connect()`. Nor is any well-formed [`DaemonResponseFrame`] proof: a
3365/// `config_mismatch`/`version_mismatch` response, a `metrics_only` snapshot
3366/// response, or a legacy pre-probe daemon that falls through to normal
3367/// dispatch on the empty `ops` string, all deserialize cleanly without being
3368/// the unambiguous "yes, alive and identity-matching" answer this check
3369/// needs — the daemon's `metrics_only` arm in particular echoes the same
3370/// `ok=true, result=None, error=None`, all-mismatch-flags-false, matching
3371/// protocol version and `served_config_id` shape as the probe-ack arm, and
3372/// is distinguished only by carrying `metrics: Some(...)`. This sends a
3373/// bounded `probe_only` frame (the same identity probe the client-side
3374/// recovery path uses, `crates/khive-mcp/src/daemon.rs::probe_daemon_identity`)
3375/// carrying this process's own `config_id`, and requires a compatible
3376/// probe-branch shape back: `ok=true`, `result=None`, `error=None`,
3377/// `metrics=None`, `request_id=None` (this probe frame never sets one), no
3378/// mismatch flags, matching protocol version, and matching
3379/// `served_config_id` — mirroring the client probe's `is_probe_ack` check so
3380/// both sides of the protocol agree on what "alive" means. Connect, write,
3381/// and read are all inside the one bounded timeout: `UnixStream::connect`
3382/// itself awaits write readiness, so a listener with a saturated accept
3383/// backlog could otherwise hold this call open past the advertised bound.
3384/// A connect that succeeds but never answers, times out, or answers with
3385/// non-protocol bytes, a mismatched identity, a `metrics_only` snapshot, or
3386/// any other non-probe-shaped response is not treated as the same khived
3387/// and falls through to the stale-socket recovery path instead.
3388#[cfg(unix)]
3389async fn socket_speaks_khived_protocol(sock: &std::path::Path, expected_config_id: &str) -> bool {
3390    let probe = DaemonRequestFrame {
3391        probe_only: true,
3392        protocol_version: PROTOCOL_VERSION,
3393        config_id: expected_config_id.to_string(),
3394        ..Default::default()
3395    };
3396    let Ok(payload) = serde_json::to_vec(&probe) else {
3397        return false;
3398    };
3399    let response = tokio::time::timeout(DUPLICATE_PROBE_TIMEOUT, async {
3400        let mut stream = UnixStream::connect(sock).await.ok()?;
3401        write_frame(&mut stream, &payload).await.ok()?;
3402        let raw = read_frame(&mut stream).await.ok()?;
3403        serde_json::from_slice::<DaemonResponseFrame>(&raw).ok()
3404    })
3405    .await
3406    .ok()
3407    .flatten();
3408
3409    let Some(resp) = response else {
3410        return false;
3411    };
3412    let is_probe_ack = resp.ok
3413        && resp.result.is_none()
3414        && resp.error.is_none()
3415        && resp.metrics.is_none()
3416        && resp.request_id.is_none();
3417    is_probe_ack
3418        && !resp.version_mismatch
3419        && !resp.namespace_mismatch
3420        && !resp.config_mismatch
3421        && resp.daemon_protocol_version == PROTOCOL_VERSION
3422        && resp
3423            .served_config_id
3424            .as_deref()
3425            .is_some_and(|served| config_ids_compatible(expected_config_id, served))
3426}
3427
3428/// Whether connecting to an existing socket path is definitely unreachable.
3429/// Timeouts and other errors remain ambiguous so cleanup fails closed.
3430#[cfg(unix)]
3431async fn socket_is_unreachable(sock: &std::path::Path) -> bool {
3432    match tokio::time::timeout(DUPLICATE_PROBE_TIMEOUT, UnixStream::connect(sock)).await {
3433        Ok(Err(error)) => matches!(
3434            error.kind(),
3435            std::io::ErrorKind::NotFound | std::io::ErrorKind::ConnectionRefused
3436        ),
3437        _ => false,
3438    }
3439}
3440
3441/// What owns the daemon PID file, from the point of view of a process that wants
3442/// to start. A live owner is never cleaned up: a draining incumbent closes its
3443/// listener before it releases writers, so an unanswered socket is ambiguous and
3444/// deleting its PID file is how two daemons end up on one store.
3445#[cfg(unix)]
3446enum Incumbent {
3447    /// A live process that answered the khived protocol on the socket.
3448    Serving(u32),
3449    /// A live PID still has an active or ambiguous rendezvous. Nothing removed.
3450    Live(u32),
3451    /// Nothing live owns the store; the socket and PID file were removed.
3452    Stale,
3453}
3454
3455/// Check whether `pid_file`/`sock` already name a live daemon and, if not,
3456/// remove the stale rendezvous files so the caller may bind fresh.
3457///
3458/// A live protocol responder, reachable socket, or live holder of the PID-file
3459/// lock means the caller must refuse to start. An unlocked PID with no listener
3460/// is a reused PID and may be reclaimed.
3461#[cfg(unix)]
3462async fn cleanup_stale_daemon(
3463    sock: &std::path::Path,
3464    pid_file: &std::path::Path,
3465    allow_same_process_incumbent: bool,
3466    expected_config_id: &str,
3467) -> Incumbent {
3468    let mut stale_pid_file_guard = None;
3469    if let Ok(pid_str) = std::fs::read_to_string(pid_file) {
3470        if let Ok(pid) = pid_str.trim().parse::<u32>() {
3471            if pid_can_name_incumbent(pid, std::process::id(), allow_same_process_incumbent)
3472                && is_process_running(pid)
3473            {
3474                if sock.exists() && socket_speaks_khived_protocol(sock, expected_config_id).await {
3475                    return Incumbent::Serving(pid);
3476                }
3477                if sock.exists() && !socket_is_unreachable(sock).await {
3478                    return Incumbent::Live(pid);
3479                }
3480                match try_acquire_pid_file_lock(pid_file) {
3481                    Ok(Some(guard)) => stale_pid_file_guard = Some(guard),
3482                    Ok(None) => return Incumbent::Live(pid),
3483                    Err(e) => {
3484                        tracing::warn!(
3485                            error = %e,
3486                            path = ?pid_file,
3487                            "cannot check daemon PID-file lock"
3488                        );
3489                        return Incumbent::Live(pid);
3490                    }
3491                }
3492            }
3493        }
3494    }
3495    if sock.exists() {
3496        if let Err(e) = std::fs::remove_file(sock) {
3497            tracing::warn!(error = %e, path = ?sock, "failed to remove stale socket");
3498        }
3499    }
3500    if pid_file.exists() {
3501        if let Err(e) = std::fs::remove_file(pid_file) {
3502            tracing::warn!(error = %e, path = ?pid_file, "failed to remove stale PID file");
3503        }
3504    }
3505    drop(stale_pid_file_guard);
3506    Incumbent::Stale
3507}
3508
3509/// Create and lock `pid_file` exclusively (`O_EXCL`) and write this process's PID.
3510///
3511/// Uses `create_new(true)` rather than `create(true).truncate(true)` so
3512/// this can never silently overwrite a PID file another process created —
3513/// the held file lock also identifies a starting or draining daemon when its
3514/// socket is not yet reachable.
3515#[cfg(unix)]
3516fn write_pid_file_exclusive(pid_file: &std::path::Path) -> std::io::Result<std::fs::File> {
3517    use std::os::unix::fs::OpenOptionsExt;
3518
3519    let mut opts = std::fs::OpenOptions::new();
3520    opts.write(true).create_new(true).mode(0o600);
3521    let mut f = opts.open(pid_file)?;
3522    // SAFETY: flock is a POSIX advisory lock with no memory side effects.
3523    let rc = unsafe { libc::flock(f.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
3524    if rc != 0 {
3525        return Err(std::io::Error::last_os_error());
3526    }
3527    f.write_all(std::process::id().to_string().as_bytes())?;
3528    Ok(f)
3529}
3530
3531/// Try to lock an existing PID file without creating it. `Some(file)` means
3532/// there is no daemon lock holder; `None` means a daemon still owns the file.
3533#[cfg(unix)]
3534fn try_acquire_pid_file_lock(pid_file: &std::path::Path) -> std::io::Result<Option<std::fs::File>> {
3535    let file = std::fs::OpenOptions::new()
3536        .read(true)
3537        .write(true)
3538        .open(pid_file)?;
3539    // SAFETY: flock is a POSIX advisory lock with no memory side effects.
3540    let rc = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) };
3541    if rc == 0 {
3542        return Ok(Some(file));
3543    }
3544    let error = std::io::Error::last_os_error();
3545    if error.kind() == std::io::ErrorKind::WouldBlock
3546        || error.raw_os_error() == Some(libc::EWOULDBLOCK)
3547    {
3548        Ok(None)
3549    } else {
3550        Err(error)
3551    }
3552}
3553
3554/// Remove the PID file on setup failure only while its path still names the
3555/// file this start attempt created.
3556#[cfg(unix)]
3557fn remove_pid_file_if_owned(pid_file: &std::path::Path, guard: &std::fs::File) {
3558    let Ok(owned) = guard.metadata() else {
3559        return;
3560    };
3561    let Ok(current) = std::fs::metadata(pid_file) else {
3562        return;
3563    };
3564    if owned.dev() == current.dev() && owned.ino() == current.ino() {
3565        if let Err(e) = std::fs::remove_file(pid_file) {
3566            tracing::warn!(error = %e, path = ?pid_file, "failed to remove unbound PID file");
3567        }
3568    }
3569}
3570
3571/// Return `true` if `pid_file` currently names an eligible live process that
3572/// still answers on `sock` — i.e. a daemon already owns this rendezvous and it
3573/// is safe to defer to it rather than treat the `AlreadyExists` PID-file
3574/// collision as a boot failure. Eligibility requires a different PID in
3575/// production; the explicit in-process harness may allow the current PID.
3576#[cfg(unix)]
3577async fn pid_file_names_a_reachable_daemon(
3578    pid_file: &std::path::Path,
3579    sock: &std::path::Path,
3580    allow_same_process_incumbent: bool,
3581    expected_config_id: &str,
3582) -> bool {
3583    let Ok(pid_str) = std::fs::read_to_string(pid_file) else {
3584        return false;
3585    };
3586    let Ok(pid) = pid_str.trim().parse::<u32>() else {
3587        return false;
3588    };
3589    pid_can_name_incumbent(pid, std::process::id(), allow_same_process_incumbent)
3590        && is_process_running(pid)
3591        && sock.exists()
3592        && socket_speaks_khived_protocol(sock, expected_config_id).await
3593}
3594
3595#[cfg(unix)]
3596async fn drain(active: &std::sync::atomic::AtomicUsize) -> bool {
3597    drain_with_timeout(active, drain_timeout()).await
3598}
3599
3600/// Voluntary retirement keeps admitted workers and rendezvous ownership alive.
3601#[cfg(unix)]
3602async fn drain_for_idle(active: &std::sync::atomic::AtomicUsize, timeout: std::time::Duration) {
3603    let deadline = tokio::time::Instant::now() + timeout;
3604    let mut warned = false;
3605    while active.load(std::sync::atomic::Ordering::SeqCst) + background_task_count() != 0 {
3606        if !warned && tokio::time::Instant::now() >= deadline {
3607            tracing::warn!(
3608                "idle drain interval elapsed; retaining workers and rendezvous until settled"
3609            );
3610            warned = true;
3611        }
3612        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3613    }
3614}
3615
3616#[cfg(unix)]
3617async fn drain_with_timeout(
3618    active: &std::sync::atomic::AtomicUsize,
3619    timeout: std::time::Duration,
3620) -> bool {
3621    use std::sync::atomic::Ordering;
3622    // One sequentially-consistent order spans connection handoff and tracked
3623    // task publication. Drain must not observe an ended connection and a
3624    // stale pre-publication background count as two simultaneous zeroes.
3625    let remaining = || active.load(Ordering::SeqCst) + background_task_count();
3626    if remaining() == 0 {
3627        return true;
3628    }
3629    let deadline = tokio::time::Instant::now() + timeout;
3630    while remaining() > 0 {
3631        if tokio::time::Instant::now() >= deadline {
3632            tracing::warn!(
3633                remaining_connections = active.load(Ordering::SeqCst),
3634                remaining_background_tasks = background_task_count(),
3635                outstanding_background_tasks = %background_task_names().join(", "),
3636                "drain timeout reached; forcing shutdown"
3637            );
3638            return false;
3639        }
3640        tokio::select! {
3641            _ = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3642            _ = tokio::time::sleep_until(deadline) => {}
3643        }
3644    }
3645    true
3646}
3647
3648#[cfg(unix)]
3649async fn finish_connection_tasks(tasks: Vec<tokio::task::JoinHandle<()>>, drained: bool) {
3650    if !drained {
3651        for task in &tasks {
3652            if !task.is_finished() {
3653                task.abort();
3654            }
3655        }
3656    }
3657    for task in tasks {
3658        let _ = task.await;
3659    }
3660}
3661
3662/// The bound `drain()` waits for tracked background tasks at daemon shutdown
3663/// (`KHIVE_DRAIN_TIMEOUT_SECS`, default 10s). Public so component supervision
3664/// can clamp per-component shutdown timeouts against it — a component timeout
3665/// longer than the drain bound could never complete its abort/state
3666/// transition before the daemon returns.
3667pub fn drain_timeout() -> std::time::Duration {
3668    let secs = khive_db::env::env_parse_or("KHIVE_DRAIN_TIMEOUT_SECS", DEFAULT_DRAIN_TIMEOUT_SECS);
3669    std::time::Duration::from_secs(secs)
3670}
3671
3672/// Returns `true` for non-empty env values that are not `"0"` or `"false"`.
3673#[cfg(unix)]
3674pub fn env_truthy(key: &str) -> bool {
3675    std::env::var(key)
3676        .map(|v| {
3677            let v = v.trim();
3678            !v.is_empty() && v != "0" && !v.eq_ignore_ascii_case("false")
3679        })
3680        .unwrap_or(false)
3681}
3682
3683include!("daemon_khive_root_tests.rs");
3684
3685/// Serve one already-admitted test connection through the production frame handler.
3686///
3687/// This seam owns no socket path, PID, boot guard, background components, or
3688/// process-wide shutdown state. The caller owns and joins the connection task.
3689/// It deliberately does not exercise listener admission or daemon lifecycle.
3690#[cfg(all(unix, any(test, feature = "test-internals")))]
3691#[doc(hidden)]
3692pub async fn serve_connection_for_test<D: DaemonDispatch>(stream: UnixStream, dispatcher: D) {
3693    handle_conn_with_shutdown(
3694        stream,
3695        dispatcher,
3696        None,
3697        tokio::time::Instant::now() + INITIAL_FRAME_READ_TIMEOUT,
3698    )
3699    .await;
3700}
3701
3702#[cfg(all(test, unix))]
3703#[path = "daemon_tests.rs"]
3704mod tests;