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}
147
148impl BootstrapConfig {
149 pub fn new(connection_file_path: impl Into<PathBuf>, port: u16) -> Self {
150 Self {
151 connection_file_path: connection_file_path.into(),
152 port,
153 daemon_ver: DAEMON_VERSION.to_owned(),
154 configured_modules: Vec::new(),
155 storage_config: None,
156 admission_facts: AdmissionFactsConfig::default(),
157 scope_authority_owners: daemon_config::default_scope_authority_owners(),
158 daemon_config_path: None,
159 configured_port: None,
160 route_bind_relay_default_ms: None,
161 reserved_capabilities: BTreeMap::new(),
162 watchdog_config: DaemonSelfWatchdogConfig::default(),
163 connection_file_source: ConnectionFileSource::Explicit,
164 cgroup_placement: CgroupPlacementConfig::default(),
165 capture_logs_dir: None,
166 terminal_journal_path: None,
167 machine_id_path: None,
168 live_children_path: None,
169 }
170 }
171
172 pub fn with_live_children_record(mut self, path: impl Into<PathBuf>) -> Self {
176 self.live_children_path = Some(path.into());
177 self
178 }
179
180 pub fn with_machine_id_path(mut self, path: impl Into<PathBuf>) -> Self {
184 self.machine_id_path = Some(path.into());
185 self
186 }
187
188 pub fn with_cgroup_placement(mut self, placement: CgroupPlacementConfig) -> Self {
190 self.cgroup_placement = placement;
191 self
192 }
193
194 pub fn with_capture_logs_dir(mut self, dir: impl Into<PathBuf>) -> Self {
198 self.capture_logs_dir = Some(dir.into());
199 self
200 }
201
202 pub fn with_terminal_journal_path(mut self, path: impl Into<PathBuf>) -> Self {
205 self.terminal_journal_path = Some(path.into());
206 self
207 }
208
209 pub fn from_env() -> Result<Self, BootstrapError> {
210 Self::from_env_with_daemon_config_path(daemon_config::default_config_path())
211 }
212
213 pub fn from_env_for_daemon_binary() -> Result<Self, BootstrapError> {
222 let run_dir = daemon_config::daemon_run_dir().map_err(BootstrapError::RunDir)?;
223 let machine_id_path =
224 crate::machine_id::default_machine_id_path().map_err(BootstrapError::MachineId)?;
225 Ok(Self::from_env()?
226 .with_capture_logs_dir(run_dir.join("logs"))
227 .with_terminal_journal_path(run_dir.join("terminals.jsonl"))
228 .with_live_children_record(crate::live_children::record_path(&run_dir))
229 .with_machine_id_path(machine_id_path))
230 }
231
232 pub fn from_env_with_daemon_config_path(
233 daemon_config_path: impl AsRef<Path>,
234 ) -> Result<Self, BootstrapError> {
235 let daemon_config_path = daemon_config_path.as_ref().to_path_buf();
236 let daemon_config =
237 daemon_config::load(&daemon_config_path).map_err(BootstrapError::DaemonConfig)?;
238 let config_port = daemon_config.as_ref().and_then(|config| config.port);
239 let storage_config = daemon_config
240 .as_ref()
241 .and_then(|config| config.storage.clone());
242 let admission_facts_carrier_module_id = daemon_config
243 .as_ref()
244 .and_then(|config| config.admission_facts_carrier_module_id.clone());
245 let admission_facts_targets = daemon_config
246 .as_ref()
247 .and_then(|config| config.admission_facts_targets.clone());
248 let scope_authority_owners = daemon_config
249 .as_ref()
250 .map(|config| config.scope_authority_owners.clone())
251 .unwrap_or_else(daemon_config::default_scope_authority_owners);
252 let route_bind_relay_default_ms = daemon_config
253 .as_ref()
254 .and_then(|config| config.route_bind_relay_timeout_ms);
255 let reserved_capabilities = daemon_config
256 .as_ref()
257 .map(|config| config.reserved_capabilities.clone())
258 .unwrap_or_default();
259 let configured_modules = daemon_config
260 .map(|config| config.modules)
261 .unwrap_or_default();
262
263 let port = match env::var(SUBC_PORT_ENV) {
264 Ok(raw) if !raw.trim().is_empty() => {
265 let port = raw
266 .parse::<u16>()
267 .map_err(|source| BootstrapError::InvalidPort { raw, source })?;
268 if let Some(config_port) = config_port {
269 info!(
270 env = SUBC_PORT_ENV,
271 env_port = port,
272 config_port,
273 "SUBC_PORT overrides daemon config port"
274 );
275 }
276 port
277 }
278 Ok(_) | Err(_) => config_port.unwrap_or(DEFAULT_SUBC_PORT),
279 };
280
281 let (connection_file_path, connection_file_source) =
282 connection_file_path_with_source(non_empty_os_var("XDG_RUNTIME_DIR"));
283 Ok(Self::new(connection_file_path, port)
284 .with_configured_modules(configured_modules)
285 .with_storage_config(storage_config)
286 .with_admission_facts_config(admission_facts_carrier_module_id, admission_facts_targets)
287 .with_scope_authority_owners(scope_authority_owners)
288 .with_route_bind_relay_default_ms(route_bind_relay_default_ms)
289 .with_reserved_capabilities(reserved_capabilities)
290 .with_daemon_config_source(daemon_config_path, config_port)
291 .with_connection_file_source(connection_file_source))
292 }
293
294 pub fn with_daemon_config_path(
295 self,
296 daemon_config_path: impl AsRef<Path>,
297 ) -> Result<Self, BootstrapError> {
298 let daemon_config_path = daemon_config_path.as_ref().to_path_buf();
299 let daemon_config =
300 daemon_config::load(&daemon_config_path).map_err(BootstrapError::DaemonConfig)?;
301 let configured_port = daemon_config.as_ref().and_then(|config| config.port);
302 let storage_config = daemon_config
303 .as_ref()
304 .and_then(|config| config.storage.clone());
305 let admission_facts_carrier_module_id = daemon_config
306 .as_ref()
307 .and_then(|config| config.admission_facts_carrier_module_id.clone());
308 let admission_facts_targets = daemon_config
309 .as_ref()
310 .and_then(|config| config.admission_facts_targets.clone());
311 let scope_authority_owners = daemon_config
312 .as_ref()
313 .map(|config| config.scope_authority_owners.clone())
314 .unwrap_or_else(daemon_config::default_scope_authority_owners);
315 let route_bind_relay_default_ms = daemon_config
316 .as_ref()
317 .and_then(|config| config.route_bind_relay_timeout_ms);
318 let reserved_capabilities = daemon_config
319 .as_ref()
320 .map(|config| config.reserved_capabilities.clone())
321 .unwrap_or_default();
322 let configured_modules = daemon_config
323 .map(|config| config.modules)
324 .unwrap_or_default();
325 Ok(self
326 .with_configured_modules(configured_modules)
327 .with_storage_config(storage_config)
328 .with_admission_facts_config(admission_facts_carrier_module_id, admission_facts_targets)
329 .with_scope_authority_owners(scope_authority_owners)
330 .with_route_bind_relay_default_ms(route_bind_relay_default_ms)
331 .with_reserved_capabilities(reserved_capabilities)
332 .with_daemon_config_source(daemon_config_path, configured_port))
333 }
334
335 pub fn with_configured_modules(
336 mut self,
337 modules: impl IntoIterator<Item = ConfiguredModule>,
338 ) -> Self {
339 self.configured_modules = modules.into_iter().collect();
340 self.configured_modules
341 .sort_by(|left, right| left.module_id.cmp(&right.module_id));
342 self
343 }
344
345 pub fn with_storage_config(
346 mut self,
347 storage_config: Option<daemon_config::StorageConfig>,
348 ) -> Self {
349 self.storage_config = storage_config;
350 self
351 }
352
353 pub fn with_admission_facts_config(
354 mut self,
355 carrier_module_id: Option<String>,
356 targets: Option<Vec<String>>,
357 ) -> Self {
358 self.admission_facts = AdmissionFactsConfig {
359 carrier_module_id,
360 targets,
361 };
362 self
363 }
364
365 pub fn with_scope_authority_owners(mut self, owners: Vec<String>) -> Self {
368 self.scope_authority_owners = owners;
369 self
370 }
371
372 pub fn with_route_bind_relay_default_ms(mut self, ms: Option<u64>) -> Self {
377 self.route_bind_relay_default_ms = ms;
378 self
379 }
380
381 pub fn with_reserved_capabilities(
382 mut self,
383 reserved_capabilities: BTreeMap<String, String>,
384 ) -> Self {
385 self.reserved_capabilities = reserved_capabilities;
386 self
387 }
388
389 fn with_daemon_config_source(
390 mut self,
391 daemon_config_path: PathBuf,
392 configured_port: Option<u16>,
393 ) -> Self {
394 self.daemon_config_path = Some(daemon_config_path);
395 self.configured_port = configured_port;
396 self
397 }
398
399 fn with_connection_file_source(mut self, source: ConnectionFileSource) -> Self {
400 self.connection_file_source = source;
401 self
402 }
403
404 pub fn with_watchdog_config(mut self, watchdog_config: DaemonSelfWatchdogConfig) -> Self {
405 self.watchdog_config = watchdog_config;
406 self
407 }
408}
409
410#[allow(clippy::large_enum_variant)]
414#[derive(Debug)]
415pub enum Outcome {
416 AlreadyRunning,
418 Bound(BoundDaemon),
421}
422
423#[derive(Debug)]
424pub struct BoundDaemon {
425 pub listeners: Vec<TcpListener>,
426 pub connection_info: ConnectionInfo,
427 pub connection_file_path: PathBuf,
428 pub connection_file_source: ConnectionFileSource,
429 pub machine_id: Option<crate::machine_id::MachineId>,
432 run_dir_lock: Option<crate::run_dir_lock::RunDirLock>,
436 publication_lock: Option<StartLock>,
442}
443
444pub fn connection_file_path() -> PathBuf {
451 connection_file_path_with_source(non_empty_os_var("XDG_RUNTIME_DIR")).0
452}
453
454fn connection_file_path_with_source(
455 runtime_dir: Option<OsString>,
456) -> (PathBuf, ConnectionFileSource) {
457 if let Some(runtime_dir) = runtime_dir.filter(|value| !value.is_empty()) {
458 return (
459 PathBuf::from(runtime_dir).join(CONNECTION_FILE_NAME),
460 ConnectionFileSource::XdgRuntimeDir,
461 );
462 }
463
464 (
465 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token())),
466 ConnectionFileSource::TempDirFallback,
467 )
468}
469
470pub async fn run() -> Result<(), BootstrapError> {
476 run_with_config(
479 BootstrapConfig::from_env_for_daemon_binary()?
480 .with_cgroup_placement(CgroupPlacementConfig::Current),
481 )
482 .await
483}
484
485pub async fn run_with_config(config: BootstrapConfig) -> Result<(), BootstrapError> {
509 let configured_modules = config.configured_modules.clone();
510 let storage_config = config.storage_config.clone();
511 let admission_facts = config.admission_facts.clone();
512 let scope_authority_owners = config.scope_authority_owners.clone();
513 let daemon_config_path = config.daemon_config_path.clone();
514 let configured_port = config.configured_port;
515 let route_bind_relay_default_ms = config.route_bind_relay_default_ms;
516 let reserved_capabilities = config.reserved_capabilities.clone();
517 let watchdog_config = config.watchdog_config.clone();
518 let cgroup_placement_config = config.cgroup_placement.clone();
519 let capture_logs_dir = config.capture_logs_dir.clone();
520 let terminal_journal_path = config.terminal_journal_path.clone();
521 let live_children_path = config.live_children_path.clone();
522 match ensure_singleton_inner(config, false).await? {
523 Outcome::AlreadyRunning => {
524 info!("subc daemon already running");
525 Ok(())
526 }
527 Outcome::Bound(bound) => {
528 #[cfg(target_os = "linux")]
529 let cgroup_placement = prepare_cgroup_placement(&cgroup_placement_config);
530 #[cfg(not(target_os = "linux"))]
531 let _ = cgroup_placement_config;
532 serve_bound_daemon(
533 bound,
534 configured_modules,
535 storage_config,
536 admission_facts,
537 scope_authority_owners,
538 daemon_config_path,
539 configured_port,
540 route_bind_relay_default_ms,
541 reserved_capabilities,
542 watchdog_config,
543 capture_logs_dir,
544 terminal_journal_path,
545 live_children_path,
546 #[cfg(target_os = "linux")]
547 cgroup_placement,
548 )
549 .await
550 }
551 }
552}
553
554#[cfg(target_os = "linux")]
555fn prepare_cgroup_placement(config: &CgroupPlacementConfig) -> Option<subc_cgroup::Placement> {
556 let result = match config {
557 CgroupPlacementConfig::Disabled => return None,
558 CgroupPlacementConfig::Current => subc_cgroup::prepare_current(),
559 CgroupPlacementConfig::Root(root) => subc_cgroup::prepare_at(root),
560 };
561
562 match result {
563 Ok(Some(placement)) => Some(placement),
564 Ok(None) => {
565 warn!(
566 placement = ?config,
567 "module cgroup placement is disabled: configured cgroup root is not delegated"
568 );
569 None
570 }
571 Err(error) => {
572 warn!(
573 placement = ?config,
574 error = %error,
575 "module cgroup placement is disabled by an unexpected cgroup probe error"
576 );
577 None
578 }
579 }
580}
581
582#[cfg(unix)]
589const NOFILE_TARGET: u64 = 65536;
590
591#[cfg(unix)]
595fn raise_nofile_limit() {
596 #[cfg(target_os = "macos")]
597 let ceiling = std::process::Command::new("/usr/sbin/sysctl")
598 .args(["-n", "kern.maxfilesperproc"])
599 .output()
600 .ok()
601 .filter(|output| output.status.success())
602 .and_then(|output| String::from_utf8(output.stdout).ok())
603 .and_then(|value| value.trim().parse::<u64>().ok())
604 .filter(|value| *value > 0);
605 #[cfg(not(target_os = "macos"))]
606 let ceiling = None;
607 raise_nofile_limit_with_ceiling(ceiling);
608}
609
610#[cfg(unix)]
611fn raise_nofile_limit_with_ceiling(kernel_ceiling: Option<u64>) {
612 match rlimit::Resource::NOFILE.get() {
613 Ok((soft, hard)) => {
614 let target = NOFILE_TARGET
618 .min(hard)
619 .min(kernel_ceiling.unwrap_or(u64::MAX));
620 if soft >= target {
621 return;
622 }
623 match rlimit::Resource::NOFILE.set(target, hard) {
624 Ok(()) => info!(
625 previous_soft = soft,
626 new_soft = target,
627 hard,
628 "raised open-file soft limit for daemon and module children"
629 ),
630 Err(err) => warn!(
631 soft,
632 hard,
633 error = %err,
634 "could not raise open-file soft limit; multi-root modules may exhaust descriptors"
635 ),
636 }
637 }
638 Err(err) => warn!(error = %err, "could not read open-file limit"),
639 }
640}
641
642#[cfg(windows)]
649fn raise_nofile_limit() {
650 const MAXSTDIO_TARGET: u32 = 8192;
651 let current = rlimit::getmaxstdio();
652 if current >= MAXSTDIO_TARGET {
653 return;
654 }
655 match rlimit::setmaxstdio(MAXSTDIO_TARGET) {
656 Ok(new_max) => info!(
657 previous = current,
658 new_max, "raised CRT stdio-stream limit for daemon"
659 ),
660 Err(err) => warn!(
661 current,
662 error = %err,
663 "could not raise CRT stdio-stream limit"
664 ),
665 }
666}
667
668#[cfg(not(any(unix, windows)))]
669fn raise_nofile_limit() {}
670
671pub async fn run_with_daemon_config_path(
672 config: BootstrapConfig,
673 daemon_config_path: impl AsRef<Path>,
674) -> Result<(), BootstrapError> {
675 run_with_config(config.with_daemon_config_path(daemon_config_path)?).await
676}
677
678#[allow(clippy::too_many_arguments)]
679async fn serve_bound_daemon(
680 bound: BoundDaemon,
681 configured_modules: Vec<ConfiguredModule>,
682 storage_config: Option<daemon_config::StorageConfig>,
683 admission_facts: AdmissionFactsConfig,
684 scope_authority_owners: Vec<String>,
685 daemon_config_path: Option<PathBuf>,
686 configured_port: Option<u16>,
687 route_bind_relay_default_ms: Option<u64>,
688 reserved_capabilities: BTreeMap<String, String>,
689 watchdog_config: DaemonSelfWatchdogConfig,
690 capture_logs_dir: Option<PathBuf>,
691 terminal_journal_path: Option<PathBuf>,
692 live_children_path: Option<PathBuf>,
693 #[cfg(target_os = "linux")] cgroup_placement: Option<subc_cgroup::Placement>,
694) -> Result<(), BootstrapError> {
695 #[cfg(unix)]
696 let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
697 .map_err(BootstrapError::Signal)?;
698 raise_nofile_limit();
701
702 info!(
703 connection_file = %bound.connection_file_path.display(),
704 connection_file_source = %bound.connection_file_source,
705 connection_file_source_reason = bound.connection_file_source.reason(),
706 endpoints = ?bound.connection_info.endpoints,
707 configured_modules = configured_modules.len(),
708 machine_id = bound.machine_id.as_ref().map(|id| id.as_str()).unwrap_or("none"),
709 "subc daemon starting"
710 );
711
712 let run_dir_lock = bound.run_dir_lock;
723 let publication_lock = bound.publication_lock;
724 if let Some(owner) = &run_dir_lock {
725 crate::live_children::sweep_orphans(
726 owner,
727 &crate::live_children::AdoptedPids::none(),
728 crate::live_children::SweepBounds::default(),
729 )
730 .await;
731 }
732
733 let registry = Arc::new(Registry::default());
734 let process_liveness = Arc::new(SupervisorProcessLiveness::new());
735 let supervisor_handle = SupervisorHandle::new();
736 let connected_clients = ConnectedClients::new();
737 let forwarding = Arc::new(ForwardingTable::default());
738 let daemon_incarnation = format!(
739 "{:032x}",
740 u128::from_be_bytes(bound.connection_info.daemon_id)
741 );
742 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
743 .with_process_liveness(process_liveness.clone())
744 .with_forwarding(Arc::clone(&forwarding))
745 .with_handle(supervisor_handle.clone())
746 .with_connection_file_path(bound.connection_file_path.clone())
747 .with_daemon_incarnation(daemon_incarnation.clone());
748 let supervisor = match terminal_journal_path {
785 Some(path) => supervisor.with_terminal_journal(path, daemon_incarnation),
786 None => supervisor,
787 };
788 let supervisor = match capture_logs_dir {
789 Some(dir) => supervisor.with_capture_logs_dir(dir),
790 None => supervisor,
791 };
792 let supervisor = match live_children_path {
793 Some(path) => supervisor.with_live_children_record(path),
794 None => supervisor,
795 };
796 #[cfg(target_os = "linux")]
797 let supervisor = supervisor.with_cgroup_placement(cgroup_placement);
798 let route_bind_relay_timeouts = configured_modules
804 .iter()
805 .filter_map(|module| {
806 module
807 .route_bind_relay_timeout_ms
808 .map(|ms| (module.module_id.clone(), Duration::from_millis(ms)))
809 })
810 .collect::<std::collections::BTreeMap<_, _>>();
811 let control_start_clock = crate::clock::StartClock::capture();
812 let mut control = ControlHandler::with_forwarding(Arc::clone(®istry), forwarding)
813 .with_process_liveness(process_liveness)
814 .with_supervisor(supervisor_handle)
815 .with_connected_clients(connected_clients.clone())
816 .with_storage_config(storage_config)
817 .with_machine_id(bound.machine_id.clone())
818 .with_admission_facts_config(admission_facts.carrier_module_id, admission_facts.targets)
819 .with_scope_authority_owners(scope_authority_owners)
820 .with_route_bind_relay_timeouts(route_bind_relay_timeouts)
821 .with_daemon_provenance(
822 bound.connection_info.pid,
823 control_start_clock.started_at_ms(),
824 std::env::current_exe().ok(),
825 normalized_build_provenance(env!("SUBC_BUILD_GIT_SHA")),
826 normalized_build_provenance(env!("SUBC_BUILD_LOCK_DIGEST")),
827 )
828 .with_daemon_start_clock(control_start_clock)
829 .with_capability_config(
830 configured_modules
831 .iter()
832 .map(|module| (module.module_id.clone(), module.enabled)),
833 reserved_capabilities,
834 );
835 if let Some(ms) = route_bind_relay_default_ms {
836 control = control.with_route_bind_relay_timeout(Duration::from_millis(ms));
839 }
840 if let Some(config_path) = daemon_config_path {
841 control = control.with_supervisor_rescan(supervisor.clone(), config_path, configured_port);
842 }
843 let control = Arc::new(control);
844 let router = Arc::new(Router::with_control_handler(Arc::clone(&control)));
845 let auth = ServerAuth::new(
846 bound.connection_info.key.clone(),
847 bound.connection_info.daemon_id,
848 bound.connection_info.daemon_ver.clone(),
849 )
850 .with_connected_clients(connected_clients);
851
852 let mut serve_task =
853 AbortOnDrop::new(tokio::spawn(serve_listeners(bound.listeners, router, auth)));
854 verify_serving(&bound.connection_info).await?;
855 write_atomic(&bound.connection_file_path, &bound.connection_info).map_err(|source| {
859 BootstrapError::ConnectionFileWrite {
860 path: bound.connection_file_path.clone(),
861 source,
862 }
863 })?;
864 drop(publication_lock);
865 let _clock_step_task = AbortOnDrop::new(crate::watchdog::spawn_clock_step_monitor());
866 #[cfg(target_os = "linux")]
868 let _kill_mode_check = AbortOnDrop::new(tokio::spawn(
869 crate::systemd_kill_mode::warn_if_kill_mode_defeats_ordered_shutdown(),
870 ));
871 let _watchdog_task = AbortOnDrop::new(
872 DaemonSelfWatchdog::new(
873 bound.connection_info.clone(),
874 bound.connection_file_path.clone(),
875 )
876 .with_config(watchdog_config)
877 .spawn(),
878 );
879
880 for configured in configured_modules {
881 let enabled = configured.enabled;
882 let health = configured.health.clone();
883 let module_id = configured.module_id.clone();
884 match supervisor.supervise_configured_with_health(
885 configured.module_spec(),
886 enabled,
887 health.clone(),
888 configured.drain_timeout_ms,
889 configured.restart,
890 ) {
891 Ok(_) => {
892 let default_threshold = HealthConfig::default().failure_threshold;
901 if enabled && health.failure_threshold > default_threshold {
902 warn!(
903 module_id = %module_id,
904 failure_threshold = health.failure_threshold,
905 default_threshold,
906 tolerance_secs = health.cadence.as_secs() * u64::from(health.failure_threshold),
907 "health failure threshold is relaxed above the default; a wedged module stays unflagged for longer"
908 );
909 }
910 info!(module_id = %module_id, enabled, "configured module supervised");
911 }
912 Err(err) => {
913 error!(module_id = %module_id, error = %err, "failed to supervise configured module; continuing daemon startup");
914 }
915 }
916 }
917
918 control.refresh_capability_requirements();
919 Arc::clone(&control).spawn_capability_deadline_loop();
920
921 #[cfg(unix)]
922 {
923 tokio::select! {
924 result = serve_task.join() => {
925 return result.map_err(BootstrapError::ServeJoin)?.map_err(BootstrapError::Serve);
926 }
927 _ = terminate.recv() => {}
928 }
929 supervisor.begin_daemon_shutdown();
934 drop(_watchdog_task);
938 drop(serve_task);
941 let escalated = tokio::select! {
942 biased;
943 _ = terminate.recv() => {
944 info!("second SIGTERM: abandoning daemon shutdown wait");
945 true
946 }
947 result = supervisor.drain_for_daemon_shutdown() => {
948 if let Err(error) = result {
949 warn!(%error, "daemon shutdown drain failed; exiting anyway");
950 }
951 false
952 }
953 };
954 supervisor
961 .end_children_for_daemon_shutdown(escalated, async {
962 terminate.recv().await;
963 })
964 .await;
965 Ok(())
966 }
967 #[cfg(not(unix))]
968 serve_task
969 .join()
970 .await
971 .map_err(BootstrapError::ServeJoin)?
972 .map_err(BootstrapError::Serve)
973}
974
975fn normalized_build_provenance(value: &str) -> Option<String> {
976 match value.trim() {
977 "" | "unavailable" => None,
978 value => Some(value.to_string()),
979 }
980}
981
982pub async fn ensure_singleton(
990 connection_file_path: impl AsRef<Path>,
991 port: u16,
992) -> Result<Outcome, BootstrapError> {
993 ensure_singleton_with_config(BootstrapConfig::new(connection_file_path.as_ref(), port)).await
994}
995
996pub async fn ensure_singleton_with_config(
997 config: BootstrapConfig,
998) -> Result<Outcome, BootstrapError> {
999 ensure_singleton_inner(config, true).await
1000}
1001
1002async fn ensure_singleton_inner(
1003 config: BootstrapConfig,
1004 publish: bool,
1005) -> Result<Outcome, BootstrapError> {
1006 let path = config.connection_file_path;
1007
1008 if matches!(probe_existing(&path).await?, Probe::Live) {
1009 return Ok(Outcome::AlreadyRunning);
1010 }
1011
1012 let lock = StartLock::acquire(&path).await?;
1013
1014 if matches!(probe_existing(&path).await?, Probe::Live) {
1018 return Ok(Outcome::AlreadyRunning);
1019 }
1020
1021 let run_dir_lock = config
1027 .live_children_path
1028 .as_deref()
1029 .map(crate::run_dir_lock::RunDirLock::acquire)
1030 .transpose()?;
1031
1032 remove_stale_connection_file_if_present(&path)?;
1033
1034 let machine_id = config
1038 .machine_id_path
1039 .as_deref()
1040 .map(crate::machine_id::load_or_mint)
1041 .transpose()
1042 .map_err(BootstrapError::MachineId)?;
1043
1044 let (listeners, endpoints) = bind_loopback(config.port).await?;
1045 let connection_info = ConnectionInfo {
1046 schema: SCHEMA_VERSION,
1047 wire_version: Some(PROTOCOL_VERSION),
1048 endpoints,
1049 key: generate_key().map_err(BootstrapError::GenerateConnectionFile)?,
1050 daemon_id: generate_daemon_id().map_err(BootstrapError::GenerateConnectionFile)?,
1051 pid: process::id(),
1052 daemon_ver: config.daemon_ver,
1053 };
1054
1055 if publish {
1059 if let Err(source) = write_atomic(&path, &connection_info) {
1060 drop(listeners);
1061 return Err(BootstrapError::ConnectionFileWrite { path, source });
1062 }
1063 }
1064
1065 Ok(Outcome::Bound(BoundDaemon {
1066 listeners,
1067 connection_info,
1068 connection_file_path: path,
1069 connection_file_source: config.connection_file_source,
1070 machine_id,
1071 run_dir_lock,
1072 publication_lock: (!publish).then_some(lock),
1073 }))
1074}
1075
1076#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1077enum Probe {
1078 Live,
1079 StaleOrAbsent,
1080}
1081
1082async fn probe_existing(path: &Path) -> Result<Probe, BootstrapError> {
1083 let info = match connection_file::read(path) {
1084 Ok(info) => info,
1085 Err(source) if is_absent_or_stale_connection_file(&source) => {
1086 return Ok(Probe::StaleOrAbsent)
1087 }
1088 Err(source) => {
1089 return Err(BootstrapError::ConnectionFileRead {
1090 path: path.to_path_buf(),
1091 source,
1092 })
1093 }
1094 };
1095
1096 for endpoint in &info.endpoints {
1097 if matches!(probe_endpoint(&info, endpoint).await, Probe::Live) {
1098 return Ok(Probe::Live);
1099 }
1100 }
1101
1102 Ok(Probe::StaleOrAbsent)
1103}
1104
1105async fn probe_endpoint(info: &ConnectionInfo, endpoint: &Endpoint) -> Probe {
1106 let Ok(ip) = endpoint.host.parse::<IpAddr>() else {
1107 return Probe::StaleOrAbsent;
1108 };
1109 if !ip.is_loopback() {
1110 return Probe::StaleOrAbsent;
1111 }
1112 let addr = SocketAddr::new(ip, endpoint.port);
1113
1114 let mut stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(addr)).await {
1115 Ok(Ok(stream)) => stream,
1116 Ok(Err(_)) | Err(_) => return Probe::StaleOrAbsent,
1117 };
1118
1119 match authenticate_client(&mut stream, info, PROBE_AUTH_DEADLINE).await {
1120 Ok(()) => Probe::Live,
1121 Err(AuthError::DaemonIdMismatch)
1122 | Err(AuthError::InvalidServerProof)
1123 | Err(AuthError::UnexpectedEof { .. })
1124 | Err(AuthError::Timeout { .. })
1125 | Err(AuthError::JsonEncode { .. })
1126 | Err(AuthError::JsonDecode { .. })
1127 | Err(AuthError::Io { .. })
1128 | Err(AuthError::MessageTooLarge { .. })
1129 | Err(AuthError::KeyTooShort { .. })
1130 | Err(AuthError::Random(_))
1131 | Err(AuthError::InvalidClientAuth) => Probe::StaleOrAbsent,
1132 }
1133}
1134
1135async fn verify_serving(info: &ConnectionInfo) -> Result<(), BootstrapError> {
1136 if let Some(endpoint) = info.endpoints.first() {
1140 if matches!(probe_endpoint(info, endpoint).await, Probe::Live) {
1141 return Ok(());
1142 }
1143 }
1144 Err(BootstrapError::StartupNotServing)
1145}
1146
1147fn is_absent_or_stale_connection_file(err: &ConnectionFileError) -> bool {
1148 match err {
1149 ConnectionFileError::Io { source, .. } if source.kind() == io::ErrorKind::NotFound => true,
1150 ConnectionFileError::JsonRead { .. }
1151 | ConnectionFileError::UnsupportedSchema { .. }
1152 | ConnectionFileError::Invalid { .. }
1153 | ConnectionFileError::KeyTooShort { .. }
1154 | ConnectionFileError::InsecurePermissions { .. } => true,
1158 ConnectionFileError::MissingParent { .. }
1159 | ConnectionFileError::MissingFileName { .. }
1160 | ConnectionFileError::InsecureParentDirectory { .. }
1164 | ConnectionFileError::Io { .. }
1165 | ConnectionFileError::JsonWrite { .. }
1166 | ConnectionFileError::Random(_)
1167 | ConnectionFileError::WireVersionMismatch { .. } => false,
1170 }
1171}
1172
1173fn remove_stale_connection_file_if_present(path: &Path) -> Result<(), BootstrapError> {
1174 match fs::remove_file(path) {
1175 Ok(()) => Ok(()),
1176 Err(err) if err.kind() == io::ErrorKind::NotFound => Ok(()),
1177 Err(source) => Err(BootstrapError::RemoveStale {
1178 path: path.to_path_buf(),
1179 source,
1180 }),
1181 }
1182}
1183
1184async fn bind_loopback(port: u16) -> Result<(Vec<TcpListener>, Vec<Endpoint>), BootstrapError> {
1185 bind_loopback_with(port, |port| TcpListener::bind((Ipv6Addr::LOCALHOST, port))).await
1186}
1187
1188async fn bind_loopback_with<F, Fut>(
1189 port: u16,
1190 mut bind_v6: F,
1191) -> Result<(Vec<TcpListener>, Vec<Endpoint>), BootstrapError>
1192where
1193 F: FnMut(u16) -> Fut,
1194 Fut: std::future::Future<Output = io::Result<TcpListener>>,
1195{
1196 let v4_host = Ipv4Addr::LOCALHOST;
1197 let v4 = TcpListener::bind((v4_host, port))
1198 .await
1199 .map_err(|source| BootstrapError::Bind {
1200 host: v4_host.to_string(),
1201 port,
1202 source,
1203 })?;
1204 let actual_port = v4
1205 .local_addr()
1206 .map_err(|source| BootstrapError::LocalAddr {
1207 host: v4_host.to_string(),
1208 source,
1209 })?
1210 .port();
1211
1212 let mut listeners = vec![v4];
1213 let mut endpoints = vec![Endpoint {
1214 host: v4_host.to_string(),
1215 port: actual_port,
1216 }];
1217
1218 let v6_host = Ipv6Addr::LOCALHOST;
1219 match bind_v6(actual_port).await {
1220 Ok(v6) => {
1221 listeners.push(v6);
1222 endpoints.push(Endpoint {
1223 host: v6_host.to_string(),
1224 port: actual_port,
1225 });
1226 }
1227 Err(err)
1228 if ipv6_loopback_unavailable(&err)
1229 || (port == 0 && err.kind() == io::ErrorKind::AddrInUse) =>
1230 {
1231 warn!(
1232 port = actual_port,
1233 error = %err,
1234 "IPv6 loopback unavailable; serving only IPv4 loopback"
1235 );
1236 }
1237 Err(source) => {
1238 drop(listeners);
1239 return Err(BootstrapError::Bind {
1240 host: v6_host.to_string(),
1241 port: actual_port,
1242 source,
1243 });
1244 }
1245 }
1246
1247 Ok((listeners, endpoints))
1248}
1249
1250fn ipv6_loopback_unavailable(err: &io::Error) -> bool {
1251 matches!(
1252 err.kind(),
1253 io::ErrorKind::AddrNotAvailable | io::ErrorKind::Unsupported
1254 ) || matches!(err.raw_os_error(), Some(47) | Some(49) | Some(97))
1255}
1256
1257struct AbortOnDrop<T> {
1258 handle: JoinHandle<T>,
1259}
1260
1261impl<T> AbortOnDrop<T> {
1262 fn new(handle: JoinHandle<T>) -> Self {
1263 Self { handle }
1264 }
1265
1266 async fn join(&mut self) -> Result<T, JoinError> {
1267 (&mut self.handle).await
1268 }
1269}
1270
1271impl<T> Drop for AbortOnDrop<T> {
1272 fn drop(&mut self) {
1273 if !self.handle.is_finished() {
1274 self.handle.abort();
1275 }
1276 }
1277}
1278
1279#[derive(Debug)]
1280struct StartLock {
1281 _file: fs::File,
1284}
1285
1286impl StartLock {
1287 async fn acquire(connection_file_path: &Path) -> Result<Self, BootstrapError> {
1288 let path = start_lock_path(connection_file_path);
1289 for _ in 0..START_LOCK_RETRIES {
1290 let file = match open_owner_only_lock(&path) {
1291 Ok(file) => file,
1292 Err(source) => return Err(BootstrapError::StartLockCreate { path, source }),
1293 };
1294 match FileExt::try_lock(&file) {
1295 Ok(()) => return Ok(Self { _file: file }),
1296 Err(TryLockError::WouldBlock) => sleep(START_LOCK_RETRY_DELAY).await,
1297 Err(TryLockError::Error(source)) => {
1298 return Err(BootstrapError::StartLockCreate { path, source });
1299 }
1300 }
1301 }
1302
1303 Err(BootstrapError::StartLockBusy {
1304 path,
1305 attempts: START_LOCK_RETRIES,
1306 })
1307 }
1308}
1309
1310pub(crate) fn open_owner_only_lock(path: &Path) -> io::Result<fs::File> {
1311 let mut options = fs::OpenOptions::new();
1312 options.read(true).write(true).create(true);
1313 #[cfg(unix)]
1314 {
1315 use std::os::unix::fs::OpenOptionsExt;
1316 options.mode(0o600);
1317 }
1318 options.open(path)
1319}
1320
1321fn start_lock_path(connection_file_path: &Path) -> PathBuf {
1322 let file_name = connection_file_path
1323 .file_name()
1324 .map(|name| name.to_string_lossy())
1325 .unwrap_or_else(|| CONNECTION_FILE_NAME.into());
1326 let lock_name = format!("{file_name}.start-lock");
1327 connection_file_path
1328 .parent()
1329 .filter(|parent| !parent.as_os_str().is_empty())
1330 .unwrap_or_else(|| Path::new("."))
1331 .join(lock_name)
1332}
1333
1334fn non_empty_os_var(key: &str) -> Option<OsString> {
1335 let value = env::var_os(key)?;
1336 if value.is_empty() {
1337 None
1338 } else {
1339 Some(value)
1340 }
1341}
1342
1343#[derive(Debug)]
1346pub enum BootstrapError {
1347 StartupNotServing,
1349 #[cfg(unix)]
1350 Signal(io::Error),
1351 InvalidPort {
1352 raw: String,
1353 source: std::num::ParseIntError,
1354 },
1355 ConnectionFileRead {
1356 path: PathBuf,
1357 source: ConnectionFileError,
1358 },
1359 ConnectionFileWrite {
1360 path: PathBuf,
1361 source: ConnectionFileError,
1362 },
1363 GenerateConnectionFile(ConnectionFileError),
1364 StartLockCreate {
1365 path: PathBuf,
1366 source: io::Error,
1367 },
1368 StartLockBusy {
1369 path: PathBuf,
1370 attempts: usize,
1371 },
1372 RunDirLockCreate {
1374 path: PathBuf,
1375 source: io::Error,
1376 },
1377 RunDirBusy {
1382 path: PathBuf,
1383 holder_pid: Option<u32>,
1384 },
1385 RemoveStale {
1386 path: PathBuf,
1387 source: io::Error,
1388 },
1389 Bind {
1390 host: String,
1391 port: u16,
1392 source: io::Error,
1393 },
1394 LocalAddr {
1395 host: String,
1396 source: io::Error,
1397 },
1398 DaemonConfig(DaemonConfigError),
1399 MachineId(crate::machine_id::MachineIdFileError),
1402 RunDir(daemon_config::DaemonRunDirError),
1405 Serve(ServerError),
1406 ServeJoin(tokio::task::JoinError),
1407}
1408
1409impl fmt::Display for BootstrapError {
1410 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1411 match self {
1412 Self::StartupNotServing => write!(f, "bound daemon did not authenticate its readiness probe; connection file was not published"),
1413 #[cfg(unix)]
1414 Self::Signal(error) => write!(f, "failed to register SIGTERM handler: {error}"),
1415 Self::InvalidPort { raw, source } => {
1416 write!(f, "invalid {SUBC_PORT_ENV} value '{raw}': {source}")
1417 }
1418 Self::ConnectionFileRead { path, source } => write!(
1419 f,
1420 "failed to read connection file {}: {source}",
1421 path.display()
1422 ),
1423 Self::ConnectionFileWrite { path, source } => write!(
1424 f,
1425 "failed to publish connection file {}: {source}",
1426 path.display()
1427 ),
1428 Self::GenerateConnectionFile(err) => {
1429 write!(f, "failed to generate connection-file auth material: {err}")
1430 }
1431 Self::StartLockCreate { path, source } => {
1432 write!(
1433 f,
1434 "failed to create start lock {}: {source}",
1435 path.display()
1436 )
1437 }
1438 Self::StartLockBusy { path, attempts } => write!(
1439 f,
1440 "start lock {} remained busy after {attempts} attempts",
1441 path.display()
1442 ),
1443 Self::RunDirLockCreate { path, source } => write!(
1444 f,
1445 "refusing to start: failed to lock run directory via {}: {source}",
1446 path.display()
1447 ),
1448 Self::RunDirBusy { path, holder_pid } => {
1449 write!(
1450 f,
1451 "refusing to start: run directory lock {} is held by another daemon",
1452 path.display()
1453 )?;
1454 if let Some(pid) = holder_pid {
1455 write!(f, " (pid {pid})")?;
1456 }
1457 write!(
1458 f,
1459 "; a live daemon owns this data home's run state (typically one started with a different XDG_RUNTIME_DIR)"
1460 )
1461 }
1462 Self::RemoveStale { path, source } => write!(
1463 f,
1464 "failed to remove stale connection file {}: {source}",
1465 path.display()
1466 ),
1467 Self::Bind { host, port, source } if source.kind() == io::ErrorKind::AddrInUse => {
1468 write!(
1469 f,
1470 "port {port} in use on loopback {host}: {source}; set the port in config"
1471 )
1472 }
1473 Self::Bind { host, port, source } => {
1474 write!(f, "failed to bind loopback TCP {host}:{port}: {source}")
1475 }
1476 Self::LocalAddr { host, source } => {
1477 write!(f, "failed to read local address for {host}: {source}")
1478 }
1479 Self::DaemonConfig(err) => write!(f, "failed to load daemon config: {err}"),
1480 Self::MachineId(err) => write!(f, "refusing to start: {err}"),
1481 Self::RunDir(err) => write!(f, "refusing to start: {err}"),
1482 Self::Serve(err) => write!(f, "daemon server failed: {err}"),
1483 Self::ServeJoin(err) => write!(f, "daemon server task failed: {err}"),
1484 }
1485 }
1486}
1487
1488impl Error for BootstrapError {
1489 fn source(&self) -> Option<&(dyn Error + 'static)> {
1490 match self {
1491 Self::StartupNotServing => None,
1492 #[cfg(unix)]
1493 Self::Signal(source) => Some(source),
1494 Self::InvalidPort { source, .. } => Some(source),
1495 Self::ConnectionFileRead { source, .. }
1496 | Self::ConnectionFileWrite { source, .. }
1497 | Self::GenerateConnectionFile(source) => Some(source),
1498 Self::StartLockCreate { source, .. }
1499 | Self::RunDirLockCreate { source, .. }
1500 | Self::RemoveStale { source, .. }
1501 | Self::Bind { source, .. }
1502 | Self::LocalAddr { source, .. } => Some(source),
1503 Self::DaemonConfig(err) => Some(err),
1504 Self::MachineId(err) => Some(err),
1505 Self::RunDir(err) => Some(err),
1506 Self::Serve(err) => Some(err),
1507 Self::ServeJoin(err) => Some(err),
1508 Self::StartLockBusy { .. } | Self::RunDirBusy { .. } => None,
1509 }
1510 }
1511}
1512
1513#[cfg(test)]
1514mod tests {
1515 use super::*;
1516 use crate::server::ServerAuth;
1517 #[cfg(target_os = "linux")]
1518 use std::collections::BTreeSet;
1519 use std::sync::Mutex;
1520 #[cfg(target_os = "linux")]
1521 use subc_control::ModuleProtocol;
1522 use subc_test_support::TestTempDir;
1523 use subc_transport::MIN_KEY_LEN;
1524 use tokio::io::AsyncReadExt;
1525 use tokio::task::JoinHandle;
1526
1527 #[cfg(unix)]
1528 use std::os::unix::fs::PermissionsExt;
1529
1530 static ENV_LOCK: Mutex<()> = Mutex::new(());
1531
1532 #[test]
1533 fn normalized_build_provenance_preserves_real_values() {
1534 assert_eq!(normalized_build_provenance("abc"), Some("abc".to_string()));
1535 }
1536
1537 #[test]
1538 fn normalized_build_provenance_omits_unavailable_and_empty_values() {
1539 assert_eq!(normalized_build_provenance("unavailable"), None);
1540 assert_eq!(normalized_build_provenance(""), None);
1541 }
1542
1543 #[cfg(target_os = "linux")]
1544 fn current_cgroup_path_for_test() -> io::Result<PathBuf> {
1545 let cgroups = fs::read_to_string("/proc/self/cgroup")?;
1546 let relative = cgroups
1547 .lines()
1548 .find_map(|line| line.strip_prefix("0::"))
1549 .ok_or_else(|| {
1550 io::Error::new(io::ErrorKind::Unsupported, "cgroup v2 is unavailable")
1551 })?;
1552 Ok(Path::new("/sys/fs/cgroup").join(relative.trim_start_matches('/')))
1553 }
1554
1555 #[cfg(target_os = "linux")]
1556 fn module_cgroup_directories() -> io::Result<BTreeSet<OsString>> {
1557 let modules = current_cgroup_path_for_test()?.join("subc-modules");
1558 let entries = match fs::read_dir(modules) {
1559 Ok(entries) => entries,
1560 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(BTreeSet::new()),
1561 Err(error) => return Err(error),
1562 };
1563 let mut directories = BTreeSet::new();
1564 for entry in entries {
1565 let entry = entry?;
1566 if entry.file_type()?.is_dir() {
1567 directories.insert(entry.file_name());
1568 }
1569 }
1570 Ok(directories)
1571 }
1572
1573 #[cfg(target_os = "linux")]
1574 async fn wait_for_path(path: &Path, task: &JoinHandle<Result<(), BootstrapError>>) {
1575 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1576 while !path.exists() && tokio::time::Instant::now() < deadline {
1577 assert!(
1578 !task.is_finished(),
1579 "daemon exited before creating {}",
1580 path.display()
1581 );
1582 sleep(Duration::from_millis(10)).await;
1583 }
1584 assert!(path.exists(), "daemon did not create {}", path.display());
1585 }
1586
1587 #[cfg(target_os = "linux")]
1592 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1593 async fn run_with_config_does_not_reconcile_the_ambient_cgroup_by_default() {
1594 let temp = unique_temp_dir("bootstrap-cgroup-default-disabled");
1595 let module_id = format!("cgroup-isolation-probe-{}", process::id());
1596 let before = module_cgroup_directories().expect("read ambient module cgroups before boot");
1597 assert!(
1598 !before.contains(&OsString::from(&module_id)),
1599 "isolation probe cgroup already exists before this daemon starts"
1600 );
1601 let capture = temp.join("logs").join(format!("{module_id}.stderr.log"));
1602 let module = ConfiguredModule {
1603 module_id,
1604 program: PathBuf::from("sh"),
1605 args: vec!["-c".to_string(), "sleep 30".to_string()],
1606 env: Vec::new(),
1607 log: None,
1608 enabled: true,
1609 reserved: false,
1610 reserved_prefixes: Vec::new(),
1611 protocol: ModuleProtocol::None,
1612 overlap: Default::default(),
1613 health: HealthConfig::default(),
1614 drain_timeout_ms: None,
1615 route_bind_relay_timeout_ms: None,
1616 restart: RestartPolicy::default(),
1617 };
1618 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1619 .with_configured_modules([module])
1620 .with_capture_logs_dir(temp.join("logs"))
1621 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1622 let task = tokio::spawn(run_with_config(config));
1623
1624 wait_for_path(&capture, &task).await;
1625 let after = module_cgroup_directories().expect("read ambient module cgroups after boot");
1626
1627 task.abort();
1628 assert!(task
1629 .await
1630 .expect_err("aborted daemon task must cancel")
1631 .is_cancelled());
1632 assert_eq!(
1633 before, after,
1634 "an in-process daemon must not create or reconcile ambient module cgroups"
1635 );
1636 }
1637
1638 #[cfg(target_os = "linux")]
1639 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1640 async fn explicit_cgroup_root_is_prepared_inside_the_fixture_tree() {
1641 let temp = unique_temp_dir("bootstrap-cgroup-explicit-root");
1642 let cgroup_root = temp.join("cgroup");
1643 fs::create_dir(&cgroup_root).expect("create scratch cgroup root");
1644 fs::write(cgroup_root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
1645 let modules = cgroup_root.join("subc-modules");
1646 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1647 .with_cgroup_placement(CgroupPlacementConfig::Root(cgroup_root))
1648 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1649 let task = tokio::spawn(run_with_config(config));
1650
1651 wait_for_path(&modules, &task).await;
1652
1653 task.abort();
1654 assert!(task
1655 .await
1656 .expect_err("aborted daemon task must cancel")
1657 .is_cancelled());
1658 assert!(
1659 fs::read_dir(&modules)
1660 .expect("read prepared modules directory")
1661 .next()
1662 .is_none(),
1663 "the delegation probe must clean up after itself"
1664 );
1665 }
1666
1667 struct EnvGuard {
1668 key: &'static str,
1669 previous: Option<OsString>,
1670 }
1671
1672 impl EnvGuard {
1673 fn set(key: &'static str, value: &Path) -> Self {
1674 let previous = env::var_os(key);
1675 env::set_var(key, value);
1676 Self { key, previous }
1677 }
1678
1679 fn set_str(key: &'static str, value: &str) -> Self {
1680 let previous = env::var_os(key);
1681 env::set_var(key, value);
1682 Self { key, previous }
1683 }
1684
1685 fn unset(key: &'static str) -> Self {
1686 let previous = env::var_os(key);
1687 env::remove_var(key);
1688 Self { key, previous }
1689 }
1690 }
1691
1692 impl Drop for EnvGuard {
1693 fn drop(&mut self) {
1694 match &self.previous {
1695 Some(value) => env::set_var(self.key, value),
1696 None => env::remove_var(self.key),
1697 }
1698 }
1699 }
1700
1701 fn unique_temp_dir(name: &str) -> TestTempDir {
1702 TestTempDir::new(name)
1703 }
1704
1705 fn temp_connection_file_path(name: &str) -> (TestTempDir, PathBuf) {
1706 let dir = unique_temp_dir(name);
1707 let path = dir.join("conn.json");
1708 (dir, path)
1709 }
1710
1711 fn auth_for(info: &ConnectionInfo) -> ServerAuth {
1712 ServerAuth::new(info.key.clone(), info.daemon_id, info.daemon_ver.clone())
1713 }
1714
1715 fn start_server(bound: BoundDaemon) -> JoinHandle<Result<(), ServerError>> {
1716 let auth = auth_for(&bound.connection_info);
1717 tokio::spawn(serve_listeners(
1718 bound.listeners,
1719 Arc::new(Router::with_default_self_handler()),
1720 auth,
1721 ))
1722 }
1723
1724 fn expect_bound(outcome: Outcome) -> BoundDaemon {
1725 match outcome {
1726 Outcome::Bound(bound) => bound,
1727 Outcome::AlreadyRunning => panic!("fresh connection file unexpectedly had a daemon"),
1728 }
1729 }
1730
1731 async fn connect_from_info(conn: &ConnectionInfo) -> io::Result<TcpStream> {
1732 let endpoint = conn
1733 .endpoints
1734 .first()
1735 .expect("test connection file should have an endpoint");
1736 let ip: IpAddr = endpoint.host.parse().unwrap();
1737 TcpStream::connect(SocketAddr::new(ip, endpoint.port)).await
1738 }
1739
1740 fn make_connection_info(port: u16) -> ConnectionInfo {
1741 ConnectionInfo {
1742 schema: SCHEMA_VERSION,
1743 wire_version: Some(PROTOCOL_VERSION),
1744 endpoints: vec![Endpoint {
1745 host: "127.0.0.1".to_owned(),
1746 port,
1747 }],
1748 key: generate_key().unwrap(),
1749 daemon_id: generate_daemon_id().unwrap(),
1750 pid: process::id(),
1751 daemon_ver: "test-subc".to_owned(),
1752 }
1753 }
1754
1755 fn write_raw_owner_only_connection_file(path: &Path, contents: &[u8]) {
1756 fs::write(path, contents).unwrap();
1757 #[cfg(unix)]
1758 fs::set_permissions(path, fs::Permissions::from_mode(0o600)).unwrap();
1759 }
1760
1761 fn assert_owner_only_connection_file(path: &Path) {
1762 #[cfg(unix)]
1765 {
1766 let mode = fs::metadata(path).unwrap().permissions().mode() & 0o777;
1767 assert_eq!(mode, 0o600);
1768 }
1769 #[cfg(not(unix))]
1770 let _ = path;
1771 }
1772
1773 #[test]
1774 fn connection_file_path_uses_xdg_runtime_dir_when_set() {
1775 let _env_lock = ENV_LOCK.lock().unwrap();
1776 let runtime_dir = unique_temp_dir("xdg-runtime");
1777 let _xdg = EnvGuard::set("XDG_RUNTIME_DIR", runtime_dir.path());
1778
1779 assert_eq!(
1780 connection_file_path(),
1781 runtime_dir.join(CONNECTION_FILE_NAME)
1782 );
1783 }
1784
1785 #[test]
1786 fn connection_file_path_source_is_xdg_runtime_dir_when_set() {
1787 let runtime_dir = OsString::from("/run/user/1000");
1788
1789 let (path, source) = connection_file_path_with_source(Some(runtime_dir));
1790
1791 assert_eq!(
1792 path,
1793 PathBuf::from("/run/user/1000").join(CONNECTION_FILE_NAME)
1794 );
1795 assert_eq!(source, ConnectionFileSource::XdgRuntimeDir);
1796 }
1797
1798 #[test]
1799 fn connection_file_path_falls_back_to_temp_dir_with_user_token_when_xdg_unset() {
1800 let _env_lock = ENV_LOCK.lock().unwrap();
1801 let _xdg = EnvGuard::unset("XDG_RUNTIME_DIR");
1802
1803 assert_eq!(
1804 connection_file_path(),
1805 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1806 );
1807 }
1808
1809 #[test]
1810 fn connection_file_path_source_is_temp_dir_when_xdg_unset() {
1811 let (path, source) = connection_file_path_with_source(None);
1812
1813 assert_eq!(
1814 path,
1815 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1816 );
1817 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1818 }
1819
1820 #[test]
1825 fn user_connection_token_is_stable_under_concurrent_callers() {
1826 let expected = user_connection_token();
1827 let workers: Vec<_> = (0..32)
1828 .map(|_| {
1829 std::thread::spawn(|| (0..40).map(|_| user_connection_token()).collect::<Vec<_>>())
1830 })
1831 .collect();
1832 for worker in workers {
1833 for token in worker.join().expect("probe thread") {
1834 assert_eq!(token, expected, "token diverged under concurrent probes");
1835 }
1836 }
1837 }
1838
1839 #[test]
1840 fn connection_file_path_source_is_temp_dir_when_xdg_empty() {
1841 let (path, source) = connection_file_path_with_source(Some(OsString::new()));
1842
1843 assert_eq!(
1844 path,
1845 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1846 );
1847 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1848 }
1849
1850 #[test]
1854 fn daemon_binary_config_journals_and_captures_into_the_run_dir() {
1855 let _env_lock = ENV_LOCK.lock().unwrap();
1856 let root = unique_temp_dir("daemon-binary-config");
1857 let data_home = root.join("data");
1858 let _data = EnvGuard::set("XDG_DATA_HOME", &data_home);
1859 let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1860 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1861
1862 let config = BootstrapConfig::from_env_for_daemon_binary().unwrap();
1863
1864 let run_dir = data_home.join("cortexkit").join("run");
1865 assert_eq!(
1866 config.terminal_journal_path,
1867 Some(run_dir.join("terminals.jsonl"))
1868 );
1869 assert_eq!(config.capture_logs_dir, Some(run_dir.join("logs")));
1870 assert_eq!(
1871 BootstrapConfig::new(root.join("connection.json"), 0).terminal_journal_path,
1872 None,
1873 "an in-process config must not journal anywhere unless asked to"
1874 );
1875 }
1876
1877 #[test]
1878 fn daemon_binary_config_refuses_a_relative_data_home() {
1879 let _env_lock = ENV_LOCK.lock().unwrap();
1880 let root = unique_temp_dir("daemon-binary-relative-data");
1881 let _data = EnvGuard::set_str("XDG_DATA_HOME", "relative-data-home");
1882 let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1883 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1884
1885 let error = BootstrapConfig::from_env_for_daemon_binary()
1886 .expect_err("a relative data home must refuse the daemon binary's config");
1887 assert!(
1888 matches!(error, BootstrapError::RunDir(_)),
1889 "expected a run-directory refusal, got {error}"
1890 );
1891 assert!(error.to_string().contains("XDG_DATA_HOME"), "{error}");
1892 }
1893
1894 #[test]
1895 fn configured_port_uses_default_config_and_env_override() {
1896 let _env_lock = ENV_LOCK.lock().unwrap();
1897 let (_dir, conn_path) = temp_connection_file_path("daemon-config-port");
1898 let config_path = conn_path.with_file_name("subc.jsonc");
1899
1900 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1901 assert_eq!(
1902 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1903 .unwrap()
1904 .port,
1905 DEFAULT_SUBC_PORT
1906 );
1907
1908 fs::write(&config_path, r#"{ "version": 1, "port": 8123 }"#).unwrap();
1909 assert_eq!(
1910 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1911 .unwrap()
1912 .port,
1913 8123
1914 );
1915
1916 let _port = EnvGuard::set_str(SUBC_PORT_ENV, "9012");
1917 assert_eq!(
1918 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1919 .unwrap()
1920 .port,
1921 9012
1922 );
1923 }
1924
1925 #[tokio::test]
1926 async fn second_singleton_probe_against_served_tcp_daemon_reports_already_running() {
1927 let (_dir, path) = temp_connection_file_path("already-running");
1928
1929 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1930 let server = start_server(bound);
1931
1932 let second = ensure_singleton(&path, 0).await.unwrap();
1933 assert!(matches!(second, Outcome::AlreadyRunning));
1934
1935 server.abort();
1936 let _ = server.await;
1937 }
1938
1939 #[tokio::test]
1940 async fn daemon_connection_file_publishes_protocol_wire_version() {
1941 let (_dir, path) = temp_connection_file_path("wire-version");
1942 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1943 assert_eq!(bound.connection_info.wire_version, Some(PROTOCOL_VERSION));
1944 assert_eq!(
1945 connection_file::read(&path).unwrap().wire_version,
1946 Some(PROTOCOL_VERSION)
1947 );
1948
1949 drop(bound.listeners);
1950 }
1951
1952 #[tokio::test]
1953 async fn stale_unbound_connection_file_is_reclaimed() {
1954 let (_dir, path) = temp_connection_file_path("stale-reclaim");
1955 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1956 let stale_port = stale.local_addr().unwrap().port();
1957 drop(stale);
1958 let stale_info = make_connection_info(stale_port);
1959 write_atomic(&path, &stale_info).unwrap();
1960
1961 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1962 assert_ne!(bound.connection_info.key, stale_info.key);
1963 drop(bound.listeners);
1964 }
1965
1966 #[cfg(unix)]
1967 #[tokio::test]
1968 async fn ensure_singleton_reclaims_insecure_connection_file() {
1969 let (_dir, path) = temp_connection_file_path("insecure-reclaim");
1970 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1971 let stale_port = stale.local_addr().unwrap().port();
1972 drop(stale);
1973 let stale_info = make_connection_info(stale_port);
1974 write_atomic(&path, &stale_info).unwrap();
1975 fs::set_permissions(&path, fs::Permissions::from_mode(0o644)).unwrap();
1976
1977 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1978 assert_ne!(bound.connection_info.key, stale_info.key);
1979 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1980 assert_owner_only_connection_file(&path);
1981
1982 drop(bound.listeners);
1983 }
1984
1985 #[tokio::test]
1986 async fn ensure_singleton_reclaims_non_loopback_connection_file() {
1987 let (_dir, path) = temp_connection_file_path("non-loopback-reclaim");
1988 let mut stale_info = make_connection_info(8757);
1989 stale_info.endpoints = vec![Endpoint {
1990 host: "192.0.2.10".to_owned(),
1991 port: 8757,
1992 }];
1993 write_atomic(&path, &stale_info).unwrap();
1994
1995 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1996 assert_ne!(bound.connection_info.key, stale_info.key);
1997 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1998 assert!(bound
1999 .connection_info
2000 .endpoints
2001 .iter()
2002 .all(|endpoint| endpoint.host.parse::<IpAddr>().unwrap().is_loopback()));
2003 assert_owner_only_connection_file(&path);
2004
2005 drop(bound.listeners);
2006 }
2007
2008 #[tokio::test]
2009 async fn ensure_singleton_reclaims_invalid_connection_file_shapes() {
2010 let mut unsupported_schema = make_connection_info(8757);
2011 unsupported_schema.schema = SCHEMA_VERSION + 1;
2012
2013 let mut empty_endpoints = make_connection_info(8757);
2014 empty_endpoints.endpoints.clear();
2015
2016 let mut short_key = make_connection_info(8757);
2017 short_key.key = vec![0x5A; MIN_KEY_LEN - 1];
2018
2019 let cases = vec![
2020 (
2021 "unsupported-schema",
2022 serde_json::to_vec(&unsupported_schema).unwrap(),
2023 Some(unsupported_schema),
2024 ),
2025 (
2026 "empty-endpoints",
2027 serde_json::to_vec(&empty_endpoints).unwrap(),
2028 Some(empty_endpoints),
2029 ),
2030 (
2031 "short-key",
2032 serde_json::to_vec(&short_key).unwrap(),
2033 Some(short_key),
2034 ),
2035 ("invalid-json", b"{not valid connection json".to_vec(), None),
2036 ];
2037
2038 for (label, contents, old_info) in cases {
2039 let (_dir, path) = temp_connection_file_path(label);
2040 write_raw_owner_only_connection_file(&path, &contents);
2041
2042 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2043 if let Some(old_info) = old_info {
2044 assert_ne!(bound.connection_info.key, old_info.key, "{label}");
2045 assert_ne!(
2046 bound.connection_info.daemon_id, old_info.daemon_id,
2047 "{label}"
2048 );
2049 }
2050 assert!(bound.connection_info.key.len() >= MIN_KEY_LEN, "{label}");
2051 assert_ne!(bound.connection_info.daemon_id, [0u8; 16], "{label}");
2052 assert_owner_only_connection_file(&path);
2053
2054 drop(bound.listeners);
2055 }
2056 }
2057
2058 #[tokio::test]
2059 async fn foreign_reused_port_connection_file_is_reclaimed_after_auth_probe_fails() {
2060 let (_dir, path) = temp_connection_file_path("foreign-reclaim");
2061 let foreign = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
2062 let foreign_port = foreign.local_addr().unwrap().port();
2063 write_atomic(&path, &make_connection_info(foreign_port)).unwrap();
2064 let foreign_task = tokio::spawn(async move {
2065 if let Ok((mut stream, _)) = foreign.accept().await {
2066 let mut buf = [0u8; 64];
2067 let _ = stream.read(&mut buf).await;
2068 }
2069 });
2070
2071 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2072 assert!(bound
2073 .connection_info
2074 .endpoints
2075 .iter()
2076 .all(|endpoint| endpoint.port != foreign_port));
2077
2078 drop(bound.listeners);
2079 let _ = foreign_task.await;
2080 }
2081
2082 #[tokio::test]
2083 async fn stale_start_lock_file_is_reclaimable() {
2084 let (_dir, path) = temp_connection_file_path("start-lock-stale-file");
2085 let lock_path = start_lock_path(&path);
2086 drop(open_owner_only_lock(&lock_path).unwrap());
2087 assert!(lock_path.is_file());
2088
2089 let lock = StartLock::acquire(&path).await.unwrap();
2090 assert!(lock_path.is_file());
2091
2092 drop(lock);
2093 assert!(lock_path.is_file());
2094 }
2095
2096 #[tokio::test]
2097 async fn held_start_lock_blocks_second_acquire_until_release() {
2098 let (_dir, path) = temp_connection_file_path("start-lock-held");
2099 let lock_path = start_lock_path(&path);
2100 let first = StartLock::acquire(&path).await.unwrap();
2101
2102 let err = match StartLock::acquire(&path).await {
2103 Ok(_) => panic!("second acquire while held must stay busy"),
2104 Err(err) => err,
2105 };
2106 assert!(matches!(
2107 err,
2108 BootstrapError::StartLockBusy {
2109 ref path,
2110 attempts: START_LOCK_RETRIES,
2111 } if path == &lock_path
2112 ));
2113
2114 drop(first);
2115
2116 let second = StartLock::acquire(&path)
2117 .await
2118 .expect("released advisory lock should be reclaimable");
2119 drop(second);
2120 }
2121
2122 #[tokio::test]
2123 async fn bind_conflict_on_fixed_port_fails_loud_without_reselecting() {
2124 let (_dir, path) = temp_connection_file_path("bind-conflict");
2125 let occupied = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
2126 let occupied_port = occupied.local_addr().unwrap().port();
2127
2128 let err = ensure_singleton(&path, occupied_port).await.unwrap_err();
2129 assert!(matches!(
2130 err,
2131 BootstrapError::Bind { ref source, .. } if source.kind() == io::ErrorKind::AddrInUse
2132 ));
2133 assert!(err.to_string().contains("set the port in config"));
2134
2135 drop(occupied);
2136 }
2137
2138 #[tokio::test]
2139 async fn ephemeral_ipv6_collision_still_serves_ipv4() {
2140 let (listeners, endpoints) = bind_loopback_with(0, |port| async move {
2141 let occupied = TcpListener::bind((Ipv6Addr::LOCALHOST, port)).await?;
2142 let result = TcpListener::bind((Ipv6Addr::LOCALHOST, port)).await;
2143 drop(occupied);
2144 result
2145 })
2146 .await
2147 .expect("an incidental IPv6 collision must not prevent an ephemeral daemon from starting");
2148 assert_eq!(endpoints.len(), 1);
2149 assert_eq!(endpoints[0].host, "127.0.0.1");
2150 TcpStream::connect((Ipv4Addr::LOCALHOST, endpoints[0].port))
2151 .await
2152 .unwrap();
2153 drop(listeners);
2154 }
2155
2156 #[tokio::test]
2157 async fn unpublished_boot_holds_singleton_lock_until_publication() {
2158 let dir = TestTempDir::new("unpublished-singleton-lock");
2159 let path = dir.join("connection.json");
2160 let bound = expect_bound(
2161 ensure_singleton_inner(BootstrapConfig::new(&path, 0), false)
2162 .await
2163 .unwrap(),
2164 );
2165 assert!(!path.exists());
2166 let contender = open_owner_only_lock(&start_lock_path(&path)).unwrap();
2167 assert!(
2168 matches!(FileExt::try_lock(&contender), Err(TryLockError::WouldBlock)),
2169 "an unpublished boot must keep its singleton claim"
2170 );
2171 drop(bound);
2172 FileExt::try_lock(&contender).unwrap();
2173 }
2174
2175 #[tokio::test]
2176 async fn readiness_probe_requires_an_authenticating_server() {
2177 let dir = TestTempDir::new("unserved-readiness");
2178 let bound = expect_bound(
2179 ensure_singleton_inner(BootstrapConfig::new(dir.join("connection.json"), 0), false)
2180 .await
2181 .unwrap(),
2182 );
2183 assert!(matches!(
2184 verify_serving(&bound.connection_info).await,
2185 Err(BootstrapError::StartupNotServing)
2186 ));
2187 }
2188
2189 #[cfg(unix)]
2190 #[tokio::test]
2191 async fn orphan_sweep_finishes_before_connection_file_publication() {
2192 use crate::live_children::{ExecutableIdentity, LiveChild};
2193 let dir = TestTempDir::new("publish-after-sweep");
2194 let path = dir.join("connection.json");
2195 let record = dir.join("live-children.json");
2196 let ready = dir.join("ready");
2197 let mut child = std::process::Command::new("/bin/sh")
2198 .args([
2199 "-c",
2200 "trap '' TERM; : > \"$0\"; read -r _",
2201 ready.to_str().unwrap(),
2202 ])
2203 .stdin(std::process::Stdio::piped())
2204 .env("XDG_DATA_HOME", dir.path())
2205 .env("XDG_RUNTIME_DIR", dir.path())
2206 .env("XDG_CONFIG_HOME", dir.path())
2207 .spawn()
2208 .unwrap();
2209 timeout(Duration::from_secs(5), async {
2210 while !ready.exists() {
2211 sleep(Duration::from_millis(5)).await;
2212 }
2213 })
2214 .await
2215 .unwrap();
2216 let observed = subc_os::Process::open(child.id())
2217 .unwrap()
2218 .unwrap()
2219 .observe()
2220 .unwrap();
2221 crate::live_children::write_record(
2222 &record,
2223 &[LiveChild {
2224 module_id: "old-sleep".into(),
2225 pid: child.id(),
2226 protocol: subc_control::ModuleProtocol::None,
2227 start_time: Some(observed.start_time),
2228 executable: observed.executable.map(ExecutableIdentity::from),
2229 cgroup_name: None,
2230 }],
2231 )
2232 .unwrap();
2233 let config = BootstrapConfig::new(&path, 0).with_live_children_record(&record);
2234 let daemon = tokio::spawn(run_with_config(config));
2235 timeout(Duration::from_secs(5), async {
2236 while !path.exists() {
2237 sleep(Duration::from_millis(5)).await;
2238 }
2239 })
2240 .await
2241 .unwrap();
2242 assert!(path.exists());
2243 let ended = child.try_wait().unwrap().is_some();
2244 if !ended {
2245 child.kill().unwrap();
2246 }
2247 child.wait().unwrap();
2248 if ended {
2249 assert!(matches!(probe_existing(&path).await.unwrap(), Probe::Live));
2250 }
2251 daemon.abort();
2252 let _ = daemon.await;
2253 assert!(
2254 ended,
2255 "the published daemon must not still owe its orphan sweep"
2256 );
2257 }
2258
2259 #[cfg(target_os = "macos")]
2260 #[test]
2261 fn nofile_raise_worker() {
2262 let Some(path) = std::env::var_os("SUBC_NOFILE_TEST_FILE") else {
2263 return;
2264 };
2265 let (_, hard) = rlimit::Resource::NOFILE.get().unwrap();
2266 rlimit::Resource::NOFILE.set(256, hard).unwrap();
2267 raise_nofile_limit_with_ceiling(Some(10240));
2268 std::fs::write(path, rlimit::Resource::NOFILE.get().unwrap().0.to_string()).unwrap();
2269 }
2270
2271 #[cfg(target_os = "macos")]
2272 #[test]
2273 fn macos_nofile_raise_clamps_to_kernel_ceiling() {
2274 let dir = TestTempDir::new("nofile-ceiling");
2275 let result = dir.join("soft");
2276 let status = std::process::Command::new(std::env::current_exe().unwrap())
2277 .args(["--exact", "bootstrap::tests::nofile_raise_worker"])
2278 .env("SUBC_NOFILE_TEST_FILE", &result)
2279 .env("XDG_DATA_HOME", dir.path())
2280 .env("XDG_RUNTIME_DIR", dir.path())
2281 .env("XDG_CONFIG_HOME", dir.path())
2282 .status()
2283 .unwrap();
2284 assert!(status.success());
2285 let soft: u64 = std::fs::read_to_string(result).unwrap().parse().unwrap();
2286 assert!(
2287 soft > 256 && soft <= 10240,
2288 "the inherited limit must be raised within the kernel ceiling; got {soft}"
2289 );
2290 }
2291
2292 #[tokio::test]
2293 async fn key_rotation_republishes_new_material_and_old_file_fails_auth() {
2294 let (_dir, path) = temp_connection_file_path("key-rotation");
2295 let first = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2296 let old_info = first.connection_info.clone();
2297 let fixed_port = old_info.endpoints[0].port;
2298 drop(first.listeners);
2299
2300 let second = expect_bound(ensure_singleton(&path, fixed_port).await.unwrap());
2301 let new_info = second.connection_info.clone();
2302 assert_ne!(old_info.key, new_info.key);
2303 assert_ne!(old_info.daemon_id, new_info.daemon_id);
2304 let server = start_server(second);
2305
2306 let mut old_stream = connect_from_info(&old_info).await.unwrap();
2307 let old_auth = authenticate_client(&mut old_stream, &old_info, PROBE_AUTH_DEADLINE).await;
2308 assert!(
2309 old_auth.is_err(),
2310 "old key must not authenticate after restart"
2311 );
2312
2313 let reread = connection_file::read(&path).unwrap();
2314 let mut new_stream = connect_from_info(&reread).await.unwrap();
2315 authenticate_client(&mut new_stream, &reread, PROBE_AUTH_DEADLINE)
2316 .await
2317 .unwrap();
2318
2319 server.abort();
2320 let _ = server.await;
2321 }
2322
2323 #[cfg(unix)]
2324 #[tokio::test]
2325 async fn published_connection_file_permissions_are_owner_only() {
2326 let (_dir, path) = temp_connection_file_path("permissions");
2327 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2328
2329 let mode = fs::metadata(&path).unwrap().permissions().mode() & 0o777;
2330 assert_eq!(mode, 0o600);
2331
2332 drop(bound.listeners);
2333 }
2334}