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 let escalated = tokio::select! {
791 biased;
792 _ = terminate.recv() => {
793 info!("second SIGTERM: abandoning daemon shutdown wait");
794 true
795 }
796 result = supervisor.drain_for_daemon_shutdown() => {
797 if let Err(error) = result {
798 warn!(%error, "daemon shutdown drain failed; exiting anyway");
799 }
800 false
801 }
802 };
803 supervisor
808 .end_children_for_daemon_shutdown(escalated, async {
809 terminate.recv().await;
810 })
811 .await;
812 Ok(())
813 }
814 #[cfg(not(unix))]
815 serve_task
816 .join()
817 .await
818 .map_err(BootstrapError::ServeJoin)?
819 .map_err(BootstrapError::Serve)
820}
821
822fn normalized_build_provenance(value: &str) -> Option<String> {
823 match value.trim() {
824 "" | "unavailable" => None,
825 value => Some(value.to_string()),
826 }
827}
828
829pub async fn ensure_singleton(
837 connection_file_path: impl AsRef<Path>,
838 port: u16,
839) -> Result<Outcome, BootstrapError> {
840 ensure_singleton_with_config(BootstrapConfig::new(connection_file_path.as_ref(), port)).await
841}
842
843pub async fn ensure_singleton_with_config(
844 config: BootstrapConfig,
845) -> Result<Outcome, BootstrapError> {
846 let path = config.connection_file_path;
847
848 if matches!(probe_existing(&path).await?, Probe::Live) {
849 return Ok(Outcome::AlreadyRunning);
850 }
851
852 let _lock = StartLock::acquire(&path).await?;
853
854 if matches!(probe_existing(&path).await?, Probe::Live) {
858 return Ok(Outcome::AlreadyRunning);
859 }
860
861 remove_stale_connection_file_if_present(&path)?;
862
863 let machine_id = config
867 .machine_id_path
868 .as_deref()
869 .map(crate::machine_id::load_or_mint)
870 .transpose()
871 .map_err(BootstrapError::MachineId)?;
872
873 let (listeners, endpoints) = bind_loopback(config.port).await?;
874 let connection_info = ConnectionInfo {
875 schema: SCHEMA_VERSION,
876 wire_version: Some(PROTOCOL_VERSION),
877 endpoints,
878 key: generate_key().map_err(BootstrapError::GenerateConnectionFile)?,
879 daemon_id: generate_daemon_id().map_err(BootstrapError::GenerateConnectionFile)?,
880 pid: process::id(),
881 daemon_ver: config.daemon_ver,
882 };
883
884 if let Err(source) = write_atomic(&path, &connection_info) {
885 drop(listeners);
886 return Err(BootstrapError::ConnectionFileWrite { path, source });
887 }
888
889 Ok(Outcome::Bound(BoundDaemon {
890 listeners,
891 connection_info,
892 connection_file_path: path,
893 connection_file_source: config.connection_file_source,
894 machine_id,
895 }))
896}
897
898#[derive(Debug, Clone, Copy, PartialEq, Eq)]
899enum Probe {
900 Live,
901 StaleOrAbsent,
902}
903
904async fn probe_existing(path: &Path) -> Result<Probe, BootstrapError> {
905 let info = match connection_file::read(path) {
906 Ok(info) => info,
907 Err(source) if is_absent_or_stale_connection_file(&source) => {
908 return Ok(Probe::StaleOrAbsent)
909 }
910 Err(source) => {
911 return Err(BootstrapError::ConnectionFileRead {
912 path: path.to_path_buf(),
913 source,
914 })
915 }
916 };
917
918 for endpoint in &info.endpoints {
919 if matches!(probe_endpoint(&info, endpoint).await, Probe::Live) {
920 return Ok(Probe::Live);
921 }
922 }
923
924 Ok(Probe::StaleOrAbsent)
925}
926
927async fn probe_endpoint(info: &ConnectionInfo, endpoint: &Endpoint) -> Probe {
928 let Ok(ip) = endpoint.host.parse::<IpAddr>() else {
929 return Probe::StaleOrAbsent;
930 };
931 if !ip.is_loopback() {
932 return Probe::StaleOrAbsent;
933 }
934 let addr = SocketAddr::new(ip, endpoint.port);
935
936 let mut stream = match timeout(CONNECT_TIMEOUT, TcpStream::connect(addr)).await {
937 Ok(Ok(stream)) => stream,
938 Ok(Err(_)) | Err(_) => return Probe::StaleOrAbsent,
939 };
940
941 match authenticate_client(&mut stream, info, PROBE_AUTH_DEADLINE).await {
942 Ok(()) => Probe::Live,
943 Err(AuthError::DaemonIdMismatch)
944 | Err(AuthError::InvalidServerProof)
945 | Err(AuthError::UnexpectedEof { .. })
946 | Err(AuthError::Timeout { .. })
947 | Err(AuthError::JsonEncode { .. })
948 | Err(AuthError::JsonDecode { .. })
949 | Err(AuthError::Io { .. })
950 | Err(AuthError::MessageTooLarge { .. })
951 | Err(AuthError::KeyTooShort { .. })
952 | Err(AuthError::Random(_))
953 | Err(AuthError::InvalidClientAuth) => Probe::StaleOrAbsent,
954 }
955}
956
957fn is_absent_or_stale_connection_file(err: &ConnectionFileError) -> bool {
958 match err {
959 ConnectionFileError::Io { source, .. } if source.kind() == io::ErrorKind::NotFound => true,
960 ConnectionFileError::JsonRead { .. }
961 | ConnectionFileError::UnsupportedSchema { .. }
962 | ConnectionFileError::Invalid { .. }
963 | ConnectionFileError::KeyTooShort { .. }
964 | ConnectionFileError::InsecurePermissions { .. } => true,
968 ConnectionFileError::MissingParent { .. }
969 | ConnectionFileError::MissingFileName { .. }
970 | ConnectionFileError::InsecureParentDirectory { .. }
974 | ConnectionFileError::Io { .. }
975 | ConnectionFileError::JsonWrite { .. }
976 | ConnectionFileError::Random(_)
977 | ConnectionFileError::WireVersionMismatch { .. } => false,
980 }
981}
982
983fn remove_stale_connection_file_if_present(path: &Path) -> Result<(), BootstrapError> {
984 match fs::remove_file(path) {
985 Ok(()) => Ok(()),
986 Err(err) if err.kind() == io::ErrorKind::NotFound => Ok(()),
987 Err(source) => Err(BootstrapError::RemoveStale {
988 path: path.to_path_buf(),
989 source,
990 }),
991 }
992}
993
994async fn bind_loopback(port: u16) -> Result<(Vec<TcpListener>, Vec<Endpoint>), BootstrapError> {
995 let v4_host = Ipv4Addr::LOCALHOST;
996 let v4 = TcpListener::bind((v4_host, port))
997 .await
998 .map_err(|source| BootstrapError::Bind {
999 host: v4_host.to_string(),
1000 port,
1001 source,
1002 })?;
1003 let actual_port = v4
1004 .local_addr()
1005 .map_err(|source| BootstrapError::LocalAddr {
1006 host: v4_host.to_string(),
1007 source,
1008 })?
1009 .port();
1010
1011 let mut listeners = vec![v4];
1012 let mut endpoints = vec![Endpoint {
1013 host: v4_host.to_string(),
1014 port: actual_port,
1015 }];
1016
1017 let v6_host = Ipv6Addr::LOCALHOST;
1018 match TcpListener::bind((v6_host, actual_port)).await {
1019 Ok(v6) => {
1020 listeners.push(v6);
1021 endpoints.push(Endpoint {
1022 host: v6_host.to_string(),
1023 port: actual_port,
1024 });
1025 }
1026 Err(err) if ipv6_loopback_unavailable(&err) => {
1027 warn!(
1028 port = actual_port,
1029 error = %err,
1030 "IPv6 loopback unavailable; serving only IPv4 loopback"
1031 );
1032 }
1033 Err(source) => {
1034 drop(listeners);
1035 return Err(BootstrapError::Bind {
1036 host: v6_host.to_string(),
1037 port: actual_port,
1038 source,
1039 });
1040 }
1041 }
1042
1043 Ok((listeners, endpoints))
1044}
1045
1046fn ipv6_loopback_unavailable(err: &io::Error) -> bool {
1047 matches!(
1048 err.kind(),
1049 io::ErrorKind::AddrNotAvailable | io::ErrorKind::Unsupported
1050 ) || matches!(err.raw_os_error(), Some(47) | Some(49) | Some(97))
1051}
1052
1053struct AbortOnDrop<T> {
1054 handle: JoinHandle<T>,
1055}
1056
1057impl<T> AbortOnDrop<T> {
1058 fn new(handle: JoinHandle<T>) -> Self {
1059 Self { handle }
1060 }
1061
1062 async fn join(&mut self) -> Result<T, JoinError> {
1063 (&mut self.handle).await
1064 }
1065}
1066
1067impl<T> Drop for AbortOnDrop<T> {
1068 fn drop(&mut self) {
1069 if !self.handle.is_finished() {
1070 self.handle.abort();
1071 }
1072 }
1073}
1074
1075struct StartLock {
1076 _file: fs::File,
1079}
1080
1081impl StartLock {
1082 async fn acquire(connection_file_path: &Path) -> Result<Self, BootstrapError> {
1083 let path = start_lock_path(connection_file_path);
1084 for _ in 0..START_LOCK_RETRIES {
1085 let file = match open_owner_only_lock(&path) {
1086 Ok(file) => file,
1087 Err(source) => return Err(BootstrapError::StartLockCreate { path, source }),
1088 };
1089 match FileExt::try_lock(&file) {
1090 Ok(()) => return Ok(Self { _file: file }),
1091 Err(TryLockError::WouldBlock) => sleep(START_LOCK_RETRY_DELAY).await,
1092 Err(TryLockError::Error(source)) => {
1093 return Err(BootstrapError::StartLockCreate { path, source });
1094 }
1095 }
1096 }
1097
1098 Err(BootstrapError::StartLockBusy {
1099 path,
1100 attempts: START_LOCK_RETRIES,
1101 })
1102 }
1103}
1104
1105fn open_owner_only_lock(path: &Path) -> io::Result<fs::File> {
1106 let mut options = fs::OpenOptions::new();
1107 options.read(true).write(true).create(true);
1108 #[cfg(unix)]
1109 {
1110 use std::os::unix::fs::OpenOptionsExt;
1111 options.mode(0o600);
1112 }
1113 options.open(path)
1114}
1115
1116fn start_lock_path(connection_file_path: &Path) -> PathBuf {
1117 let file_name = connection_file_path
1118 .file_name()
1119 .map(|name| name.to_string_lossy())
1120 .unwrap_or_else(|| CONNECTION_FILE_NAME.into());
1121 let lock_name = format!("{file_name}.start-lock");
1122 connection_file_path
1123 .parent()
1124 .filter(|parent| !parent.as_os_str().is_empty())
1125 .unwrap_or_else(|| Path::new("."))
1126 .join(lock_name)
1127}
1128
1129fn non_empty_os_var(key: &str) -> Option<OsString> {
1130 let value = env::var_os(key)?;
1131 if value.is_empty() {
1132 None
1133 } else {
1134 Some(value)
1135 }
1136}
1137
1138#[derive(Debug)]
1141pub enum BootstrapError {
1142 #[cfg(unix)]
1143 Signal(io::Error),
1144 InvalidPort {
1145 raw: String,
1146 source: std::num::ParseIntError,
1147 },
1148 ConnectionFileRead {
1149 path: PathBuf,
1150 source: ConnectionFileError,
1151 },
1152 ConnectionFileWrite {
1153 path: PathBuf,
1154 source: ConnectionFileError,
1155 },
1156 GenerateConnectionFile(ConnectionFileError),
1157 StartLockCreate {
1158 path: PathBuf,
1159 source: io::Error,
1160 },
1161 StartLockBusy {
1162 path: PathBuf,
1163 attempts: usize,
1164 },
1165 RemoveStale {
1166 path: PathBuf,
1167 source: io::Error,
1168 },
1169 Bind {
1170 host: String,
1171 port: u16,
1172 source: io::Error,
1173 },
1174 LocalAddr {
1175 host: String,
1176 source: io::Error,
1177 },
1178 DaemonConfig(DaemonConfigError),
1179 MachineId(crate::machine_id::MachineIdFileError),
1182 Serve(ServerError),
1183 ServeJoin(tokio::task::JoinError),
1184}
1185
1186impl fmt::Display for BootstrapError {
1187 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1188 match self {
1189 #[cfg(unix)]
1190 Self::Signal(error) => write!(f, "failed to register SIGTERM handler: {error}"),
1191 Self::InvalidPort { raw, source } => {
1192 write!(f, "invalid {SUBC_PORT_ENV} value '{raw}': {source}")
1193 }
1194 Self::ConnectionFileRead { path, source } => write!(
1195 f,
1196 "failed to read connection file {}: {source}",
1197 path.display()
1198 ),
1199 Self::ConnectionFileWrite { path, source } => write!(
1200 f,
1201 "failed to publish connection file {}: {source}",
1202 path.display()
1203 ),
1204 Self::GenerateConnectionFile(err) => {
1205 write!(f, "failed to generate connection-file auth material: {err}")
1206 }
1207 Self::StartLockCreate { path, source } => {
1208 write!(
1209 f,
1210 "failed to create start lock {}: {source}",
1211 path.display()
1212 )
1213 }
1214 Self::StartLockBusy { path, attempts } => write!(
1215 f,
1216 "start lock {} remained busy after {attempts} attempts",
1217 path.display()
1218 ),
1219 Self::RemoveStale { path, source } => write!(
1220 f,
1221 "failed to remove stale connection file {}: {source}",
1222 path.display()
1223 ),
1224 Self::Bind { host, port, source } if source.kind() == io::ErrorKind::AddrInUse => {
1225 write!(
1226 f,
1227 "port {port} in use on loopback {host}: {source}; set the port in config"
1228 )
1229 }
1230 Self::Bind { host, port, source } => {
1231 write!(f, "failed to bind loopback TCP {host}:{port}: {source}")
1232 }
1233 Self::LocalAddr { host, source } => {
1234 write!(f, "failed to read local address for {host}: {source}")
1235 }
1236 Self::DaemonConfig(err) => write!(f, "failed to load daemon config: {err}"),
1237 Self::MachineId(err) => write!(f, "refusing to start: {err}"),
1238 Self::Serve(err) => write!(f, "daemon server failed: {err}"),
1239 Self::ServeJoin(err) => write!(f, "daemon server task failed: {err}"),
1240 }
1241 }
1242}
1243
1244impl Error for BootstrapError {
1245 fn source(&self) -> Option<&(dyn Error + 'static)> {
1246 match self {
1247 #[cfg(unix)]
1248 Self::Signal(source) => Some(source),
1249 Self::InvalidPort { source, .. } => Some(source),
1250 Self::ConnectionFileRead { source, .. }
1251 | Self::ConnectionFileWrite { source, .. }
1252 | Self::GenerateConnectionFile(source) => Some(source),
1253 Self::StartLockCreate { source, .. }
1254 | Self::RemoveStale { source, .. }
1255 | Self::Bind { source, .. }
1256 | Self::LocalAddr { source, .. } => Some(source),
1257 Self::DaemonConfig(err) => Some(err),
1258 Self::MachineId(err) => Some(err),
1259 Self::Serve(err) => Some(err),
1260 Self::ServeJoin(err) => Some(err),
1261 Self::StartLockBusy { .. } => None,
1262 }
1263 }
1264}
1265
1266#[cfg(test)]
1267mod tests {
1268 use super::*;
1269 use crate::server::ServerAuth;
1270 use crate::test_support::TestTempDir;
1271 #[cfg(target_os = "linux")]
1272 use std::collections::BTreeSet;
1273 use std::sync::Mutex;
1274 #[cfg(target_os = "linux")]
1275 use subc_control::ModuleProtocol;
1276 use subc_transport::MIN_KEY_LEN;
1277 use tokio::io::AsyncReadExt;
1278 use tokio::task::JoinHandle;
1279
1280 #[cfg(unix)]
1281 use std::os::unix::fs::PermissionsExt;
1282
1283 static ENV_LOCK: Mutex<()> = Mutex::new(());
1284
1285 #[test]
1286 fn normalized_build_provenance_preserves_real_values() {
1287 assert_eq!(normalized_build_provenance("abc"), Some("abc".to_string()));
1288 }
1289
1290 #[test]
1291 fn normalized_build_provenance_omits_unavailable_and_empty_values() {
1292 assert_eq!(normalized_build_provenance("unavailable"), None);
1293 assert_eq!(normalized_build_provenance(""), None);
1294 }
1295
1296 #[cfg(target_os = "linux")]
1297 fn current_cgroup_path_for_test() -> io::Result<PathBuf> {
1298 let cgroups = fs::read_to_string("/proc/self/cgroup")?;
1299 let relative = cgroups
1300 .lines()
1301 .find_map(|line| line.strip_prefix("0::"))
1302 .ok_or_else(|| {
1303 io::Error::new(io::ErrorKind::Unsupported, "cgroup v2 is unavailable")
1304 })?;
1305 Ok(Path::new("/sys/fs/cgroup").join(relative.trim_start_matches('/')))
1306 }
1307
1308 #[cfg(target_os = "linux")]
1309 fn module_cgroup_directories() -> io::Result<BTreeSet<OsString>> {
1310 let modules = current_cgroup_path_for_test()?.join("subc-modules");
1311 let entries = match fs::read_dir(modules) {
1312 Ok(entries) => entries,
1313 Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(BTreeSet::new()),
1314 Err(error) => return Err(error),
1315 };
1316 let mut directories = BTreeSet::new();
1317 for entry in entries {
1318 let entry = entry?;
1319 if entry.file_type()?.is_dir() {
1320 directories.insert(entry.file_name());
1321 }
1322 }
1323 Ok(directories)
1324 }
1325
1326 #[cfg(target_os = "linux")]
1327 async fn wait_for_path(path: &Path, task: &JoinHandle<Result<(), BootstrapError>>) {
1328 let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1329 while !path.exists() && tokio::time::Instant::now() < deadline {
1330 assert!(
1331 !task.is_finished(),
1332 "daemon exited before creating {}",
1333 path.display()
1334 );
1335 sleep(Duration::from_millis(10)).await;
1336 }
1337 assert!(path.exists(), "daemon did not create {}", path.display());
1338 }
1339
1340 #[cfg(target_os = "linux")]
1341 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1342 async fn run_with_config_does_not_reconcile_the_ambient_cgroup_by_default() {
1343 let temp = unique_temp_dir("bootstrap-cgroup-default-disabled");
1344 let module_id = format!("cgroup-isolation-probe-{}", process::id());
1345 let before = module_cgroup_directories().expect("read ambient module cgroups before boot");
1346 assert!(
1347 !before.contains(&OsString::from(&module_id)),
1348 "isolation probe cgroup already exists before this daemon starts"
1349 );
1350 let capture = temp.join("logs").join(format!("{module_id}.stderr.log"));
1351 let module = ConfiguredModule {
1352 module_id,
1353 program: PathBuf::from("sh"),
1354 args: vec!["-c".to_string(), "sleep 30".to_string()],
1355 env: Vec::new(),
1356 log: None,
1357 enabled: true,
1358 reserved: false,
1359 reserved_prefixes: Vec::new(),
1360 protocol: ModuleProtocol::None,
1361 overlap: Default::default(),
1362 health: HealthConfig::default(),
1363 drain_timeout_ms: None,
1364 route_bind_relay_timeout_ms: None,
1365 restart: RestartPolicy::default(),
1366 };
1367 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1368 .with_configured_modules([module])
1369 .with_capture_logs_dir(temp.join("logs"))
1370 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1371 let task = tokio::spawn(run_with_config(config));
1372
1373 wait_for_path(&capture, &task).await;
1374 let after = module_cgroup_directories().expect("read ambient module cgroups after boot");
1375
1376 task.abort();
1377 assert!(task
1378 .await
1379 .expect_err("aborted daemon task must cancel")
1380 .is_cancelled());
1381 assert_eq!(
1382 before, after,
1383 "an in-process daemon must not create or reconcile ambient module cgroups"
1384 );
1385 }
1386
1387 #[cfg(target_os = "linux")]
1388 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1389 async fn explicit_cgroup_root_is_prepared_inside_the_fixture_tree() {
1390 let temp = unique_temp_dir("bootstrap-cgroup-explicit-root");
1391 let cgroup_root = temp.join("cgroup");
1392 fs::create_dir(&cgroup_root).expect("create scratch cgroup root");
1393 fs::write(cgroup_root.join("cgroup.procs"), b"").expect("write scratch cgroup marker");
1394 let modules = cgroup_root.join("subc-modules");
1395 let config = BootstrapConfig::new(temp.join("connection.json"), 0)
1396 .with_cgroup_placement(CgroupPlacementConfig::Root(cgroup_root))
1397 .with_terminal_journal_path(temp.join("terminals.jsonl"));
1398 let task = tokio::spawn(run_with_config(config));
1399
1400 wait_for_path(&modules, &task).await;
1401
1402 task.abort();
1403 assert!(task
1404 .await
1405 .expect_err("aborted daemon task must cancel")
1406 .is_cancelled());
1407 assert!(
1408 fs::read_dir(&modules)
1409 .expect("read prepared modules directory")
1410 .next()
1411 .is_none(),
1412 "the delegation probe must clean up after itself"
1413 );
1414 }
1415
1416 struct EnvGuard {
1417 key: &'static str,
1418 previous: Option<OsString>,
1419 }
1420
1421 impl EnvGuard {
1422 fn set(key: &'static str, value: &Path) -> Self {
1423 let previous = env::var_os(key);
1424 env::set_var(key, value);
1425 Self { key, previous }
1426 }
1427
1428 fn set_str(key: &'static str, value: &str) -> Self {
1429 let previous = env::var_os(key);
1430 env::set_var(key, value);
1431 Self { key, previous }
1432 }
1433
1434 fn unset(key: &'static str) -> Self {
1435 let previous = env::var_os(key);
1436 env::remove_var(key);
1437 Self { key, previous }
1438 }
1439 }
1440
1441 impl Drop for EnvGuard {
1442 fn drop(&mut self) {
1443 match &self.previous {
1444 Some(value) => env::set_var(self.key, value),
1445 None => env::remove_var(self.key),
1446 }
1447 }
1448 }
1449
1450 fn unique_temp_dir(name: &str) -> TestTempDir {
1451 TestTempDir::new(name)
1452 }
1453
1454 fn temp_connection_file_path(name: &str) -> (TestTempDir, PathBuf) {
1455 let dir = unique_temp_dir(name);
1456 let path = dir.join("conn.json");
1457 (dir, path)
1458 }
1459
1460 fn auth_for(info: &ConnectionInfo) -> ServerAuth {
1461 ServerAuth::new(info.key.clone(), info.daemon_id, info.daemon_ver.clone())
1462 }
1463
1464 fn start_server(bound: BoundDaemon) -> JoinHandle<Result<(), ServerError>> {
1465 let auth = auth_for(&bound.connection_info);
1466 tokio::spawn(serve_listeners(
1467 bound.listeners,
1468 Arc::new(Router::with_default_self_handler()),
1469 auth,
1470 ))
1471 }
1472
1473 fn expect_bound(outcome: Outcome) -> BoundDaemon {
1474 match outcome {
1475 Outcome::Bound(bound) => bound,
1476 Outcome::AlreadyRunning => panic!("fresh connection file unexpectedly had a daemon"),
1477 }
1478 }
1479
1480 async fn connect_from_info(conn: &ConnectionInfo) -> io::Result<TcpStream> {
1481 let endpoint = conn
1482 .endpoints
1483 .first()
1484 .expect("test connection file should have an endpoint");
1485 let ip: IpAddr = endpoint.host.parse().unwrap();
1486 TcpStream::connect(SocketAddr::new(ip, endpoint.port)).await
1487 }
1488
1489 fn make_connection_info(port: u16) -> ConnectionInfo {
1490 ConnectionInfo {
1491 schema: SCHEMA_VERSION,
1492 wire_version: Some(PROTOCOL_VERSION),
1493 endpoints: vec![Endpoint {
1494 host: "127.0.0.1".to_owned(),
1495 port,
1496 }],
1497 key: generate_key().unwrap(),
1498 daemon_id: generate_daemon_id().unwrap(),
1499 pid: process::id(),
1500 daemon_ver: "test-subc".to_owned(),
1501 }
1502 }
1503
1504 fn write_raw_owner_only_connection_file(path: &Path, contents: &[u8]) {
1505 fs::write(path, contents).unwrap();
1506 #[cfg(unix)]
1507 fs::set_permissions(path, fs::Permissions::from_mode(0o600)).unwrap();
1508 }
1509
1510 fn assert_owner_only_connection_file(path: &Path) {
1511 #[cfg(unix)]
1514 {
1515 let mode = fs::metadata(path).unwrap().permissions().mode() & 0o777;
1516 assert_eq!(mode, 0o600);
1517 }
1518 #[cfg(not(unix))]
1519 let _ = path;
1520 }
1521
1522 #[test]
1523 fn connection_file_path_uses_xdg_runtime_dir_when_set() {
1524 let _env_lock = ENV_LOCK.lock().unwrap();
1525 let runtime_dir = unique_temp_dir("xdg-runtime");
1526 let _xdg = EnvGuard::set("XDG_RUNTIME_DIR", runtime_dir.path());
1527
1528 assert_eq!(
1529 connection_file_path(),
1530 runtime_dir.join(CONNECTION_FILE_NAME)
1531 );
1532 }
1533
1534 #[test]
1535 fn connection_file_path_source_is_xdg_runtime_dir_when_set() {
1536 let runtime_dir = OsString::from("/run/user/1000");
1537
1538 let (path, source) = connection_file_path_with_source(Some(runtime_dir));
1539
1540 assert_eq!(
1541 path,
1542 PathBuf::from("/run/user/1000").join(CONNECTION_FILE_NAME)
1543 );
1544 assert_eq!(source, ConnectionFileSource::XdgRuntimeDir);
1545 }
1546
1547 #[test]
1548 fn connection_file_path_falls_back_to_temp_dir_with_user_token_when_xdg_unset() {
1549 let _env_lock = ENV_LOCK.lock().unwrap();
1550 let _xdg = EnvGuard::unset("XDG_RUNTIME_DIR");
1551
1552 assert_eq!(
1553 connection_file_path(),
1554 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1555 );
1556 }
1557
1558 #[test]
1559 fn connection_file_path_source_is_temp_dir_when_xdg_unset() {
1560 let (path, source) = connection_file_path_with_source(None);
1561
1562 assert_eq!(
1563 path,
1564 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1565 );
1566 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1567 }
1568
1569 #[test]
1574 fn user_connection_token_is_stable_under_concurrent_callers() {
1575 let expected = user_connection_token();
1576 let workers: Vec<_> = (0..32)
1577 .map(|_| {
1578 std::thread::spawn(|| (0..40).map(|_| user_connection_token()).collect::<Vec<_>>())
1579 })
1580 .collect();
1581 for worker in workers {
1582 for token in worker.join().expect("probe thread") {
1583 assert_eq!(token, expected, "token diverged under concurrent probes");
1584 }
1585 }
1586 }
1587
1588 #[test]
1589 fn connection_file_path_source_is_temp_dir_when_xdg_empty() {
1590 let (path, source) = connection_file_path_with_source(Some(OsString::new()));
1591
1592 assert_eq!(
1593 path,
1594 env::temp_dir().join(format!("subc-{}.connection.json", user_connection_token()))
1595 );
1596 assert_eq!(source, ConnectionFileSource::TempDirFallback);
1597 }
1598
1599 #[test]
1600 fn configured_port_uses_default_config_and_env_override() {
1601 let _env_lock = ENV_LOCK.lock().unwrap();
1602 let (_dir, conn_path) = temp_connection_file_path("daemon-config-port");
1603 let config_path = conn_path.with_file_name("subc.jsonc");
1604
1605 let _port = EnvGuard::unset(SUBC_PORT_ENV);
1606 assert_eq!(
1607 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1608 .unwrap()
1609 .port,
1610 DEFAULT_SUBC_PORT
1611 );
1612
1613 fs::write(&config_path, r#"{ "version": 1, "port": 8123 }"#).unwrap();
1614 assert_eq!(
1615 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1616 .unwrap()
1617 .port,
1618 8123
1619 );
1620
1621 let _port = EnvGuard::set_str(SUBC_PORT_ENV, "9012");
1622 assert_eq!(
1623 BootstrapConfig::from_env_with_daemon_config_path(&config_path)
1624 .unwrap()
1625 .port,
1626 9012
1627 );
1628 }
1629
1630 #[tokio::test]
1631 async fn second_singleton_probe_against_served_tcp_daemon_reports_already_running() {
1632 let (_dir, path) = temp_connection_file_path("already-running");
1633
1634 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1635 let server = start_server(bound);
1636
1637 let second = ensure_singleton(&path, 0).await.unwrap();
1638 assert!(matches!(second, Outcome::AlreadyRunning));
1639
1640 server.abort();
1641 let _ = server.await;
1642 }
1643
1644 #[tokio::test]
1645 async fn daemon_connection_file_publishes_protocol_wire_version() {
1646 let (_dir, path) = temp_connection_file_path("wire-version");
1647 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1648 assert_eq!(bound.connection_info.wire_version, Some(PROTOCOL_VERSION));
1649 assert_eq!(
1650 connection_file::read(&path).unwrap().wire_version,
1651 Some(PROTOCOL_VERSION)
1652 );
1653
1654 drop(bound.listeners);
1655 }
1656
1657 #[tokio::test]
1658 async fn stale_unbound_connection_file_is_reclaimed() {
1659 let (_dir, path) = temp_connection_file_path("stale-reclaim");
1660 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1661 let stale_port = stale.local_addr().unwrap().port();
1662 drop(stale);
1663 let stale_info = make_connection_info(stale_port);
1664 write_atomic(&path, &stale_info).unwrap();
1665
1666 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1667 assert_ne!(bound.connection_info.key, stale_info.key);
1668 drop(bound.listeners);
1669 }
1670
1671 #[cfg(unix)]
1672 #[tokio::test]
1673 async fn ensure_singleton_reclaims_insecure_connection_file() {
1674 let (_dir, path) = temp_connection_file_path("insecure-reclaim");
1675 let stale = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1676 let stale_port = stale.local_addr().unwrap().port();
1677 drop(stale);
1678 let stale_info = make_connection_info(stale_port);
1679 write_atomic(&path, &stale_info).unwrap();
1680 fs::set_permissions(&path, fs::Permissions::from_mode(0o644)).unwrap();
1681
1682 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1683 assert_ne!(bound.connection_info.key, stale_info.key);
1684 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1685 assert_owner_only_connection_file(&path);
1686
1687 drop(bound.listeners);
1688 }
1689
1690 #[tokio::test]
1691 async fn ensure_singleton_reclaims_non_loopback_connection_file() {
1692 let (_dir, path) = temp_connection_file_path("non-loopback-reclaim");
1693 let mut stale_info = make_connection_info(8757);
1694 stale_info.endpoints = vec![Endpoint {
1695 host: "192.0.2.10".to_owned(),
1696 port: 8757,
1697 }];
1698 write_atomic(&path, &stale_info).unwrap();
1699
1700 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1701 assert_ne!(bound.connection_info.key, stale_info.key);
1702 assert_ne!(bound.connection_info.daemon_id, stale_info.daemon_id);
1703 assert!(bound
1704 .connection_info
1705 .endpoints
1706 .iter()
1707 .all(|endpoint| endpoint.host.parse::<IpAddr>().unwrap().is_loopback()));
1708 assert_owner_only_connection_file(&path);
1709
1710 drop(bound.listeners);
1711 }
1712
1713 #[tokio::test]
1714 async fn ensure_singleton_reclaims_invalid_connection_file_shapes() {
1715 let mut unsupported_schema = make_connection_info(8757);
1716 unsupported_schema.schema = SCHEMA_VERSION + 1;
1717
1718 let mut empty_endpoints = make_connection_info(8757);
1719 empty_endpoints.endpoints.clear();
1720
1721 let mut short_key = make_connection_info(8757);
1722 short_key.key = vec![0x5A; MIN_KEY_LEN - 1];
1723
1724 let cases = vec![
1725 (
1726 "unsupported-schema",
1727 serde_json::to_vec(&unsupported_schema).unwrap(),
1728 Some(unsupported_schema),
1729 ),
1730 (
1731 "empty-endpoints",
1732 serde_json::to_vec(&empty_endpoints).unwrap(),
1733 Some(empty_endpoints),
1734 ),
1735 (
1736 "short-key",
1737 serde_json::to_vec(&short_key).unwrap(),
1738 Some(short_key),
1739 ),
1740 ("invalid-json", b"{not valid connection json".to_vec(), None),
1741 ];
1742
1743 for (label, contents, old_info) in cases {
1744 let (_dir, path) = temp_connection_file_path(label);
1745 write_raw_owner_only_connection_file(&path, &contents);
1746
1747 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1748 if let Some(old_info) = old_info {
1749 assert_ne!(bound.connection_info.key, old_info.key, "{label}");
1750 assert_ne!(
1751 bound.connection_info.daemon_id, old_info.daemon_id,
1752 "{label}"
1753 );
1754 }
1755 assert!(bound.connection_info.key.len() >= MIN_KEY_LEN, "{label}");
1756 assert_ne!(bound.connection_info.daemon_id, [0u8; 16], "{label}");
1757 assert_owner_only_connection_file(&path);
1758
1759 drop(bound.listeners);
1760 }
1761 }
1762
1763 #[tokio::test]
1764 async fn foreign_reused_port_connection_file_is_reclaimed_after_auth_probe_fails() {
1765 let (_dir, path) = temp_connection_file_path("foreign-reclaim");
1766 let foreign = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1767 let foreign_port = foreign.local_addr().unwrap().port();
1768 write_atomic(&path, &make_connection_info(foreign_port)).unwrap();
1769 let foreign_task = tokio::spawn(async move {
1770 if let Ok((mut stream, _)) = foreign.accept().await {
1771 let mut buf = [0u8; 64];
1772 let _ = stream.read(&mut buf).await;
1773 }
1774 });
1775
1776 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1777 assert!(bound
1778 .connection_info
1779 .endpoints
1780 .iter()
1781 .all(|endpoint| endpoint.port != foreign_port));
1782
1783 drop(bound.listeners);
1784 let _ = foreign_task.await;
1785 }
1786
1787 #[tokio::test]
1788 async fn stale_start_lock_file_is_reclaimable() {
1789 let (_dir, path) = temp_connection_file_path("start-lock-stale-file");
1790 let lock_path = start_lock_path(&path);
1791 drop(open_owner_only_lock(&lock_path).unwrap());
1792 assert!(lock_path.is_file());
1793
1794 let lock = StartLock::acquire(&path).await.unwrap();
1795 assert!(lock_path.is_file());
1796
1797 drop(lock);
1798 assert!(lock_path.is_file());
1799 }
1800
1801 #[tokio::test]
1802 async fn held_start_lock_blocks_second_acquire_until_release() {
1803 let (_dir, path) = temp_connection_file_path("start-lock-held");
1804 let lock_path = start_lock_path(&path);
1805 let first = StartLock::acquire(&path).await.unwrap();
1806
1807 let err = match StartLock::acquire(&path).await {
1808 Ok(_) => panic!("second acquire while held must stay busy"),
1809 Err(err) => err,
1810 };
1811 assert!(matches!(
1812 err,
1813 BootstrapError::StartLockBusy {
1814 ref path,
1815 attempts: START_LOCK_RETRIES,
1816 } if path == &lock_path
1817 ));
1818
1819 drop(first);
1820
1821 let second = StartLock::acquire(&path)
1822 .await
1823 .expect("released advisory lock should be reclaimable");
1824 drop(second);
1825 }
1826
1827 #[tokio::test]
1828 async fn bind_conflict_on_fixed_port_fails_loud_without_reselecting() {
1829 let (_dir, path) = temp_connection_file_path("bind-conflict");
1830 let occupied = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap();
1831 let occupied_port = occupied.local_addr().unwrap().port();
1832
1833 let err = ensure_singleton(&path, occupied_port).await.unwrap_err();
1834 assert!(matches!(
1835 err,
1836 BootstrapError::Bind { ref source, .. } if source.kind() == io::ErrorKind::AddrInUse
1837 ));
1838 assert!(err.to_string().contains("set the port in config"));
1839
1840 drop(occupied);
1841 }
1842
1843 #[tokio::test]
1844 async fn key_rotation_republishes_new_material_and_old_file_fails_auth() {
1845 let (_dir, path) = temp_connection_file_path("key-rotation");
1846 let first = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1847 let old_info = first.connection_info.clone();
1848 let fixed_port = old_info.endpoints[0].port;
1849 drop(first.listeners);
1850
1851 let second = expect_bound(ensure_singleton(&path, fixed_port).await.unwrap());
1852 let new_info = second.connection_info.clone();
1853 assert_ne!(old_info.key, new_info.key);
1854 assert_ne!(old_info.daemon_id, new_info.daemon_id);
1855 let server = start_server(second);
1856
1857 let mut old_stream = connect_from_info(&old_info).await.unwrap();
1858 let old_auth = authenticate_client(&mut old_stream, &old_info, PROBE_AUTH_DEADLINE).await;
1859 assert!(
1860 old_auth.is_err(),
1861 "old key must not authenticate after restart"
1862 );
1863
1864 let reread = connection_file::read(&path).unwrap();
1865 let mut new_stream = connect_from_info(&reread).await.unwrap();
1866 authenticate_client(&mut new_stream, &reread, PROBE_AUTH_DEADLINE)
1867 .await
1868 .unwrap();
1869
1870 server.abort();
1871 let _ = server.await;
1872 }
1873
1874 #[cfg(unix)]
1875 #[tokio::test]
1876 async fn published_connection_file_permissions_are_owner_only() {
1877 let (_dir, path) = temp_connection_file_path("permissions");
1878 let bound = expect_bound(ensure_singleton(&path, 0).await.unwrap());
1879
1880 let mode = fs::metadata(&path).unwrap().permissions().mode() & 0o777;
1881 assert_eq!(mode, 0o600);
1882
1883 drop(bound.listeners);
1884 }
1885}