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