Skip to main content

subc_daemon/
bootstrap.rs

1use std::{
2    collections::BTreeMap,
3    env,
4    error::Error,
5    ffi::OsString,
6    fmt, fs, io,
7    net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr},
8    path::{Path, PathBuf},
9    process,
10    time::Duration,
11};
12
13use fs4::{FileExt, TryLockError};
14use subc_protocol::PROTOCOL_VERSION;
15pub use subc_transport::user_connection_token;
16use subc_transport::{
17    authenticate_client, connection_file, generate_daemon_id, generate_key, write_atomic,
18    AuthError, ConnectionFileError, ConnectionInfo, Endpoint, SCHEMA_VERSION,
19};
20use tokio::{
21    net::{TcpListener, TcpStream},
22    task::{JoinError, JoinHandle},
23    time::{sleep, timeout},
24};
25use tracing::{error, info, warn};
26
27use crate::{
28    daemon_config::{self, ConfiguredModule, DaemonConfigError},
29    server::{serve_listeners, ServerAuth, ServerError},
30    supervise::HealthConfig,
31    ConnectedClients, ControlHandler, DaemonSelfWatchdog, DaemonSelfWatchdogConfig,
32    ForwardingTable, Registry, RestartPolicy, Router, Supervisor, SupervisorHandle,
33    SupervisorProcessLiveness,
34};
35use std::sync::Arc;
36
37pub const DEFAULT_SUBC_PORT: u16 = 8757;
38pub const SUBC_PORT_ENV: &str = "SUBC_PORT";
39use subc_transport::CONNECTION_FILE_NAME;
40const DAEMON_VERSION: &str = env!("CARGO_PKG_VERSION");
41const CONNECT_TIMEOUT: Duration = Duration::from_secs(2);
42const PROBE_AUTH_DEADLINE: Duration = Duration::from_secs(2);
43const START_LOCK_RETRIES: usize = 40;
44const START_LOCK_RETRY_DELAY: Duration = Duration::from_millis(25);
45
46#[derive(Debug, Clone, Copy, PartialEq, Eq)]
47pub enum ConnectionFileSource {
48    XdgRuntimeDir,
49    TempDirFallback,
50    Explicit,
51}
52
53impl ConnectionFileSource {
54    fn reason(self) -> &'static str {
55        match self {
56            Self::XdgRuntimeDir => "XDG_RUNTIME_DIR set and non-empty",
57            Self::TempDirFallback => "XDG_RUNTIME_DIR unset or empty",
58            Self::Explicit => "configured path",
59        }
60    }
61}
62
63impl fmt::Display for ConnectionFileSource {
64    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
65        formatter.write_str(match self {
66            Self::XdgRuntimeDir => "xdg_runtime_dir",
67            Self::TempDirFallback => "temp_dir_fallback",
68            Self::Explicit => "explicit",
69        })
70    }
71}
72
73/// Runtime bootstrap configuration. Production uses the default fixed port and
74/// optional daemon-config override; tests pass port 0 to let the OS assign a free
75/// loopback port and discover it from the connection file.
76#[derive(Debug, Clone, Default)]
77struct AdmissionFactsConfig {
78    carrier_module_id: Option<String>,
79    targets: Option<Vec<String>>,
80}
81
82/// Controls where module cgroups are prepared.
83///
84/// In-process daemons default to [`Self::Disabled`] so they never derive a
85/// production location from the host process. The shipped daemon explicitly uses
86/// [`Self::Current`]. The `ck-subc` binary accepts
87/// `SUBC_CGROUP_PLACEMENT=disabled` for isolated test processes; an unset
88/// variable retains `Current`, and every other value is rejected at startup.
89/// Tests that exercise placement can own a [`Self::Root`].
90#[derive(Debug, Clone, Default, PartialEq, Eq)]
91pub enum CgroupPlacementConfig {
92    #[default]
93    Disabled,
94    Current,
95    Root(PathBuf),
96}
97
98#[derive(Debug, Clone)]
99pub struct BootstrapConfig {
100    pub operator_queue_wait: Duration,
101    pub operator_timeout: Duration,
102    pub operator_stuck_grace: Duration,
103    #[cfg(feature = "test-support")]
104    pub operator_script: Option<PathBuf>,
105    #[cfg(feature = "test-support")]
106    pub operator_events: Option<PathBuf>,
107    pub connection_file_path: PathBuf,
108    pub port: u16,
109    pub daemon_ver: String,
110    configured_modules: Vec<ConfiguredModule>,
111    storage_config: Option<daemon_config::StorageConfig>,
112    admission_facts: AdmissionFactsConfig,
113    /// Module ids whose scopes may carry the attributes that grant authority.
114    scope_authority_owners: Vec<String>,
115    daemon_config_path: Option<PathBuf>,
116    configured_port: Option<u16>,
117    /// Daemon-wide route.bind relay budget in milliseconds (the fallback for
118    /// any module without a per-module override). `None` = built-in default
119    /// (12s — see `control::DEFAULT_ROUTE_BIND_RELAY_TIMEOUT`).
120    route_bind_relay_default_ms: Option<u64>,
121    reserved_capabilities: BTreeMap<String, String>,
122    watchdog_config: DaemonSelfWatchdogConfig,
123    connection_file_source: ConnectionFileSource,
124    /// Where module cgroups are prepared. Disabled unless a caller explicitly
125    /// opts in, because `Current` derives a host location from `/proc/self/cgroup`.
126    cgroup_placement: CgroupPlacementConfig,
127    /// Directory the supervisor writes per-module stdout/stderr capture files
128    /// into. `None` disables capture; the shipped binary supplies its real run
129    /// directory explicitly.
130    capture_logs_dir: Option<PathBuf>,
131    /// File the supervisor appends every module's terminal exits to, so exit
132    /// history survives a daemon restart. `None` keeps terminal history in
133    /// memory only (each module's ring). Absent by default for the same reason
134    /// as `capture_logs_dir`: an in-process daemon booted by a test must not
135    /// append its exits to the operator's real `terminals.jsonl`, where they
136    /// would show up in `ck module terminals`. The shipped binary supplies
137    /// `<run dir>/terminals.jsonl` explicitly.
138    terminal_journal_path: Option<PathBuf>,
139    /// Where the machine id is read from, or minted into when absent. `None`
140    /// serves no machine id. Absent by default for the same reason as
141    /// `capture_logs_dir`: an in-process daemon booted by a test must never
142    /// derive the operator's real data home and mint into it. The shipped binary
143    /// supplies `<data home>/cortexkit/machine-id` explicitly.
144    machine_id_path: Option<PathBuf>,
145    /// The live-children record: every supervised process this daemon has
146    /// running, kept so the next daemon can end the ones a crash left behind.
147    /// At startup, before any module is spawned, the previous daemon's record
148    /// here is swept. `None` keeps no record and sweeps nothing, for the same
149    /// reason as `capture_logs_dir`: an in-process daemon booted by a test
150    /// must never signal processes listed in the operator's real run
151    /// directory. The shipped binary supplies `<run dir>/live-children.json`.
152    live_children_path: Option<PathBuf>,
153    privacy_trampoline: Option<PathBuf>,
154}
155
156impl BootstrapConfig {
157    pub fn new(connection_file_path: impl Into<PathBuf>, port: u16) -> Self {
158        Self {
159            connection_file_path: connection_file_path.into(),
160            operator_queue_wait: Duration::from_secs(150),
161            operator_timeout: Duration::from_secs(120),
162            operator_stuck_grace: Duration::from_secs(10),
163            #[cfg(feature = "test-support")]
164            operator_script: None,
165            #[cfg(feature = "test-support")]
166            operator_events: None,
167            port,
168            daemon_ver: DAEMON_VERSION.to_owned(),
169            configured_modules: Vec::new(),
170            storage_config: None,
171            admission_facts: AdmissionFactsConfig::default(),
172            scope_authority_owners: daemon_config::default_scope_authority_owners(),
173            daemon_config_path: None,
174            configured_port: None,
175            route_bind_relay_default_ms: None,
176            reserved_capabilities: BTreeMap::new(),
177            watchdog_config: DaemonSelfWatchdogConfig::default(),
178            connection_file_source: ConnectionFileSource::Explicit,
179            cgroup_placement: CgroupPlacementConfig::default(),
180            capture_logs_dir: None,
181            terminal_journal_path: None,
182            machine_id_path: None,
183            live_children_path: None,
184            privacy_trampoline: None,
185        }
186    }
187
188    /// Keep the live-children record at `path`, and at startup end the
189    /// processes a previous daemon recorded there that are still running.
190    /// Embedding daemons and tests pass a path inside their own fixture tree.
191    pub fn with_live_children_record(mut self, path: impl Into<PathBuf>) -> Self {
192        self.live_children_path = Some(path.into());
193        self
194    }
195
196    /// macOS module launches require an explicit executable implementing the
197    /// subc-os trampoline entry point. Other systems ignore this setting.
198    pub fn with_privacy_trampoline(mut self, path: impl Into<PathBuf>) -> Self {
199        self.privacy_trampoline = Some(path.into());
200        self
201    }
202
203    /// Serve the machine id stored at `path`, minting it there at startup when
204    /// the file is absent. Embedding daemons and tests pass a path inside their
205    /// own fixture tree.
206    pub fn with_machine_id_path(mut self, path: impl Into<PathBuf>) -> Self {
207        self.machine_id_path = Some(path.into());
208        self
209    }
210
211    /// Selects module cgroup placement. The default is disabled.
212    pub fn with_cgroup_placement(mut self, placement: CgroupPlacementConfig) -> Self {
213        self.cgroup_placement = placement;
214        self
215    }
216
217    /// Redirects per-module stdout/stderr capture files out of the real run
218    /// directory. Tests that start an in-process daemon must call this with a
219    /// path inside their fixture tree.
220    pub fn with_capture_logs_dir(mut self, dir: impl Into<PathBuf>) -> Self {
221        self.capture_logs_dir = Some(dir.into());
222        self
223    }
224
225    /// Redirects the daemon-private journal, allowing embedded daemons and tests
226    /// to keep their observations out of the operator's live run directory.
227    pub fn with_terminal_journal_path(mut self, path: impl Into<PathBuf>) -> Self {
228        self.terminal_journal_path = Some(path.into());
229        self
230    }
231
232    pub fn from_env() -> Result<Self, BootstrapError> {
233        Self::from_env_with_daemon_config_path(daemon_config::default_config_path())
234    }
235
236    /// The SHIPPED BINARY's config, which is the only caller that should capture
237    /// child output into the operator's real run directory.
238    ///
239    /// Kept separate from `from_env` deliberately: see `capture_logs_dir` and
240    /// `terminal_journal_path` on this struct for why an absent value must mean
241    /// NO CAPTURE and NO JOURNAL rather than the operator's real run directory.
242    /// A run directory that cannot be resolved (a relative data home) refuses
243    /// startup instead of landing under the working directory.
244    pub fn from_env_for_daemon_binary() -> Result<Self, BootstrapError> {
245        let run_dir = daemon_config::daemon_run_dir().map_err(BootstrapError::RunDir)?;
246        let machine_id_path =
247            crate::machine_id::default_machine_id_path().map_err(BootstrapError::MachineId)?;
248        let config = Self::from_env()?
249            .with_capture_logs_dir(run_dir.join("logs"))
250            .with_terminal_journal_path(run_dir.join("terminals.jsonl"))
251            .with_live_children_record(crate::live_children::record_path(&run_dir))
252            .with_machine_id_path(machine_id_path);
253        #[cfg(feature = "test-support")]
254        let config = {
255            let mut config = config;
256            config.operator_script =
257                std::env::var_os("SUBC_TEST_OPERATOR_SCRIPT").map(PathBuf::from);
258            config.operator_events =
259                std::env::var_os("SUBC_TEST_OPERATOR_EVENTS").map(PathBuf::from);
260            for (name, value) in [
261                (
262                    "SUBC_TEST_OPERATOR_QUEUE_WAIT_MS",
263                    &mut config.operator_queue_wait,
264                ),
265                (
266                    "SUBC_TEST_OPERATOR_TIMEOUT_MS",
267                    &mut config.operator_timeout,
268                ),
269                (
270                    "SUBC_TEST_OPERATOR_STUCK_MS",
271                    &mut config.operator_stuck_grace,
272                ),
273            ] {
274                if let Ok(ms) = std::env::var(name) {
275                    if let Ok(ms) = ms.parse::<u64>() {
276                        *value = Duration::from_millis(ms);
277                    }
278                }
279            }
280            config
281        };
282        Ok(config)
283    }
284
285    pub fn from_env_with_daemon_config_path(
286        daemon_config_path: impl AsRef<Path>,
287    ) -> Result<Self, BootstrapError> {
288        let daemon_config_path = daemon_config_path.as_ref().to_path_buf();
289        let daemon_config =
290            daemon_config::load(&daemon_config_path).map_err(BootstrapError::DaemonConfig)?;
291        let config_port = daemon_config.as_ref().and_then(|config| config.port);
292        let storage_config = daemon_config
293            .as_ref()
294            .and_then(|config| config.storage.clone());
295        let admission_facts_carrier_module_id = daemon_config
296            .as_ref()
297            .and_then(|config| config.admission_facts_carrier_module_id.clone());
298        let admission_facts_targets = daemon_config
299            .as_ref()
300            .and_then(|config| config.admission_facts_targets.clone());
301        let scope_authority_owners = daemon_config
302            .as_ref()
303            .map(|config| config.scope_authority_owners.clone())
304            .unwrap_or_else(daemon_config::default_scope_authority_owners);
305        let route_bind_relay_default_ms = daemon_config
306            .as_ref()
307            .and_then(|config| config.route_bind_relay_timeout_ms);
308        let reserved_capabilities = daemon_config
309            .as_ref()
310            .map(|config| config.reserved_capabilities.clone())
311            .unwrap_or_default();
312        let configured_modules = daemon_config
313            .map(|config| config.modules)
314            .unwrap_or_default();
315
316        let port = match env::var(SUBC_PORT_ENV) {
317            Ok(raw) if !raw.trim().is_empty() => {
318                let port = raw
319                    .parse::<u16>()
320                    .map_err(|source| BootstrapError::InvalidPort { raw, source })?;
321                if let Some(config_port) = config_port {
322                    info!(
323                        env = SUBC_PORT_ENV,
324                        env_port = port,
325                        config_port,
326                        "SUBC_PORT overrides daemon config port"
327                    );
328                }
329                port
330            }
331            Ok(_) | Err(_) => config_port.unwrap_or(DEFAULT_SUBC_PORT),
332        };
333
334        let (connection_file_path, connection_file_source) =
335            connection_file_path_with_source(non_empty_os_var("XDG_RUNTIME_DIR"));
336        Ok(Self::new(connection_file_path, port)
337            .with_configured_modules(configured_modules)
338            .with_storage_config(storage_config)
339            .with_admission_facts_config(admission_facts_carrier_module_id, admission_facts_targets)
340            .with_scope_authority_owners(scope_authority_owners)
341            .with_route_bind_relay_default_ms(route_bind_relay_default_ms)
342            .with_reserved_capabilities(reserved_capabilities)
343            .with_daemon_config_source(daemon_config_path, config_port)
344            .with_connection_file_source(connection_file_source))
345    }
346
347    pub fn with_daemon_config_path(
348        self,
349        daemon_config_path: impl AsRef<Path>,
350    ) -> Result<Self, BootstrapError> {
351        let daemon_config_path = daemon_config_path.as_ref().to_path_buf();
352        let daemon_config =
353            daemon_config::load(&daemon_config_path).map_err(BootstrapError::DaemonConfig)?;
354        let configured_port = daemon_config.as_ref().and_then(|config| config.port);
355        let storage_config = daemon_config
356            .as_ref()
357            .and_then(|config| config.storage.clone());
358        let admission_facts_carrier_module_id = daemon_config
359            .as_ref()
360            .and_then(|config| config.admission_facts_carrier_module_id.clone());
361        let admission_facts_targets = daemon_config
362            .as_ref()
363            .and_then(|config| config.admission_facts_targets.clone());
364        let scope_authority_owners = daemon_config
365            .as_ref()
366            .map(|config| config.scope_authority_owners.clone())
367            .unwrap_or_else(daemon_config::default_scope_authority_owners);
368        let route_bind_relay_default_ms = daemon_config
369            .as_ref()
370            .and_then(|config| config.route_bind_relay_timeout_ms);
371        let reserved_capabilities = daemon_config
372            .as_ref()
373            .map(|config| config.reserved_capabilities.clone())
374            .unwrap_or_default();
375        let configured_modules = daemon_config
376            .map(|config| config.modules)
377            .unwrap_or_default();
378        Ok(self
379            .with_configured_modules(configured_modules)
380            .with_storage_config(storage_config)
381            .with_admission_facts_config(admission_facts_carrier_module_id, admission_facts_targets)
382            .with_scope_authority_owners(scope_authority_owners)
383            .with_route_bind_relay_default_ms(route_bind_relay_default_ms)
384            .with_reserved_capabilities(reserved_capabilities)
385            .with_daemon_config_source(daemon_config_path, configured_port))
386    }
387
388    pub fn with_configured_modules(
389        mut self,
390        modules: impl IntoIterator<Item = ConfiguredModule>,
391    ) -> Self {
392        self.configured_modules = modules.into_iter().collect();
393        self.configured_modules
394            .sort_by(|left, right| left.module_id.cmp(&right.module_id));
395        self
396    }
397
398    pub fn with_storage_config(
399        mut self,
400        storage_config: Option<daemon_config::StorageConfig>,
401    ) -> Self {
402        self.storage_config = storage_config;
403        self
404    }
405
406    pub fn with_admission_facts_config(
407        mut self,
408        carrier_module_id: Option<String>,
409        targets: Option<Vec<String>>,
410    ) -> Self {
411        self.admission_facts = AdmissionFactsConfig {
412            carrier_module_id,
413            targets,
414        };
415        self
416    }
417
418    /// Set the module ids whose scopes may carry `agent_id` and `delegates`.
419    /// Absent from a config, this is `daemon_config::default_scope_authority_owners`.
420    pub fn with_scope_authority_owners(mut self, owners: Vec<String>) -> Self {
421        self.scope_authority_owners = owners;
422        self
423    }
424
425    /// Set the daemon-wide route.bind relay default (the fallback for any
426    /// module without a per-module override). `None` preserves the built-in
427    /// default (12s). `serve_bound_daemon` reads this at startup and threads
428    /// it into the control handler's daemon-wide field.
429    pub fn with_route_bind_relay_default_ms(mut self, ms: Option<u64>) -> Self {
430        self.route_bind_relay_default_ms = ms;
431        self
432    }
433
434    pub fn with_reserved_capabilities(
435        mut self,
436        reserved_capabilities: BTreeMap<String, String>,
437    ) -> Self {
438        self.reserved_capabilities = reserved_capabilities;
439        self
440    }
441
442    fn with_daemon_config_source(
443        mut self,
444        daemon_config_path: PathBuf,
445        configured_port: Option<u16>,
446    ) -> Self {
447        self.daemon_config_path = Some(daemon_config_path);
448        self.configured_port = configured_port;
449        self
450    }
451
452    fn with_connection_file_source(mut self, source: ConnectionFileSource) -> Self {
453        self.connection_file_source = source;
454        self
455    }
456
457    pub fn with_watchdog_config(mut self, watchdog_config: DaemonSelfWatchdogConfig) -> Self {
458        self.watchdog_config = watchdog_config;
459        self
460    }
461}
462
463/// Result of singleton discovery.
464// Built once per daemon start and immediately matched, so the size gap between
465// the variants costs nothing worth a box on the public shape.
466#[allow(clippy::large_enum_variant)]
467#[derive(Debug)]
468pub enum Outcome {
469    /// A live daemon authenticated from the connection file; this invocation should exit 0.
470    AlreadyRunning,
471    /// This process won the singleton race, owns bound loopback listener(s), and
472    /// has published a fresh connection file.
473    Bound(BoundDaemon),
474}
475
476#[derive(Debug)]
477pub struct BoundDaemon {
478    pub listeners: Vec<TcpListener>,
479    pub connection_info: ConnectionInfo,
480    pub connection_file_path: PathBuf,
481    pub connection_file_source: ConnectionFileSource,
482    /// The machine id this daemon serves, established before any connection is
483    /// accepted. `None` when the config named no machine id path.
484    pub machine_id: Option<crate::machine_id::MachineId>,
485    /// Ownership of the run directory, held until the daemon exits so no other
486    /// daemon sweeps or writes this one's run state. `None` when the config
487    /// keeps no live-children record, and so has no run directory to own.
488    run_dir_lock: Option<crate::run_dir_lock::RunDirLock>,
489    /// The singleton start lock, kept when binding deferred publishing the
490    /// connection file. It stays held until the daemon has verified it is serving
491    /// and written the connection file, so a second daemon starting meanwhile
492    /// cannot claim the singleton too. `None` when the connection
493    /// file was already written at bind time.
494    publication_lock: Option<StartLock>,
495}
496
497/// Resolve subc's per-user TCP connection-file path.
498///
499/// `$XDG_RUNTIME_DIR/subc-connection.json` is preferred because the runtime
500/// directory is already per-user on Unix desktops. Without it, subc falls back
501/// to the system temp dir with a per-user token in the filename so different OS
502/// users do not collide on shared temp directories.
503pub fn connection_file_path() -> PathBuf {
504    connection_file_path_with_source(non_empty_os_var("XDG_RUNTIME_DIR")).0
505}
506
507fn connection_file_path_with_source(
508    runtime_dir: Option<OsString>,
509) -> (PathBuf, ConnectionFileSource) {
510    if let Some(runtime_dir) = runtime_dir.filter(|value| !value.is_empty()) {
511        return (
512            PathBuf::from(runtime_dir).join(CONNECTION_FILE_NAME),
513            ConnectionFileSource::XdgRuntimeDir,
514        );
515    }
516
517    (
518        env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token())),
519        ConnectionFileSource::TempDirFallback,
520    )
521}
522
523/// Resolve, claim, and serve the per-user daemon singleton.
524///
525/// A second invocation is successful: if a live daemon authenticates from the
526/// existing connection file, this returns `Ok(())` after logging and the caller
527/// exits with status 0.
528pub async fn run() -> Result<(), BootstrapError> {
529    // `run` is a binary entry point, so it opts into the production cgroup and
530    // child-log locations. In-process callers keep both features disabled by default.
531    run_with_config(
532        BootstrapConfig::from_env_for_daemon_binary()?
533            .with_cgroup_placement(CgroupPlacementConfig::Current),
534    )
535    .await
536}
537
538/// Serve a daemon from an explicit config. This is the entry point the twelve
539/// sibling repos use to boot an in-process daemon in their integration tests.
540///
541/// # THIS INSTALLS NO TRACING SUBSCRIBER, SO THE DAEMON IS SILENT BY DEFAULT
542///
543/// The daemon's own diagnostics go through `tracing`, and `tracing` DISCARDS
544/// every event when no subscriber is installed. The shipped binary installs one
545/// in `main` (`init_tracing`); this function deliberately does not, because a
546/// library that installs a global subscriber fights with whatever the host
547/// process already set up.
548///
549/// The consequence for a test harness is not "less verbose": it is that the
550/// daemon has NOTHING TO SAY about any failure, and a missing instrument reads
551/// exactly like a clean one. PLEX found this on 2026-09-18 while trying to
552/// capture daemon logs beside an intermittent bind failure, and discovered the
553/// daemon had been silent in every conformance run that repo had ever done --
554/// so the one client-side error string was all the evidence that could exist,
555/// and they had spent a real investigation on a failure whose second source was
556/// never being recorded.
557///
558/// Install one in the harness before calling this, and assert it did something
559/// (they measured 236 daemon lines with the subscriber installed, 0 with the
560/// call commented out) -- otherwise the fix is itself unverified.
561pub async fn run_with_config(config: BootstrapConfig) -> Result<(), BootstrapError> {
562    let forwarding = Arc::new(ForwardingTable::default());
563    forwarding.operator_confirms().configure_bootstrap(&config);
564    let configured_modules = config.configured_modules.clone();
565    let storage_config = config.storage_config.clone();
566    let admission_facts = config.admission_facts.clone();
567    let scope_authority_owners = config.scope_authority_owners.clone();
568    let daemon_config_path = config.daemon_config_path.clone();
569    let configured_port = config.configured_port;
570    let route_bind_relay_default_ms = config.route_bind_relay_default_ms;
571    let reserved_capabilities = config.reserved_capabilities.clone();
572    let watchdog_config = config.watchdog_config.clone();
573    let cgroup_placement_config = config.cgroup_placement.clone();
574    let capture_logs_dir = config.capture_logs_dir.clone();
575    let terminal_journal_path = config.terminal_journal_path.clone();
576    let live_children_path = config.live_children_path.clone();
577    let privacy_trampoline = config.privacy_trampoline.clone();
578    match ensure_singleton_inner(config, false).await? {
579        Outcome::AlreadyRunning => {
580            info!("subc daemon already running");
581            Ok(())
582        }
583        Outcome::Bound(bound) => {
584            #[cfg(target_os = "linux")]
585            let cgroup_placement = prepare_cgroup_placement(&cgroup_placement_config);
586            #[cfg(not(target_os = "linux"))]
587            let _ = cgroup_placement_config;
588            serve_bound_daemon(
589                bound,
590                forwarding,
591                configured_modules,
592                storage_config,
593                admission_facts,
594                scope_authority_owners,
595                daemon_config_path,
596                configured_port,
597                route_bind_relay_default_ms,
598                reserved_capabilities,
599                watchdog_config,
600                capture_logs_dir,
601                terminal_journal_path,
602                live_children_path,
603                privacy_trampoline,
604                #[cfg(target_os = "linux")]
605                cgroup_placement,
606            )
607            .await
608        }
609    }
610}
611
612#[cfg(target_os = "linux")]
613fn prepare_cgroup_placement(config: &CgroupPlacementConfig) -> Option<subc_cgroup::Placement> {
614    let result = match config {
615        CgroupPlacementConfig::Disabled => return None,
616        CgroupPlacementConfig::Current => subc_cgroup::prepare_current(),
617        CgroupPlacementConfig::Root(root) => subc_cgroup::prepare_at(root),
618    };
619
620    match result {
621        Ok(Some(placement)) => Some(placement),
622        Ok(None) => {
623            warn!(
624                placement = ?config,
625                "module cgroup placement is disabled: configured cgroup root is not delegated"
626            );
627            None
628        }
629        Err(error) => {
630            warn!(
631                placement = ?config,
632                error = %error,
633                "module cgroup placement is disabled by an unexpected cgroup probe error"
634            );
635            None
636        }
637    }
638}
639
640/// Target soft limit for open file descriptors, applied to the daemon before any
641/// module is spawned so children inherit it. Multi-root modules (one process
642/// aggregating every project root's sqlite stores, index caches, watchers, and
643/// LSP pipes) trivially exceed the macOS default soft limit of 256; a launchd
644/// user agent does not pass login-shell ulimits through, so the raise must
645/// happen in-process.
646#[cfg(unix)]
647const NOFILE_TARGET: u64 = 65536;
648
649/// Raise RLIMIT_NOFILE to `NOFILE_TARGET` (clamped to the hard limit).
650/// Best-effort: failure is logged and never fatal, since the daemon can run
651/// under the inherited limit — modules with few roots just have less headroom.
652#[cfg(unix)]
653fn raise_nofile_limit() {
654    #[cfg(target_os = "macos")]
655    let ceiling = std::process::Command::new("/usr/sbin/sysctl")
656        .args(["-n", "kern.maxfilesperproc"])
657        .output()
658        .ok()
659        .filter(|output| output.status.success())
660        .and_then(|output| String::from_utf8(output.stdout).ok())
661        .and_then(|value| value.trim().parse::<u64>().ok())
662        .filter(|value| *value > 0);
663    #[cfg(not(target_os = "macos"))]
664    let ceiling = None;
665    raise_nofile_limit_with_ceiling(ceiling);
666}
667
668#[cfg(unix)]
669fn raise_nofile_limit_with_ceiling(kernel_ceiling: Option<u64>) {
670    match rlimit::Resource::NOFILE.get() {
671        Ok((soft, hard)) => {
672            // On XNU an infinite hard limit does not remove the separate
673            // per-process kernel ceiling. Exceeding it makes setrlimit fail
674            // without raising the inherited (often 256) soft limit at all.
675            let target = NOFILE_TARGET
676                .min(hard)
677                .min(kernel_ceiling.unwrap_or(u64::MAX));
678            if soft >= target {
679                return;
680            }
681            match rlimit::Resource::NOFILE.set(target, hard) {
682                Ok(()) => info!(
683                    previous_soft = soft,
684                    new_soft = target,
685                    hard,
686                    "raised open-file soft limit for daemon and module children"
687                ),
688                Err(err) => warn!(
689                    soft,
690                    hard,
691                    error = %err,
692                    "could not raise open-file soft limit; multi-root modules may exhaust descriptors"
693                ),
694            }
695        }
696        Err(err) => warn!(error = %err, "could not read open-file limit"),
697    }
698}
699
700/// CRT stdio-stream target on Windows (the `_setmaxstdio` maximum). Win32
701/// HANDLEs — what Rust `File`, tokio sockets, and SQLite's Win32 VFS actually
702/// consume — have a per-process quota in the millions and need no raise; the
703/// C-runtime stream table (default 512) is the only low ceiling, and it is
704/// per-process rather than inherited, so supervised modules linking the CRT
705/// must raise their own. Raising it here covers the daemon itself.
706#[cfg(windows)]
707fn raise_nofile_limit() {
708    const MAXSTDIO_TARGET: u32 = 8192;
709    let current = rlimit::getmaxstdio();
710    if current >= MAXSTDIO_TARGET {
711        return;
712    }
713    match rlimit::setmaxstdio(MAXSTDIO_TARGET) {
714        Ok(new_max) => info!(
715            previous = current,
716            new_max, "raised CRT stdio-stream limit for daemon"
717        ),
718        Err(err) => warn!(
719            current,
720            error = %err,
721            "could not raise CRT stdio-stream limit"
722        ),
723    }
724}
725
726#[cfg(not(any(unix, windows)))]
727fn raise_nofile_limit() {}
728
729pub async fn run_with_daemon_config_path(
730    config: BootstrapConfig,
731    daemon_config_path: impl AsRef<Path>,
732) -> Result<(), BootstrapError> {
733    run_with_config(config.with_daemon_config_path(daemon_config_path)?).await
734}
735
736#[allow(clippy::too_many_arguments)]
737async fn serve_bound_daemon(
738    bound: BoundDaemon,
739    forwarding: Arc<ForwardingTable>,
740    configured_modules: Vec<ConfiguredModule>,
741    storage_config: Option<daemon_config::StorageConfig>,
742    admission_facts: AdmissionFactsConfig,
743    scope_authority_owners: Vec<String>,
744    daemon_config_path: Option<PathBuf>,
745    configured_port: Option<u16>,
746    route_bind_relay_default_ms: Option<u64>,
747    reserved_capabilities: BTreeMap<String, String>,
748    watchdog_config: DaemonSelfWatchdogConfig,
749    capture_logs_dir: Option<PathBuf>,
750    terminal_journal_path: Option<PathBuf>,
751    live_children_path: Option<PathBuf>,
752    privacy_trampoline: Option<PathBuf>,
753    #[cfg(target_os = "linux")] cgroup_placement: Option<subc_cgroup::Placement>,
754) -> Result<(), BootstrapError> {
755    #[cfg(unix)]
756    let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
757        .map_err(BootstrapError::Signal)?;
758    // Windows has no SIGTERM. Ctrl-C is a different event, not an equivalent
759    // service-stop contract, so this shutdown handler is intentionally Unix-only.
760    raise_nofile_limit();
761
762    info!(
763        connection_file = %bound.connection_file_path.display(),
764        connection_file_source = %bound.connection_file_source,
765        connection_file_source_reason = bound.connection_file_source.reason(),
766        endpoints = ?bound.connection_info.endpoints,
767        configured_modules = configured_modules.len(),
768        machine_id = bound.machine_id.as_ref().map(|id| id.as_str()).unwrap_or("none"),
769        "subc daemon starting"
770    );
771
772    // Before anything can spawn a module (the configured modules below, or a
773    // client's start request once the listeners are served): a previous daemon
774    // that died without its shutdown stop may have left children running, and
775    // a fresh copy beside one would fight it for its port and stores. The
776    // record is not a running daemon's because this daemon holds the lock on
777    // the run directory the record lives in, and a running daemon holds that
778    // lock for its whole life. (Claiming the singleton is not enough: it is
779    // keyed on the connection file, which can live in a different runtime
780    // directory from the run directory.) The lock stays held until this
781    // function returns, which is when the daemon exits.
782    let run_dir_lock = bound.run_dir_lock;
783    let publication_lock = bound.publication_lock;
784    if let Some(owner) = &run_dir_lock {
785        crate::live_children::sweep_orphans(
786            owner,
787            &crate::live_children::AdoptedPids::none(),
788            crate::live_children::SweepBounds::default(),
789        )
790        .await;
791    }
792
793    let registry = Arc::new(Registry::default());
794    let process_liveness = Arc::new(SupervisorProcessLiveness::new());
795    let supervisor_handle = SupervisorHandle::new();
796    let connected_clients = ConnectedClients::new();
797    let daemon_incarnation = format!(
798        "{:032x}",
799        u128::from_be_bytes(bound.connection_info.daemon_id)
800    );
801    let supervisor = Supervisor::new(Arc::clone(&registry), RestartPolicy::default())
802        .with_process_liveness(process_liveness.clone())
803        .with_forwarding(Arc::clone(&forwarding))
804        .with_handle(supervisor_handle.clone())
805        .with_connection_file_path(bound.connection_file_path.clone())
806        .with_daemon_incarnation(daemon_incarnation.clone());
807    let supervisor = match privacy_trampoline {
808        Some(path) => supervisor.with_privacy_trampoline(path),
809        None => supervisor,
810    };
811    // ABSENT MEANS NO CAPTURE AND NO JOURNAL, NOT "THE REAL RUN DIRECTORY", and
812    // the difference is a production-corruption hazard rather than a preference.
813    // Both fields below follow the same rule: `None` for `capture_logs_dir`
814    // means supervised output is not captured, and `None` for
815    // `terminal_journal_path` means terminal history lives only in each
816    // module's in-memory ring (`supervisor.terminals` still answers, with the
817    // journal counters at zero).
818    //
819    // The terminal journal used to fall back to
820    // `daemon_run_dir().join("terminals.jsonl")`, so an in-process test daemon
821    // appended its fixture exits to the operator's journal, where they then
822    // appeared in `ck module terminals`.
823    //
824    // The capture line used to be `unwrap_or_else(|| daemon_run_dir().join("logs"))`,
825    // so ANY caller that did not set the field captured supervised children into
826    // the operator's live `~/.local/share/cortexkit/run/logs/`. That is twelve
827    // sibling repos whose integration tests boot an in-process daemon through
828    // `run_with_config` -- none of which asked for it, and none of which can see
829    // it from their side.
830    //
831    // Harmless while fixture module ids are fixture-shaped: this host carries 20
832    // zero-byte files from subc's own tests (good-aft, missing-aft,
833    // preview-consumer...). THE HAZARD IS A COLLISION. A fixture named "broca"
834    // or "aft" appends to a PRODUCTION capture file that operators read
835    // forensically and that placement gates count lines in -- with no residue to
836    // notice, because the file legitimately exists and legitimately grows.
837    //
838    // Found by BROCA (2026-09-19) from the other side: their rigs spawn the
839    // SHIPPED ck-subc and set XDG_CONFIG_HOME + XDG_RUNTIME_DIR but not
840    // XDG_DATA_HOME, so every local rig run supervised a module named "broca"
841    // and captured it into production's broca.stderr.log -- the same file I
842    // count seal lines in before and after placing their binaries.
843    //
844    // The supervisor already treats `None` as no-capture and no-journal, so
845    // this only removes invented defaults. The binary keeps capturing and
846    // journaling via `BootstrapConfig::from_env_for_daemon_binary`.
847    let supervisor = match terminal_journal_path {
848        Some(path) => supervisor.with_terminal_journal(path, daemon_incarnation),
849        None => supervisor,
850    };
851    let supervisor = match capture_logs_dir {
852        Some(dir) => supervisor.with_capture_logs_dir(dir),
853        None => supervisor,
854    };
855    let supervisor = match live_children_path {
856        Some(path) => supervisor.with_live_children_record(path),
857        None => supervisor,
858    };
859    #[cfg(target_os = "linux")]
860    let supervisor = supervisor.with_cgroup_placement(cgroup_placement);
861    // Collect per-module route.bind relay overrides BEFORE handing the
862    // `configured_modules` vector to the supervisor (which only needs each
863    // module's `drain_timeout_ms`). Each entry was filled in by parse-time
864    // resolution (per-module > daemon-wide > absent), so modules with no
865    // override are absent from this map and the daemon-wide default applies.
866    let route_bind_relay_timeouts = configured_modules
867        .iter()
868        .filter_map(|module| {
869            module
870                .route_bind_relay_timeout_ms
871                .map(|ms| (module.module_id.clone(), Duration::from_millis(ms)))
872        })
873        .collect::<std::collections::BTreeMap<_, _>>();
874    let control_start_clock = crate::clock::StartClock::capture();
875    let mut control = ControlHandler::with_forwarding(Arc::clone(&registry), forwarding)
876        .with_process_liveness(process_liveness)
877        .with_supervisor(supervisor_handle)
878        .with_connected_clients(connected_clients.clone())
879        .with_storage_config(storage_config)
880        .with_machine_id(bound.machine_id.clone())
881        .with_admission_facts_config(admission_facts.carrier_module_id, admission_facts.targets)
882        .with_scope_authority_owners(scope_authority_owners)
883        .with_route_bind_relay_timeouts(route_bind_relay_timeouts)
884        .with_daemon_provenance(
885            bound.connection_info.pid,
886            control_start_clock.started_at_ms(),
887            std::env::current_exe().ok(),
888            normalized_build_provenance(env!("SUBC_BUILD_GIT_SHA")),
889            normalized_build_provenance(env!("SUBC_BUILD_LOCK_DIGEST")),
890        )
891        .with_daemon_start_clock(control_start_clock)
892        .with_capability_config(
893            configured_modules
894                .iter()
895                .map(|module| (module.module_id.clone(), module.enabled)),
896            reserved_capabilities,
897        );
898    if let Some(ms) = route_bind_relay_default_ms {
899        // A daemon-wide config value overrides the built-in default; a
900        // `None` here leaves the ControlHandler's 12s default in place.
901        control = control.with_route_bind_relay_timeout(Duration::from_millis(ms));
902    }
903    if let Some(config_path) = daemon_config_path {
904        control = control.with_supervisor_rescan(supervisor.clone(), config_path, configured_port);
905    }
906    let control = Arc::new(control);
907    let router = Arc::new(Router::with_control_handler(Arc::clone(&control)));
908    let auth = ServerAuth::new(
909        bound.connection_info.key.clone(),
910        bound.connection_info.daemon_id,
911        bound.connection_info.daemon_ver.clone(),
912    )
913    .with_connected_clients(connected_clients);
914
915    let mut serve_task =
916        AbortOnDrop::new(tokio::spawn(serve_listeners(bound.listeners, router, auth)));
917    verify_serving(&bound.connection_info).await?;
918    // A published endpoint promises that authentication can complete now.
919    // Orphan cleanup and control-plane construction can take longer than an
920    // existing-daemon probe's deadline, so finish both before discovery.
921    write_atomic(&bound.connection_file_path, &bound.connection_info).map_err(|source| {
922        BootstrapError::ConnectionFileWrite {
923            path: bound.connection_file_path.clone(),
924            source,
925        }
926    })?;
927    drop(publication_lock);
928    let _clock_step_task = AbortOnDrop::new(crate::watchdog::spawn_clock_step_monitor());
929    // Off the startup path: it runs `systemctl`, and only ever logs.
930    #[cfg(target_os = "linux")]
931    let _kill_mode_check = AbortOnDrop::new(tokio::spawn(
932        crate::systemd_kill_mode::warn_if_kill_mode_defeats_ordered_shutdown(),
933    ));
934    let _watchdog_task = AbortOnDrop::new(
935        DaemonSelfWatchdog::new(
936            bound.connection_info.clone(),
937            bound.connection_file_path.clone(),
938        )
939        .with_config(watchdog_config)
940        .spawn(),
941    );
942
943    for configured in configured_modules {
944        let enabled = configured.enabled;
945        let health = configured.health.clone();
946        let module_id = configured.module_id.clone();
947        match supervisor.supervise_configured_with_health(
948            configured.module_spec(),
949            enabled,
950            health.clone(),
951            configured.drain_timeout_ms,
952            configured.restart,
953        ) {
954            Ok(_) => {
955                // A raised failure threshold is normally a temporary allowance for a
956                // drive that deliberately stops a module, and it widens the window in
957                // which a genuinely wedged module looks fine. It is only ever noticed
958                // when someone thinks to re-read the config, so a relaxation outlives
959                // its reason silently: a rig ran five days at 240s of tolerance against
960                // a 90s default because a comment promising a revert was mistaken for
961                // the revert. Saying so on every boot costs one line and removes the
962                // need for anyone to remember.
963                let default_threshold = HealthConfig::default().failure_threshold;
964                if enabled && health.failure_threshold > default_threshold {
965                    warn!(
966                        module_id = %module_id,
967                        failure_threshold = health.failure_threshold,
968                        default_threshold,
969                        tolerance_secs = health.cadence.as_secs() * u64::from(health.failure_threshold),
970                        "health failure threshold is relaxed above the default; a wedged module stays unflagged for longer"
971                    );
972                }
973                info!(module_id = %module_id, enabled, "configured module supervised");
974            }
975            Err(err) => {
976                error!(module_id = %module_id, error = %err, "failed to supervise configured module; continuing daemon startup");
977            }
978        }
979    }
980
981    control.refresh_capability_requirements();
982    Arc::clone(&control).spawn_capability_deadline_loop();
983    Arc::clone(&control).spawn_scope_expiry_loop();
984
985    #[cfg(unix)]
986    {
987        tokio::select! {
988            result = serve_task.join() => {
989                return result.map_err(BootstrapError::ServeJoin)?.map_err(BootstrapError::Serve);
990            }
991            _ = terminate.recv() => {}
992        }
993        // First, before the notice, the drain, or any connection close: from
994        // here on a module exit is recorded as `daemon_shutdown` and never
995        // respawned. Also before allowing a second signal to cut the bounded
996        // wait short, so the journal marker is always written.
997        supervisor.begin_daemon_shutdown();
998        // Stop the self-watchdog before closing the listener. Its next tick would
999        // connect to that listener, fail, and log an ERROR indistinguishable from
1000        // a wedged daemon, once for every interval a planned stop lasts.
1001        drop(_watchdog_task);
1002        // Dropping the listener stops new accepts, not established connections:
1003        // their detached tasks must remain live throughout notice and drain.
1004        drop(serve_task);
1005        let escalated = tokio::select! {
1006            biased;
1007            _ = terminate.recv() => {
1008                info!("second SIGTERM: abandoning daemon shutdown wait");
1009                true
1010            }
1011            result = supervisor.drain_for_daemon_shutdown() => {
1012                if let Err(error) = result {
1013                    warn!(%error, "daemon shutdown drain failed; exiting anyway");
1014                }
1015                false
1016            }
1017        };
1018        // Supervised modules lead their own process groups, so a service
1019        // manager's kill of this process's group does not reach them. (A
1020        // systemd unit with KillMode=control-group kills by cgroup instead and
1021        // does reach them; see `systemd_kill_mode`.) The daemon ends them
1022        // itself: EOF (or SIGTERM for a protocol none child) first, then
1023        // signals at each child's own drain deadline.
1024        supervisor
1025            .end_children_for_daemon_shutdown(escalated, async {
1026                terminate.recv().await;
1027            })
1028            .await;
1029        Ok(())
1030    }
1031    #[cfg(not(unix))]
1032    serve_task
1033        .join()
1034        .await
1035        .map_err(BootstrapError::ServeJoin)?
1036        .map_err(BootstrapError::Serve)
1037}
1038
1039fn normalized_build_provenance(value: &str) -> Option<String> {
1040    match value.trim() {
1041        "" | "unavailable" => None,
1042        value => Some(value.to_string()),
1043    }
1044}
1045
1046/// Find an existing daemon or atomically bind loopback TCP for this daemon.
1047///
1048/// The algorithm is intentionally connect-first: an endpoint from the connection
1049/// file is treated as live only after the TCP+key server-proof authenticates for
1050/// that file's key and daemon_id. Stale or foreign connection files are reclaimed
1051/// only while holding the per-user start lock; the TCP port is never the
1052/// singleton primitive.
1053pub async fn ensure_singleton(
1054    connection_file_path: impl AsRef<Path>,
1055    port: u16,
1056) -> Result<Outcome, BootstrapError> {
1057    ensure_singleton_with_config(BootstrapConfig::new(connection_file_path.as_ref(), port)).await
1058}
1059
1060pub async fn ensure_singleton_with_config(
1061    config: BootstrapConfig,
1062) -> Result<Outcome, BootstrapError> {
1063    ensure_singleton_inner(config, true).await
1064}
1065
1066async fn ensure_singleton_inner(
1067    config: BootstrapConfig,
1068    publish: bool,
1069) -> Result<Outcome, BootstrapError> {
1070    let path = config.connection_file_path;
1071
1072    if matches!(probe_existing(&path).await?, Probe::Live) {
1073        return Ok(Outcome::AlreadyRunning);
1074    }
1075
1076    let lock = StartLock::acquire(&path).await?;
1077
1078    // Re-probe after acquiring the start lock so a peer that won the race between
1079    // our first failed probe and the lock acquisition is observed instead of
1080    // overwritten.
1081    if matches!(probe_existing(&path).await?, Probe::Live) {
1082        return Ok(Outcome::AlreadyRunning);
1083    }
1084
1085    // The start lock and the probe above only rule out a daemon using this
1086    // connection file. A daemon started with another runtime directory but
1087    // the same data home shares the run directory while passing both, so the
1088    // run directory needs its own owner: refuse at once if a live daemon holds
1089    // it, before anything below reads or writes run or data-home state.
1090    let run_dir_lock = config
1091        .live_children_path
1092        .as_deref()
1093        .map(crate::run_dir_lock::RunDirLock::acquire)
1094        .transpose()?;
1095
1096    remove_stale_connection_file_if_present(&path)?;
1097
1098    // Established under the start lock, and under the run-directory lock when
1099    // there is one, and before binding: a corrupt file stops boot before
1100    // anything is published, and two daemons racing to start cannot both mint.
1101    let machine_id = config
1102        .machine_id_path
1103        .as_deref()
1104        .map(crate::machine_id::load_or_mint)
1105        .transpose()
1106        .map_err(BootstrapError::MachineId)?;
1107
1108    let (listeners, endpoints) = bind_loopback(config.port).await?;
1109    let connection_info = ConnectionInfo {
1110        schema: SCHEMA_VERSION,
1111        wire_version: Some(PROTOCOL_VERSION),
1112        endpoints,
1113        key: generate_key().map_err(BootstrapError::GenerateConnectionFile)?,
1114        daemon_id: generate_daemon_id().map_err(BootstrapError::GenerateConnectionFile)?,
1115        pid: process::id(),
1116        daemon_ver: config.daemon_ver,
1117    };
1118
1119    // The low-level binding API historically publishes for callers that serve
1120    // these listeners themselves. The daemon entry point defers publication
1121    // until its own authentication server is running.
1122    if publish {
1123        if let Err(source) = write_atomic(&path, &connection_info) {
1124            drop(listeners);
1125            return Err(BootstrapError::ConnectionFileWrite { path, source });
1126        }
1127    }
1128
1129    Ok(Outcome::Bound(BoundDaemon {
1130        listeners,
1131        connection_info,
1132        connection_file_path: path,
1133        connection_file_source: config.connection_file_source,
1134        machine_id,
1135        run_dir_lock,
1136        publication_lock: (!publish).then_some(lock),
1137    }))
1138}
1139
1140#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1141enum Probe {
1142    Live,
1143    StaleOrAbsent,
1144}
1145
1146async fn probe_existing(path: &Path) -> Result<Probe, BootstrapError> {
1147    let info = match connection_file::read(path) {
1148        Ok(info) => info,
1149        Err(source) if is_absent_or_stale_connection_file(&source) => {
1150            return Ok(Probe::StaleOrAbsent)
1151        }
1152        Err(source) => {
1153            return Err(BootstrapError::ConnectionFileRead {
1154                path: path.to_path_buf(),
1155                source,
1156            })
1157        }
1158    };
1159
1160    for endpoint in &info.endpoints {
1161        if matches!(probe_endpoint(&info, endpoint).await, Probe::Live) {
1162            return Ok(Probe::Live);
1163        }
1164    }
1165
1166    Ok(Probe::StaleOrAbsent)
1167}
1168
1169async fn probe_endpoint(info: &ConnectionInfo, endpoint: &Endpoint) -> Probe {
1170    let Ok(ip) = endpoint.host.parse::<IpAddr>() else {
1171        return Probe::StaleOrAbsent;
1172    };
1173    if !ip.is_loopback() {
1174        return Probe::StaleOrAbsent;
1175    }
1176    let addr = SocketAddr::new(ip, endpoint.port);
1177
1178    let mut stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(addr)).await {
1179        Ok(Ok(stream)) => stream,
1180        Ok(Err(_)) | Err(_) => return Probe::StaleOrAbsent,
1181    };
1182
1183    match authenticate_client(&mut stream, info, PROBE_AUTH_DEADLINE).await {
1184        Ok(()) => Probe::Live,
1185        Err(AuthError::DaemonIdMismatch)
1186        | Err(AuthError::InvalidServerProof)
1187        | Err(AuthError::UnexpectedEof { .. })
1188        | Err(AuthError::Timeout { .. })
1189        | Err(AuthError::JsonEncode { .. })
1190        | Err(AuthError::JsonDecode { .. })
1191        | Err(AuthError::Io { .. })
1192        | Err(AuthError::MessageTooLarge { .. })
1193        | Err(AuthError::KeyTooShort { .. })
1194        | Err(AuthError::Random(_))
1195        | Err(AuthError::InvalidClientAuth) => Probe::StaleOrAbsent,
1196    }
1197}
1198
1199async fn verify_serving(info: &ConnectionInfo) -> Result<(), BootstrapError> {
1200    // Scheduling a server task is not a readiness barrier. Authenticate once
1201    // against the mandatory IPv4 listener before exposing discovery to peers.
1202    // This uses the normal probe deadline, with no retries or widened budget.
1203    if let Some(endpoint) = info.endpoints.first() {
1204        if matches!(probe_endpoint(info, endpoint).await, Probe::Live) {
1205            return Ok(());
1206        }
1207    }
1208    Err(BootstrapError::StartupNotServing)
1209}
1210
1211fn is_absent_or_stale_connection_file(err: &ConnectionFileError) -> bool {
1212    match err {
1213        ConnectionFileError::Io { source, .. } if source.kind() == io::ErrorKind::NotFound => true,
1214        ConnectionFileError::JsonRead { .. }
1215        | ConnectionFileError::UnsupportedSchema { .. }
1216        | ConnectionFileError::Invalid { .. }
1217        | ConnectionFileError::KeyTooShort { .. }
1218        // A live daemon always publishes the file owner-only (0600), so a file
1219        // with insecure permissions is never a daemon we should defer to: treat it
1220        // as stale and take over (which republishes a correct 0600 file).
1221        | ConnectionFileError::InsecurePermissions { .. } => true,
1222        ConnectionFileError::MissingParent { .. }
1223        | ConnectionFileError::MissingFileName { .. }
1224        // A writable ancestor is an operator misconfiguration, never evidence
1225        // about whether a daemon is live. Reclaiming the file would republish key
1226        // material into the same directory the refusal is about.
1227        | ConnectionFileError::InsecureParentDirectory { .. }
1228        | ConnectionFileError::Io { .. }
1229        | ConnectionFileError::JsonWrite { .. }
1230        | ConnectionFileError::Random(_)
1231        // A wire mismatch may identify a newer live daemon, so never reclaim its
1232        // connection file merely because this binary cannot speak its envelope.
1233        | ConnectionFileError::WireVersionMismatch { .. } => false,
1234    }
1235}
1236
1237fn remove_stale_connection_file_if_present(path: &Path) -> Result<(), BootstrapError> {
1238    match fs::remove_file(path) {
1239        Ok(()) => Ok(()),
1240        Err(err) if err.kind() == io::ErrorKind::NotFound => Ok(()),
1241        Err(source) => Err(BootstrapError::RemoveStale {
1242            path: path.to_path_buf(),
1243            source,
1244        }),
1245    }
1246}
1247
1248async fn bind_loopback(port: u16) -> Result<(Vec<TcpListener>, Vec<Endpoint>), BootstrapError> {
1249    bind_loopback_with(port, |port| TcpListener::bind((Ipv6Addr::LOCALHOST, port))).await
1250}
1251
1252async fn bind_loopback_with<F, Fut>(
1253    port: u16,
1254    mut bind_v6: F,
1255) -> Result<(Vec<TcpListener>, Vec<Endpoint>), BootstrapError>
1256where
1257    F: FnMut(u16) -> Fut,
1258    Fut: std::future::Future<Output = io::Result<TcpListener>>,
1259{
1260    let v4_host = Ipv4Addr::LOCALHOST;
1261    let v4 = TcpListener::bind((v4_host, port))
1262        .await
1263        .map_err(|source| BootstrapError::Bind {
1264            host: v4_host.to_string(),
1265            port,
1266            source,
1267        })?;
1268    let actual_port = v4
1269        .local_addr()
1270        .map_err(|source| BootstrapError::LocalAddr {
1271            host: v4_host.to_string(),
1272            source,
1273        })?
1274        .port();
1275
1276    let mut listeners = vec![v4];
1277    let mut endpoints = vec![Endpoint {
1278        host: v4_host.to_string(),
1279        port: actual_port,
1280    }];
1281
1282    let v6_host = Ipv6Addr::LOCALHOST;
1283    match bind_v6(actual_port).await {
1284        Ok(v6) => {
1285            listeners.push(v6);
1286            endpoints.push(Endpoint {
1287                host: v6_host.to_string(),
1288                port: actual_port,
1289            });
1290        }
1291        Err(err)
1292            if ipv6_loopback_unavailable(&err)
1293                || (port == 0 && err.kind() == io::ErrorKind::AddrInUse) =>
1294        {
1295            warn!(
1296                port = actual_port,
1297                error = %err,
1298                "IPv6 loopback unavailable; serving only IPv4 loopback"
1299            );
1300        }
1301        Err(source) => {
1302            drop(listeners);
1303            return Err(BootstrapError::Bind {
1304                host: v6_host.to_string(),
1305                port: actual_port,
1306                source,
1307            });
1308        }
1309    }
1310
1311    Ok((listeners, endpoints))
1312}
1313
1314fn ipv6_loopback_unavailable(err: &io::Error) -> bool {
1315    matches!(
1316        err.kind(),
1317        io::ErrorKind::AddrNotAvailable | io::ErrorKind::Unsupported
1318    ) || matches!(err.raw_os_error(), Some(47) | Some(49) | Some(97))
1319}
1320
1321struct AbortOnDrop<T> {
1322    handle: JoinHandle<T>,
1323}
1324
1325impl<T> AbortOnDrop<T> {
1326    fn new(handle: JoinHandle<T>) -> Self {
1327        Self { handle }
1328    }
1329
1330    async fn join(&mut self) -> Result<T, JoinError> {
1331        (&mut self.handle).await
1332    }
1333}
1334
1335impl<T> Drop for AbortOnDrop<T> {
1336    fn drop(&mut self) {
1337        if !self.handle.is_finished() {
1338            self.handle.abort();
1339        }
1340    }
1341}
1342
1343#[derive(Debug)]
1344struct StartLock {
1345    // Keep the locked file handle alive for the duration of bootstrap; closing
1346    // it releases the advisory lock while leaving the stable path in place.
1347    _file: fs::File,
1348}
1349
1350impl StartLock {
1351    async fn acquire(connection_file_path: &Path) -> Result<Self, BootstrapError> {
1352        let path = start_lock_path(connection_file_path);
1353        for _ in 0..START_LOCK_RETRIES {
1354            let file = match open_owner_only_lock(&path) {
1355                Ok(file) => file,
1356                Err(source) => return Err(BootstrapError::StartLockCreate { path, source }),
1357            };
1358            match FileExt::try_lock(&file) {
1359                Ok(()) => return Ok(Self { _file: file }),
1360                Err(TryLockError::WouldBlock) => sleep(START_LOCK_RETRY_DELAY).await,
1361                Err(TryLockError::Error(source)) => {
1362                    return Err(BootstrapError::StartLockCreate { path, source });
1363                }
1364            }
1365        }
1366
1367        Err(BootstrapError::StartLockBusy {
1368            path,
1369            attempts: START_LOCK_RETRIES,
1370        })
1371    }
1372}
1373
1374pub(crate) fn open_owner_only_lock(path: &Path) -> io::Result<fs::File> {
1375    let mut options = fs::OpenOptions::new();
1376    options.read(true).write(true).create(true);
1377    #[cfg(unix)]
1378    {
1379        use std::os::unix::fs::OpenOptionsExt;
1380        options.mode(0o600);
1381    }
1382    options.open(path)
1383}
1384
1385fn start_lock_path(connection_file_path: &Path) -> PathBuf {
1386    let file_name = connection_file_path
1387        .file_name()
1388        .map(|name| name.to_string_lossy())
1389        .unwrap_or_else(|| CONNECTION_FILE_NAME.into());
1390    let lock_name = format!("{file_name}.start-lock");
1391    connection_file_path
1392        .parent()
1393        .filter(|parent| !parent.as_os_str().is_empty())
1394        .unwrap_or_else(|| Path::new("."))
1395        .join(lock_name)
1396}
1397
1398fn non_empty_os_var(key: &str) -> Option<OsString> {
1399    let value = env::var_os(key)?;
1400    if value.is_empty() {
1401        None
1402    } else {
1403        Some(value)
1404    }
1405}
1406
1407/// Bootstrap-layer errors are deliberately typed so startup never panics for
1408/// ordinary daemon-discovery races or stale filesystem state.
1409#[derive(Debug)]
1410pub enum BootstrapError {
1411    /// The bound server could not authenticate its own readiness probe.
1412    StartupNotServing,
1413    #[cfg(unix)]
1414    Signal(io::Error),
1415    InvalidPort {
1416        raw: String,
1417        source: std::num::ParseIntError,
1418    },
1419    ConnectionFileRead {
1420        path: PathBuf,
1421        source: ConnectionFileError,
1422    },
1423    ConnectionFileWrite {
1424        path: PathBuf,
1425        source: ConnectionFileError,
1426    },
1427    GenerateConnectionFile(ConnectionFileError),
1428    StartLockCreate {
1429        path: PathBuf,
1430        source: io::Error,
1431    },
1432    StartLockBusy {
1433        path: PathBuf,
1434        attempts: usize,
1435    },
1436    /// The run-directory lock file could not be created, opened or locked.
1437    RunDirLockCreate {
1438        path: PathBuf,
1439        source: io::Error,
1440    },
1441    /// Another live process, almost certainly a daemon started with a
1442    /// different runtime directory over the same data home, owns the run
1443    /// directory. The daemon does not start. `holder_pid` is read from the
1444    /// lock file when possible and is informational only.
1445    RunDirBusy {
1446        path: PathBuf,
1447        holder_pid: Option<u32>,
1448    },
1449    RemoveStale {
1450        path: PathBuf,
1451        source: io::Error,
1452    },
1453    Bind {
1454        host: String,
1455        port: u16,
1456        source: io::Error,
1457    },
1458    LocalAddr {
1459        host: String,
1460        source: io::Error,
1461    },
1462    DaemonConfig(DaemonConfigError),
1463    /// The machine id could not be established: its file is corrupt, unreadable
1464    /// or unwritable, or the data home is relative. The daemon does not start.
1465    MachineId(crate::machine_id::MachineIdFileError),
1466    /// The daemon run directory could not be resolved because the data home is
1467    /// relative. The daemon does not start.
1468    RunDir(daemon_config::DaemonRunDirError),
1469    Serve(ServerError),
1470    ServeJoin(tokio::task::JoinError),
1471}
1472
1473impl fmt::Display for BootstrapError {
1474    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1475        match self {
1476            Self::StartupNotServing => write!(f, "bound daemon did not authenticate its readiness probe; connection file was not published"),
1477            #[cfg(unix)]
1478            Self::Signal(error) => write!(f, "failed to register SIGTERM handler: {error}"),
1479            Self::InvalidPort { raw, source } => {
1480                write!(f, "invalid {SUBC_PORT_ENV} value '{raw}': {source}")
1481            }
1482            Self::ConnectionFileRead { path, source } => write!(
1483                f,
1484                "failed to read connection file {}: {source}",
1485                path.display()
1486            ),
1487            Self::ConnectionFileWrite { path, source } => write!(
1488                f,
1489                "failed to publish connection file {}: {source}",
1490                path.display()
1491            ),
1492            Self::GenerateConnectionFile(err) => {
1493                write!(f, "failed to generate connection-file auth material: {err}")
1494            }
1495            Self::StartLockCreate { path, source } => {
1496                write!(
1497                    f,
1498                    "failed to create start lock {}: {source}",
1499                    path.display()
1500                )
1501            }
1502            Self::StartLockBusy { path, attempts } => write!(
1503                f,
1504                "start lock {} remained busy after {attempts} attempts",
1505                path.display()
1506            ),
1507            Self::RunDirLockCreate { path, source } => write!(
1508                f,
1509                "refusing to start: failed to lock run directory via {}: {source}",
1510                path.display()
1511            ),
1512            Self::RunDirBusy { path, holder_pid } => {
1513                write!(
1514                    f,
1515                    "refusing to start: run directory lock {} is held by another daemon",
1516                    path.display()
1517                )?;
1518                if let Some(pid) = holder_pid {
1519                    write!(f, " (pid {pid})")?;
1520                }
1521                write!(
1522                    f,
1523                    "; a live daemon owns this data home's run state (typically one started with a different XDG_RUNTIME_DIR)"
1524                )
1525            }
1526            Self::RemoveStale { path, source } => write!(
1527                f,
1528                "failed to remove stale connection file {}: {source}",
1529                path.display()
1530            ),
1531            Self::Bind { host, port, source } if source.kind() == io::ErrorKind::AddrInUse => {
1532                write!(
1533                    f,
1534                    "port {port} in use on loopback {host}: {source}; set the port in config"
1535                )
1536            }
1537            Self::Bind { host, port, source } => {
1538                write!(f, "failed to bind loopback TCP {host}:{port}: {source}")
1539            }
1540            Self::LocalAddr { host, source } => {
1541                write!(f, "failed to read local address for {host}: {source}")
1542            }
1543            Self::DaemonConfig(err) => write!(f, "failed to load daemon config: {err}"),
1544            Self::MachineId(err) => write!(f, "refusing to start: {err}"),
1545            Self::RunDir(err) => write!(f, "refusing to start: {err}"),
1546            Self::Serve(err) => write!(f, "daemon server failed: {err}"),
1547            Self::ServeJoin(err) => write!(f, "daemon server task failed: {err}"),
1548        }
1549    }
1550}
1551
1552impl Error for BootstrapError {
1553    fn source(&self) -> Option<&(dyn Error + 'static)> {
1554        match self {
1555            Self::StartupNotServing => None,
1556            #[cfg(unix)]
1557            Self::Signal(source) => Some(source),
1558            Self::InvalidPort { source, .. } => Some(source),
1559            Self::ConnectionFileRead { source, .. }
1560            | Self::ConnectionFileWrite { source, .. }
1561            | Self::GenerateConnectionFile(source) => Some(source),
1562            Self::StartLockCreate { source, .. }
1563            | Self::RunDirLockCreate { source, .. }
1564            | Self::RemoveStale { source, .. }
1565            | Self::Bind { source, .. }
1566            | Self::LocalAddr { source, .. } => Some(source),
1567            Self::DaemonConfig(err) => Some(err),
1568            Self::MachineId(err) => Some(err),
1569            Self::RunDir(err) => Some(err),
1570            Self::Serve(err) => Some(err),
1571            Self::ServeJoin(err) => Some(err),
1572            Self::StartLockBusy { .. } | Self::RunDirBusy { .. } => None,
1573        }
1574    }
1575}
1576
1577#[cfg(test)]
1578mod tests {
1579    use super::*;
1580    use crate::server::ServerAuth;
1581    use cortexkit_test_support::ScratchDir;
1582    #[cfg(target_os = "linux")]
1583    use std::collections::BTreeSet;
1584    use std::sync::Mutex;
1585    #[cfg(target_os = "linux")]
1586    use subc_control::ModuleProtocol;
1587    use subc_transport::MIN_KEY_LEN;
1588    use tokio::io::AsyncReadExt;
1589    use tokio::task::JoinHandle;
1590
1591    #[cfg(unix)]
1592    use std::os::unix::fs::PermissionsExt;
1593
1594    static ENV_LOCK: Mutex<()> = Mutex::new(());
1595
1596    #[test]
1597    fn normalized_build_provenance_preserves_real_values() {
1598        assert_eq!(normalized_build_provenance("abc"), Some("abc".to_string()));
1599    }
1600
1601    #[test]
1602    fn normalized_build_provenance_omits_unavailable_and_empty_values() {
1603        assert_eq!(normalized_build_provenance("unavailable"), None);
1604        assert_eq!(normalized_build_provenance(""), None);
1605    }
1606
1607    #[cfg(target_os = "linux")]
1608    fn current_cgroup_path_for_test() -> io::Result<PathBuf> {
1609        let cgroups = fs::read_to_string("/proc/self/cgroup")?;
1610        let relative = cgroups
1611            .lines()
1612            .find_map(|line| line.strip_prefix("0::"))
1613            .ok_or_else(|| {
1614                io::Error::new(io::ErrorKind::Unsupported, "cgroup v2 is unavailable")
1615            })?;
1616        Ok(Path::new("/sys/fs/cgroup").join(relative.trim_start_matches('/')))
1617    }
1618
1619    #[cfg(target_os = "linux")]
1620    fn module_cgroup_directories() -> io::Result<BTreeSet<OsString>> {
1621        let modules = current_cgroup_path_for_test()?.join("subc-modules");
1622        let entries = match fs::read_dir(modules) {
1623            Ok(entries) => entries,
1624            Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(BTreeSet::new()),
1625            Err(error) => return Err(error),
1626        };
1627        let mut directories = BTreeSet::new();
1628        for entry in entries {
1629            let entry = entry?;
1630            if entry.file_type()?.is_dir() {
1631                directories.insert(entry.file_name());
1632            }
1633        }
1634        Ok(directories)
1635    }
1636
1637    #[cfg(target_os = "linux")]
1638    async fn wait_for_path(path: &Path, task: &JoinHandle<Result<(), BootstrapError>>) {
1639        let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1640        while !path.exists() && tokio::time::Instant::now() < deadline {
1641            assert!(
1642                !task.is_finished(),
1643                "daemon exited before creating {}",
1644                path.display()
1645            );
1646            sleep(Duration::from_millis(10)).await;
1647        }
1648        assert!(path.exists(), "daemon did not create {}", path.display());
1649    }
1650
1651    /// An in-process daemon with the default (disabled) cgroup placement must
1652    /// leave the host's live module cgroup tree exactly as it found it.
1653    /// Red-checking this test (making it fail) creates a
1654    /// `cgroup-isolation-probe-*` directory in the live cgroup tree, which must be removed by hand.
1655    #[cfg(target_os = "linux")]
1656    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1657    async fn run_with_config_does_not_reconcile_the_ambient_cgroup_by_default() {
1658        let temp = unique_temp_dir("bootstrap-cgroup-default-disabled");
1659        let module_id = format!("cgroup-isolation-probe-{}", process::id());
1660        let before = module_cgroup_directories().expect("read ambient module cgroups before boot");
1661        assert!(
1662            !before.contains(&OsString::from(&module_id)),
1663            "isolation probe cgroup already exists before this daemon starts"
1664        );
1665        let capture = temp.join("logs").join(format!("{module_id}.stderr.log"));
1666        let module = ConfiguredModule {
1667            module_id,
1668            program: PathBuf::from("sh"),
1669            args: vec!["-c".to_string(), "sleep 30".to_string()],
1670            env: Vec::new(),
1671            log: None,
1672            enabled: true,
1673            reserved: false,
1674            reserved_prefixes: Vec::new(),
1675            protocol: ModuleProtocol::None,
1676            overlap: Default::default(),
1677            health: HealthConfig::default(),
1678            drain_timeout_ms: None,
1679            route_bind_relay_timeout_ms: None,
1680            restart: RestartPolicy::default(),
1681        };
1682        let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1683            .with_configured_modules([module])
1684            .with_capture_logs_dir(temp.join("logs"))
1685            .with_terminal_journal_path(temp.join("terminals.jsonl"));
1686        let task = tokio::spawn(run_with_config(config));
1687
1688        wait_for_path(&capture, &task).await;
1689        let after = module_cgroup_directories().expect("read ambient module cgroups after boot");
1690
1691        task.abort();
1692        assert!(task
1693            .await
1694            .expect_err("aborted daemon task must cancel")
1695            .is_cancelled());
1696        assert_eq!(
1697            before, after,
1698            "an in-process daemon must not create or reconcile ambient module cgroups"
1699        );
1700    }
1701
1702    #[cfg(target_os = "linux")]
1703    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1704    async fn explicit_cgroup_root_is_prepared_inside_the_fixture_tree() {
1705        let temp = unique_temp_dir("bootstrap-cgroup-explicit-root");
1706        let cgroup_root = temp.join("cgroup");
1707        fs::create_dir(&cgroup_root).expect("create scratch cgroup root");
1708        fs::write(cgroup_root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
1709        let modules = cgroup_root.join("subc-modules");
1710        let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1711            .with_cgroup_placement(CgroupPlacementConfig::Root(cgroup_root))
1712            .with_terminal_journal_path(temp.join("terminals.jsonl"));
1713        let task = tokio::spawn(run_with_config(config));
1714
1715        wait_for_path(&modules, &task).await;
1716
1717        task.abort();
1718        assert!(task
1719            .await
1720            .expect_err("aborted daemon task must cancel")
1721            .is_cancelled());
1722        assert!(
1723            fs::read_dir(&modules)
1724                .expect("read prepared modules directory")
1725                .next()
1726                .is_none(),
1727            "the delegation probe must clean up after itself"
1728        );
1729    }
1730
1731    struct EnvGuard {
1732        key: &'static str,
1733        previous: Option<OsString>,
1734    }
1735
1736    impl EnvGuard {
1737        fn set(key: &'static str, value: &Path) -> Self {
1738            let previous = env::var_os(key);
1739            env::set_var(key, value);
1740            Self { key, previous }
1741        }
1742
1743        fn set_str(key: &'static str, value: &str) -> Self {
1744            let previous = env::var_os(key);
1745            env::set_var(key, value);
1746            Self { key, previous }
1747        }
1748
1749        fn unset(key: &'static str) -> Self {
1750            let previous = env::var_os(key);
1751            env::remove_var(key);
1752            Self { key, previous }
1753        }
1754    }
1755
1756    impl Drop for EnvGuard {
1757        fn drop(&mut self) {
1758            match &self.previous {
1759                Some(value) => env::set_var(self.key, value),
1760                None => env::remove_var(self.key),
1761            }
1762        }
1763    }
1764
1765    fn unique_temp_dir(name: &str) -> ScratchDir {
1766        ScratchDir::new(name)
1767    }
1768
1769    fn temp_connection_file_path(name: &str) -> (ScratchDir, PathBuf) {
1770        let dir = unique_temp_dir(name);
1771        let path = dir.join("conn.json");
1772        (dir, path)
1773    }
1774
1775    fn auth_for(info: &ConnectionInfo) -> ServerAuth {
1776        ServerAuth::new(info.key.clone(), info.daemon_id, info.daemon_ver.clone())
1777    }
1778
1779    fn start_server(bound: BoundDaemon) -> JoinHandle<Result<(), ServerError>> {
1780        let auth = auth_for(&bound.connection_info);
1781        tokio::spawn(serve_listeners(
1782            bound.listeners,
1783            Arc::new(Router::with_default_self_handler()),
1784            auth,
1785        ))
1786    }
1787
1788    fn expect_bound(outcome: Outcome) -> BoundDaemon {
1789        match outcome {
1790            Outcome::Bound(bound) => bound,
1791            Outcome::AlreadyRunning => panic!("fresh connection file unexpectedly had a daemon"),
1792        }
1793    }
1794
1795    async fn connect_from_info(conn: &ConnectionInfo) -> io::Result<TcpStream> {
1796        let endpoint = conn
1797            .endpoints
1798            .first()
1799            .expect("test connection file should have an endpoint");
1800        let ip: IpAddr = endpoint.host.parse().unwrap();
1801        TcpStream::connect(SocketAddr::new(ip, endpoint.port)).await
1802    }
1803
1804    fn make_connection_info(port: u16) -> ConnectionInfo {
1805        ConnectionInfo {
1806            schema: SCHEMA_VERSION,
1807            wire_version: Some(PROTOCOL_VERSION),
1808            endpoints: vec![Endpoint {
1809                host: "127.0.0.1".to_owned(),
1810                port,
1811            }],
1812            key: generate_key().unwrap(),
1813            daemon_id: generate_daemon_id().unwrap(),
1814            pid: process::id(),
1815            daemon_ver: "test-subc".to_owned(),
1816        }
1817    }
1818
1819    fn write_raw_owner_only_connection_file(path: &Path, contents: &[u8]) {
1820        fs::write(path, contents).unwrap();
1821        #[cfg(unix)]
1822        fs::set_permissions(path, fs::Permissions::from_mode(0o600)).unwrap();
1823    }
1824
1825    fn assert_owner_only_connection_file(path: &Path) {
1826        // `path` is only inspected on Unix (mode bits); on Windows the owner-only
1827        // guarantee comes from the inherited %TEMP% ACL, nothing to assert here.
1828        #[cfg(unix)]
1829        {
1830            let mode = fs::metadata(path).unwrap().permissions().mode() & 0o777;
1831            assert_eq!(mode, 0o600);
1832        }
1833        #[cfg(not(unix))]
1834        let _ = path;
1835    }
1836
1837    #[test]
1838    fn connection_file_path_uses_xdg_runtime_dir_when_set() {
1839        let _env_lock = ENV_LOCK.lock().unwrap();
1840        let runtime_dir = unique_temp_dir("xdg-runtime");
1841        let _xdg = EnvGuard::set("XDG_RUNTIME_DIR", runtime_dir.path());
1842
1843        assert_eq!(
1844            connection_file_path(),
1845            runtime_dir.join(CONNECTION_FILE_NAME)
1846        );
1847    }
1848
1849    #[test]
1850    fn connection_file_path_source_is_xdg_runtime_dir_when_set() {
1851        let runtime_dir = OsString::from("/run/user/1000");
1852
1853        let (path, source) = connection_file_path_with_source(Some(runtime_dir));
1854
1855        assert_eq!(
1856            path,
1857            PathBuf::from("/run/user/1000").join(CONNECTION_FILE_NAME)
1858        );
1859        assert_eq!(source, ConnectionFileSource::XdgRuntimeDir);
1860    }
1861
1862    #[test]
1863    fn connection_file_path_falls_back_to_temp_dir_with_user_token_when_xdg_unset() {
1864        let _env_lock = ENV_LOCK.lock().unwrap();
1865        let _xdg = EnvGuard::unset("XDG_RUNTIME_DIR");
1866
1867        assert_eq!(
1868            connection_file_path(),
1869            env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1870        );
1871    }
1872
1873    #[test]
1874    fn connection_file_path_source_is_temp_dir_when_xdg_unset() {
1875        let (path, source) = connection_file_path_with_source(None);
1876
1877        assert_eq!(
1878            path,
1879            env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1880        );
1881        assert_eq!(source, ConnectionFileSource::TempDirFallback);
1882    }
1883
1884    /// Concurrent callers must derive one token. The former temp-file uid probe
1885    /// could fail transiently (same-tick name collision, fd exhaustion) and send
1886    /// the loser down the env-derived fallback with a different identity for the
1887    /// same user; this fence keeps identity independent of filesystem luck.
1888    #[test]
1889    fn user_connection_token_is_stable_under_concurrent_callers() {
1890        let expected = user_connection_token();
1891        let workers: Vec<_> = (0..32)
1892            .map(|_| {
1893                std::thread::spawn(|| (0..40).map(|_| user_connection_token()).collect::<Vec<_>>())
1894            })
1895            .collect();
1896        for worker in workers {
1897            for token in worker.join().expect("probe thread") {
1898                assert_eq!(token, expected, "token diverged under concurrent probes");
1899            }
1900        }
1901    }
1902
1903    #[test]
1904    fn connection_file_path_source_is_temp_dir_when_xdg_empty() {
1905        let (path, source) = connection_file_path_with_source(Some(OsString::new()));
1906
1907        assert_eq!(
1908            path,
1909            env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1910        );
1911        assert_eq!(source, ConnectionFileSource::TempDirFallback);
1912    }
1913
1914    /// The shipped binary is the only caller that journals terminal exits into
1915    /// the real run directory, so its constructor must supply that path itself;
1916    /// the in-process default leaves it unset.
1917    #[test]
1918    fn daemon_binary_config_journals_and_captures_into_the_run_dir() {
1919        let _env_lock = ENV_LOCK.lock().unwrap();
1920        let root = unique_temp_dir("daemon-binary-config");
1921        let data_home = root.join("data");
1922        let _data = EnvGuard::set("XDG_DATA_HOME", &data_home);
1923        let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1924        let _port = EnvGuard::unset(SUBC_PORT_ENV);
1925
1926        let config = BootstrapConfig::from_env_for_daemon_binary().unwrap();
1927
1928        let run_dir = data_home.join("cortexkit").join("run");
1929        assert_eq!(
1930            config.terminal_journal_path,
1931            Some(run_dir.join("terminals.jsonl"))
1932        );
1933        assert_eq!(config.capture_logs_dir, Some(run_dir.join("logs")));
1934        assert_eq!(
1935            BootstrapConfig::new(root.join("connection.json"), 0).terminal_journal_path,
1936            None,
1937            "an in-process config must not journal anywhere unless asked to"
1938        );
1939    }
1940
1941    #[test]
1942    fn daemon_binary_config_refuses_a_relative_data_home() {
1943        let _env_lock = ENV_LOCK.lock().unwrap();
1944        let root = unique_temp_dir("daemon-binary-relative-data");
1945        let _data = EnvGuard::set_str("XDG_DATA_HOME", "relative-data-home");
1946        let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1947        let _port = EnvGuard::unset(SUBC_PORT_ENV);
1948
1949        let error = BootstrapConfig::from_env_for_daemon_binary()
1950            .expect_err("a relative data home must refuse the daemon binary's config");
1951        assert!(
1952            matches!(error, BootstrapError::RunDir(_)),
1953            "expected a run-directory refusal, got {error}"
1954        );
1955        assert!(error.to_string().contains("XDG_DATA_HOME"), "{error}");
1956    }
1957
1958    #[test]
1959    fn configured_port_uses_default_config_and_env_override() {
1960        let _env_lock = ENV_LOCK.lock().unwrap();
1961        let (_dir, conn_path) = temp_connection_file_path("daemon-config-port");
1962        let config_path = conn_path.with_file_name("subc.jsonc");
1963
1964        let _port = EnvGuard::unset(SUBC_PORT_ENV);
1965        assert_eq!(
1966            BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1967                .unwrap()
1968                .port,
1969            DEFAULT_SUBC_PORT
1970        );
1971
1972        fs::write(&config_path, r#"{ "version": 1, "port": 8123 }"#).unwrap();
1973        assert_eq!(
1974            BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1975                .unwrap()
1976                .port,
1977            8123
1978        );
1979
1980        let _port = EnvGuard::set_str(SUBC_PORT_ENV, "9012");
1981        assert_eq!(
1982            BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1983                .unwrap()
1984                .port,
1985            9012
1986        );
1987    }
1988
1989    #[tokio::test]
1990    async fn second_singleton_probe_against_served_tcp_daemon_reports_already_running() {
1991        let (_dir, path) = temp_connection_file_path("already-running");
1992
1993        let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1994        let server = start_server(bound);
1995
1996        let second = ensure_singleton(&path, 0).await.unwrap();
1997        assert!(matches!(second, Outcome::AlreadyRunning));
1998
1999        server.abort();
2000        let _ = server.await;
2001    }
2002
2003    #[tokio::test]
2004    async fn daemon_connection_file_publishes_protocol_wire_version() {
2005        let (_dir, path) = temp_connection_file_path("wire-version");
2006        let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2007        assert_eq!(bound.connection_info.wire_version, Some(PROTOCOL_VERSION));
2008        assert_eq!(
2009            connection_file::read(&path).unwrap().wire_version,
2010            Some(PROTOCOL_VERSION)
2011        );
2012
2013        drop(bound.listeners);
2014    }
2015
2016    #[tokio::test]
2017    async fn stale_unbound_connection_file_is_reclaimed() {
2018        let (_dir, path) = temp_connection_file_path("stale-reclaim");
2019        let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
2020        let stale_port = stale.local_addr().unwrap().port();
2021        drop(stale);
2022        let stale_info = make_connection_info(stale_port);
2023        write_atomic(&path, &stale_info).unwrap();
2024
2025        let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2026        assert_ne!(bound.connection_info.key, stale_info.key);
2027        drop(bound.listeners);
2028    }
2029
2030    #[cfg(unix)]
2031    #[tokio::test]
2032    async fn ensure_singleton_reclaims_insecure_connection_file() {
2033        let (_dir, path) = temp_connection_file_path("insecure-reclaim");
2034        let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
2035        let stale_port = stale.local_addr().unwrap().port();
2036        drop(stale);
2037        let stale_info = make_connection_info(stale_port);
2038        write_atomic(&path, &stale_info).unwrap();
2039        fs::set_permissions(&path, fs::Permissions::from_mode(0o644)).unwrap();
2040
2041        let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2042        assert_ne!(bound.connection_info.key, stale_info.key);
2043        assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
2044        assert_owner_only_connection_file(&path);
2045
2046        drop(bound.listeners);
2047    }
2048
2049    #[tokio::test]
2050    async fn ensure_singleton_reclaims_non_loopback_connection_file() {
2051        let (_dir, path) = temp_connection_file_path("non-loopback-reclaim");
2052        let mut stale_info = make_connection_info(8757);
2053        stale_info.endpoints = vec![Endpoint {
2054            host: "192.0.2.10".to_owned(),
2055            port: 8757,
2056        }];
2057        write_atomic(&path, &stale_info).unwrap();
2058
2059        let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2060        assert_ne!(bound.connection_info.key, stale_info.key);
2061        assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
2062        assert!(bound
2063            .connection_info
2064            .endpoints
2065            .iter()
2066            .all(|endpoint| endpoint.host.parse::<IpAddr>().unwrap().is_loopback()));
2067        assert_owner_only_connection_file(&path);
2068
2069        drop(bound.listeners);
2070    }
2071
2072    #[tokio::test]
2073    async fn ensure_singleton_reclaims_invalid_connection_file_shapes() {
2074        let mut unsupported_schema = make_connection_info(8757);
2075        unsupported_schema.schema = SCHEMA_VERSION + 1;
2076
2077        let mut empty_endpoints = make_connection_info(8757);
2078        empty_endpoints.endpoints.clear();
2079
2080        let mut short_key = make_connection_info(8757);
2081        short_key.key = vec![0x5A; MIN_KEY_LEN - 1];
2082
2083        let cases = vec![
2084            (
2085                "unsupported-schema",
2086                serde_json::to_vec(&unsupported_schema).unwrap(),
2087                Some(unsupported_schema),
2088            ),
2089            (
2090                "empty-endpoints",
2091                serde_json::to_vec(&empty_endpoints).unwrap(),
2092                Some(empty_endpoints),
2093            ),
2094            (
2095                "short-key",
2096                serde_json::to_vec(&short_key).unwrap(),
2097                Some(short_key),
2098            ),
2099            ("invalid-json", b"{not valid connection json".to_vec(), None),
2100        ];
2101
2102        for (label, contents, old_info) in cases {
2103            let (_dir, path) = temp_connection_file_path(label);
2104            write_raw_owner_only_connection_file(&path, &contents);
2105
2106            let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2107            if let Some(old_info) = old_info {
2108                assert_ne!(bound.connection_info.key, old_info.key, "{label}");
2109                assert_ne!(
2110                    bound.connection_info.daemon_id, old_info.daemon_id,
2111                    "{label}"
2112                );
2113            }
2114            assert!(bound.connection_info.key.len() >= MIN_KEY_LEN, "{label}");
2115            assert_ne!(bound.connection_info.daemon_id, [0u8; 16], "{label}");
2116            assert_owner_only_connection_file(&path);
2117
2118            drop(bound.listeners);
2119        }
2120    }
2121
2122    #[tokio::test]
2123    async fn foreign_reused_port_connection_file_is_reclaimed_after_auth_probe_fails() {
2124        let (_dir, path) = temp_connection_file_path("foreign-reclaim");
2125        let foreign = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
2126        let foreign_port = foreign.local_addr().unwrap().port();
2127        write_atomic(&path, &make_connection_info(foreign_port)).unwrap();
2128        let foreign_task = tokio::spawn(async move {
2129            if let Ok((mut stream, _)) = foreign.accept().await {
2130                let mut buf = [0u8; 64];
2131                let _ = stream.read(&mut buf).await;
2132            }
2133        });
2134
2135        let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2136        assert!(bound
2137            .connection_info
2138            .endpoints
2139            .iter()
2140            .all(|endpoint| endpoint.port != foreign_port));
2141
2142        drop(bound.listeners);
2143        let _ = foreign_task.await;
2144    }
2145
2146    #[tokio::test]
2147    async fn stale_start_lock_file_is_reclaimable() {
2148        let (_dir, path) = temp_connection_file_path("start-lock-stale-file");
2149        let lock_path = start_lock_path(&path);
2150        drop(open_owner_only_lock(&lock_path).unwrap());
2151        assert!(lock_path.is_file());
2152
2153        let lock = StartLock::acquire(&path).await.unwrap();
2154        assert!(lock_path.is_file());
2155
2156        drop(lock);
2157        assert!(lock_path.is_file());
2158    }
2159
2160    #[tokio::test]
2161    async fn held_start_lock_blocks_second_acquire_until_release() {
2162        let (_dir, path) = temp_connection_file_path("start-lock-held");
2163        let lock_path = start_lock_path(&path);
2164        let first = StartLock::acquire(&path).await.unwrap();
2165
2166        let err = match StartLock::acquire(&path).await {
2167            Ok(_) => panic!("second acquire while held must stay busy"),
2168            Err(err) => err,
2169        };
2170        assert!(matches!(
2171            err,
2172            BootstrapError::StartLockBusy {
2173                ref path,
2174                attempts: START_LOCK_RETRIES,
2175            } if path == &lock_path
2176        ));
2177
2178        drop(first);
2179
2180        let second = StartLock::acquire(&path)
2181            .await
2182            .expect("released advisory lock should be reclaimable");
2183        drop(second);
2184    }
2185
2186    #[tokio::test]
2187    async fn bind_conflict_on_fixed_port_fails_loud_without_reselecting() {
2188        let (_dir, path) = temp_connection_file_path("bind-conflict");
2189        let occupied = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
2190        let occupied_port = occupied.local_addr().unwrap().port();
2191
2192        let err = ensure_singleton(&path, occupied_port).await.unwrap_err();
2193        assert!(matches!(
2194            err,
2195            BootstrapError::Bind { ref source, .. } if source.kind() == io::ErrorKind::AddrInUse
2196        ));
2197        assert!(err.to_string().contains("set the port in config"));
2198
2199        drop(occupied);
2200    }
2201
2202    #[tokio::test]
2203    async fn ephemeral_ipv6_collision_still_serves_ipv4() {
2204        let (listeners, endpoints) = bind_loopback_with(0, |port| async move {
2205            let occupied = TcpListener::bind((Ipv6Addr::LOCALHOST, port)).await?;
2206            let result = TcpListener::bind((Ipv6Addr::LOCALHOST, port)).await;
2207            drop(occupied);
2208            result
2209        })
2210        .await
2211        .expect("an incidental IPv6 collision must not prevent an ephemeral daemon from starting");
2212        assert_eq!(endpoints.len(), 1);
2213        assert_eq!(endpoints[0].host, "127.0.0.1");
2214        TcpStream::connect((Ipv4Addr::LOCALHOST, endpoints[0].port))
2215            .await
2216            .unwrap();
2217        drop(listeners);
2218    }
2219
2220    #[tokio::test]
2221    async fn unpublished_boot_holds_singleton_lock_until_publication() {
2222        let dir = ScratchDir::new("unpublished-singleton-lock");
2223        let path = dir.join("connection.json");
2224        let bound = expect_bound(
2225            ensure_singleton_inner(BootstrapConfig::new(&path, 0), false)
2226                .await
2227                .unwrap(),
2228        );
2229        assert!(!path.exists());
2230        let contender = open_owner_only_lock(&start_lock_path(&path)).unwrap();
2231        assert!(
2232            matches!(FileExt::try_lock(&contender), Err(TryLockError::WouldBlock)),
2233            "an unpublished boot must keep its singleton claim"
2234        );
2235        drop(bound);
2236        FileExt::try_lock(&contender).unwrap();
2237    }
2238
2239    #[tokio::test]
2240    async fn readiness_probe_requires_an_authenticating_server() {
2241        let dir = ScratchDir::new("unserved-readiness");
2242        let bound = expect_bound(
2243            ensure_singleton_inner(BootstrapConfig::new(dir.join("connection.json"), 0), false)
2244                .await
2245                .unwrap(),
2246        );
2247        assert!(matches!(
2248            verify_serving(&bound.connection_info).await,
2249            Err(BootstrapError::StartupNotServing)
2250        ));
2251    }
2252
2253    #[cfg(unix)]
2254    #[tokio::test]
2255    async fn orphan_sweep_finishes_before_connection_file_publication() {
2256        use crate::live_children::{ExecutableIdentity, LiveChild};
2257        let dir = ScratchDir::new("publish-after-sweep");
2258        let path = dir.join("connection.json");
2259        let record = dir.join("live-children.json");
2260        let ready = dir.join("ready");
2261        let mut child = std::process::Command::new("/bin/sh")
2262            .args([
2263                "-c",
2264                "trap '' TERM; : > \"$0\"; read -r _",
2265                ready.to_str().unwrap(),
2266            ])
2267            .stdin(std::process::Stdio::piped())
2268            .env("XDG_DATA_HOME", dir.path())
2269            .env("XDG_RUNTIME_DIR", dir.path())
2270            .env("XDG_CONFIG_HOME", dir.path())
2271            .spawn()
2272            .unwrap();
2273        timeout(Duration::from_secs(5), async {
2274            while !ready.exists() {
2275                sleep(Duration::from_millis(5)).await;
2276            }
2277        })
2278        .await
2279        .unwrap();
2280        let observed = subc_os::Process::open(child.id())
2281            .unwrap()
2282            .unwrap()
2283            .observe()
2284            .unwrap();
2285        crate::live_children::write_record(
2286            &record,
2287            &[LiveChild {
2288                module_id: "old-sleep".into(),
2289                pid: child.id(),
2290                protocol: subc_control::ModuleProtocol::None,
2291                start_time: Some(observed.start_time),
2292                executable: observed.executable.map(ExecutableIdentity::from),
2293                cgroup_name: None,
2294            }],
2295        )
2296        .unwrap();
2297        let config = BootstrapConfig::new(&path, 0).with_live_children_record(&record);
2298        let daemon = tokio::spawn(run_with_config(config));
2299        timeout(Duration::from_secs(5), async {
2300            while !path.exists() {
2301                sleep(Duration::from_millis(5)).await;
2302            }
2303        })
2304        .await
2305        .unwrap();
2306        assert!(path.exists());
2307        let ended = child.try_wait().unwrap().is_some();
2308        if !ended {
2309            child.kill().unwrap();
2310        }
2311        child.wait().unwrap();
2312        if ended {
2313            assert!(matches!(probe_existing(&path).await.unwrap(), Probe::Live));
2314        }
2315        daemon.abort();
2316        let _ = daemon.await;
2317        assert!(
2318            ended,
2319            "the published daemon must not still owe its orphan sweep"
2320        );
2321    }
2322
2323    #[cfg(target_os = "macos")]
2324    #[test]
2325    fn nofile_raise_worker() {
2326        let Some(path) = std::env::var_os("SUBC_NOFILE_TEST_FILE") else {
2327            return;
2328        };
2329        let (_, hard) = rlimit::Resource::NOFILE.get().unwrap();
2330        rlimit::Resource::NOFILE.set(256, hard).unwrap();
2331        raise_nofile_limit_with_ceiling(Some(10240));
2332        std::fs::write(path, rlimit::Resource::NOFILE.get().unwrap().0.to_string()).unwrap();
2333    }
2334
2335    #[cfg(target_os = "macos")]
2336    #[test]
2337    fn macos_nofile_raise_clamps_to_kernel_ceiling() {
2338        let dir = ScratchDir::new("nofile-ceiling");
2339        let result = dir.join("soft");
2340        let status = std::process::Command::new(std::env::current_exe().unwrap())
2341            .args(["--exact", "bootstrap::tests::nofile_raise_worker"])
2342            .env("SUBC_NOFILE_TEST_FILE", &result)
2343            .env("XDG_DATA_HOME", dir.path())
2344            .env("XDG_RUNTIME_DIR", dir.path())
2345            .env("XDG_CONFIG_HOME", dir.path())
2346            .status()
2347            .unwrap();
2348        assert!(status.success());
2349        let soft: u64 = std::fs::read_to_string(result).unwrap().parse().unwrap();
2350        assert!(
2351            soft > 256 && soft <= 10240,
2352            "the inherited limit must be raised within the kernel ceiling; got {soft}"
2353        );
2354    }
2355
2356    #[tokio::test]
2357    async fn key_rotation_republishes_new_material_and_old_file_fails_auth() {
2358        let (_dir, path) = temp_connection_file_path("key-rotation");
2359        let first = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2360        let old_info = first.connection_info.clone();
2361        let fixed_port = old_info.endpoints[0].port;
2362        drop(first.listeners);
2363
2364        let second = expect_bound(ensure_singleton(&path, fixed_port).await.unwrap());
2365        let new_info = second.connection_info.clone();
2366        assert_ne!(old_info.key, new_info.key);
2367        assert_ne!(old_info.daemon_id, new_info.daemon_id);
2368        let server = start_server(second);
2369
2370        let mut old_stream = connect_from_info(&old_info).await.unwrap();
2371        let old_auth = authenticate_client(&mut old_stream, &old_info, PROBE_AUTH_DEADLINE).await;
2372        assert!(
2373            old_auth.is_err(),
2374            "old key must not authenticate after restart"
2375        );
2376
2377        let reread = connection_file::read(&path).unwrap();
2378        let mut new_stream = connect_from_info(&reread).await.unwrap();
2379        authenticate_client(&mut new_stream, &reread, PROBE_AUTH_DEADLINE)
2380            .await
2381            .unwrap();
2382
2383        server.abort();
2384        let _ = server.await;
2385    }
2386
2387    #[cfg(unix)]
2388    #[tokio::test]
2389    async fn published_connection_file_permissions_are_owner_only() {
2390        let (_dir, path) = temp_connection_file_path("permissions");
2391        let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2392
2393        let mode = fs::metadata(&path).unwrap().permissions().mode() & 0o777;
2394        assert_eq!(mode, 0o600);
2395
2396        drop(bound.listeners);
2397    }
2398}