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}
437
438pub fn connection_file_path() -> PathBuf {
445 connection_file_path_with_source(non_empty_os_var("XDG_RUNTIME_DIR")).0
446}
447
448fn connection_file_path_with_source(
449 runtime_dir: Option<OsString>,
450) -> (PathBuf, ConnectionFileSource) {
451 if let Some(runtime_dir) = runtime_dir.filter(|value| !value.is_empty()) {
452 return (
453 PathBuf::from(runtime_dir).join(CONNECTION_FILE_NAME),
454 ConnectionFileSource::XdgRuntimeDir,
455 );
456 }
457
458 (
459 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token())),
460 ConnectionFileSource::TempDirFallback,
461 )
462}
463
464pub async fn run() -> Result<(), BootstrapError> {
470 run_with_config(
473 BootstrapConfig::from_env_for_daemon_binary()?
474 .with_cgroup_placement(CgroupPlacementConfig::Current),
475 )
476 .await
477}
478
479pub async fn run_with_config(config: BootstrapConfig) -> Result<(), BootstrapError> {
503 let configured_modules = config.configured_modules.clone();
504 let storage_config = config.storage_config.clone();
505 let admission_facts = config.admission_facts.clone();
506 let scope_authority_owners = config.scope_authority_owners.clone();
507 let daemon_config_path = config.daemon_config_path.clone();
508 let configured_port = config.configured_port;
509 let route_bind_relay_default_ms = config.route_bind_relay_default_ms;
510 let reserved_capabilities = config.reserved_capabilities.clone();
511 let watchdog_config = config.watchdog_config.clone();
512 let cgroup_placement_config = config.cgroup_placement.clone();
513 let capture_logs_dir = config.capture_logs_dir.clone();
514 let terminal_journal_path = config.terminal_journal_path.clone();
515 let live_children_path = config.live_children_path.clone();
516 match ensure_singleton_with_config(config).await? {
517 Outcome::AlreadyRunning => {
518 info!("subc daemon already running");
519 Ok(())
520 }
521 Outcome::Bound(bound) => {
522 #[cfg(target_os = "linux")]
523 let cgroup_placement = prepare_cgroup_placement(&cgroup_placement_config);
524 #[cfg(not(target_os = "linux"))]
525 let _ = cgroup_placement_config;
526 serve_bound_daemon(
527 bound,
528 configured_modules,
529 storage_config,
530 admission_facts,
531 scope_authority_owners,
532 daemon_config_path,
533 configured_port,
534 route_bind_relay_default_ms,
535 reserved_capabilities,
536 watchdog_config,
537 capture_logs_dir,
538 terminal_journal_path,
539 live_children_path,
540 #[cfg(target_os = "linux")]
541 cgroup_placement,
542 )
543 .await
544 }
545 }
546}
547
548#[cfg(target_os = "linux")]
549fn prepare_cgroup_placement(config: &CgroupPlacementConfig) -> Option<subc_cgroup::Placement> {
550 let result = match config {
551 CgroupPlacementConfig::Disabled => return None,
552 CgroupPlacementConfig::Current => subc_cgroup::prepare_current(),
553 CgroupPlacementConfig::Root(root) => subc_cgroup::prepare_at(root),
554 };
555
556 match result {
557 Ok(Some(placement)) => Some(placement),
558 Ok(None) => {
559 warn!(
560 placement = ?config,
561 "module cgroup placement is disabled: configured cgroup root is not delegated"
562 );
563 None
564 }
565 Err(error) => {
566 warn!(
567 placement = ?config,
568 error = %error,
569 "module cgroup placement is disabled by an unexpected cgroup probe error"
570 );
571 None
572 }
573 }
574}
575
576#[cfg(unix)]
583const NOFILE_TARGET: u64 = 65536;
584
585#[cfg(unix)]
589fn raise_nofile_limit() {
590 match rlimit::Resource::NOFILE.get() {
591 Ok((soft, hard)) => {
592 if soft >= NOFILE_TARGET {
593 return;
594 }
595 let target = NOFILE_TARGET.min(hard);
596 match rlimit::Resource::NOFILE.set(target, hard) {
597 Ok(()) => info!(
598 previous_soft = soft,
599 new_soft = target,
600 hard,
601 "raised open-file soft limit for daemon and module children"
602 ),
603 Err(err) => warn!(
604 soft,
605 hard,
606 error = %err,
607 "could not raise open-file soft limit; multi-root modules may exhaust descriptors"
608 ),
609 }
610 }
611 Err(err) => warn!(error = %err, "could not read open-file limit"),
612 }
613}
614
615#[cfg(windows)]
622fn raise_nofile_limit() {
623 const MAXSTDIO_TARGET: u32 = 8192;
624 let current = rlimit::getmaxstdio();
625 if current >= MAXSTDIO_TARGET {
626 return;
627 }
628 match rlimit::setmaxstdio(MAXSTDIO_TARGET) {
629 Ok(new_max) => info!(
630 previous = current,
631 new_max, "raised CRT stdio-stream limit for daemon"
632 ),
633 Err(err) => warn!(
634 current,
635 error = %err,
636 "could not raise CRT stdio-stream limit"
637 ),
638 }
639}
640
641#[cfg(not(any(unix, windows)))]
642fn raise_nofile_limit() {}
643
644pub async fn run_with_daemon_config_path(
645 config: BootstrapConfig,
646 daemon_config_path: impl AsRef<Path>,
647) -> Result<(), BootstrapError> {
648 run_with_config(config.with_daemon_config_path(daemon_config_path)?).await
649}
650
651#[allow(clippy::too_many_arguments)]
652async fn serve_bound_daemon(
653 bound: BoundDaemon,
654 configured_modules: Vec<ConfiguredModule>,
655 storage_config: Option<daemon_config::StorageConfig>,
656 admission_facts: AdmissionFactsConfig,
657 scope_authority_owners: Vec<String>,
658 daemon_config_path: Option<PathBuf>,
659 configured_port: Option<u16>,
660 route_bind_relay_default_ms: Option<u64>,
661 reserved_capabilities: BTreeMap<String, String>,
662 watchdog_config: DaemonSelfWatchdogConfig,
663 capture_logs_dir: Option<PathBuf>,
664 terminal_journal_path: Option<PathBuf>,
665 live_children_path: Option<PathBuf>,
666 #[cfg(target_os = "linux")] cgroup_placement: Option<subc_cgroup::Placement>,
667) -> Result<(), BootstrapError> {
668 #[cfg(unix)]
669 let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
670 .map_err(BootstrapError::Signal)?;
671 raise_nofile_limit();
674
675 info!(
676 connection_file = %bound.connection_file_path.display(),
677 connection_file_source = %bound.connection_file_source,
678 connection_file_source_reason = bound.connection_file_source.reason(),
679 endpoints = ?bound.connection_info.endpoints,
680 configured_modules = configured_modules.len(),
681 machine_id = bound.machine_id.as_ref().map(|id| id.as_str()).unwrap_or("none"),
682 "subc daemon starting"
683 );
684
685 let run_dir_lock = bound.run_dir_lock;
696 if let Some(owner) = &run_dir_lock {
697 crate::live_children::sweep_orphans(
698 owner,
699 &crate::live_children::AdoptedPids::none(),
700 crate::live_children::SweepBounds::default(),
701 )
702 .await;
703 }
704
705 let registry = Arc::new(Registry::default());
706 let process_liveness = Arc::new(SupervisorProcessLiveness::new());
707 let supervisor_handle = SupervisorHandle::new();
708 let connected_clients = ConnectedClients::new();
709 let forwarding = Arc::new(ForwardingTable::default());
710 let daemon_incarnation = format!(
711 "{:032x}",
712 u128::from_be_bytes(bound.connection_info.daemon_id)
713 );
714 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
715 .with_process_liveness(process_liveness.clone())
716 .with_forwarding(Arc::clone(&forwarding))
717 .with_handle(supervisor_handle.clone())
718 .with_connection_file_path(bound.connection_file_path.clone())
719 .with_daemon_incarnation(daemon_incarnation.clone());
720 let supervisor = match terminal_journal_path {
757 Some(path) => supervisor.with_terminal_journal(path, daemon_incarnation),
758 None => supervisor,
759 };
760 let supervisor = match capture_logs_dir {
761 Some(dir) => supervisor.with_capture_logs_dir(dir),
762 None => supervisor,
763 };
764 let supervisor = match live_children_path {
765 Some(path) => supervisor.with_live_children_record(path),
766 None => supervisor,
767 };
768 #[cfg(target_os = "linux")]
769 let supervisor = supervisor.with_cgroup_placement(cgroup_placement);
770 let route_bind_relay_timeouts = configured_modules
776 .iter()
777 .filter_map(|module| {
778 module
779 .route_bind_relay_timeout_ms
780 .map(|ms| (module.module_id.clone(), Duration::from_millis(ms)))
781 })
782 .collect::<std::collections::BTreeMap<_, _>>();
783 let control_start_clock = crate::clock::StartClock::capture();
784 let mut control = ControlHandler::with_forwarding(Arc::clone(®istry), forwarding)
785 .with_process_liveness(process_liveness)
786 .with_supervisor(supervisor_handle)
787 .with_connected_clients(connected_clients.clone())
788 .with_storage_config(storage_config)
789 .with_machine_id(bound.machine_id.clone())
790 .with_admission_facts_config(admission_facts.carrier_module_id, admission_facts.targets)
791 .with_scope_authority_owners(scope_authority_owners)
792 .with_route_bind_relay_timeouts(route_bind_relay_timeouts)
793 .with_daemon_provenance(
794 bound.connection_info.pid,
795 control_start_clock.started_at_ms(),
796 std::env::current_exe().ok(),
797 normalized_build_provenance(env!("SUBC_BUILD_GIT_SHA")),
798 normalized_build_provenance(env!("SUBC_BUILD_LOCK_DIGEST")),
799 )
800 .with_daemon_start_clock(control_start_clock)
801 .with_capability_config(
802 configured_modules
803 .iter()
804 .map(|module| (module.module_id.clone(), module.enabled)),
805 reserved_capabilities,
806 );
807 if let Some(ms) = route_bind_relay_default_ms {
808 control = control.with_route_bind_relay_timeout(Duration::from_millis(ms));
811 }
812 if let Some(config_path) = daemon_config_path {
813 control = control.with_supervisor_rescan(supervisor.clone(), config_path, configured_port);
814 }
815 let control = Arc::new(control);
816 let router = Arc::new(Router::with_control_handler(Arc::clone(&control)));
817 let auth = ServerAuth::new(
818 bound.connection_info.key.clone(),
819 bound.connection_info.daemon_id,
820 bound.connection_info.daemon_ver.clone(),
821 )
822 .with_connected_clients(connected_clients);
823
824 let mut serve_task =
825 AbortOnDrop::new(tokio::spawn(serve_listeners(bound.listeners, router, auth)));
826 tokio::task::yield_now().await;
827 let _clock_step_task = AbortOnDrop::new(crate::watchdog::spawn_clock_step_monitor());
828 #[cfg(target_os = "linux")]
830 let _kill_mode_check = AbortOnDrop::new(tokio::spawn(
831 crate::systemd_kill_mode::warn_if_kill_mode_defeats_ordered_shutdown(),
832 ));
833 let _watchdog_task = AbortOnDrop::new(
834 DaemonSelfWatchdog::new(
835 bound.connection_info.clone(),
836 bound.connection_file_path.clone(),
837 )
838 .with_config(watchdog_config)
839 .spawn(),
840 );
841
842 for configured in configured_modules {
843 let enabled = configured.enabled;
844 let health = configured.health;
845 let module_id = configured.module_id.clone();
846 match supervisor.supervise_configured_with_health(
847 configured.module_spec(),
848 enabled,
849 health,
850 configured.drain_timeout_ms,
851 configured.restart,
852 ) {
853 Ok(_) => {
854 let default_threshold = HealthConfig::default().failure_threshold;
863 if enabled && health.failure_threshold > default_threshold {
864 warn!(
865 module_id = %module_id,
866 failure_threshold = health.failure_threshold,
867 default_threshold,
868 tolerance_secs = health.cadence.as_secs() * u64::from(health.failure_threshold),
869 "health failure threshold is relaxed above the default; a wedged module stays unflagged for longer"
870 );
871 }
872 info!(module_id = %module_id, enabled, "configured module supervised");
873 }
874 Err(err) => {
875 error!(module_id = %module_id, error = %err, "failed to supervise configured module; continuing daemon startup");
876 }
877 }
878 }
879
880 control.refresh_capability_requirements();
881 Arc::clone(&control).spawn_capability_deadline_loop();
882
883 #[cfg(unix)]
884 {
885 tokio::select! {
886 result = serve_task.join() => {
887 return result.map_err(BootstrapError::ServeJoin)?.map_err(BootstrapError::Serve);
888 }
889 _ = terminate.recv() => {}
890 }
891 supervisor.begin_daemon_shutdown();
896 drop(_watchdog_task);
900 drop(serve_task);
903 let escalated = tokio::select! {
904 biased;
905 _ = terminate.recv() => {
906 info!("second SIGTERM: abandoning daemon shutdown wait");
907 true
908 }
909 result = supervisor.drain_for_daemon_shutdown() => {
910 if let Err(error) = result {
911 warn!(%error, "daemon shutdown drain failed; exiting anyway");
912 }
913 false
914 }
915 };
916 supervisor
923 .end_children_for_daemon_shutdown(escalated, async {
924 terminate.recv().await;
925 })
926 .await;
927 Ok(())
928 }
929 #[cfg(not(unix))]
930 serve_task
931 .join()
932 .await
933 .map_err(BootstrapError::ServeJoin)?
934 .map_err(BootstrapError::Serve)
935}
936
937fn normalized_build_provenance(value: &str) -> Option<String> {
938 match value.trim() {
939 "" | "unavailable" => None,
940 value => Some(value.to_string()),
941 }
942}
943
944pub async fn ensure_singleton(
952 connection_file_path: impl AsRef<Path>,
953 port: u16,
954) -> Result<Outcome, BootstrapError> {
955 ensure_singleton_with_config(BootstrapConfig::new(connection_file_path.as_ref(), port)).await
956}
957
958pub async fn ensure_singleton_with_config(
959 config: BootstrapConfig,
960) -> Result<Outcome, BootstrapError> {
961 let path = config.connection_file_path;
962
963 if matches!(probe_existing(&path).await?, Probe::Live) {
964 return Ok(Outcome::AlreadyRunning);
965 }
966
967 let _lock = StartLock::acquire(&path).await?;
968
969 if matches!(probe_existing(&path).await?, Probe::Live) {
973 return Ok(Outcome::AlreadyRunning);
974 }
975
976 let run_dir_lock = config
982 .live_children_path
983 .as_deref()
984 .map(crate::run_dir_lock::RunDirLock::acquire)
985 .transpose()?;
986
987 remove_stale_connection_file_if_present(&path)?;
988
989 let machine_id = config
993 .machine_id_path
994 .as_deref()
995 .map(crate::machine_id::load_or_mint)
996 .transpose()
997 .map_err(BootstrapError::MachineId)?;
998
999 let (listeners, endpoints) = bind_loopback(config.port).await?;
1000 let connection_info = ConnectionInfo {
1001 schema: SCHEMA_VERSION,
1002 wire_version: Some(PROTOCOL_VERSION),
1003 endpoints,
1004 key: generate_key().map_err(BootstrapError::GenerateConnectionFile)?,
1005 daemon_id: generate_daemon_id().map_err(BootstrapError::GenerateConnectionFile)?,
1006 pid: process::id(),
1007 daemon_ver: config.daemon_ver,
1008 };
1009
1010 if let Err(source) = write_atomic(&path, &connection_info) {
1011 drop(listeners);
1012 return Err(BootstrapError::ConnectionFileWrite { path, source });
1013 }
1014
1015 Ok(Outcome::Bound(BoundDaemon {
1016 listeners,
1017 connection_info,
1018 connection_file_path: path,
1019 connection_file_source: config.connection_file_source,
1020 machine_id,
1021 run_dir_lock,
1022 }))
1023}
1024
1025#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1026enum Probe {
1027 Live,
1028 StaleOrAbsent,
1029}
1030
1031async fn probe_existing(path: &Path) -> Result<Probe, BootstrapError> {
1032 let info = match connection_file::read(path) {
1033 Ok(info) => info,
1034 Err(source) if is_absent_or_stale_connection_file(&source) => {
1035 return Ok(Probe::StaleOrAbsent)
1036 }
1037 Err(source) => {
1038 return Err(BootstrapError::ConnectionFileRead {
1039 path: path.to_path_buf(),
1040 source,
1041 })
1042 }
1043 };
1044
1045 for endpoint in &info.endpoints {
1046 if matches!(probe_endpoint(&info, endpoint).await, Probe::Live) {
1047 return Ok(Probe::Live);
1048 }
1049 }
1050
1051 Ok(Probe::StaleOrAbsent)
1052}
1053
1054async fn probe_endpoint(info: &ConnectionInfo, endpoint: &Endpoint) -> Probe {
1055 let Ok(ip) = endpoint.host.parse::<IpAddr>() else {
1056 return Probe::StaleOrAbsent;
1057 };
1058 if !ip.is_loopback() {
1059 return Probe::StaleOrAbsent;
1060 }
1061 let addr = SocketAddr::new(ip, endpoint.port);
1062
1063 let mut stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(addr)).await {
1064 Ok(Ok(stream)) => stream,
1065 Ok(Err(_)) | Err(_) => return Probe::StaleOrAbsent,
1066 };
1067
1068 match authenticate_client(&mut stream, info, PROBE_AUTH_DEADLINE).await {
1069 Ok(()) => Probe::Live,
1070 Err(AuthError::DaemonIdMismatch)
1071 | Err(AuthError::InvalidServerProof)
1072 | Err(AuthError::UnexpectedEof { .. })
1073 | Err(AuthError::Timeout { .. })
1074 | Err(AuthError::JsonEncode { .. })
1075 | Err(AuthError::JsonDecode { .. })
1076 | Err(AuthError::Io { .. })
1077 | Err(AuthError::MessageTooLarge { .. })
1078 | Err(AuthError::KeyTooShort { .. })
1079 | Err(AuthError::Random(_))
1080 | Err(AuthError::InvalidClientAuth) => Probe::StaleOrAbsent,
1081 }
1082}
1083
1084fn is_absent_or_stale_connection_file(err: &ConnectionFileError) -> bool {
1085 match err {
1086 ConnectionFileError::Io { source, .. } if source.kind() == io::ErrorKind::NotFound => true,
1087 ConnectionFileError::JsonRead { .. }
1088 | ConnectionFileError::UnsupportedSchema { .. }
1089 | ConnectionFileError::Invalid { .. }
1090 | ConnectionFileError::KeyTooShort { .. }
1091 | ConnectionFileError::InsecurePermissions { .. } => true,
1095 ConnectionFileError::MissingParent { .. }
1096 | ConnectionFileError::MissingFileName { .. }
1097 | ConnectionFileError::InsecureParentDirectory { .. }
1101 | ConnectionFileError::Io { .. }
1102 | ConnectionFileError::JsonWrite { .. }
1103 | ConnectionFileError::Random(_)
1104 | ConnectionFileError::WireVersionMismatch { .. } => false,
1107 }
1108}
1109
1110fn remove_stale_connection_file_if_present(path: &Path) -> Result<(), BootstrapError> {
1111 match fs::remove_file(path) {
1112 Ok(()) => Ok(()),
1113 Err(err) if err.kind() == io::ErrorKind::NotFound => Ok(()),
1114 Err(source) => Err(BootstrapError::RemoveStale {
1115 path: path.to_path_buf(),
1116 source,
1117 }),
1118 }
1119}
1120
1121async fn bind_loopback(port: u16) -> Result<(Vec<TcpListener>, Vec<Endpoint>), BootstrapError> {
1122 let v4_host = Ipv4Addr::LOCALHOST;
1123 let v4 = TcpListener::bind((v4_host, port))
1124 .await
1125 .map_err(|source| BootstrapError::Bind {
1126 host: v4_host.to_string(),
1127 port,
1128 source,
1129 })?;
1130 let actual_port = v4
1131 .local_addr()
1132 .map_err(|source| BootstrapError::LocalAddr {
1133 host: v4_host.to_string(),
1134 source,
1135 })?
1136 .port();
1137
1138 let mut listeners = vec![v4];
1139 let mut endpoints = vec![Endpoint {
1140 host: v4_host.to_string(),
1141 port: actual_port,
1142 }];
1143
1144 let v6_host = Ipv6Addr::LOCALHOST;
1145 match TcpListener::bind((v6_host, actual_port)).await {
1146 Ok(v6) => {
1147 listeners.push(v6);
1148 endpoints.push(Endpoint {
1149 host: v6_host.to_string(),
1150 port: actual_port,
1151 });
1152 }
1153 Err(err) if ipv6_loopback_unavailable(&err) => {
1154 warn!(
1155 port = actual_port,
1156 error = %err,
1157 "IPv6 loopback unavailable; serving only IPv4 loopback"
1158 );
1159 }
1160 Err(source) => {
1161 drop(listeners);
1162 return Err(BootstrapError::Bind {
1163 host: v6_host.to_string(),
1164 port: actual_port,
1165 source,
1166 });
1167 }
1168 }
1169
1170 Ok((listeners, endpoints))
1171}
1172
1173fn ipv6_loopback_unavailable(err: &io::Error) -> bool {
1174 matches!(
1175 err.kind(),
1176 io::ErrorKind::AddrNotAvailable | io::ErrorKind::Unsupported
1177 ) || matches!(err.raw_os_error(), Some(47) | Some(49) | Some(97))
1178}
1179
1180struct AbortOnDrop<T> {
1181 handle: JoinHandle<T>,
1182}
1183
1184impl<T> AbortOnDrop<T> {
1185 fn new(handle: JoinHandle<T>) -> Self {
1186 Self { handle }
1187 }
1188
1189 async fn join(&mut self) -> Result<T, JoinError> {
1190 (&mut self.handle).await
1191 }
1192}
1193
1194impl<T> Drop for AbortOnDrop<T> {
1195 fn drop(&mut self) {
1196 if !self.handle.is_finished() {
1197 self.handle.abort();
1198 }
1199 }
1200}
1201
1202struct StartLock {
1203 _file: fs::File,
1206}
1207
1208impl StartLock {
1209 async fn acquire(connection_file_path: &Path) -> Result<Self, BootstrapError> {
1210 let path = start_lock_path(connection_file_path);
1211 for _ in 0..START_LOCK_RETRIES {
1212 let file = match open_owner_only_lock(&path) {
1213 Ok(file) => file,
1214 Err(source) => return Err(BootstrapError::StartLockCreate { path, source }),
1215 };
1216 match FileExt::try_lock(&file) {
1217 Ok(()) => return Ok(Self { _file: file }),
1218 Err(TryLockError::WouldBlock) => sleep(START_LOCK_RETRY_DELAY).await,
1219 Err(TryLockError::Error(source)) => {
1220 return Err(BootstrapError::StartLockCreate { path, source });
1221 }
1222 }
1223 }
1224
1225 Err(BootstrapError::StartLockBusy {
1226 path,
1227 attempts: START_LOCK_RETRIES,
1228 })
1229 }
1230}
1231
1232pub(crate) fn open_owner_only_lock(path: &Path) -> io::Result<fs::File> {
1233 let mut options = fs::OpenOptions::new();
1234 options.read(true).write(true).create(true);
1235 #[cfg(unix)]
1236 {
1237 use std::os::unix::fs::OpenOptionsExt;
1238 options.mode(0o600);
1239 }
1240 options.open(path)
1241}
1242
1243fn start_lock_path(connection_file_path: &Path) -> PathBuf {
1244 let file_name = connection_file_path
1245 .file_name()
1246 .map(|name| name.to_string_lossy())
1247 .unwrap_or_else(|| CONNECTION_FILE_NAME.into());
1248 let lock_name = format!("{file_name}.start-lock");
1249 connection_file_path
1250 .parent()
1251 .filter(|parent| !parent.as_os_str().is_empty())
1252 .unwrap_or_else(|| Path::new("."))
1253 .join(lock_name)
1254}
1255
1256fn non_empty_os_var(key: &str) -> Option<OsString> {
1257 let value = env::var_os(key)?;
1258 if value.is_empty() {
1259 None
1260 } else {
1261 Some(value)
1262 }
1263}
1264
1265#[derive(Debug)]
1268pub enum BootstrapError {
1269 #[cfg(unix)]
1270 Signal(io::Error),
1271 InvalidPort {
1272 raw: String,
1273 source: std::num::ParseIntError,
1274 },
1275 ConnectionFileRead {
1276 path: PathBuf,
1277 source: ConnectionFileError,
1278 },
1279 ConnectionFileWrite {
1280 path: PathBuf,
1281 source: ConnectionFileError,
1282 },
1283 GenerateConnectionFile(ConnectionFileError),
1284 StartLockCreate {
1285 path: PathBuf,
1286 source: io::Error,
1287 },
1288 StartLockBusy {
1289 path: PathBuf,
1290 attempts: usize,
1291 },
1292 RunDirLockCreate {
1294 path: PathBuf,
1295 source: io::Error,
1296 },
1297 RunDirBusy {
1302 path: PathBuf,
1303 holder_pid: Option<u32>,
1304 },
1305 RemoveStale {
1306 path: PathBuf,
1307 source: io::Error,
1308 },
1309 Bind {
1310 host: String,
1311 port: u16,
1312 source: io::Error,
1313 },
1314 LocalAddr {
1315 host: String,
1316 source: io::Error,
1317 },
1318 DaemonConfig(DaemonConfigError),
1319 MachineId(crate::machine_id::MachineIdFileError),
1322 RunDir(daemon_config::DaemonRunDirError),
1325 Serve(ServerError),
1326 ServeJoin(tokio::task::JoinError),
1327}
1328
1329impl fmt::Display for BootstrapError {
1330 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1331 match self {
1332 #[cfg(unix)]
1333 Self::Signal(error) => write!(f, "failed to register SIGTERM handler: {error}"),
1334 Self::InvalidPort { raw, source } => {
1335 write!(f, "invalid {SUBC_PORT_ENV} value '{raw}': {source}")
1336 }
1337 Self::ConnectionFileRead { path, source } => write!(
1338 f,
1339 "failed to read connection file {}: {source}",
1340 path.display()
1341 ),
1342 Self::ConnectionFileWrite { path, source } => write!(
1343 f,
1344 "failed to publish connection file {}: {source}",
1345 path.display()
1346 ),
1347 Self::GenerateConnectionFile(err) => {
1348 write!(f, "failed to generate connection-file auth material: {err}")
1349 }
1350 Self::StartLockCreate { path, source } => {
1351 write!(
1352 f,
1353 "failed to create start lock {}: {source}",
1354 path.display()
1355 )
1356 }
1357 Self::StartLockBusy { path, attempts } => write!(
1358 f,
1359 "start lock {} remained busy after {attempts} attempts",
1360 path.display()
1361 ),
1362 Self::RunDirLockCreate { path, source } => write!(
1363 f,
1364 "refusing to start: failed to lock run directory via {}: {source}",
1365 path.display()
1366 ),
1367 Self::RunDirBusy { path, holder_pid } => {
1368 write!(
1369 f,
1370 "refusing to start: run directory lock {} is held by another daemon",
1371 path.display()
1372 )?;
1373 if let Some(pid) = holder_pid {
1374 write!(f, " (pid {pid})")?;
1375 }
1376 write!(
1377 f,
1378 "; a live daemon owns this data home's run state (typically one started with a different XDG_RUNTIME_DIR)"
1379 )
1380 }
1381 Self::RemoveStale { path, source } => write!(
1382 f,
1383 "failed to remove stale connection file {}: {source}",
1384 path.display()
1385 ),
1386 Self::Bind { host, port, source } if source.kind() == io::ErrorKind::AddrInUse => {
1387 write!(
1388 f,
1389 "port {port} in use on loopback {host}: {source}; set the port in config"
1390 )
1391 }
1392 Self::Bind { host, port, source } => {
1393 write!(f, "failed to bind loopback TCP {host}:{port}: {source}")
1394 }
1395 Self::LocalAddr { host, source } => {
1396 write!(f, "failed to read local address for {host}: {source}")
1397 }
1398 Self::DaemonConfig(err) => write!(f, "failed to load daemon config: {err}"),
1399 Self::MachineId(err) => write!(f, "refusing to start: {err}"),
1400 Self::RunDir(err) => write!(f, "refusing to start: {err}"),
1401 Self::Serve(err) => write!(f, "daemon server failed: {err}"),
1402 Self::ServeJoin(err) => write!(f, "daemon server task failed: {err}"),
1403 }
1404 }
1405}
1406
1407impl Error for BootstrapError {
1408 fn source(&self) -> Option<&(dyn Error + 'static)> {
1409 match self {
1410 #[cfg(unix)]
1411 Self::Signal(source) => Some(source),
1412 Self::InvalidPort { source, .. } => Some(source),
1413 Self::ConnectionFileRead { source, .. }
1414 | Self::ConnectionFileWrite { source, .. }
1415 | Self::GenerateConnectionFile(source) => Some(source),
1416 Self::StartLockCreate { source, .. }
1417 | Self::RunDirLockCreate { source, .. }
1418 | Self::RemoveStale { source, .. }
1419 | Self::Bind { source, .. }
1420 | Self::LocalAddr { source, .. } => Some(source),
1421 Self::DaemonConfig(err) => Some(err),
1422 Self::MachineId(err) => Some(err),
1423 Self::RunDir(err) => Some(err),
1424 Self::Serve(err) => Some(err),
1425 Self::ServeJoin(err) => Some(err),
1426 Self::StartLockBusy { .. } | Self::RunDirBusy { .. } => None,
1427 }
1428 }
1429}
1430
1431#[cfg(test)]
1432mod tests {
1433 use super::*;
1434 use crate::server::ServerAuth;
1435 #[cfg(target_os = "linux")]
1436 use std::collections::BTreeSet;
1437 use std::sync::Mutex;
1438 #[cfg(target_os = "linux")]
1439 use subc_control::ModuleProtocol;
1440 use subc_test_support::TestTempDir;
1441 use subc_transport::MIN_KEY_LEN;
1442 use tokio::io::AsyncReadExt;
1443 use tokio::task::JoinHandle;
1444
1445 #[cfg(unix)]
1446 use std::os::unix::fs::PermissionsExt;
1447
1448 static ENV_LOCK: Mutex<()> = Mutex::new(());
1449
1450 #[test]
1451 fn normalized_build_provenance_preserves_real_values() {
1452 assert_eq!(normalized_build_provenance("abc"), Some("abc".to_string()));
1453 }
1454
1455 #[test]
1456 fn normalized_build_provenance_omits_unavailable_and_empty_values() {
1457 assert_eq!(normalized_build_provenance("unavailable"), None);
1458 assert_eq!(normalized_build_provenance(""), None);
1459 }
1460
1461 #[cfg(target_os = "linux")]
1462 fn current_cgroup_path_for_test() -> io::Result<PathBuf> {
1463 let cgroups = fs::read_to_string("/proc/self/cgroup")?;
1464 let relative = cgroups
1465 .lines()
1466 .find_map(|line| line.strip_prefix("0::"))
1467 .ok_or_else(|| {
1468 io::Error::new(io::ErrorKind::Unsupported, "cgroup v2 is unavailable")
1469 })?;
1470 Ok(Path::new("/sys/fs/cgroup").join(relative.trim_start_matches('/')))
1471 }
1472
1473 #[cfg(target_os = "linux")]
1474 fn module_cgroup_directories() -> io::Result<BTreeSet<OsString>> {
1475 let modules = current_cgroup_path_for_test()?.join("subc-modules");
1476 let entries = match fs::read_dir(modules) {
1477 Ok(entries) => entries,
1478 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(BTreeSet::new()),
1479 Err(error) => return Err(error),
1480 };
1481 let mut directories = BTreeSet::new();
1482 for entry in entries {
1483 let entry = entry?;
1484 if entry.file_type()?.is_dir() {
1485 directories.insert(entry.file_name());
1486 }
1487 }
1488 Ok(directories)
1489 }
1490
1491 #[cfg(target_os = "linux")]
1492 async fn wait_for_path(path: &Path, task: &JoinHandle<Result<(), BootstrapError>>) {
1493 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1494 while !path.exists() && tokio::time::Instant::now() < deadline {
1495 assert!(
1496 !task.is_finished(),
1497 "daemon exited before creating {}",
1498 path.display()
1499 );
1500 sleep(Duration::from_millis(10)).await;
1501 }
1502 assert!(path.exists(), "daemon did not create {}", path.display());
1503 }
1504
1505 #[cfg(target_os = "linux")]
1510 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1511 async fn run_with_config_does_not_reconcile_the_ambient_cgroup_by_default() {
1512 let temp = unique_temp_dir("bootstrap-cgroup-default-disabled");
1513 let module_id = format!("cgroup-isolation-probe-{}", process::id());
1514 let before = module_cgroup_directories().expect("read ambient module cgroups before boot");
1515 assert!(
1516 !before.contains(&OsString::from(&module_id)),
1517 "isolation probe cgroup already exists before this daemon starts"
1518 );
1519 let capture = temp.join("logs").join(format!("{module_id}.stderr.log"));
1520 let module = ConfiguredModule {
1521 module_id,
1522 program: PathBuf::from("sh"),
1523 args: vec!["-c".to_string(), "sleep 30".to_string()],
1524 env: Vec::new(),
1525 log: None,
1526 enabled: true,
1527 reserved: false,
1528 reserved_prefixes: Vec::new(),
1529 protocol: ModuleProtocol::None,
1530 overlap: Default::default(),
1531 health: HealthConfig::default(),
1532 drain_timeout_ms: None,
1533 route_bind_relay_timeout_ms: None,
1534 restart: RestartPolicy::default(),
1535 };
1536 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1537 .with_configured_modules([module])
1538 .with_capture_logs_dir(temp.join("logs"))
1539 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1540 let task = tokio::spawn(run_with_config(config));
1541
1542 wait_for_path(&capture, &task).await;
1543 let after = module_cgroup_directories().expect("read ambient module cgroups after boot");
1544
1545 task.abort();
1546 assert!(task
1547 .await
1548 .expect_err("aborted daemon task must cancel")
1549 .is_cancelled());
1550 assert_eq!(
1551 before, after,
1552 "an in-process daemon must not create or reconcile ambient module cgroups"
1553 );
1554 }
1555
1556 #[cfg(target_os = "linux")]
1557 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1558 async fn explicit_cgroup_root_is_prepared_inside_the_fixture_tree() {
1559 let temp = unique_temp_dir("bootstrap-cgroup-explicit-root");
1560 let cgroup_root = temp.join("cgroup");
1561 fs::create_dir(&cgroup_root).expect("create scratch cgroup root");
1562 fs::write(cgroup_root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
1563 let modules = cgroup_root.join("subc-modules");
1564 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1565 .with_cgroup_placement(CgroupPlacementConfig::Root(cgroup_root))
1566 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1567 let task = tokio::spawn(run_with_config(config));
1568
1569 wait_for_path(&modules, &task).await;
1570
1571 task.abort();
1572 assert!(task
1573 .await
1574 .expect_err("aborted daemon task must cancel")
1575 .is_cancelled());
1576 assert!(
1577 fs::read_dir(&modules)
1578 .expect("read prepared modules directory")
1579 .next()
1580 .is_none(),
1581 "the delegation probe must clean up after itself"
1582 );
1583 }
1584
1585 struct EnvGuard {
1586 key: &'static str,
1587 previous: Option<OsString>,
1588 }
1589
1590 impl EnvGuard {
1591 fn set(key: &'static str, value: &Path) -> Self {
1592 let previous = env::var_os(key);
1593 env::set_var(key, value);
1594 Self { key, previous }
1595 }
1596
1597 fn set_str(key: &'static str, value: &str) -> Self {
1598 let previous = env::var_os(key);
1599 env::set_var(key, value);
1600 Self { key, previous }
1601 }
1602
1603 fn unset(key: &'static str) -> Self {
1604 let previous = env::var_os(key);
1605 env::remove_var(key);
1606 Self { key, previous }
1607 }
1608 }
1609
1610 impl Drop for EnvGuard {
1611 fn drop(&mut self) {
1612 match &self.previous {
1613 Some(value) => env::set_var(self.key, value),
1614 None => env::remove_var(self.key),
1615 }
1616 }
1617 }
1618
1619 fn unique_temp_dir(name: &str) -> TestTempDir {
1620 TestTempDir::new(name)
1621 }
1622
1623 fn temp_connection_file_path(name: &str) -> (TestTempDir, PathBuf) {
1624 let dir = unique_temp_dir(name);
1625 let path = dir.join("conn.json");
1626 (dir, path)
1627 }
1628
1629 fn auth_for(info: &ConnectionInfo) -> ServerAuth {
1630 ServerAuth::new(info.key.clone(), info.daemon_id, info.daemon_ver.clone())
1631 }
1632
1633 fn start_server(bound: BoundDaemon) -> JoinHandle<Result<(), ServerError>> {
1634 let auth = auth_for(&bound.connection_info);
1635 tokio::spawn(serve_listeners(
1636 bound.listeners,
1637 Arc::new(Router::with_default_self_handler()),
1638 auth,
1639 ))
1640 }
1641
1642 fn expect_bound(outcome: Outcome) -> BoundDaemon {
1643 match outcome {
1644 Outcome::Bound(bound) => bound,
1645 Outcome::AlreadyRunning => panic!("fresh connection file unexpectedly had a daemon"),
1646 }
1647 }
1648
1649 async fn connect_from_info(conn: &ConnectionInfo) -> io::Result<TcpStream> {
1650 let endpoint = conn
1651 .endpoints
1652 .first()
1653 .expect("test connection file should have an endpoint");
1654 let ip: IpAddr = endpoint.host.parse().unwrap();
1655 TcpStream::connect(SocketAddr::new(ip, endpoint.port)).await
1656 }
1657
1658 fn make_connection_info(port: u16) -> ConnectionInfo {
1659 ConnectionInfo {
1660 schema: SCHEMA_VERSION,
1661 wire_version: Some(PROTOCOL_VERSION),
1662 endpoints: vec![Endpoint {
1663 host: "127.0.0.1".to_owned(),
1664 port,
1665 }],
1666 key: generate_key().unwrap(),
1667 daemon_id: generate_daemon_id().unwrap(),
1668 pid: process::id(),
1669 daemon_ver: "test-subc".to_owned(),
1670 }
1671 }
1672
1673 fn write_raw_owner_only_connection_file(path: &Path, contents: &[u8]) {
1674 fs::write(path, contents).unwrap();
1675 #[cfg(unix)]
1676 fs::set_permissions(path, fs::Permissions::from_mode(0o600)).unwrap();
1677 }
1678
1679 fn assert_owner_only_connection_file(path: &Path) {
1680 #[cfg(unix)]
1683 {
1684 let mode = fs::metadata(path).unwrap().permissions().mode() & 0o777;
1685 assert_eq!(mode, 0o600);
1686 }
1687 #[cfg(not(unix))]
1688 let _ = path;
1689 }
1690
1691 #[test]
1692 fn connection_file_path_uses_xdg_runtime_dir_when_set() {
1693 let _env_lock = ENV_LOCK.lock().unwrap();
1694 let runtime_dir = unique_temp_dir("xdg-runtime");
1695 let _xdg = EnvGuard::set("XDG_RUNTIME_DIR", runtime_dir.path());
1696
1697 assert_eq!(
1698 connection_file_path(),
1699 runtime_dir.join(CONNECTION_FILE_NAME)
1700 );
1701 }
1702
1703 #[test]
1704 fn connection_file_path_source_is_xdg_runtime_dir_when_set() {
1705 let runtime_dir = OsString::from("/run/user/1000");
1706
1707 let (path, source) = connection_file_path_with_source(Some(runtime_dir));
1708
1709 assert_eq!(
1710 path,
1711 PathBuf::from("/run/user/1000").join(CONNECTION_FILE_NAME)
1712 );
1713 assert_eq!(source, ConnectionFileSource::XdgRuntimeDir);
1714 }
1715
1716 #[test]
1717 fn connection_file_path_falls_back_to_temp_dir_with_user_token_when_xdg_unset() {
1718 let _env_lock = ENV_LOCK.lock().unwrap();
1719 let _xdg = EnvGuard::unset("XDG_RUNTIME_DIR");
1720
1721 assert_eq!(
1722 connection_file_path(),
1723 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1724 );
1725 }
1726
1727 #[test]
1728 fn connection_file_path_source_is_temp_dir_when_xdg_unset() {
1729 let (path, source) = connection_file_path_with_source(None);
1730
1731 assert_eq!(
1732 path,
1733 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1734 );
1735 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1736 }
1737
1738 #[test]
1743 fn user_connection_token_is_stable_under_concurrent_callers() {
1744 let expected = user_connection_token();
1745 let workers: Vec<_> = (0..32)
1746 .map(|_| {
1747 std::thread::spawn(|| (0..40).map(|_| user_connection_token()).collect::<Vec<_>>())
1748 })
1749 .collect();
1750 for worker in workers {
1751 for token in worker.join().expect("probe thread") {
1752 assert_eq!(token, expected, "token diverged under concurrent probes");
1753 }
1754 }
1755 }
1756
1757 #[test]
1758 fn connection_file_path_source_is_temp_dir_when_xdg_empty() {
1759 let (path, source) = connection_file_path_with_source(Some(OsString::new()));
1760
1761 assert_eq!(
1762 path,
1763 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1764 );
1765 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1766 }
1767
1768 #[test]
1772 fn daemon_binary_config_journals_and_captures_into_the_run_dir() {
1773 let _env_lock = ENV_LOCK.lock().unwrap();
1774 let root = unique_temp_dir("daemon-binary-config");
1775 let data_home = root.join("data");
1776 let _data = EnvGuard::set("XDG_DATA_HOME", &data_home);
1777 let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1778 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1779
1780 let config = BootstrapConfig::from_env_for_daemon_binary().unwrap();
1781
1782 let run_dir = data_home.join("cortexkit").join("run");
1783 assert_eq!(
1784 config.terminal_journal_path,
1785 Some(run_dir.join("terminals.jsonl"))
1786 );
1787 assert_eq!(config.capture_logs_dir, Some(run_dir.join("logs")));
1788 assert_eq!(
1789 BootstrapConfig::new(root.join("connection.json"), 0).terminal_journal_path,
1790 None,
1791 "an in-process config must not journal anywhere unless asked to"
1792 );
1793 }
1794
1795 #[test]
1796 fn daemon_binary_config_refuses_a_relative_data_home() {
1797 let _env_lock = ENV_LOCK.lock().unwrap();
1798 let root = unique_temp_dir("daemon-binary-relative-data");
1799 let _data = EnvGuard::set_str("XDG_DATA_HOME", "relative-data-home");
1800 let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1801 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1802
1803 let error = BootstrapConfig::from_env_for_daemon_binary()
1804 .expect_err("a relative data home must refuse the daemon binary's config");
1805 assert!(
1806 matches!(error, BootstrapError::RunDir(_)),
1807 "expected a run-directory refusal, got {error}"
1808 );
1809 assert!(error.to_string().contains("XDG_DATA_HOME"), "{error}");
1810 }
1811
1812 #[test]
1813 fn configured_port_uses_default_config_and_env_override() {
1814 let _env_lock = ENV_LOCK.lock().unwrap();
1815 let (_dir, conn_path) = temp_connection_file_path("daemon-config-port");
1816 let config_path = conn_path.with_file_name("subc.jsonc");
1817
1818 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1819 assert_eq!(
1820 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1821 .unwrap()
1822 .port,
1823 DEFAULT_SUBC_PORT
1824 );
1825
1826 fs::write(&config_path, r#"{ "version": 1, "port": 8123 }"#).unwrap();
1827 assert_eq!(
1828 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1829 .unwrap()
1830 .port,
1831 8123
1832 );
1833
1834 let _port = EnvGuard::set_str(SUBC_PORT_ENV, "9012");
1835 assert_eq!(
1836 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1837 .unwrap()
1838 .port,
1839 9012
1840 );
1841 }
1842
1843 #[tokio::test]
1844 async fn second_singleton_probe_against_served_tcp_daemon_reports_already_running() {
1845 let (_dir, path) = temp_connection_file_path("already-running");
1846
1847 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1848 let server = start_server(bound);
1849
1850 let second = ensure_singleton(&path, 0).await.unwrap();
1851 assert!(matches!(second, Outcome::AlreadyRunning));
1852
1853 server.abort();
1854 let _ = server.await;
1855 }
1856
1857 #[tokio::test]
1858 async fn daemon_connection_file_publishes_protocol_wire_version() {
1859 let (_dir, path) = temp_connection_file_path("wire-version");
1860 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1861 assert_eq!(bound.connection_info.wire_version, Some(PROTOCOL_VERSION));
1862 assert_eq!(
1863 connection_file::read(&path).unwrap().wire_version,
1864 Some(PROTOCOL_VERSION)
1865 );
1866
1867 drop(bound.listeners);
1868 }
1869
1870 #[tokio::test]
1871 async fn stale_unbound_connection_file_is_reclaimed() {
1872 let (_dir, path) = temp_connection_file_path("stale-reclaim");
1873 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1874 let stale_port = stale.local_addr().unwrap().port();
1875 drop(stale);
1876 let stale_info = make_connection_info(stale_port);
1877 write_atomic(&path, &stale_info).unwrap();
1878
1879 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1880 assert_ne!(bound.connection_info.key, stale_info.key);
1881 drop(bound.listeners);
1882 }
1883
1884 #[cfg(unix)]
1885 #[tokio::test]
1886 async fn ensure_singleton_reclaims_insecure_connection_file() {
1887 let (_dir, path) = temp_connection_file_path("insecure-reclaim");
1888 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1889 let stale_port = stale.local_addr().unwrap().port();
1890 drop(stale);
1891 let stale_info = make_connection_info(stale_port);
1892 write_atomic(&path, &stale_info).unwrap();
1893 fs::set_permissions(&path, fs::Permissions::from_mode(0o644)).unwrap();
1894
1895 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1896 assert_ne!(bound.connection_info.key, stale_info.key);
1897 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1898 assert_owner_only_connection_file(&path);
1899
1900 drop(bound.listeners);
1901 }
1902
1903 #[tokio::test]
1904 async fn ensure_singleton_reclaims_non_loopback_connection_file() {
1905 let (_dir, path) = temp_connection_file_path("non-loopback-reclaim");
1906 let mut stale_info = make_connection_info(8757);
1907 stale_info.endpoints = vec![Endpoint {
1908 host: "192.0.2.10".to_owned(),
1909 port: 8757,
1910 }];
1911 write_atomic(&path, &stale_info).unwrap();
1912
1913 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1914 assert_ne!(bound.connection_info.key, stale_info.key);
1915 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1916 assert!(bound
1917 .connection_info
1918 .endpoints
1919 .iter()
1920 .all(|endpoint| endpoint.host.parse::<IpAddr>().unwrap().is_loopback()));
1921 assert_owner_only_connection_file(&path);
1922
1923 drop(bound.listeners);
1924 }
1925
1926 #[tokio::test]
1927 async fn ensure_singleton_reclaims_invalid_connection_file_shapes() {
1928 let mut unsupported_schema = make_connection_info(8757);
1929 unsupported_schema.schema = SCHEMA_VERSION + 1;
1930
1931 let mut empty_endpoints = make_connection_info(8757);
1932 empty_endpoints.endpoints.clear();
1933
1934 let mut short_key = make_connection_info(8757);
1935 short_key.key = vec![0x5A; MIN_KEY_LEN - 1];
1936
1937 let cases = vec![
1938 (
1939 "unsupported-schema",
1940 serde_json::to_vec(&unsupported_schema).unwrap(),
1941 Some(unsupported_schema),
1942 ),
1943 (
1944 "empty-endpoints",
1945 serde_json::to_vec(&empty_endpoints).unwrap(),
1946 Some(empty_endpoints),
1947 ),
1948 (
1949 "short-key",
1950 serde_json::to_vec(&short_key).unwrap(),
1951 Some(short_key),
1952 ),
1953 ("invalid-json", b"{not valid connection json".to_vec(), None),
1954 ];
1955
1956 for (label, contents, old_info) in cases {
1957 let (_dir, path) = temp_connection_file_path(label);
1958 write_raw_owner_only_connection_file(&path, &contents);
1959
1960 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1961 if let Some(old_info) = old_info {
1962 assert_ne!(bound.connection_info.key, old_info.key, "{label}");
1963 assert_ne!(
1964 bound.connection_info.daemon_id, old_info.daemon_id,
1965 "{label}"
1966 );
1967 }
1968 assert!(bound.connection_info.key.len() >= MIN_KEY_LEN, "{label}");
1969 assert_ne!(bound.connection_info.daemon_id, [0u8; 16], "{label}");
1970 assert_owner_only_connection_file(&path);
1971
1972 drop(bound.listeners);
1973 }
1974 }
1975
1976 #[tokio::test]
1977 async fn foreign_reused_port_connection_file_is_reclaimed_after_auth_probe_fails() {
1978 let (_dir, path) = temp_connection_file_path("foreign-reclaim");
1979 let foreign = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1980 let foreign_port = foreign.local_addr().unwrap().port();
1981 write_atomic(&path, &make_connection_info(foreign_port)).unwrap();
1982 let foreign_task = tokio::spawn(async move {
1983 if let Ok((mut stream, _)) = foreign.accept().await {
1984 let mut buf = [0u8; 64];
1985 let _ = stream.read(&mut buf).await;
1986 }
1987 });
1988
1989 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1990 assert!(bound
1991 .connection_info
1992 .endpoints
1993 .iter()
1994 .all(|endpoint| endpoint.port != foreign_port));
1995
1996 drop(bound.listeners);
1997 let _ = foreign_task.await;
1998 }
1999
2000 #[tokio::test]
2001 async fn stale_start_lock_file_is_reclaimable() {
2002 let (_dir, path) = temp_connection_file_path("start-lock-stale-file");
2003 let lock_path = start_lock_path(&path);
2004 drop(open_owner_only_lock(&lock_path).unwrap());
2005 assert!(lock_path.is_file());
2006
2007 let lock = StartLock::acquire(&path).await.unwrap();
2008 assert!(lock_path.is_file());
2009
2010 drop(lock);
2011 assert!(lock_path.is_file());
2012 }
2013
2014 #[tokio::test]
2015 async fn held_start_lock_blocks_second_acquire_until_release() {
2016 let (_dir, path) = temp_connection_file_path("start-lock-held");
2017 let lock_path = start_lock_path(&path);
2018 let first = StartLock::acquire(&path).await.unwrap();
2019
2020 let err = match StartLock::acquire(&path).await {
2021 Ok(_) => panic!("second acquire while held must stay busy"),
2022 Err(err) => err,
2023 };
2024 assert!(matches!(
2025 err,
2026 BootstrapError::StartLockBusy {
2027 ref path,
2028 attempts: START_LOCK_RETRIES,
2029 } if path == &lock_path
2030 ));
2031
2032 drop(first);
2033
2034 let second = StartLock::acquire(&path)
2035 .await
2036 .expect("released advisory lock should be reclaimable");
2037 drop(second);
2038 }
2039
2040 #[tokio::test]
2041 async fn bind_conflict_on_fixed_port_fails_loud_without_reselecting() {
2042 let (_dir, path) = temp_connection_file_path("bind-conflict");
2043 let occupied = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
2044 let occupied_port = occupied.local_addr().unwrap().port();
2045
2046 let err = ensure_singleton(&path, occupied_port).await.unwrap_err();
2047 assert!(matches!(
2048 err,
2049 BootstrapError::Bind { ref source, .. } if source.kind() == io::ErrorKind::AddrInUse
2050 ));
2051 assert!(err.to_string().contains("set the port in config"));
2052
2053 drop(occupied);
2054 }
2055
2056 #[tokio::test]
2057 async fn key_rotation_republishes_new_material_and_old_file_fails_auth() {
2058 let (_dir, path) = temp_connection_file_path("key-rotation");
2059 let first = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2060 let old_info = first.connection_info.clone();
2061 let fixed_port = old_info.endpoints[0].port;
2062 drop(first.listeners);
2063
2064 let second = expect_bound(ensure_singleton(&path, fixed_port).await.unwrap());
2065 let new_info = second.connection_info.clone();
2066 assert_ne!(old_info.key, new_info.key);
2067 assert_ne!(old_info.daemon_id, new_info.daemon_id);
2068 let server = start_server(second);
2069
2070 let mut old_stream = connect_from_info(&old_info).await.unwrap();
2071 let old_auth = authenticate_client(&mut old_stream, &old_info, PROBE_AUTH_DEADLINE).await;
2072 assert!(
2073 old_auth.is_err(),
2074 "old key must not authenticate after restart"
2075 );
2076
2077 let reread = connection_file::read(&path).unwrap();
2078 let mut new_stream = connect_from_info(&reread).await.unwrap();
2079 authenticate_client(&mut new_stream, &reread, PROBE_AUTH_DEADLINE)
2080 .await
2081 .unwrap();
2082
2083 server.abort();
2084 let _ = server.await;
2085 }
2086
2087 #[cfg(unix)]
2088 #[tokio::test]
2089 async fn published_connection_file_permissions_are_owner_only() {
2090 let (_dir, path) = temp_connection_file_path("permissions");
2091 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2092
2093 let mode = fs::metadata(&path).unwrap().permissions().mode() & 0o777;
2094 assert_eq!(mode, 0o600);
2095
2096 drop(bound.listeners);
2097 }
2098}