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