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 daemon_config_path: Option<PathBuf>,
107 configured_port: Option<u16>,
108 route_bind_relay_default_ms: Option<u64>,
112 reserved_capabilities: BTreeMap<String, String>,
113 watchdog_config: DaemonSelfWatchdogConfig,
114 connection_file_source: ConnectionFileSource,
115 cgroup_placement: CgroupPlacementConfig,
118 capture_logs_dir: Option<PathBuf>,
122 terminal_journal_path: Option<PathBuf>,
130 machine_id_path: Option<PathBuf>,
136 live_children_path: Option<PathBuf>,
144}
145
146impl BootstrapConfig {
147 pub fn new(connection_file_path: impl Into<PathBuf>, port: u16) -> Self {
148 Self {
149 connection_file_path: connection_file_path.into(),
150 port,
151 daemon_ver: DAEMON_VERSION.to_owned(),
152 configured_modules: Vec::new(),
153 storage_config: None,
154 admission_facts: AdmissionFactsConfig::default(),
155 daemon_config_path: None,
156 configured_port: None,
157 route_bind_relay_default_ms: None,
158 reserved_capabilities: BTreeMap::new(),
159 watchdog_config: DaemonSelfWatchdogConfig::default(),
160 connection_file_source: ConnectionFileSource::Explicit,
161 cgroup_placement: CgroupPlacementConfig::default(),
162 capture_logs_dir: None,
163 terminal_journal_path: None,
164 machine_id_path: None,
165 live_children_path: None,
166 }
167 }
168
169 pub fn with_live_children_record(mut self, path: impl Into<PathBuf>) -> Self {
173 self.live_children_path = Some(path.into());
174 self
175 }
176
177 pub fn with_machine_id_path(mut self, path: impl Into<PathBuf>) -> Self {
181 self.machine_id_path = Some(path.into());
182 self
183 }
184
185 pub fn with_cgroup_placement(mut self, placement: CgroupPlacementConfig) -> Self {
187 self.cgroup_placement = placement;
188 self
189 }
190
191 pub fn with_capture_logs_dir(mut self, dir: impl Into<PathBuf>) -> Self {
195 self.capture_logs_dir = Some(dir.into());
196 self
197 }
198
199 pub fn with_terminal_journal_path(mut self, path: impl Into<PathBuf>) -> Self {
202 self.terminal_journal_path = Some(path.into());
203 self
204 }
205
206 pub fn from_env() -> Result<Self, BootstrapError> {
207 Self::from_env_with_daemon_config_path(daemon_config::default_config_path())
208 }
209
210 pub fn from_env_for_daemon_binary() -> Result<Self, BootstrapError> {
219 let run_dir = daemon_config::daemon_run_dir().map_err(BootstrapError::RunDir)?;
220 let machine_id_path =
221 crate::machine_id::default_machine_id_path().map_err(BootstrapError::MachineId)?;
222 Ok(Self::from_env()?
223 .with_capture_logs_dir(run_dir.join("logs"))
224 .with_terminal_journal_path(run_dir.join("terminals.jsonl"))
225 .with_live_children_record(crate::live_children::record_path(&run_dir))
226 .with_machine_id_path(machine_id_path))
227 }
228
229 pub fn from_env_with_daemon_config_path(
230 daemon_config_path: impl AsRef<Path>,
231 ) -> Result<Self, BootstrapError> {
232 let daemon_config_path = daemon_config_path.as_ref().to_path_buf();
233 let daemon_config =
234 daemon_config::load(&daemon_config_path).map_err(BootstrapError::DaemonConfig)?;
235 let config_port = daemon_config.as_ref().and_then(|config| config.port);
236 let storage_config = daemon_config
237 .as_ref()
238 .and_then(|config| config.storage.clone());
239 let admission_facts_carrier_module_id = daemon_config
240 .as_ref()
241 .and_then(|config| config.admission_facts_carrier_module_id.clone());
242 let admission_facts_targets = daemon_config
243 .as_ref()
244 .and_then(|config| config.admission_facts_targets.clone());
245 let route_bind_relay_default_ms = daemon_config
246 .as_ref()
247 .and_then(|config| config.route_bind_relay_timeout_ms);
248 let reserved_capabilities = daemon_config
249 .as_ref()
250 .map(|config| config.reserved_capabilities.clone())
251 .unwrap_or_default();
252 let configured_modules = daemon_config
253 .map(|config| config.modules)
254 .unwrap_or_default();
255
256 let port = match env::var(SUBC_PORT_ENV) {
257 Ok(raw) if !raw.trim().is_empty() => {
258 let port = raw
259 .parse::<u16>()
260 .map_err(|source| BootstrapError::InvalidPort { raw, source })?;
261 if let Some(config_port) = config_port {
262 info!(
263 env = SUBC_PORT_ENV,
264 env_port = port,
265 config_port,
266 "SUBC_PORT overrides daemon config port"
267 );
268 }
269 port
270 }
271 Ok(_) | Err(_) => config_port.unwrap_or(DEFAULT_SUBC_PORT),
272 };
273
274 let (connection_file_path, connection_file_source) =
275 connection_file_path_with_source(non_empty_os_var("XDG_RUNTIME_DIR"));
276 Ok(Self::new(connection_file_path, port)
277 .with_configured_modules(configured_modules)
278 .with_storage_config(storage_config)
279 .with_admission_facts_config(admission_facts_carrier_module_id, admission_facts_targets)
280 .with_route_bind_relay_default_ms(route_bind_relay_default_ms)
281 .with_reserved_capabilities(reserved_capabilities)
282 .with_daemon_config_source(daemon_config_path, config_port)
283 .with_connection_file_source(connection_file_source))
284 }
285
286 pub fn with_daemon_config_path(
287 self,
288 daemon_config_path: impl AsRef<Path>,
289 ) -> Result<Self, BootstrapError> {
290 let daemon_config_path = daemon_config_path.as_ref().to_path_buf();
291 let daemon_config =
292 daemon_config::load(&daemon_config_path).map_err(BootstrapError::DaemonConfig)?;
293 let configured_port = daemon_config.as_ref().and_then(|config| config.port);
294 let storage_config = daemon_config
295 .as_ref()
296 .and_then(|config| config.storage.clone());
297 let admission_facts_carrier_module_id = daemon_config
298 .as_ref()
299 .and_then(|config| config.admission_facts_carrier_module_id.clone());
300 let admission_facts_targets = daemon_config
301 .as_ref()
302 .and_then(|config| config.admission_facts_targets.clone());
303 let route_bind_relay_default_ms = daemon_config
304 .as_ref()
305 .and_then(|config| config.route_bind_relay_timeout_ms);
306 let reserved_capabilities = daemon_config
307 .as_ref()
308 .map(|config| config.reserved_capabilities.clone())
309 .unwrap_or_default();
310 let configured_modules = daemon_config
311 .map(|config| config.modules)
312 .unwrap_or_default();
313 Ok(self
314 .with_configured_modules(configured_modules)
315 .with_storage_config(storage_config)
316 .with_admission_facts_config(admission_facts_carrier_module_id, admission_facts_targets)
317 .with_route_bind_relay_default_ms(route_bind_relay_default_ms)
318 .with_reserved_capabilities(reserved_capabilities)
319 .with_daemon_config_source(daemon_config_path, configured_port))
320 }
321
322 pub fn with_configured_modules(
323 mut self,
324 modules: impl IntoIterator<Item = ConfiguredModule>,
325 ) -> Self {
326 self.configured_modules = modules.into_iter().collect();
327 self.configured_modules
328 .sort_by(|left, right| left.module_id.cmp(&right.module_id));
329 self
330 }
331
332 pub fn with_storage_config(
333 mut self,
334 storage_config: Option<daemon_config::StorageConfig>,
335 ) -> Self {
336 self.storage_config = storage_config;
337 self
338 }
339
340 pub fn with_admission_facts_config(
341 mut self,
342 carrier_module_id: Option<String>,
343 targets: Option<Vec<String>>,
344 ) -> Self {
345 self.admission_facts = AdmissionFactsConfig {
346 carrier_module_id,
347 targets,
348 };
349 self
350 }
351
352 pub fn with_route_bind_relay_default_ms(mut self, ms: Option<u64>) -> Self {
357 self.route_bind_relay_default_ms = ms;
358 self
359 }
360
361 pub fn with_reserved_capabilities(
362 mut self,
363 reserved_capabilities: BTreeMap<String, String>,
364 ) -> Self {
365 self.reserved_capabilities = reserved_capabilities;
366 self
367 }
368
369 fn with_daemon_config_source(
370 mut self,
371 daemon_config_path: PathBuf,
372 configured_port: Option<u16>,
373 ) -> Self {
374 self.daemon_config_path = Some(daemon_config_path);
375 self.configured_port = configured_port;
376 self
377 }
378
379 fn with_connection_file_source(mut self, source: ConnectionFileSource) -> Self {
380 self.connection_file_source = source;
381 self
382 }
383
384 pub fn with_watchdog_config(mut self, watchdog_config: DaemonSelfWatchdogConfig) -> Self {
385 self.watchdog_config = watchdog_config;
386 self
387 }
388}
389
390#[derive(Debug)]
392pub enum Outcome {
393 AlreadyRunning,
395 Bound(BoundDaemon),
398}
399
400#[derive(Debug)]
401pub struct BoundDaemon {
402 pub listeners: Vec<TcpListener>,
403 pub connection_info: ConnectionInfo,
404 pub connection_file_path: PathBuf,
405 pub connection_file_source: ConnectionFileSource,
406 pub machine_id: Option<crate::machine_id::MachineId>,
409}
410
411pub fn connection_file_path() -> PathBuf {
418 connection_file_path_with_source(non_empty_os_var("XDG_RUNTIME_DIR")).0
419}
420
421fn connection_file_path_with_source(
422 runtime_dir: Option<OsString>,
423) -> (PathBuf, ConnectionFileSource) {
424 if let Some(runtime_dir) = runtime_dir.filter(|value| !value.is_empty()) {
425 return (
426 PathBuf::from(runtime_dir).join(CONNECTION_FILE_NAME),
427 ConnectionFileSource::XdgRuntimeDir,
428 );
429 }
430
431 (
432 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token())),
433 ConnectionFileSource::TempDirFallback,
434 )
435}
436
437pub async fn run() -> Result<(), BootstrapError> {
443 run_with_config(
446 BootstrapConfig::from_env_for_daemon_binary()?
447 .with_cgroup_placement(CgroupPlacementConfig::Current),
448 )
449 .await
450}
451
452pub async fn run_with_config(config: BootstrapConfig) -> Result<(), BootstrapError> {
476 let configured_modules = config.configured_modules.clone();
477 let storage_config = config.storage_config.clone();
478 let admission_facts = config.admission_facts.clone();
479 let daemon_config_path = config.daemon_config_path.clone();
480 let configured_port = config.configured_port;
481 let route_bind_relay_default_ms = config.route_bind_relay_default_ms;
482 let reserved_capabilities = config.reserved_capabilities.clone();
483 let watchdog_config = config.watchdog_config.clone();
484 let cgroup_placement_config = config.cgroup_placement.clone();
485 let capture_logs_dir = config.capture_logs_dir.clone();
486 let terminal_journal_path = config.terminal_journal_path.clone();
487 let live_children_path = config.live_children_path.clone();
488 match ensure_singleton_with_config(config).await? {
489 Outcome::AlreadyRunning => {
490 info!("subc daemon already running");
491 Ok(())
492 }
493 Outcome::Bound(bound) => {
494 #[cfg(target_os = "linux")]
495 let cgroup_placement = prepare_cgroup_placement(&cgroup_placement_config);
496 #[cfg(not(target_os = "linux"))]
497 let _ = cgroup_placement_config;
498 serve_bound_daemon(
499 bound,
500 configured_modules,
501 storage_config,
502 admission_facts,
503 daemon_config_path,
504 configured_port,
505 route_bind_relay_default_ms,
506 reserved_capabilities,
507 watchdog_config,
508 capture_logs_dir,
509 terminal_journal_path,
510 live_children_path,
511 #[cfg(target_os = "linux")]
512 cgroup_placement,
513 )
514 .await
515 }
516 }
517}
518
519#[cfg(target_os = "linux")]
520fn prepare_cgroup_placement(config: &CgroupPlacementConfig) -> Option<subc_cgroup::Placement> {
521 let result = match config {
522 CgroupPlacementConfig::Disabled => return None,
523 CgroupPlacementConfig::Current => subc_cgroup::prepare_current(),
524 CgroupPlacementConfig::Root(root) => subc_cgroup::prepare_at(root),
525 };
526
527 match result {
528 Ok(Some(placement)) => Some(placement),
529 Ok(None) => {
530 warn!(
531 placement = ?config,
532 "module cgroup placement is disabled: configured cgroup root is not delegated"
533 );
534 None
535 }
536 Err(error) => {
537 warn!(
538 placement = ?config,
539 error = %error,
540 "module cgroup placement is disabled by an unexpected cgroup probe error"
541 );
542 None
543 }
544 }
545}
546
547#[cfg(unix)]
554const NOFILE_TARGET: u64 = 65536;
555
556#[cfg(unix)]
560fn raise_nofile_limit() {
561 match rlimit::Resource::NOFILE.get() {
562 Ok((soft, hard)) => {
563 if soft >= NOFILE_TARGET {
564 return;
565 }
566 let target = NOFILE_TARGET.min(hard);
567 match rlimit::Resource::NOFILE.set(target, hard) {
568 Ok(()) => info!(
569 previous_soft = soft,
570 new_soft = target,
571 hard,
572 "raised open-file soft limit for daemon and module children"
573 ),
574 Err(err) => warn!(
575 soft,
576 hard,
577 error = %err,
578 "could not raise open-file soft limit; multi-root modules may exhaust descriptors"
579 ),
580 }
581 }
582 Err(err) => warn!(error = %err, "could not read open-file limit"),
583 }
584}
585
586#[cfg(windows)]
593fn raise_nofile_limit() {
594 const MAXSTDIO_TARGET: u32 = 8192;
595 let current = rlimit::getmaxstdio();
596 if current >= MAXSTDIO_TARGET {
597 return;
598 }
599 match rlimit::setmaxstdio(MAXSTDIO_TARGET) {
600 Ok(new_max) => info!(
601 previous = current,
602 new_max, "raised CRT stdio-stream limit for daemon"
603 ),
604 Err(err) => warn!(
605 current,
606 error = %err,
607 "could not raise CRT stdio-stream limit"
608 ),
609 }
610}
611
612#[cfg(not(any(unix, windows)))]
613fn raise_nofile_limit() {}
614
615pub async fn run_with_daemon_config_path(
616 config: BootstrapConfig,
617 daemon_config_path: impl AsRef<Path>,
618) -> Result<(), BootstrapError> {
619 run_with_config(config.with_daemon_config_path(daemon_config_path)?).await
620}
621
622#[allow(clippy::too_many_arguments)]
623async fn serve_bound_daemon(
624 bound: BoundDaemon,
625 configured_modules: Vec<ConfiguredModule>,
626 storage_config: Option<daemon_config::StorageConfig>,
627 admission_facts: AdmissionFactsConfig,
628 daemon_config_path: Option<PathBuf>,
629 configured_port: Option<u16>,
630 route_bind_relay_default_ms: Option<u64>,
631 reserved_capabilities: BTreeMap<String, String>,
632 watchdog_config: DaemonSelfWatchdogConfig,
633 capture_logs_dir: Option<PathBuf>,
634 terminal_journal_path: Option<PathBuf>,
635 live_children_path: Option<PathBuf>,
636 #[cfg(target_os = "linux")] cgroup_placement: Option<subc_cgroup::Placement>,
637) -> Result<(), BootstrapError> {
638 #[cfg(unix)]
639 let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
640 .map_err(BootstrapError::Signal)?;
641 raise_nofile_limit();
644
645 info!(
646 connection_file = %bound.connection_file_path.display(),
647 connection_file_source = %bound.connection_file_source,
648 connection_file_source_reason = bound.connection_file_source.reason(),
649 endpoints = ?bound.connection_info.endpoints,
650 configured_modules = configured_modules.len(),
651 machine_id = bound.machine_id.as_ref().map(|id| id.as_str()).unwrap_or("none"),
652 "subc daemon starting"
653 );
654
655 if let Some(path) = &live_children_path {
662 crate::live_children::sweep_orphans(
663 path,
664 &crate::live_children::AdoptedPids::none(),
665 crate::live_children::SweepBounds::default(),
666 )
667 .await;
668 }
669
670 let registry = Arc::new(Registry::default());
671 let process_liveness = Arc::new(SupervisorProcessLiveness::new());
672 let supervisor_handle = SupervisorHandle::new();
673 let connected_clients = ConnectedClients::new();
674 let forwarding = Arc::new(ForwardingTable::default());
675 let daemon_incarnation = format!(
676 "{:032x}",
677 u128::from_be_bytes(bound.connection_info.daemon_id)
678 );
679 let supervisor = Supervisor::new(Arc::clone(®istry), RestartPolicy::default())
680 .with_process_liveness(process_liveness.clone())
681 .with_forwarding(Arc::clone(&forwarding))
682 .with_handle(supervisor_handle.clone())
683 .with_connection_file_path(bound.connection_file_path.clone())
684 .with_daemon_incarnation(daemon_incarnation.clone());
685 let supervisor = match terminal_journal_path {
722 Some(path) => supervisor.with_terminal_journal(path, daemon_incarnation),
723 None => supervisor,
724 };
725 let supervisor = match capture_logs_dir {
726 Some(dir) => supervisor.with_capture_logs_dir(dir),
727 None => supervisor,
728 };
729 let supervisor = match live_children_path {
730 Some(path) => supervisor.with_live_children_record(path),
731 None => supervisor,
732 };
733 #[cfg(target_os = "linux")]
734 let supervisor = supervisor.with_cgroup_placement(cgroup_placement);
735 let route_bind_relay_timeouts = configured_modules
741 .iter()
742 .filter_map(|module| {
743 module
744 .route_bind_relay_timeout_ms
745 .map(|ms| (module.module_id.clone(), Duration::from_millis(ms)))
746 })
747 .collect::<std::collections::BTreeMap<_, _>>();
748 let control_start_clock = crate::clock::StartClock::capture();
749 let mut control = ControlHandler::with_forwarding(Arc::clone(®istry), forwarding)
750 .with_process_liveness(process_liveness)
751 .with_supervisor(supervisor_handle)
752 .with_connected_clients(connected_clients.clone())
753 .with_storage_config(storage_config)
754 .with_machine_id(bound.machine_id.clone())
755 .with_admission_facts_config(admission_facts.carrier_module_id, admission_facts.targets)
756 .with_route_bind_relay_timeouts(route_bind_relay_timeouts)
757 .with_daemon_provenance(
758 bound.connection_info.pid,
759 control_start_clock.started_at_ms(),
760 std::env::current_exe().ok(),
761 normalized_build_provenance(env!("SUBC_BUILD_GIT_SHA")),
762 normalized_build_provenance(env!("SUBC_BUILD_LOCK_DIGEST")),
763 )
764 .with_daemon_start_clock(control_start_clock)
765 .with_capability_config(
766 configured_modules
767 .iter()
768 .map(|module| (module.module_id.clone(), module.enabled)),
769 reserved_capabilities,
770 );
771 if let Some(ms) = route_bind_relay_default_ms {
772 control = control.with_route_bind_relay_timeout(Duration::from_millis(ms));
775 }
776 if let Some(config_path) = daemon_config_path {
777 control = control.with_supervisor_rescan(supervisor.clone(), config_path, configured_port);
778 }
779 let control = Arc::new(control);
780 let router = Arc::new(Router::with_control_handler(Arc::clone(&control)));
781 let auth = ServerAuth::new(
782 bound.connection_info.key.clone(),
783 bound.connection_info.daemon_id,
784 bound.connection_info.daemon_ver.clone(),
785 )
786 .with_connected_clients(connected_clients);
787
788 let mut serve_task =
789 AbortOnDrop::new(tokio::spawn(serve_listeners(bound.listeners, router, auth)));
790 tokio::task::yield_now().await;
791 let _clock_step_task = AbortOnDrop::new(crate::watchdog::spawn_clock_step_monitor());
792 #[cfg(target_os = "linux")]
794 let _kill_mode_check = AbortOnDrop::new(tokio::spawn(
795 crate::systemd_kill_mode::warn_if_kill_mode_defeats_ordered_shutdown(),
796 ));
797 let _watchdog_task = AbortOnDrop::new(
798 DaemonSelfWatchdog::new(
799 bound.connection_info.clone(),
800 bound.connection_file_path.clone(),
801 )
802 .with_config(watchdog_config)
803 .spawn(),
804 );
805
806 for configured in configured_modules {
807 let enabled = configured.enabled;
808 let health = configured.health;
809 let module_id = configured.module_id.clone();
810 match supervisor.supervise_configured_with_health(
811 configured.module_spec(),
812 enabled,
813 health,
814 configured.drain_timeout_ms,
815 configured.restart,
816 ) {
817 Ok(_) => {
818 let default_threshold = HealthConfig::default().failure_threshold;
827 if enabled && health.failure_threshold > default_threshold {
828 warn!(
829 module_id = %module_id,
830 failure_threshold = health.failure_threshold,
831 default_threshold,
832 tolerance_secs = health.cadence.as_secs() * u64::from(health.failure_threshold),
833 "health failure threshold is relaxed above the default; a wedged module stays unflagged for longer"
834 );
835 }
836 info!(module_id = %module_id, enabled, "configured module supervised");
837 }
838 Err(err) => {
839 error!(module_id = %module_id, error = %err, "failed to supervise configured module; continuing daemon startup");
840 }
841 }
842 }
843
844 control.refresh_capability_requirements();
845 Arc::clone(&control).spawn_capability_deadline_loop();
846
847 #[cfg(unix)]
848 {
849 tokio::select! {
850 result = serve_task.join() => {
851 return result.map_err(BootstrapError::ServeJoin)?.map_err(BootstrapError::Serve);
852 }
853 _ = terminate.recv() => {}
854 }
855 supervisor.begin_daemon_shutdown();
860 drop(_watchdog_task);
864 drop(serve_task);
867 let escalated = tokio::select! {
868 biased;
869 _ = terminate.recv() => {
870 info!("second SIGTERM: abandoning daemon shutdown wait");
871 true
872 }
873 result = supervisor.drain_for_daemon_shutdown() => {
874 if let Err(error) = result {
875 warn!(%error, "daemon shutdown drain failed; exiting anyway");
876 }
877 false
878 }
879 };
880 supervisor
887 .end_children_for_daemon_shutdown(escalated, async {
888 terminate.recv().await;
889 })
890 .await;
891 Ok(())
892 }
893 #[cfg(not(unix))]
894 serve_task
895 .join()
896 .await
897 .map_err(BootstrapError::ServeJoin)?
898 .map_err(BootstrapError::Serve)
899}
900
901fn normalized_build_provenance(value: &str) -> Option<String> {
902 match value.trim() {
903 "" | "unavailable" => None,
904 value => Some(value.to_string()),
905 }
906}
907
908pub async fn ensure_singleton(
916 connection_file_path: impl AsRef<Path>,
917 port: u16,
918) -> Result<Outcome, BootstrapError> {
919 ensure_singleton_with_config(BootstrapConfig::new(connection_file_path.as_ref(), port)).await
920}
921
922pub async fn ensure_singleton_with_config(
923 config: BootstrapConfig,
924) -> Result<Outcome, BootstrapError> {
925 let path = config.connection_file_path;
926
927 if matches!(probe_existing(&path).await?, Probe::Live) {
928 return Ok(Outcome::AlreadyRunning);
929 }
930
931 let _lock = StartLock::acquire(&path).await?;
932
933 if matches!(probe_existing(&path).await?, Probe::Live) {
937 return Ok(Outcome::AlreadyRunning);
938 }
939
940 remove_stale_connection_file_if_present(&path)?;
941
942 let machine_id = config
946 .machine_id_path
947 .as_deref()
948 .map(crate::machine_id::load_or_mint)
949 .transpose()
950 .map_err(BootstrapError::MachineId)?;
951
952 let (listeners, endpoints) = bind_loopback(config.port).await?;
953 let connection_info = ConnectionInfo {
954 schema: SCHEMA_VERSION,
955 wire_version: Some(PROTOCOL_VERSION),
956 endpoints,
957 key: generate_key().map_err(BootstrapError::GenerateConnectionFile)?,
958 daemon_id: generate_daemon_id().map_err(BootstrapError::GenerateConnectionFile)?,
959 pid: process::id(),
960 daemon_ver: config.daemon_ver,
961 };
962
963 if let Err(source) = write_atomic(&path, &connection_info) {
964 drop(listeners);
965 return Err(BootstrapError::ConnectionFileWrite { path, source });
966 }
967
968 Ok(Outcome::Bound(BoundDaemon {
969 listeners,
970 connection_info,
971 connection_file_path: path,
972 connection_file_source: config.connection_file_source,
973 machine_id,
974 }))
975}
976
977#[derive(Debug, Clone, Copy, PartialEq, Eq)]
978enum Probe {
979 Live,
980 StaleOrAbsent,
981}
982
983async fn probe_existing(path: &Path) -> Result<Probe, BootstrapError> {
984 let info = match connection_file::read(path) {
985 Ok(info) => info,
986 Err(source) if is_absent_or_stale_connection_file(&source) => {
987 return Ok(Probe::StaleOrAbsent)
988 }
989 Err(source) => {
990 return Err(BootstrapError::ConnectionFileRead {
991 path: path.to_path_buf(),
992 source,
993 })
994 }
995 };
996
997 for endpoint in &info.endpoints {
998 if matches!(probe_endpoint(&info, endpoint).await, Probe::Live) {
999 return Ok(Probe::Live);
1000 }
1001 }
1002
1003 Ok(Probe::StaleOrAbsent)
1004}
1005
1006async fn probe_endpoint(info: &ConnectionInfo, endpoint: &Endpoint) -> Probe {
1007 let Ok(ip) = endpoint.host.parse::<IpAddr>() else {
1008 return Probe::StaleOrAbsent;
1009 };
1010 if !ip.is_loopback() {
1011 return Probe::StaleOrAbsent;
1012 }
1013 let addr = SocketAddr::new(ip, endpoint.port);
1014
1015 let mut stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(addr)).await {
1016 Ok(Ok(stream)) => stream,
1017 Ok(Err(_)) | Err(_) => return Probe::StaleOrAbsent,
1018 };
1019
1020 match authenticate_client(&mut stream, info, PROBE_AUTH_DEADLINE).await {
1021 Ok(()) => Probe::Live,
1022 Err(AuthError::DaemonIdMismatch)
1023 | Err(AuthError::InvalidServerProof)
1024 | Err(AuthError::UnexpectedEof { .. })
1025 | Err(AuthError::Timeout { .. })
1026 | Err(AuthError::JsonEncode { .. })
1027 | Err(AuthError::JsonDecode { .. })
1028 | Err(AuthError::Io { .. })
1029 | Err(AuthError::MessageTooLarge { .. })
1030 | Err(AuthError::KeyTooShort { .. })
1031 | Err(AuthError::Random(_))
1032 | Err(AuthError::InvalidClientAuth) => Probe::StaleOrAbsent,
1033 }
1034}
1035
1036fn is_absent_or_stale_connection_file(err: &ConnectionFileError) -> bool {
1037 match err {
1038 ConnectionFileError::Io { source, .. } if source.kind() == io::ErrorKind::NotFound => true,
1039 ConnectionFileError::JsonRead { .. }
1040 | ConnectionFileError::UnsupportedSchema { .. }
1041 | ConnectionFileError::Invalid { .. }
1042 | ConnectionFileError::KeyTooShort { .. }
1043 | ConnectionFileError::InsecurePermissions { .. } => true,
1047 ConnectionFileError::MissingParent { .. }
1048 | ConnectionFileError::MissingFileName { .. }
1049 | ConnectionFileError::InsecureParentDirectory { .. }
1053 | ConnectionFileError::Io { .. }
1054 | ConnectionFileError::JsonWrite { .. }
1055 | ConnectionFileError::Random(_)
1056 | ConnectionFileError::WireVersionMismatch { .. } => false,
1059 }
1060}
1061
1062fn remove_stale_connection_file_if_present(path: &Path) -> Result<(), BootstrapError> {
1063 match fs::remove_file(path) {
1064 Ok(()) => Ok(()),
1065 Err(err) if err.kind() == io::ErrorKind::NotFound => Ok(()),
1066 Err(source) => Err(BootstrapError::RemoveStale {
1067 path: path.to_path_buf(),
1068 source,
1069 }),
1070 }
1071}
1072
1073async fn bind_loopback(port: u16) -> Result<(Vec<TcpListener>, Vec<Endpoint>), BootstrapError> {
1074 let v4_host = Ipv4Addr::LOCALHOST;
1075 let v4 = TcpListener::bind((v4_host, port))
1076 .await
1077 .map_err(|source| BootstrapError::Bind {
1078 host: v4_host.to_string(),
1079 port,
1080 source,
1081 })?;
1082 let actual_port = v4
1083 .local_addr()
1084 .map_err(|source| BootstrapError::LocalAddr {
1085 host: v4_host.to_string(),
1086 source,
1087 })?
1088 .port();
1089
1090 let mut listeners = vec![v4];
1091 let mut endpoints = vec![Endpoint {
1092 host: v4_host.to_string(),
1093 port: actual_port,
1094 }];
1095
1096 let v6_host = Ipv6Addr::LOCALHOST;
1097 match TcpListener::bind((v6_host, actual_port)).await {
1098 Ok(v6) => {
1099 listeners.push(v6);
1100 endpoints.push(Endpoint {
1101 host: v6_host.to_string(),
1102 port: actual_port,
1103 });
1104 }
1105 Err(err) if ipv6_loopback_unavailable(&err) => {
1106 warn!(
1107 port = actual_port,
1108 error = %err,
1109 "IPv6 loopback unavailable; serving only IPv4 loopback"
1110 );
1111 }
1112 Err(source) => {
1113 drop(listeners);
1114 return Err(BootstrapError::Bind {
1115 host: v6_host.to_string(),
1116 port: actual_port,
1117 source,
1118 });
1119 }
1120 }
1121
1122 Ok((listeners, endpoints))
1123}
1124
1125fn ipv6_loopback_unavailable(err: &io::Error) -> bool {
1126 matches!(
1127 err.kind(),
1128 io::ErrorKind::AddrNotAvailable | io::ErrorKind::Unsupported
1129 ) || matches!(err.raw_os_error(), Some(47) | Some(49) | Some(97))
1130}
1131
1132struct AbortOnDrop<T> {
1133 handle: JoinHandle<T>,
1134}
1135
1136impl<T> AbortOnDrop<T> {
1137 fn new(handle: JoinHandle<T>) -> Self {
1138 Self { handle }
1139 }
1140
1141 async fn join(&mut self) -> Result<T, JoinError> {
1142 (&mut self.handle).await
1143 }
1144}
1145
1146impl<T> Drop for AbortOnDrop<T> {
1147 fn drop(&mut self) {
1148 if !self.handle.is_finished() {
1149 self.handle.abort();
1150 }
1151 }
1152}
1153
1154struct StartLock {
1155 _file: fs::File,
1158}
1159
1160impl StartLock {
1161 async fn acquire(connection_file_path: &Path) -> Result<Self, BootstrapError> {
1162 let path = start_lock_path(connection_file_path);
1163 for _ in 0..START_LOCK_RETRIES {
1164 let file = match open_owner_only_lock(&path) {
1165 Ok(file) => file,
1166 Err(source) => return Err(BootstrapError::StartLockCreate { path, source }),
1167 };
1168 match FileExt::try_lock(&file) {
1169 Ok(()) => return Ok(Self { _file: file }),
1170 Err(TryLockError::WouldBlock) => sleep(START_LOCK_RETRY_DELAY).await,
1171 Err(TryLockError::Error(source)) => {
1172 return Err(BootstrapError::StartLockCreate { path, source });
1173 }
1174 }
1175 }
1176
1177 Err(BootstrapError::StartLockBusy {
1178 path,
1179 attempts: START_LOCK_RETRIES,
1180 })
1181 }
1182}
1183
1184fn open_owner_only_lock(path: &Path) -> io::Result<fs::File> {
1185 let mut options = fs::OpenOptions::new();
1186 options.read(true).write(true).create(true);
1187 #[cfg(unix)]
1188 {
1189 use std::os::unix::fs::OpenOptionsExt;
1190 options.mode(0o600);
1191 }
1192 options.open(path)
1193}
1194
1195fn start_lock_path(connection_file_path: &Path) -> PathBuf {
1196 let file_name = connection_file_path
1197 .file_name()
1198 .map(|name| name.to_string_lossy())
1199 .unwrap_or_else(|| CONNECTION_FILE_NAME.into());
1200 let lock_name = format!("{file_name}.start-lock");
1201 connection_file_path
1202 .parent()
1203 .filter(|parent| !parent.as_os_str().is_empty())
1204 .unwrap_or_else(|| Path::new("."))
1205 .join(lock_name)
1206}
1207
1208fn non_empty_os_var(key: &str) -> Option<OsString> {
1209 let value = env::var_os(key)?;
1210 if value.is_empty() {
1211 None
1212 } else {
1213 Some(value)
1214 }
1215}
1216
1217#[derive(Debug)]
1220pub enum BootstrapError {
1221 #[cfg(unix)]
1222 Signal(io::Error),
1223 InvalidPort {
1224 raw: String,
1225 source: std::num::ParseIntError,
1226 },
1227 ConnectionFileRead {
1228 path: PathBuf,
1229 source: ConnectionFileError,
1230 },
1231 ConnectionFileWrite {
1232 path: PathBuf,
1233 source: ConnectionFileError,
1234 },
1235 GenerateConnectionFile(ConnectionFileError),
1236 StartLockCreate {
1237 path: PathBuf,
1238 source: io::Error,
1239 },
1240 StartLockBusy {
1241 path: PathBuf,
1242 attempts: usize,
1243 },
1244 RemoveStale {
1245 path: PathBuf,
1246 source: io::Error,
1247 },
1248 Bind {
1249 host: String,
1250 port: u16,
1251 source: io::Error,
1252 },
1253 LocalAddr {
1254 host: String,
1255 source: io::Error,
1256 },
1257 DaemonConfig(DaemonConfigError),
1258 MachineId(crate::machine_id::MachineIdFileError),
1261 RunDir(daemon_config::DaemonRunDirError),
1264 Serve(ServerError),
1265 ServeJoin(tokio::task::JoinError),
1266}
1267
1268impl fmt::Display for BootstrapError {
1269 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1270 match self {
1271 #[cfg(unix)]
1272 Self::Signal(error) => write!(f, "failed to register SIGTERM handler: {error}"),
1273 Self::InvalidPort { raw, source } => {
1274 write!(f, "invalid {SUBC_PORT_ENV} value '{raw}': {source}")
1275 }
1276 Self::ConnectionFileRead { path, source } => write!(
1277 f,
1278 "failed to read connection file {}: {source}",
1279 path.display()
1280 ),
1281 Self::ConnectionFileWrite { path, source } => write!(
1282 f,
1283 "failed to publish connection file {}: {source}",
1284 path.display()
1285 ),
1286 Self::GenerateConnectionFile(err) => {
1287 write!(f, "failed to generate connection-file auth material: {err}")
1288 }
1289 Self::StartLockCreate { path, source } => {
1290 write!(
1291 f,
1292 "failed to create start lock {}: {source}",
1293 path.display()
1294 )
1295 }
1296 Self::StartLockBusy { path, attempts } => write!(
1297 f,
1298 "start lock {} remained busy after {attempts} attempts",
1299 path.display()
1300 ),
1301 Self::RemoveStale { path, source } => write!(
1302 f,
1303 "failed to remove stale connection file {}: {source}",
1304 path.display()
1305 ),
1306 Self::Bind { host, port, source } if source.kind() == io::ErrorKind::AddrInUse => {
1307 write!(
1308 f,
1309 "port {port} in use on loopback {host}: {source}; set the port in config"
1310 )
1311 }
1312 Self::Bind { host, port, source } => {
1313 write!(f, "failed to bind loopback TCP {host}:{port}: {source}")
1314 }
1315 Self::LocalAddr { host, source } => {
1316 write!(f, "failed to read local address for {host}: {source}")
1317 }
1318 Self::DaemonConfig(err) => write!(f, "failed to load daemon config: {err}"),
1319 Self::MachineId(err) => write!(f, "refusing to start: {err}"),
1320 Self::RunDir(err) => write!(f, "refusing to start: {err}"),
1321 Self::Serve(err) => write!(f, "daemon server failed: {err}"),
1322 Self::ServeJoin(err) => write!(f, "daemon server task failed: {err}"),
1323 }
1324 }
1325}
1326
1327impl Error for BootstrapError {
1328 fn source(&self) -> Option<&(dyn Error + 'static)> {
1329 match self {
1330 #[cfg(unix)]
1331 Self::Signal(source) => Some(source),
1332 Self::InvalidPort { source, .. } => Some(source),
1333 Self::ConnectionFileRead { source, .. }
1334 | Self::ConnectionFileWrite { source, .. }
1335 | Self::GenerateConnectionFile(source) => Some(source),
1336 Self::StartLockCreate { source, .. }
1337 | Self::RemoveStale { source, .. }
1338 | Self::Bind { source, .. }
1339 | Self::LocalAddr { source, .. } => Some(source),
1340 Self::DaemonConfig(err) => Some(err),
1341 Self::MachineId(err) => Some(err),
1342 Self::RunDir(err) => Some(err),
1343 Self::Serve(err) => Some(err),
1344 Self::ServeJoin(err) => Some(err),
1345 Self::StartLockBusy { .. } => None,
1346 }
1347 }
1348}
1349
1350#[cfg(test)]
1351mod tests {
1352 use super::*;
1353 use crate::server::ServerAuth;
1354 #[cfg(target_os = "linux")]
1355 use std::collections::BTreeSet;
1356 use std::sync::Mutex;
1357 #[cfg(target_os = "linux")]
1358 use subc_control::ModuleProtocol;
1359 use subc_test_support::TestTempDir;
1360 use subc_transport::MIN_KEY_LEN;
1361 use tokio::io::AsyncReadExt;
1362 use tokio::task::JoinHandle;
1363
1364 #[cfg(unix)]
1365 use std::os::unix::fs::PermissionsExt;
1366
1367 static ENV_LOCK: Mutex<()> = Mutex::new(());
1368
1369 #[test]
1370 fn normalized_build_provenance_preserves_real_values() {
1371 assert_eq!(normalized_build_provenance("abc"), Some("abc".to_string()));
1372 }
1373
1374 #[test]
1375 fn normalized_build_provenance_omits_unavailable_and_empty_values() {
1376 assert_eq!(normalized_build_provenance("unavailable"), None);
1377 assert_eq!(normalized_build_provenance(""), None);
1378 }
1379
1380 #[cfg(target_os = "linux")]
1381 fn current_cgroup_path_for_test() -> io::Result<PathBuf> {
1382 let cgroups = fs::read_to_string("/proc/self/cgroup")?;
1383 let relative = cgroups
1384 .lines()
1385 .find_map(|line| line.strip_prefix("0::"))
1386 .ok_or_else(|| {
1387 io::Error::new(io::ErrorKind::Unsupported, "cgroup v2 is unavailable")
1388 })?;
1389 Ok(Path::new("/sys/fs/cgroup").join(relative.trim_start_matches('/')))
1390 }
1391
1392 #[cfg(target_os = "linux")]
1393 fn module_cgroup_directories() -> io::Result<BTreeSet<OsString>> {
1394 let modules = current_cgroup_path_for_test()?.join("subc-modules");
1395 let entries = match fs::read_dir(modules) {
1396 Ok(entries) => entries,
1397 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(BTreeSet::new()),
1398 Err(error) => return Err(error),
1399 };
1400 let mut directories = BTreeSet::new();
1401 for entry in entries {
1402 let entry = entry?;
1403 if entry.file_type()?.is_dir() {
1404 directories.insert(entry.file_name());
1405 }
1406 }
1407 Ok(directories)
1408 }
1409
1410 #[cfg(target_os = "linux")]
1411 async fn wait_for_path(path: &Path, task: &JoinHandle<Result<(), BootstrapError>>) {
1412 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1413 while !path.exists() && tokio::time::Instant::now() < deadline {
1414 assert!(
1415 !task.is_finished(),
1416 "daemon exited before creating {}",
1417 path.display()
1418 );
1419 sleep(Duration::from_millis(10)).await;
1420 }
1421 assert!(path.exists(), "daemon did not create {}", path.display());
1422 }
1423
1424 #[cfg(target_os = "linux")]
1429 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1430 async fn run_with_config_does_not_reconcile_the_ambient_cgroup_by_default() {
1431 let temp = unique_temp_dir("bootstrap-cgroup-default-disabled");
1432 let module_id = format!("cgroup-isolation-probe-{}", process::id());
1433 let before = module_cgroup_directories().expect("read ambient module cgroups before boot");
1434 assert!(
1435 !before.contains(&OsString::from(&module_id)),
1436 "isolation probe cgroup already exists before this daemon starts"
1437 );
1438 let capture = temp.join("logs").join(format!("{module_id}.stderr.log"));
1439 let module = ConfiguredModule {
1440 module_id,
1441 program: PathBuf::from("sh"),
1442 args: vec!["-c".to_string(), "sleep 30".to_string()],
1443 env: Vec::new(),
1444 log: None,
1445 enabled: true,
1446 reserved: false,
1447 reserved_prefixes: Vec::new(),
1448 protocol: ModuleProtocol::None,
1449 overlap: Default::default(),
1450 health: HealthConfig::default(),
1451 drain_timeout_ms: None,
1452 route_bind_relay_timeout_ms: None,
1453 restart: RestartPolicy::default(),
1454 };
1455 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1456 .with_configured_modules([module])
1457 .with_capture_logs_dir(temp.join("logs"))
1458 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1459 let task = tokio::spawn(run_with_config(config));
1460
1461 wait_for_path(&capture, &task).await;
1462 let after = module_cgroup_directories().expect("read ambient module cgroups after boot");
1463
1464 task.abort();
1465 assert!(task
1466 .await
1467 .expect_err("aborted daemon task must cancel")
1468 .is_cancelled());
1469 assert_eq!(
1470 before, after,
1471 "an in-process daemon must not create or reconcile ambient module cgroups"
1472 );
1473 }
1474
1475 #[cfg(target_os = "linux")]
1476 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1477 async fn explicit_cgroup_root_is_prepared_inside_the_fixture_tree() {
1478 let temp = unique_temp_dir("bootstrap-cgroup-explicit-root");
1479 let cgroup_root = temp.join("cgroup");
1480 fs::create_dir(&cgroup_root).expect("create scratch cgroup root");
1481 fs::write(cgroup_root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
1482 let modules = cgroup_root.join("subc-modules");
1483 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1484 .with_cgroup_placement(CgroupPlacementConfig::Root(cgroup_root))
1485 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1486 let task = tokio::spawn(run_with_config(config));
1487
1488 wait_for_path(&modules, &task).await;
1489
1490 task.abort();
1491 assert!(task
1492 .await
1493 .expect_err("aborted daemon task must cancel")
1494 .is_cancelled());
1495 assert!(
1496 fs::read_dir(&modules)
1497 .expect("read prepared modules directory")
1498 .next()
1499 .is_none(),
1500 "the delegation probe must clean up after itself"
1501 );
1502 }
1503
1504 struct EnvGuard {
1505 key: &'static str,
1506 previous: Option<OsString>,
1507 }
1508
1509 impl EnvGuard {
1510 fn set(key: &'static str, value: &Path) -> Self {
1511 let previous = env::var_os(key);
1512 env::set_var(key, value);
1513 Self { key, previous }
1514 }
1515
1516 fn set_str(key: &'static str, value: &str) -> Self {
1517 let previous = env::var_os(key);
1518 env::set_var(key, value);
1519 Self { key, previous }
1520 }
1521
1522 fn unset(key: &'static str) -> Self {
1523 let previous = env::var_os(key);
1524 env::remove_var(key);
1525 Self { key, previous }
1526 }
1527 }
1528
1529 impl Drop for EnvGuard {
1530 fn drop(&mut self) {
1531 match &self.previous {
1532 Some(value) => env::set_var(self.key, value),
1533 None => env::remove_var(self.key),
1534 }
1535 }
1536 }
1537
1538 fn unique_temp_dir(name: &str) -> TestTempDir {
1539 TestTempDir::new(name)
1540 }
1541
1542 fn temp_connection_file_path(name: &str) -> (TestTempDir, PathBuf) {
1543 let dir = unique_temp_dir(name);
1544 let path = dir.join("conn.json");
1545 (dir, path)
1546 }
1547
1548 fn auth_for(info: &ConnectionInfo) -> ServerAuth {
1549 ServerAuth::new(info.key.clone(), info.daemon_id, info.daemon_ver.clone())
1550 }
1551
1552 fn start_server(bound: BoundDaemon) -> JoinHandle<Result<(), ServerError>> {
1553 let auth = auth_for(&bound.connection_info);
1554 tokio::spawn(serve_listeners(
1555 bound.listeners,
1556 Arc::new(Router::with_default_self_handler()),
1557 auth,
1558 ))
1559 }
1560
1561 fn expect_bound(outcome: Outcome) -> BoundDaemon {
1562 match outcome {
1563 Outcome::Bound(bound) => bound,
1564 Outcome::AlreadyRunning => panic!("fresh connection file unexpectedly had a daemon"),
1565 }
1566 }
1567
1568 async fn connect_from_info(conn: &ConnectionInfo) -> io::Result<TcpStream> {
1569 let endpoint = conn
1570 .endpoints
1571 .first()
1572 .expect("test connection file should have an endpoint");
1573 let ip: IpAddr = endpoint.host.parse().unwrap();
1574 TcpStream::connect(SocketAddr::new(ip, endpoint.port)).await
1575 }
1576
1577 fn make_connection_info(port: u16) -> ConnectionInfo {
1578 ConnectionInfo {
1579 schema: SCHEMA_VERSION,
1580 wire_version: Some(PROTOCOL_VERSION),
1581 endpoints: vec![Endpoint {
1582 host: "127.0.0.1".to_owned(),
1583 port,
1584 }],
1585 key: generate_key().unwrap(),
1586 daemon_id: generate_daemon_id().unwrap(),
1587 pid: process::id(),
1588 daemon_ver: "test-subc".to_owned(),
1589 }
1590 }
1591
1592 fn write_raw_owner_only_connection_file(path: &Path, contents: &[u8]) {
1593 fs::write(path, contents).unwrap();
1594 #[cfg(unix)]
1595 fs::set_permissions(path, fs::Permissions::from_mode(0o600)).unwrap();
1596 }
1597
1598 fn assert_owner_only_connection_file(path: &Path) {
1599 #[cfg(unix)]
1602 {
1603 let mode = fs::metadata(path).unwrap().permissions().mode() & 0o777;
1604 assert_eq!(mode, 0o600);
1605 }
1606 #[cfg(not(unix))]
1607 let _ = path;
1608 }
1609
1610 #[test]
1611 fn connection_file_path_uses_xdg_runtime_dir_when_set() {
1612 let _env_lock = ENV_LOCK.lock().unwrap();
1613 let runtime_dir = unique_temp_dir("xdg-runtime");
1614 let _xdg = EnvGuard::set("XDG_RUNTIME_DIR", runtime_dir.path());
1615
1616 assert_eq!(
1617 connection_file_path(),
1618 runtime_dir.join(CONNECTION_FILE_NAME)
1619 );
1620 }
1621
1622 #[test]
1623 fn connection_file_path_source_is_xdg_runtime_dir_when_set() {
1624 let runtime_dir = OsString::from("/run/user/1000");
1625
1626 let (path, source) = connection_file_path_with_source(Some(runtime_dir));
1627
1628 assert_eq!(
1629 path,
1630 PathBuf::from("/run/user/1000").join(CONNECTION_FILE_NAME)
1631 );
1632 assert_eq!(source, ConnectionFileSource::XdgRuntimeDir);
1633 }
1634
1635 #[test]
1636 fn connection_file_path_falls_back_to_temp_dir_with_user_token_when_xdg_unset() {
1637 let _env_lock = ENV_LOCK.lock().unwrap();
1638 let _xdg = EnvGuard::unset("XDG_RUNTIME_DIR");
1639
1640 assert_eq!(
1641 connection_file_path(),
1642 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1643 );
1644 }
1645
1646 #[test]
1647 fn connection_file_path_source_is_temp_dir_when_xdg_unset() {
1648 let (path, source) = connection_file_path_with_source(None);
1649
1650 assert_eq!(
1651 path,
1652 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1653 );
1654 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1655 }
1656
1657 #[test]
1662 fn user_connection_token_is_stable_under_concurrent_callers() {
1663 let expected = user_connection_token();
1664 let workers: Vec<_> = (0..32)
1665 .map(|_| {
1666 std::thread::spawn(|| (0..40).map(|_| user_connection_token()).collect::<Vec<_>>())
1667 })
1668 .collect();
1669 for worker in workers {
1670 for token in worker.join().expect("probe thread") {
1671 assert_eq!(token, expected, "token diverged under concurrent probes");
1672 }
1673 }
1674 }
1675
1676 #[test]
1677 fn connection_file_path_source_is_temp_dir_when_xdg_empty() {
1678 let (path, source) = connection_file_path_with_source(Some(OsString::new()));
1679
1680 assert_eq!(
1681 path,
1682 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1683 );
1684 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1685 }
1686
1687 #[test]
1691 fn daemon_binary_config_journals_and_captures_into_the_run_dir() {
1692 let _env_lock = ENV_LOCK.lock().unwrap();
1693 let root = unique_temp_dir("daemon-binary-config");
1694 let data_home = root.join("data");
1695 let _data = EnvGuard::set("XDG_DATA_HOME", &data_home);
1696 let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1697 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1698
1699 let config = BootstrapConfig::from_env_for_daemon_binary().unwrap();
1700
1701 let run_dir = data_home.join("cortexkit").join("run");
1702 assert_eq!(
1703 config.terminal_journal_path,
1704 Some(run_dir.join("terminals.jsonl"))
1705 );
1706 assert_eq!(config.capture_logs_dir, Some(run_dir.join("logs")));
1707 assert_eq!(
1708 BootstrapConfig::new(root.join("connection.json"), 0).terminal_journal_path,
1709 None,
1710 "an in-process config must not journal anywhere unless asked to"
1711 );
1712 }
1713
1714 #[test]
1715 fn daemon_binary_config_refuses_a_relative_data_home() {
1716 let _env_lock = ENV_LOCK.lock().unwrap();
1717 let root = unique_temp_dir("daemon-binary-relative-data");
1718 let _data = EnvGuard::set_str("XDG_DATA_HOME", "relative-data-home");
1719 let _config = EnvGuard::set("XDG_CONFIG_HOME", &root.join("config"));
1720 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1721
1722 let error = BootstrapConfig::from_env_for_daemon_binary()
1723 .expect_err("a relative data home must refuse the daemon binary's config");
1724 assert!(
1725 matches!(error, BootstrapError::RunDir(_)),
1726 "expected a run-directory refusal, got {error}"
1727 );
1728 assert!(error.to_string().contains("XDG_DATA_HOME"), "{error}");
1729 }
1730
1731 #[test]
1732 fn configured_port_uses_default_config_and_env_override() {
1733 let _env_lock = ENV_LOCK.lock().unwrap();
1734 let (_dir, conn_path) = temp_connection_file_path("daemon-config-port");
1735 let config_path = conn_path.with_file_name("subc.jsonc");
1736
1737 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1738 assert_eq!(
1739 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1740 .unwrap()
1741 .port,
1742 DEFAULT_SUBC_PORT
1743 );
1744
1745 fs::write(&config_path, r#"{ "version": 1, "port": 8123 }"#).unwrap();
1746 assert_eq!(
1747 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1748 .unwrap()
1749 .port,
1750 8123
1751 );
1752
1753 let _port = EnvGuard::set_str(SUBC_PORT_ENV, "9012");
1754 assert_eq!(
1755 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1756 .unwrap()
1757 .port,
1758 9012
1759 );
1760 }
1761
1762 #[tokio::test]
1763 async fn second_singleton_probe_against_served_tcp_daemon_reports_already_running() {
1764 let (_dir, path) = temp_connection_file_path("already-running");
1765
1766 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1767 let server = start_server(bound);
1768
1769 let second = ensure_singleton(&path, 0).await.unwrap();
1770 assert!(matches!(second, Outcome::AlreadyRunning));
1771
1772 server.abort();
1773 let _ = server.await;
1774 }
1775
1776 #[tokio::test]
1777 async fn daemon_connection_file_publishes_protocol_wire_version() {
1778 let (_dir, path) = temp_connection_file_path("wire-version");
1779 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1780 assert_eq!(bound.connection_info.wire_version, Some(PROTOCOL_VERSION));
1781 assert_eq!(
1782 connection_file::read(&path).unwrap().wire_version,
1783 Some(PROTOCOL_VERSION)
1784 );
1785
1786 drop(bound.listeners);
1787 }
1788
1789 #[tokio::test]
1790 async fn stale_unbound_connection_file_is_reclaimed() {
1791 let (_dir, path) = temp_connection_file_path("stale-reclaim");
1792 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1793 let stale_port = stale.local_addr().unwrap().port();
1794 drop(stale);
1795 let stale_info = make_connection_info(stale_port);
1796 write_atomic(&path, &stale_info).unwrap();
1797
1798 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1799 assert_ne!(bound.connection_info.key, stale_info.key);
1800 drop(bound.listeners);
1801 }
1802
1803 #[cfg(unix)]
1804 #[tokio::test]
1805 async fn ensure_singleton_reclaims_insecure_connection_file() {
1806 let (_dir, path) = temp_connection_file_path("insecure-reclaim");
1807 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1808 let stale_port = stale.local_addr().unwrap().port();
1809 drop(stale);
1810 let stale_info = make_connection_info(stale_port);
1811 write_atomic(&path, &stale_info).unwrap();
1812 fs::set_permissions(&path, fs::Permissions::from_mode(0o644)).unwrap();
1813
1814 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1815 assert_ne!(bound.connection_info.key, stale_info.key);
1816 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1817 assert_owner_only_connection_file(&path);
1818
1819 drop(bound.listeners);
1820 }
1821
1822 #[tokio::test]
1823 async fn ensure_singleton_reclaims_non_loopback_connection_file() {
1824 let (_dir, path) = temp_connection_file_path("non-loopback-reclaim");
1825 let mut stale_info = make_connection_info(8757);
1826 stale_info.endpoints = vec![Endpoint {
1827 host: "192.0.2.10".to_owned(),
1828 port: 8757,
1829 }];
1830 write_atomic(&path, &stale_info).unwrap();
1831
1832 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1833 assert_ne!(bound.connection_info.key, stale_info.key);
1834 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1835 assert!(bound
1836 .connection_info
1837 .endpoints
1838 .iter()
1839 .all(|endpoint| endpoint.host.parse::<IpAddr>().unwrap().is_loopback()));
1840 assert_owner_only_connection_file(&path);
1841
1842 drop(bound.listeners);
1843 }
1844
1845 #[tokio::test]
1846 async fn ensure_singleton_reclaims_invalid_connection_file_shapes() {
1847 let mut unsupported_schema = make_connection_info(8757);
1848 unsupported_schema.schema = SCHEMA_VERSION + 1;
1849
1850 let mut empty_endpoints = make_connection_info(8757);
1851 empty_endpoints.endpoints.clear();
1852
1853 let mut short_key = make_connection_info(8757);
1854 short_key.key = vec![0x5A; MIN_KEY_LEN - 1];
1855
1856 let cases = vec![
1857 (
1858 "unsupported-schema",
1859 serde_json::to_vec(&unsupported_schema).unwrap(),
1860 Some(unsupported_schema),
1861 ),
1862 (
1863 "empty-endpoints",
1864 serde_json::to_vec(&empty_endpoints).unwrap(),
1865 Some(empty_endpoints),
1866 ),
1867 (
1868 "short-key",
1869 serde_json::to_vec(&short_key).unwrap(),
1870 Some(short_key),
1871 ),
1872 ("invalid-json", b"{not valid connection json".to_vec(), None),
1873 ];
1874
1875 for (label, contents, old_info) in cases {
1876 let (_dir, path) = temp_connection_file_path(label);
1877 write_raw_owner_only_connection_file(&path, &contents);
1878
1879 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1880 if let Some(old_info) = old_info {
1881 assert_ne!(bound.connection_info.key, old_info.key, "{label}");
1882 assert_ne!(
1883 bound.connection_info.daemon_id, old_info.daemon_id,
1884 "{label}"
1885 );
1886 }
1887 assert!(bound.connection_info.key.len() >= MIN_KEY_LEN, "{label}");
1888 assert_ne!(bound.connection_info.daemon_id, [0u8; 16], "{label}");
1889 assert_owner_only_connection_file(&path);
1890
1891 drop(bound.listeners);
1892 }
1893 }
1894
1895 #[tokio::test]
1896 async fn foreign_reused_port_connection_file_is_reclaimed_after_auth_probe_fails() {
1897 let (_dir, path) = temp_connection_file_path("foreign-reclaim");
1898 let foreign = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1899 let foreign_port = foreign.local_addr().unwrap().port();
1900 write_atomic(&path, &make_connection_info(foreign_port)).unwrap();
1901 let foreign_task = tokio::spawn(async move {
1902 if let Ok((mut stream, _)) = foreign.accept().await {
1903 let mut buf = [0u8; 64];
1904 let _ = stream.read(&mut buf).await;
1905 }
1906 });
1907
1908 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1909 assert!(bound
1910 .connection_info
1911 .endpoints
1912 .iter()
1913 .all(|endpoint| endpoint.port != foreign_port));
1914
1915 drop(bound.listeners);
1916 let _ = foreign_task.await;
1917 }
1918
1919 #[tokio::test]
1920 async fn stale_start_lock_file_is_reclaimable() {
1921 let (_dir, path) = temp_connection_file_path("start-lock-stale-file");
1922 let lock_path = start_lock_path(&path);
1923 drop(open_owner_only_lock(&lock_path).unwrap());
1924 assert!(lock_path.is_file());
1925
1926 let lock = StartLock::acquire(&path).await.unwrap();
1927 assert!(lock_path.is_file());
1928
1929 drop(lock);
1930 assert!(lock_path.is_file());
1931 }
1932
1933 #[tokio::test]
1934 async fn held_start_lock_blocks_second_acquire_until_release() {
1935 let (_dir, path) = temp_connection_file_path("start-lock-held");
1936 let lock_path = start_lock_path(&path);
1937 let first = StartLock::acquire(&path).await.unwrap();
1938
1939 let err = match StartLock::acquire(&path).await {
1940 Ok(_) => panic!("second acquire while held must stay busy"),
1941 Err(err) => err,
1942 };
1943 assert!(matches!(
1944 err,
1945 BootstrapError::StartLockBusy {
1946 ref path,
1947 attempts: START_LOCK_RETRIES,
1948 } if path == &lock_path
1949 ));
1950
1951 drop(first);
1952
1953 let second = StartLock::acquire(&path)
1954 .await
1955 .expect("released advisory lock should be reclaimable");
1956 drop(second);
1957 }
1958
1959 #[tokio::test]
1960 async fn bind_conflict_on_fixed_port_fails_loud_without_reselecting() {
1961 let (_dir, path) = temp_connection_file_path("bind-conflict");
1962 let occupied = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1963 let occupied_port = occupied.local_addr().unwrap().port();
1964
1965 let err = ensure_singleton(&path, occupied_port).await.unwrap_err();
1966 assert!(matches!(
1967 err,
1968 BootstrapError::Bind { ref source, .. } if source.kind() == io::ErrorKind::AddrInUse
1969 ));
1970 assert!(err.to_string().contains("set the port in config"));
1971
1972 drop(occupied);
1973 }
1974
1975 #[tokio::test]
1976 async fn key_rotation_republishes_new_material_and_old_file_fails_auth() {
1977 let (_dir, path) = temp_connection_file_path("key-rotation");
1978 let first = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1979 let old_info = first.connection_info.clone();
1980 let fixed_port = old_info.endpoints[0].port;
1981 drop(first.listeners);
1982
1983 let second = expect_bound(ensure_singleton(&path, fixed_port).await.unwrap());
1984 let new_info = second.connection_info.clone();
1985 assert_ne!(old_info.key, new_info.key);
1986 assert_ne!(old_info.daemon_id, new_info.daemon_id);
1987 let server = start_server(second);
1988
1989 let mut old_stream = connect_from_info(&old_info).await.unwrap();
1990 let old_auth = authenticate_client(&mut old_stream, &old_info, PROBE_AUTH_DEADLINE).await;
1991 assert!(
1992 old_auth.is_err(),
1993 "old key must not authenticate after restart"
1994 );
1995
1996 let reread = connection_file::read(&path).unwrap();
1997 let mut new_stream = connect_from_info(&reread).await.unwrap();
1998 authenticate_client(&mut new_stream, &reread, PROBE_AUTH_DEADLINE)
1999 .await
2000 .unwrap();
2001
2002 server.abort();
2003 let _ = server.await;
2004 }
2005
2006 #[cfg(unix)]
2007 #[tokio::test]
2008 async fn published_connection_file_permissions_are_owner_only() {
2009 let (_dir, path) = temp_connection_file_path("permissions");
2010 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
2011
2012 let mode = fs::metadata(&path).unwrap().permissions().mode() & 0o777;
2013 assert_eq!(mode, 0o600);
2014
2015 drop(bound.listeners);
2016 }
2017}