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