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#[derive(Debug, Clone, Default)]
77struct AdmissionFactsConfig {
78 carrier_module_id: Option<String>,
79 targets: Option<Vec<String>>,
80}
81
82#[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 scope_authority_owners: Vec<String>,
108 daemon_config_path: Option<PathBuf>,
109 configured_port: Option<u16>,
110 route_bind_relay_default_ms: Option<u64>,
114 reserved_capabilities: BTreeMap<String, String>,
115 watchdog_config: DaemonSelfWatchdogConfig,
116 connection_file_source: ConnectionFileSource,
117 cgroup_placement: CgroupPlacementConfig,
120 capture_logs_dir: Option<PathBuf>,
124 terminal_journal_path: Option<PathBuf>,
132 machine_id_path: Option<PathBuf>,
138 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 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 pub fn with_privacy_trampoline(mut self, path: impl Into<PathBuf>) -> Self {
185 self.privacy_trampoline = Some(path.into());
186 self
187 }
188
189 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 pub fn with_cgroup_placement(mut self, placement: CgroupPlacementConfig) -> Self {
199 self.cgroup_placement = placement;
200 self
201 }
202
203 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 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 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 pub fn with_scope_authority_owners(mut self, owners: Vec<String>) -> Self {
377 self.scope_authority_owners = owners;
378 self
379 }
380
381 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#[allow(clippy::large_enum_variant)]
423#[derive(Debug)]
424pub enum Outcome {
425 AlreadyRunning,
427 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 pub machine_id: Option<crate::machine_id::MachineId>,
441 run_dir_lock: Option<crate::run_dir_lock::RunDirLock>,
445 publication_lock: Option<StartLock>,
451}
452
453pub 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
479pub async fn run() -> Result<(), BootstrapError> {
485 run_with_config(
488 BootstrapConfig::from_env_for_daemon_binary()?
489 .with_cgroup_placement(CgroupPlacementConfig::Current),
490 )
491 .await
492}
493
494pub 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#[cfg(unix)]
600const NOFILE_TARGET: u64 = 65536;
601
602#[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 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#[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 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 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(®istry), 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 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 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(®istry), 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 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 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 #[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 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 supervisor.begin_daemon_shutdown();
950 drop(_watchdog_task);
954 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 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
998pub 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 if matches!(probe_existing(&path).await?, Probe::Live) {
1034 return Ok(Outcome::AlreadyRunning);
1035 }
1036
1037 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 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 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 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 | ConnectionFileError::InsecurePermissions { .. } => true,
1174 ConnectionFileError::MissingParent { .. }
1175 | ConnectionFileError::MissingFileName { .. }
1176 | ConnectionFileError::InsecureParentDirectory { .. }
1180 | ConnectionFileError::Io { .. }
1181 | ConnectionFileError::JsonWrite { .. }
1182 | ConnectionFileError::Random(_)
1183 | 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 _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#[derive(Debug)]
1362pub enum BootstrapError {
1363 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 RunDirLockCreate {
1390 path: PathBuf,
1391 source: io::Error,
1392 },
1393 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 MachineId(crate::machine_id::MachineIdFileError),
1418 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 #[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 #[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 #[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 #[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}