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