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