diff --git a/crates/hearthspace-session/src/main.rs b/crates/hearthspace-session/src/main.rs index 5aa6a8f..8b9639e 100644 --- a/crates/hearthspace-session/src/main.rs +++ b/crates/hearthspace-session/src/main.rs @@ -1,6 +1,6 @@ use std::{ env, fs, - io::{Read, Write}, + io::{self, Read, Write}, net::Shutdown, os::unix::fs::PermissionsExt, os::unix::net::UnixStream, @@ -102,10 +102,19 @@ enum CompositorBackend { struct ManagedProcess { spec: ProcessSpec, child: Child, + output_activity: Arc>, restarts: usize, max_restarts: usize, } +#[derive(Debug, Clone, Copy)] +enum ChildOutputStream { + Stdout, + Stderr, +} + +const STARTUP_OUTPUT_IDLE_TIMEOUT: Duration = Duration::from_secs(120); + struct Supervisor { mode: RunMode, paths: SessionPaths, @@ -397,8 +406,7 @@ impl Supervisor { fn wait_for_settingsd_ready(&mut self) -> Result<(), Box> { info!(socket = %self.paths.settings_socket.display(), "waiting for settingsd readiness"); - let deadline = Instant::now() + Duration::from_secs(120); - while Instant::now() < deadline { + loop { if let Some(status) = poll_process(&mut self.settingsd)? { return Err(format!("settingsd exited before readiness: {status}").into()); } @@ -406,9 +414,17 @@ impl Supervisor { info!("settingsd ready"); return Ok(()); } + if process_output_idle_duration(&self.settingsd) + .is_some_and(|idle| idle >= STARTUP_OUTPUT_IDLE_TIMEOUT) + { + return Err(format!( + "timed out waiting for settingsd readiness after {}s without child output", + STARTUP_OUTPUT_IDLE_TIMEOUT.as_secs() + ) + .into()); + } thread::sleep(Duration::from_millis(50)); } - Err("timed out waiting for settingsd readiness".into()) } fn wait_for_compositor_ready(&mut self) -> Result<(), Box> { @@ -417,8 +433,7 @@ impl Supervisor { .session_dir .join(&self.session_env.wayland_display); info!(socket = %socket.display(), "waiting for compositor readiness"); - let deadline = Instant::now() + Duration::from_secs(120); - while Instant::now() < deadline { + loop { if let Some(status) = poll_process(&mut self.compositor)? { return Err(format!("compositor exited before readiness: {status}").into()); } @@ -426,9 +441,17 @@ impl Supervisor { info!("compositor ready"); return Ok(()); } + if process_output_idle_duration(&self.compositor) + .is_some_and(|idle| idle >= STARTUP_OUTPUT_IDLE_TIMEOUT) + { + return Err(format!( + "timed out waiting for compositor readiness after {}s without child output", + STARTUP_OUTPUT_IDLE_TIMEOUT.as_secs() + ) + .into()); + } thread::sleep(Duration::from_millis(50)); } - Err("timed out waiting for compositor readiness".into()) } fn monitor_until_exit(&mut self) -> Result<(), Box> { @@ -683,10 +706,12 @@ fn cargo_package_exists(package: &str) -> bool { impl ManagedProcess { fn spawn(spec: ProcessSpec, max_restarts: usize) -> Result> { - let child = spawn_child(&spec)?; + let output_activity = Arc::new(std::sync::Mutex::new(Instant::now())); + let child = spawn_child(&spec, Arc::clone(&output_activity))?; Ok(Self { spec, child, + output_activity, restarts: 0, max_restarts, }) @@ -694,12 +719,16 @@ impl ManagedProcess { fn restart(&mut self) -> Result<(), Box> { self.restarts += 1; - self.child = spawn_child(&self.spec)?; + record_output_activity(&self.output_activity); + self.child = spawn_child(&self.spec, Arc::clone(&self.output_activity))?; Ok(()) } } -fn spawn_child(spec: &ProcessSpec) -> Result> { +fn spawn_child( + spec: &ProcessSpec, + output_activity: Arc>, +) -> Result> { info!(name = spec.name, program = %spec.program, args = ?spec.args, "starting process"); let mut command = Command::new(&spec.program); for key in &spec.env_removals { @@ -709,12 +738,88 @@ fn spawn_child(spec: &ProcessSpec) -> Result> .args(&spec.args) .envs(spec.envs.iter().map(|(key, value)| (key, value))) .stdin(Stdio::null()) - .stdout(Stdio::inherit()) - .stderr(Stdio::inherit()); + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); if let Some(current_dir) = &spec.current_dir { command.current_dir(current_dir); } - Ok(command.spawn()?) + let mut child = command.spawn()?; + if let Some(stdout) = child.stdout.take() { + forward_child_output( + spec.name, + ChildOutputStream::Stdout, + Box::new(stdout), + Arc::clone(&output_activity), + ); + } + if let Some(stderr) = child.stderr.take() { + forward_child_output( + spec.name, + ChildOutputStream::Stderr, + Box::new(stderr), + output_activity, + ); + } + Ok(child) +} + +fn forward_child_output( + process_name: &'static str, + stream: ChildOutputStream, + mut reader: Box, + output_activity: Arc>, +) { + thread::spawn(move || { + let mut buffer = [0; 8192]; + loop { + match reader.read(&mut buffer) { + Ok(0) => return, + Ok(read) => { + record_output_activity(&output_activity); + if let Err(error) = write_child_output(stream, &buffer[..read]) { + warn!(process_name, ?stream, %error, "failed to forward child output"); + return; + } + } + Err(error) if error.kind() == io::ErrorKind::Interrupted => continue, + Err(error) => { + warn!(process_name, ?stream, %error, "failed to read child output"); + return; + } + } + } + }); +} + +fn write_child_output(stream: ChildOutputStream, bytes: &[u8]) -> io::Result<()> { + match stream { + ChildOutputStream::Stdout => { + let mut stdout = io::stdout().lock(); + stdout.write_all(bytes)?; + stdout.flush() + } + ChildOutputStream::Stderr => { + let mut stderr = io::stderr().lock(); + stderr.write_all(bytes)?; + stderr.flush() + } + } +} + +fn record_output_activity(output_activity: &Arc>) { + match output_activity.lock() { + Ok(mut activity) => *activity = Instant::now(), + Err(error) => *error.into_inner() = Instant::now(), + } +} + +fn process_output_idle_duration(process: &Option) -> Option { + let process = process.as_ref()?; + let last_output = match process.output_activity.lock() { + Ok(activity) => *activity, + Err(error) => *error.into_inner(), + }; + Some(Instant::now().duration_since(last_output)) } fn poll_process(