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