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