diff --git a/.gitignore b/.gitignore index 17da560..e4d3eb0 100644 --- a/.gitignore +++ b/.gitignore @@ -1,4 +1,5 @@ /target +/.devloop/ /examples/*/state.json **/.#* diff --git a/CHANGELOG.md b/CHANGELOG.md index eb252f6..a42493f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,14 @@ All notable changes to `devloop` will be recorded in this file. ## [Unreleased] +## [0.10.0] - 2026-07-23 + +### Added + +- `devloop run` now persists a unique session log under the state-file + directory's `logs/` subdirectory (by default `.devloop/logs/`), including + engine, process, and hook output even when terminal inheritance is disabled. + ## [0.9.3] - 2026-07-22 ### Changed diff --git a/Cargo.lock b/Cargo.lock index 60a19ed..a837fbc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -235,7 +235,7 @@ checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" [[package]] name = "devloop" -version = "0.9.3" +version = "0.10.0" dependencies = [ "anyhow", "axum", diff --git a/Cargo.toml b/Cargo.toml index 0e30176..af1edfc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "devloop" -version = "0.9.3" +version = "0.10.0" edition = "2024" [dependencies] diff --git a/README.md b/README.md index 30ae8f6..27defb3 100644 --- a/README.md +++ b/README.md @@ -130,6 +130,12 @@ The session state file is owned by `devloop` while it is running. External edits to that file are not merged back into the live session; restart the supervisor if you need to seed a different initial state. +Each `devloop run` also writes a durable, per-session log beside that +state file. With the default state-file location, logs live under +`.devloop/logs/`; add `.devloop/` to the client repository's `.gitignore`. +The [Behavior Reference](docs/behavior.md#session-logs) defines what the log +captures and how `devloop` handles it. + ## Example use case Used as the primary local development workflow for diff --git a/docs/README.md b/docs/README.md index 5e618e4..cc8c733 100644 --- a/docs/README.md +++ b/docs/README.md @@ -7,4 +7,6 @@ This directory holds detailed reference material for `devloop`. Keep the top-level `README.md` focused on purpose, installation, and the -main development loop; put configuration detail here. +main development loop; put configuration detail here. Session-log behaviour +is documented in the Behavior Reference, and its path follows the state-file +configuration in the Configuration Reference. diff --git a/docs/behavior.md b/docs/behavior.md index c995f39..8b192c3 100644 --- a/docs/behavior.md +++ b/docs/behavior.md @@ -202,6 +202,26 @@ server for browser listeners. first. - When output color is enabled, labels are colorized per source. +### Session logs + +Every `devloop run` also persists a per-session log under the `logs/` +directory beside its state file. With the default state file, that is +`.devloop/logs/`. `devloop` reports the selected path at startup. + +The persistent log contains `tracing` output and labeled managed-process and +hook output. It records that child output even when `output.inherit` hides it +from the terminal. `devloop` creates the directory and log before starting the +runtime; a creation failure stops startup rather than running without durable +evidence. + +Devloop ignores the active session-log file when classifying watch events, so +broad patterns such as `**/*` do not retrigger workflows on that log write. +Other files in the same `logs/` directory remain normal watched files. + +`devloop` does not rotate or delete session logs. The client owns retention; +add `.devloop/` to the client repository's `.gitignore` when using the default +state-file location. + ### Color rules Colorized output is enabled when stdout is a terminal and `NO_COLOR` is diff --git a/docs/configuration.md b/docs/configuration.md index d7909cc..cf994ef 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -26,6 +26,16 @@ startup_workflows = ["startup"] - `startup_workflows`: workflows to run after autostart processes have been started. +## State and session logs + +Each `devloop run` creates a unique durable log beside the state file, under +`/logs/`. The default layout is therefore +`/.devloop/logs/`. Add `.devloop/` to the client repository's +`.gitignore`; `devloop` owns both the state and log files there. Session logs +need no configuration and are created only by `devloop run`, not by +`devloop validate` or `devloop docs`. See [Behavior Reference](behavior.md) +for capture, failure, and retention rules. + Optional watcher backend config: ```toml diff --git a/src/engine.rs b/src/engine.rs index c2b4406..9df3380 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -1,5 +1,5 @@ use std::collections::{BTreeMap, BTreeSet}; -use std::path::Path; +use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; @@ -21,10 +21,12 @@ use crate::config::{CompiledWatchGroup, CompiledWatchTarget, Config, LogStyle, W use crate::core::{RuntimeEffect, RuntimeEvent, RuntimeMachine, WorkflowEffect, WorkflowMachine}; use crate::external_events::{ExternalEventMessage, ExternalEventServer}; use crate::processes::ProcessManager; +use crate::session_log::SessionLog; use crate::state::SessionState; pub struct Engine { config: Config, + session_log: SessionLog, } trait WorkflowEffectAdapter { @@ -79,8 +81,11 @@ struct LiveRuntimeAdapter<'a, 'b> { } impl Engine { - pub fn new(config: Config) -> Self { - Self { config } + pub fn new(config: Config, session_log: SessionLog) -> Self { + Self { + config, + session_log, + } } pub async fn run(self) -> Result<()> { @@ -90,16 +95,24 @@ impl Engine { .clone() .ok_or_else(|| anyhow!("state file missing after config load"))?, )?; - let mut processes = ProcessManager::new(&self.config); + let mut processes = + ProcessManager::new(&self.config).with_session_log(self.session_log.clone()); let watch_groups = self.config.compiled_watchers()?; let watched_targets = self.config.compiled_watch_targets(); + let ignored_watch_paths = vec![self.session_log.path().to_path_buf()]; let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel(); let (external_event_tx, mut external_event_rx) = tokio::sync::mpsc::unbounded_channel(); let tx_watcher = tx.clone(); + let watcher_ignored_paths = ignored_watch_paths.clone(); let watcher_shutdown = Arc::new(AtomicBool::new(false)); let watcher_shutdown_callback = watcher_shutdown.clone(); let mut watcher = create_watcher(&self.config, move |result| { - forward_watcher_event(&tx_watcher, &watcher_shutdown_callback, result); + forward_watcher_event( + &tx_watcher, + &watcher_shutdown_callback, + result, + &watcher_ignored_paths, + ); })?; let mut maintain_tick = tokio::time::interval(Duration::from_secs(1)); let mut runtime = RuntimeMachine::new(&self.config); @@ -158,8 +171,12 @@ impl Engine { tokio::time::sleep_until(deadline).await; } }, if watch_deadline.is_some() => { - let workflows = - classify_events(&self.config.root, &watch_groups, &pending_watch_events); + let workflows = classify_events( + &self.config.root, + &watch_groups, + &pending_watch_events, + &ignored_watch_paths, + ); pending_watch_events.clear(); watch_deadline = None; if !workflows.is_empty() { @@ -425,8 +442,19 @@ async fn execute_runtime_effects( fn forward_watcher_event( tx: &tokio::sync::mpsc::UnboundedSender>, shutting_down: &AtomicBool, - result: notify::Result, + mut result: notify::Result, + ignored_paths: &[PathBuf], ) { + if let Ok(event) = &mut result { + event.paths.retain(|path| { + !ignored_paths + .iter() + .any(|ignored| path_is_equivalent_to_ignored_path(path, ignored)) + }); + if event.paths.is_empty() { + return; + } + } if let Err(error) = tx.send(result) && !shutting_down.load(Ordering::Relaxed) { @@ -600,6 +628,7 @@ fn classify_events( root: &Path, watch_groups: &[CompiledWatchGroup], events: &[Event], + ignored_paths: &[PathBuf], ) -> BTreeMap> { let mut grouped: BTreeMap> = BTreeMap::new(); for event in events { @@ -607,6 +636,12 @@ fn classify_events( continue; } for path in &event.paths { + if ignored_paths + .iter() + .any(|ignored_path| path_is_equivalent_to_ignored_path(path, ignored_path)) + { + continue; + } let Some(relative) = relativize_event_path(root, path) else { continue; }; @@ -632,11 +667,36 @@ fn relativize_event_path<'a>(root: &'a Path, path: &'a Path) -> Option<&'a Path> .or_else(|| strip_private_prefix_variant(root, path)) } +fn path_is_equivalent_to_ignored_path(path: &Path, ignored_path: &Path) -> bool { + if path == ignored_path { + return true; + } + if let Some(private_path) = private_path_variant(ignored_path) + && path == private_path + { + return true; + } + if let Some(public_path) = public_path_variant(ignored_path) + && path == public_path + { + return true; + } + false +} + fn strip_private_prefix_variant<'a>(root: &'a Path, path: &'a Path) -> Option<&'a Path> { - let private_root = Path::new("/private").join(root.strip_prefix("/").ok()?); + let private_root = private_path_variant(root)?; path.strip_prefix(&private_root).ok() } +fn private_path_variant(path: &Path) -> Option { + Some(Path::new("/private").join(path.strip_prefix("/").ok()?)) +} + +fn public_path_variant(path: &Path) -> Option { + Some(Path::new("/").join(path.strip_prefix("/private").ok()?)) +} + fn normalize_path(path: &Path) -> String { path.components() .map(|component| component.as_os_str().to_string_lossy()) @@ -697,11 +757,52 @@ mod tests { attrs: Default::default(), }, ]; - let grouped = classify_events(&root, &groups, &events); + let grouped = classify_events(&root, &groups, &events, &[]); assert_eq!(grouped["server"], vec!["src/main.rs"]); assert_eq!(grouped["content"], vec!["content/posts/example.md"]); } + #[test] + fn classify_events_ignores_active_session_log_file_even_for_broad_patterns() { + let root = PathBuf::from("/tmp/example"); + let groups = vec![CompiledWatchGroup::for_test(&["**/*"], "all").expect("watch group")]; + let ignored_paths = vec![root.join(".devloop/logs/session-1.log")]; + let events = vec![ + Event { + kind: EventKind::Modify(ModifyKind::Any), + paths: vec![root.join(".devloop/logs/session-1.log")], + attrs: Default::default(), + }, + Event { + kind: EventKind::Modify(ModifyKind::Any), + paths: vec![root.join(".devloop/logs/user-owned.log")], + attrs: Default::default(), + }, + ]; + + let grouped = classify_events(&root, &groups, &events, &ignored_paths); + + assert_eq!(grouped["all"], vec![".devloop/logs/user-owned.log"]); + } + + #[test] + fn classify_events_ignores_private_path_variant_of_active_session_log_file() { + let root = PathBuf::from("/tmp/example"); + let groups = vec![CompiledWatchGroup::for_test(&["**/*"], "all").expect("watch group")]; + let ignored_paths = vec![root.join(".devloop/logs/session-1.log")]; + let events = vec![Event { + kind: EventKind::Modify(ModifyKind::Any), + paths: vec![PathBuf::from( + "/private/tmp/example/.devloop/logs/session-1.log", + )], + attrs: Default::default(), + }]; + + let grouped = classify_events(&root, &groups, &events, &ignored_paths); + + assert!(grouped.is_empty()); + } + #[test] fn resolve_watch_registration_keeps_existing_poll_file_exact() { let dir = tempdir().expect("tempdir"); @@ -838,7 +939,7 @@ mod tests { attrs: Default::default(), }]; - let grouped = classify_events(&root, &groups, &events); + let grouped = classify_events(&root, &groups, &events, &[]); assert_eq!(grouped["content"], vec!["watched.txt"]); } @@ -866,7 +967,7 @@ mod tests { }, ]; - let grouped = classify_events(&root, &groups, &events); + let grouped = classify_events(&root, &groups, &events, &[]); assert_eq!(grouped["css"], vec!["tailwind.css"]); } @@ -882,7 +983,7 @@ mod tests { attrs: Default::default(), }]; - let grouped = classify_events(&root, &groups, &events); + let grouped = classify_events(&root, &groups, &events, &[]); assert!(grouped.is_empty()); } @@ -1645,9 +1746,30 @@ mod tests { paths: vec![PathBuf::from("content/layout.html")], attrs: Default::default(), }), + &[], ); } + #[test] + fn forward_watcher_event_drops_ignored_paths_before_queueing() { + let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel(); + let shutdown = AtomicBool::new(false); + let ignored_paths = vec![PathBuf::from("/tmp/example/.devloop/logs/session-1.log")]; + + forward_watcher_event( + &tx, + &shutdown, + Ok(Event { + kind: EventKind::Modify(ModifyKind::Any), + paths: vec![PathBuf::from("/tmp/example/.devloop/logs/session-1.log")], + attrs: Default::default(), + }), + &ignored_paths, + ); + + assert!(rx.try_recv().is_err()); + } + #[tokio::test] async fn runtime_machine_does_not_run_observed_workflow_when_hook_state_is_unchanged() { let mut config = Config { diff --git a/src/main.rs b/src/main.rs index d4ba9c6..9d8dbcf 100644 --- a/src/main.rs +++ b/src/main.rs @@ -6,20 +6,24 @@ mod env_expand; mod external_events; mod output; mod processes; +mod session_log; mod state; #[cfg(test)] mod test_support; +use std::io::{self, Write}; use std::path::PathBuf; +use std::time::Duration; -use anyhow::{Result, anyhow}; +use anyhow::{Context, Result, anyhow}; use clap::{Parser, Subcommand, ValueEnum}; use pulldown_cmark::{ CodeBlockKind, Event as MarkdownEvent, HeadingLevel, Parser as MarkdownParser, Tag, TagEnd, }; -use tracing::{Event, Subscriber}; +use tracing::{Event, Subscriber, error}; use tracing_subscriber::EnvFilter; use tracing_subscriber::fmt::FmtContext; +use tracing_subscriber::fmt::MakeWriter; use tracing_subscriber::fmt::format::{FormatEvent, FormatFields, Writer}; use tracing_subscriber::fmt::time::{FormatTime, SystemTime}; use tracing_subscriber::registry::LookupSpan; @@ -27,6 +31,9 @@ use tracing_subscriber::registry::LookupSpan; use crate::config::Config; use crate::engine::Engine; use crate::output::{format_output_prefix, normalize_internal_log_label, should_colorize_output}; +use crate::session_log::SessionLog; + +const SESSION_LOG_SHUTDOWN_FLUSH_TIMEOUT: Duration = Duration::from_secs(5); #[derive(Debug, Parser)] #[command( @@ -71,15 +78,10 @@ enum DocsTopic { #[tokio::main] async fn main() -> Result<()> { - tracing_subscriber::fmt() - .with_env_filter(EnvFilter::new(default_rust_log())) - .event_format(DevloopLogFormatter::default()) - .with_writer(std::io::stderr) - .init(); - let cli = Cli::parse(); match cli.command { Command::Validate { config } => { + init_logging(None); let config = resolve_config_path(config)?; let config = Config::load(&config)?; config.validate()?; @@ -88,7 +90,19 @@ async fn main() -> Result<()> { let config = resolve_config_path(config)?; let config = Config::load(&config)?; config.validate()?; - Engine::new(config).run().await?; + let state_file = config + .state_file + .as_deref() + .ok_or_else(|| anyhow!("state file missing after config load"))?; + let session_log = SessionLog::create(state_file)?; + init_logging(Some(session_log.clone())); + announce_session_log_path(&session_log).await?; + if let Err(error) = Engine::new(config, session_log.clone()).run().await { + error!(error = %format!("{error:#}"), "devloop run failed"); + flush_session_log_before_exit(&session_log).await; + return Err(error); + } + flush_session_log_before_exit(&session_log).await; } Command::Docs { topic } => { print!("{}", render_docs_text(topic)); @@ -97,6 +111,139 @@ async fn main() -> Result<()> { Ok(()) } +async fn flush_session_log_before_exit(session_log: &SessionLog) { + match tokio::time::timeout( + SESSION_LOG_SHUTDOWN_FLUSH_TIMEOUT, + session_log.flush_queued(), + ) + .await + { + Ok(Ok(())) => {} + Ok(Err(error)) => eprintln!("devloop: failed to flush session log before exit: {error}"), + Err(_) => eprintln!("devloop: timed out flushing session log before exit"), + } +} + +async fn announce_session_log_path(session_log: &SessionLog) -> Result<()> { + let message = format!("writing session log: {}", session_log.path().display()); + session_log + .write_labeled_line("devloop", message.as_bytes()) + .with_context(|| "failed to persist session log path before runtime start")?; + session_log + .flush_queued() + .await + .with_context(|| "failed to flush session log path before runtime start")?; + let _ = writeln!(io::stderr(), "devloop: {message}"); + Ok(()) +} + +fn init_logging(session_log: Option) { + match session_log { + Some(session_log) => tracing_subscriber::fmt() + .with_env_filter(EnvFilter::new(default_rust_log())) + .event_format(DevloopLogFormatter::default()) + .with_writer(SessionLogMakeWriter { session_log }) + .init(), + None => tracing_subscriber::fmt() + .with_env_filter(EnvFilter::new(default_rust_log())) + .event_format(DevloopLogFormatter::default()) + .with_writer(std::io::stderr) + .init(), + } +} + +#[derive(Clone)] +struct SessionLogMakeWriter { + session_log: SessionLog, +} + +impl<'a> MakeWriter<'a> for SessionLogMakeWriter { + type Writer = SessionLogTeeWriter; + + fn make_writer(&'a self) -> Self::Writer { + SessionLogTeeWriter { + session_log: self.session_log.clone(), + buffer: Vec::new(), + } + } +} + +struct SessionLogTeeWriter { + session_log: SessionLog, + buffer: Vec, +} + +impl Write for SessionLogTeeWriter { + fn write(&mut self, bytes: &[u8]) -> io::Result { + self.buffer.extend_from_slice(bytes); + Ok(bytes.len()) + } + + fn flush(&mut self) -> io::Result<()> { + self.flush_buffer_to_log_and_terminal() + } +} + +impl SessionLogTeeWriter { + fn flush_buffer_to_log_and_terminal(&mut self) -> io::Result<()> { + if self.buffer.is_empty() { + return io::stderr().flush(); + } + let log_bytes = strip_ansi_csi_sequences(&self.buffer); + let log_result = self.session_log.queue_raw(log_bytes); + let terminal_result = io::stderr().write_all(&self.buffer); + if let Err(error) = log_result { + eprintln!("devloop: failed to persist tracing output: {error}"); + } + terminal_result?; + self.buffer.clear(); + io::stderr().flush() + } +} + +impl Drop for SessionLogTeeWriter { + fn drop(&mut self) { + if let Err(error) = self.flush_buffer_to_log_and_terminal() { + eprintln!("devloop: failed to flush tracing output: {error}"); + } + } +} + +#[cfg(test)] +impl SessionLogTeeWriter { + fn write_for_test(session_log: SessionLog, chunks: &[&[u8]]) -> io::Result<()> { + let mut writer = Self { + session_log, + buffer: Vec::new(), + }; + for chunk in chunks { + writer.write_all(chunk)?; + } + writer.flush() + } +} + +fn strip_ansi_csi_sequences(bytes: &[u8]) -> Vec { + let mut stripped = Vec::with_capacity(bytes.len()); + let mut index = 0; + while index < bytes.len() { + if bytes[index] == 0x1b && bytes.get(index + 1) == Some(&b'[') { + index += 2; + while index < bytes.len() { + let byte = bytes[index]; + index += 1; + if (0x40..=0x7e).contains(&byte) { + break; + } + } + continue; + } + stripped.push(bytes[index]); + index += 1; + } + stripped +} + fn default_rust_log() -> String { std::env::var("RUST_LOG").unwrap_or_else(|_| "info".to_string()) } @@ -343,14 +490,17 @@ fn format_heading(text: &str, level: HeadingLevel) -> String { #[cfg(test)] mod tests { use super::{ - Cli, DocsTopic, default_rust_log, docs_text, format_tracing_prefix, render_docs_text, - render_markdown_for_terminal, + Cli, DocsTopic, SessionLogTeeWriter, announce_session_log_path, default_rust_log, + docs_text, format_tracing_prefix, render_docs_text, render_markdown_for_terminal, + strip_ansi_csi_sequences, }; use crate::output::{ format_output_prefix, normalize_internal_log_label, normalize_source_label, }; + use crate::session_log::SessionLog; use crate::test_support::RustLogGuard; use clap::Parser; + use tempfile::tempdir; #[test] fn default_rust_log_uses_info_when_unset() { @@ -372,6 +522,67 @@ mod tests { ); } + #[tokio::test] + async fn session_log_tee_writer_appends_one_buffered_record() { + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + + SessionLogTeeWriter::write_for_test( + log.clone(), + &[b"first ".as_slice(), b"second\n".as_slice()], + ) + .expect("write tracing event"); + log.flush_queued() + .await + .expect("flush queued tracing event"); + + assert_eq!( + std::fs::read_to_string(log.path()).expect("read log"), + "first second\n" + ); + } + + #[tokio::test] + async fn announce_session_log_path_reports_persistence_failure() { + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + log.fail_for_test( + std::io::ErrorKind::BrokenPipe, + "simulated session log failure", + ); + + let error = announce_session_log_path(&log) + .await + .expect_err("announcement should report log failure"); + + assert!(format!("{error:#}").contains("failed to persist session log path")); + } + + #[tokio::test] + async fn session_log_tee_writer_strips_ansi_sequences_from_file_copy() { + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + + SessionLogTeeWriter::write_for_test( + log.clone(), + &[b"\x1b[31mred".as_slice(), b"\x1b[0m plain\n".as_slice()], + ) + .expect("write tracing event"); + log.flush_queued() + .await + .expect("flush queued tracing event"); + + assert_eq!( + std::fs::read_to_string(log.path()).expect("read log"), + "red plain\n" + ); + } + + #[test] + fn strip_ansi_csi_sequences_removes_color_codes() { + assert_eq!(strip_ansi_csi_sequences(b"a\x1b[2;31mb\x1b[0mc"), b"abc"); + } + #[test] fn tracing_prefix_wraps_dependency_targets_under_devloop() { assert_eq!( @@ -407,6 +618,17 @@ mod tests { assert!(rendered.starts_with("# Configuration Reference")); assert!(rendered.contains("startup_workflows")); + assert!(rendered.contains("## State and session logs")); + assert!(rendered.contains(".devloop/logs/")); + } + + #[test] + fn docs_text_uses_embedded_session_log_behavior_reference() { + let rendered = docs_text(DocsTopic::Behavior); + + assert!(rendered.starts_with("# Behavior Reference")); + assert!(rendered.contains("### Session logs")); + assert!(rendered.contains("output.inherit")); } #[test] diff --git a/src/processes.rs b/src/processes.rs index 503f074..30817a8 100644 --- a/src/processes.rs +++ b/src/processes.rs @@ -1,7 +1,7 @@ use std::collections::{BTreeMap, VecDeque}; use std::path::{Path, PathBuf}; use std::process::Stdio; -use std::sync::Arc; +use std::sync::{Arc, Mutex as StdMutex}; use std::time::{Duration, Instant}; use anyhow::{Context, Result, anyhow}; @@ -10,7 +10,8 @@ use rustix::io::Errno; use rustix::process::{Pid, Signal, kill_process_group}; use tokio::io::{AsyncReadExt, AsyncWriteExt, Stderr, Stdout}; use tokio::process::{Child, Command}; -use tokio::sync::Mutex; +use tokio::sync::{Mutex, mpsc}; +use tokio::task::{JoinHandle, JoinSet}; use tokio::time::{sleep, timeout}; use tracing::{info, warn}; @@ -25,6 +26,7 @@ use crate::external_events::{ExternalEventEnvironment, apply_external_event_env} use crate::output::{ dim_start, format_output_prefix_with_style, should_colorize_output, style_reset, }; +use crate::session_log::SessionLog; use crate::state::SessionState; pub struct ProcessManager<'a> { @@ -35,16 +37,35 @@ pub struct ProcessManager<'a> { stdout: Arc>, stderr: Arc>, supervisor: ProcessSupervisor, + output_cleanup_tasks: JoinSet>, + output_generations: BTreeMap>>, clock_start: Instant, external_event_env: Option, browser_reload_env: Option, + session_log: Option, } struct ManagedProcess { child: Child, process_group: Pid, + output_tasks: Vec, } +struct OutputTask { + handle: JoinHandle<()>, +} + +#[derive(Clone)] +struct OutputStateGeneration { + current: Arc>, + value: u64, +} + +const OUTPUT_DRAIN_TOTAL_TIMEOUT: Duration = Duration::from_secs(30); +const OUTPUT_DRAIN_ABORT_TIMEOUT: Duration = Duration::from_secs(1); +const TERMINAL_OUTPUT_QUEUE_CAPACITY: usize = 256; +const SESSION_LOG_FLUSH_TIMEOUT: Duration = Duration::from_secs(5); + struct CommandContext<'a> { env: &'a BTreeMap, external_event_env: Option<&'a ExternalEventEnvironment>, @@ -65,12 +86,20 @@ impl<'a> ProcessManager<'a> { stdout: Arc::new(Mutex::new(tokio::io::stdout())), stderr: Arc::new(Mutex::new(tokio::io::stderr())), supervisor: ProcessSupervisor::new(config), + output_cleanup_tasks: JoinSet::new(), + output_generations: BTreeMap::new(), clock_start: Instant::now(), external_event_env: None, browser_reload_env: None, + session_log: None, } } + pub fn with_session_log(mut self, session_log: SessionLog) -> Self { + self.session_log = Some(session_log); + self + } + pub fn set_external_event_env(&mut self, external_event_env: Option) { self.external_event_env = external_event_env; } @@ -104,6 +133,11 @@ impl<'a> ProcessManager<'a> { }; terminate_child(name, &mut child.child, child.process_group).await?; self.supervisor.on_process_stopped(name); + if let Err(error) = + wait_for_output_tasks(name, child.output_tasks, self.session_log.clone()).await + { + warn!("output cleanup failed while stopping {name}: {error:#}"); + } Ok(()) } @@ -152,20 +186,48 @@ impl<'a> ProcessManager<'a> { workflow, }, )?; - let output = command - .output() - .await - .with_context(|| format!("failed to run hook '{name}'"))?; let source_label = process_output_source_label(name, &spec.command); - self.render_hook_output(&source_label, &spec.output, &output.stdout, &output.stderr) + command.stdin(Stdio::null()); + command.stdout(Stdio::piped()); + command.stderr(Stdio::piped()); + let mut child = command + .spawn() + .with_context(|| format!("failed to run hook '{name}'"))?; + let stdout = child + .stdout + .take() + .ok_or_else(|| anyhow!("failed to capture stdout for hook '{name}'"))?; + let stderr = child + .stderr + .take() + .ok_or_else(|| anyhow!("failed to capture stderr for hook '{name}'"))?; + let stdout_task = collect_hook_output( + stdout, + self.session_log.clone(), + source_label.clone(), + name.to_owned(), + ); + let stderr_task = collect_hook_output( + stderr, + self.session_log.clone(), + source_label.clone(), + name.to_owned(), + ); + let child_status = async { child.wait().await.map_err(anyhow::Error::from) }; + let (status, stdout, stderr) = tokio::try_join!(child_status, stdout_task, stderr_task) + .with_context(|| format!("failed to run hook '{name}'"))?; + if let Some(session_log) = &self.session_log + && let Err(error) = flush_session_log_with_timeout(session_log).await + { + eprintln!("devloop: failed to flush persisted output for hook {name}: {error}"); + } + + self.render_hook_output(&source_label, &spec.output, &stdout, &stderr) .await; - if !output.status.success() { - return Err(anyhow!( - "hook '{name}' failed with status {}", - output.status - )); + if !status.success() { + return Err(anyhow!("hook '{name}' failed with status {}", status)); } - let stdout = String::from_utf8(output.stdout) + let stdout = String::from_utf8(stdout) .with_context(|| format!("hook '{name}' produced non-utf8 stdout"))?; apply_hook_capture(spec, stdout.trim(), state) } @@ -185,10 +247,32 @@ impl<'a> ProcessManager<'a> { pub async fn stop_all(&mut self, state: &SessionState) -> Result<()> { self.initiate_shutdown(); + let mut cleanup_errors = Vec::new(); for effect in self.supervisor.on_shutdown() { - self.apply_process_effect(effect, state).await?; + match effect { + ProcessEffect::StopProcess { process } => { + if let Err(error) = self.stop_named(&process).await { + cleanup_errors.push(error.context(format!( + "failed to stop process '{process}' during shutdown" + ))); + } + } + effect => { + if let Err(error) = self.apply_process_effect(effect, state).await { + cleanup_errors + .push(error.context("failed to apply shutdown process effect")); + } + } + } + } + if let Err(error) = self.finish_output_cleanup_tasks().await { + cleanup_errors.push(error.context("failed to finish output cleanup during shutdown")); + } + if cleanup_errors.is_empty() { + Ok(()) + } else { + Err(join_cleanup_errors(cleanup_errors)) } - Ok(()) } pub async fn maintain(&mut self, state: &SessionState) -> Result<()> { @@ -205,10 +289,16 @@ impl<'a> ProcessManager<'a> { if let Some(status) = exited { warn!("process {} exited with {}", name, status); - self.children.remove(&name); + let managed = self + .children + .remove(&name) + .ok_or_else(|| anyhow!("missing exited process '{name}'"))?; + signal_process_group(&name, managed.process_group, Signal::KILL)?; + self.spawn_output_cleanup(name.clone(), managed.output_tasks); exits.push((name, status.success())); } } + self.reap_output_cleanup_tasks(); let now_ms = self.clock_start.elapsed().as_millis() as u64; for effect in self.supervisor.on_tick(self.config, now_ms, exits) { self.apply_process_effect(effect, state).await?; @@ -216,6 +306,59 @@ impl<'a> ProcessManager<'a> { Ok(()) } + fn spawn_output_cleanup(&mut self, name: String, output_tasks: Vec) { + let session_log = self.session_log.clone(); + self.output_cleanup_tasks + .spawn(async move { wait_for_output_tasks(&name, output_tasks, session_log).await }); + } + + fn reap_output_cleanup_tasks(&mut self) { + while let Some(result) = self.output_cleanup_tasks.try_join_next() { + match result { + Ok(Ok(())) => {} + Ok(Err(error)) => warn!("output cleanup failed: {error:#}"), + Err(error) => warn!("output cleanup task panicked: {error}"), + } + } + } + + fn next_output_state_generation( + &mut self, + name: &str, + rules: &[OutputRule], + state: &SessionState, + ) -> Result { + let current = self + .output_generations + .entry(name.to_owned()) + .or_insert_with(|| Arc::new(StdMutex::new(0))) + .clone(); + let mut generation = current + .lock() + .map_err(|_| anyhow!("output generation mutex for process '{name}' was poisoned"))?; + *generation += 1; + let value = *generation; + clear_output_state_keys(rules, state)?; + drop(generation); + Ok(OutputStateGeneration { current, value }) + } + + async fn finish_output_cleanup_tasks(&mut self) -> Result<()> { + let mut cleanup_errors = Vec::new(); + while let Some(result) = self.output_cleanup_tasks.join_next().await { + match result { + Ok(Ok(())) => {} + Ok(Err(error)) => cleanup_errors.push(error), + Err(error) => cleanup_errors.push(anyhow!("output cleanup task panicked: {error}")), + } + } + if cleanup_errors.is_empty() { + Ok(()) + } else { + Err(join_cleanup_errors(cleanup_errors)) + } + } + async fn start(&mut self, name: &str, spec: &ProcessSpec, state: &SessionState) -> Result<()> { if self.children.contains_key(name) { return Ok(()); @@ -240,7 +383,8 @@ impl<'a> ProcessManager<'a> { .spawn() .with_context(|| format!("failed to start process '{name}'"))?; let process_group = child_process_group(name, &child)?; - clear_output_state_keys(&spec.output.rules, state)?; + let output_generation = + self.next_output_state_generation(name, &spec.output.rules, state)?; let process_name = name.to_owned(); let source_label = process_output_source_label(name, &spec.command); let inherit_output = spec.output.inherit; @@ -248,39 +392,49 @@ impl<'a> ProcessManager<'a> { let rules = compile_output_rules(&spec.output.rules)?; let stdout_sink = OutputSink::Stdout(self.stdout.clone()); let stderr_sink = OutputSink::Stderr(self.stderr.clone()); + let mut output_tasks = Vec::new(); if let Some(stdout) = child.stdout.take() { - tokio::spawn(forward_output_lines( - stdout, - ForwardOutputConfig { - output: stdout_sink, - source_label: source_label.clone(), - inherit_output, - body_style, - }, - process_name.clone(), - rules.clone(), - state.clone(), - )); + output_tasks.push(OutputTask { + handle: tokio::spawn(forward_output_lines( + stdout, + ForwardOutputConfig { + output: stdout_sink, + source_label: source_label.clone(), + inherit_output, + body_style, + session_log: self.session_log.clone(), + }, + process_name.clone(), + rules.clone(), + state.clone(), + output_generation.clone(), + )), + }); } if let Some(stderr) = child.stderr.take() { - tokio::spawn(forward_output_lines( - stderr, - ForwardOutputConfig { - output: stderr_sink, - source_label, - inherit_output, - body_style, - }, - process_name, - rules, - state.clone(), - )); + output_tasks.push(OutputTask { + handle: tokio::spawn(forward_output_lines( + stderr, + ForwardOutputConfig { + output: stderr_sink, + source_label, + inherit_output, + body_style, + session_log: self.session_log.clone(), + }, + process_name, + rules, + state.clone(), + output_generation, + )), + }); } self.children.insert( name.to_owned(), ManagedProcess { child, process_group, + output_tasks, }, ); self.supervisor.on_process_started(name); @@ -377,6 +531,158 @@ impl<'a> ProcessManager<'a> { } } +async fn wait_for_output_tasks( + name: &str, + output_tasks: Vec, + session_log: Option, +) -> Result<()> { + let mut drain_tasks = JoinSet::new(); + for output_task in output_tasks { + let name = name.to_owned(); + let session_log = session_log.clone(); + drain_tasks + .spawn(async move { wait_for_output_task(name, output_task, session_log).await }); + } + + while let Some(result) = drain_tasks.join_next().await { + result.with_context(|| format!("output drain task for process '{name}' panicked"))??; + } + if let Some(session_log) = session_log { + flush_session_log_with_timeout(&session_log).await?; + } + Ok(()) +} + +async fn flush_session_log_with_timeout(session_log: &SessionLog) -> Result<()> { + timeout(SESSION_LOG_FLUSH_TIMEOUT, session_log.flush_queued()) + .await + .map_err(|_| anyhow!("timed out flushing session log"))? + .map_err(anyhow::Error::from) +} + +async fn wait_for_output_task( + name: String, + output_task: OutputTask, + session_log: Option, +) -> Result<()> { + wait_for_output_task_with_deadline(name, output_task, session_log, OUTPUT_DRAIN_TOTAL_TIMEOUT) + .await +} + +async fn wait_for_output_task_with_deadline( + name: String, + mut output_task: OutputTask, + session_log: Option, + drain_timeout: Duration, +) -> Result<()> { + let drain_timeout_ms = drain_timeout.as_millis(); + let drain_timeout = sleep(drain_timeout); + tokio::pin!(drain_timeout); + + tokio::select! { + result = &mut output_task.handle => { + result.with_context(|| { + format!("output forwarding task for process '{name}' panicked") + })?; + Ok(()) + } + () = &mut drain_timeout => { + output_task.handle.abort(); + if timeout(OUTPUT_DRAIN_ABORT_TIMEOUT, &mut output_task.handle) + .await + .is_err() + { + warn!( + "timed out waiting for aborted output forwarding task for process '{name}'" + ); + } + report_abandoned_output_drain( + &name, + &format!( + "process output may be truncated: output for {} did not finish draining within {} ms; abandoned a forwarding task", + name, + drain_timeout_ms + ), + session_log.as_ref(), + ); + Ok(()) + } + } +} + +fn report_abandoned_output_drain(name: &str, message: &str, session_log: Option<&SessionLog>) { + if let Some(session_log) = session_log + && let Err(error) = session_log.write_labeled_line("devloop", message.as_bytes()) + { + eprintln!("devloop: failed to persist output truncation marker for {name}: {error}"); + } + eprintln!("devloop: {message}"); +} + +fn join_cleanup_errors(errors: Vec) -> anyhow::Error { + let messages = errors + .into_iter() + .map(|error| format!("{error:#}")) + .collect::>() + .join("; "); + anyhow!("shutdown cleanup failed: {messages}") +} + +async fn collect_hook_output( + reader: T, + session_log: Option, + source_label: String, + hook_name: String, +) -> Result> +where + T: tokio::io::AsyncRead + Unpin, +{ + let mut reader = reader; + let mut chunk = [0_u8; 4096]; + let mut output = Vec::new(); + let mut session_log_line = Vec::new(); + let mut session_log_last_was_carriage_return = false; + let mut session_log = session_log; + + loop { + let bytes_read = reader + .read(&mut chunk) + .await + .with_context(|| format!("failed to read output for hook '{hook_name}'"))?; + if bytes_read == 0 { + break; + } + + output.extend_from_slice(&chunk[..bytes_read]); + for &byte in &chunk[..bytes_read] { + if let Some(log) = &session_log + && let Err(error) = persist_output_byte_blocking( + log, + &source_label, + byte, + &mut session_log_line, + &mut session_log_last_was_carriage_return, + ) + .await + { + eprintln!("devloop: failed to persist output for hook {hook_name}: {error}"); + session_log = None; + } + } + } + + if let Some(log) = &session_log + && !session_log_line.is_empty() + && let Err(error) = log + .queue_labeled_line(&source_label, session_log_line) + .await + { + eprintln!("devloop: failed to persist output for hook {hook_name}: {error}"); + } + + Ok(output) +} + #[derive(Clone)] struct CompiledOutputRule { regex: Option, @@ -390,6 +696,7 @@ struct ForwardOutputConfig { source_label: String, inherit_output: bool, body_style: OutputBodyStyle, + session_log: Option, } async fn forward_output_lines( @@ -398,6 +705,7 @@ async fn forward_output_lines( process_name: String, rules: Vec, state: SessionState, + output_generation: OutputStateGeneration, ) where T: tokio::io::AsyncRead + Unpin, { @@ -406,8 +714,13 @@ async fn forward_output_lines( let mut chunk = [0_u8; 4096]; let mut line_buffer = Vec::new(); let mut render_state = OutputRenderState::default(); + let mut terminal_output = config + .inherit_output + .then(|| TerminalOutput::start(config.output.clone())); let mut last_was_carriage_return = false; - let mut output_failed = false; + let mut session_log_line = Vec::new(); + let mut session_log_last_was_carriage_return = false; + let mut session_log = config.session_log; loop { let bytes_read = match reader.read(&mut chunk).await { @@ -420,41 +733,120 @@ async fn forward_output_lines( }; for &byte in &chunk[..bytes_read] { - if config.inherit_output - && !output_failed - && let Err(error) = write_output_byte( - &config.output, + if let Some(log) = &session_log + && let Err(error) = persist_output_byte_blocking( + log, + &config.source_label, + byte, + &mut session_log_line, + &mut session_log_last_was_carriage_return, + ) + .await + { + eprintln!("devloop: failed to persist output for {process_name}: {error}"); + session_log = None; + } + if let Some(terminal_output) = &mut terminal_output + && let Some(record) = prepare_output_byte( &config.source_label, byte, colorize, config.body_style, &mut render_state, ) - .await { - warn!("failed to write output for {}: {}", process_name, error); - output_failed = true; + terminal_output.queue(record, &process_name).await; } - process_output_byte_for_rules( + process_output_byte_for_rules_guarded( &process_name, byte, &mut line_buffer, &mut last_was_carriage_return, &rules, &state, + Some(&output_generation), ); } } if !line_buffer.is_empty() { - process_output_line(&process_name, &line_buffer, &rules, &state); + process_output_line_guarded( + &process_name, + &line_buffer, + &rules, + &state, + Some(&output_generation), + ); } - if config.inherit_output - && !output_failed - && let Err(error) = flush_rendered_output(&config.output, &mut render_state, false).await + if let Some(log) = &session_log + && !session_log_line.is_empty() + && let Err(error) = log + .queue_labeled_line(&config.source_label, session_log_line) + .await { - warn!("failed to flush output for {}: {}", process_name, error); + eprintln!("devloop: failed to persist output for {process_name}: {error}"); + } + if let Some(mut terminal_output) = terminal_output { + if let Some(record) = prepare_flush_rendered_output(&mut render_state, false) { + terminal_output.queue(record, &process_name).await; + } + terminal_output.finish(&process_name).await; + } +} + +#[cfg(test)] +fn persist_output_byte( + session_log: &SessionLog, + source_label: &str, + byte: u8, + line: &mut Vec, + last_was_carriage_return: &mut bool, +) -> std::io::Result<()> { + if byte == b'\r' { + session_log.write_labeled_line(source_label, line)?; + line.clear(); + *last_was_carriage_return = true; + return Ok(()); + } + if byte == b'\n' { + if !*last_was_carriage_return { + session_log.write_labeled_line(source_label, line)?; + line.clear(); + } + *last_was_carriage_return = false; + return Ok(()); } + *last_was_carriage_return = false; + line.push(byte); + Ok(()) +} + +async fn persist_output_byte_blocking( + session_log: &SessionLog, + source_label: &str, + byte: u8, + line: &mut Vec, + last_was_carriage_return: &mut bool, +) -> std::io::Result<()> { + if byte == b'\r' { + session_log + .queue_labeled_line(source_label, std::mem::take(line)) + .await?; + *last_was_carriage_return = true; + return Ok(()); + } + if byte == b'\n' { + if !*last_was_carriage_return { + session_log + .queue_labeled_line(source_label, std::mem::take(line)) + .await?; + } + *last_was_carriage_return = false; + return Ok(()); + } + *last_was_carriage_return = false; + line.push(byte); + Ok(()) } #[cfg(test)] @@ -517,99 +909,115 @@ enum OutputSink { Stderr(Arc>), } -async fn write_output_byte( - output: &OutputSink, - source_label: &str, - byte: u8, - colorize: bool, - body_style: OutputBodyStyle, - render_state: &mut OutputRenderState, -) -> std::io::Result<()> { - match output { - OutputSink::Stdout(writer) => { - write_output_byte_to_writer( - writer, - source_label, - byte, - colorize, - body_style, - render_state, - ) - .await +struct TerminalOutput { + records: Option>>, + handle: Option>>, +} + +impl TerminalOutput { + fn start(output: OutputSink) -> Self { + let (records, mut received_records) = + mpsc::channel::>(TERMINAL_OUTPUT_QUEUE_CAPACITY); + let handle = tokio::spawn(async move { + while let Some(record) = received_records.recv().await { + write_output_record(&output, &record).await?; + } + Ok(()) + }); + Self { + records: Some(records), + handle: Some(handle), } - OutputSink::Stderr(writer) => { - write_output_byte_to_writer( - writer, - source_label, - byte, - colorize, - body_style, - render_state, - ) - .await + } + + async fn queue(&mut self, record: Vec, process_name: &str) { + let Some(records) = &self.records else { + return; + }; + match records.try_send(record) { + Ok(()) => {} + Err(mpsc::error::TrySendError::Full(record)) => { + self.queue_after_backpressure(record, process_name).await + } + Err(mpsc::error::TrySendError::Closed(_)) => { + self.records = None; + warn!( + "terminal output writer for {process_name} stopped; suppressing inherited terminal output while persistent logging continues" + ); + } } } -} -async fn flush_rendered_output( - output: &OutputSink, - render_state: &mut OutputRenderState, - add_newline: bool, -) -> std::io::Result<()> { - match output { - OutputSink::Stdout(writer) => { - flush_rendered_output_to_writer(writer, render_state, add_newline).await + async fn queue_after_backpressure(&mut self, record: Vec, process_name: &str) { + tokio::task::yield_now().await; + let Some(records) = &self.records else { + return; + }; + if records.send(record).await.is_err() { + self.records = None; + warn!( + "terminal output writer for {process_name} stopped; suppressing inherited terminal output while persistent logging continues" + ); } - OutputSink::Stderr(writer) => { - flush_rendered_output_to_writer(writer, render_state, add_newline).await + } + + async fn finish(&mut self, process_name: &str) { + self.records = None; + let Some(handle) = self.handle.take() else { + return; + }; + match handle.await { + Ok(Ok(())) => {} + Ok(Err(error)) => { + warn!("failed to write output for {}: {}", process_name, error); + } + Err(error) => { + warn!("terminal output task for {process_name} panicked: {error}"); + } } } } -async fn write_captured_output_to_writer( +impl Drop for TerminalOutput { + fn drop(&mut self) { + if let Some(handle) = &self.handle { + handle.abort(); + } + } +} + +async fn write_output_record(output: &OutputSink, record: &[u8]) -> std::io::Result<()> { + match output { + OutputSink::Stdout(writer) => write_output_record_to_writer(writer, record).await, + OutputSink::Stderr(writer) => write_output_record_to_writer(writer, record).await, + } +} + +async fn write_output_record_to_writer( writer: &Arc>, - source_label: &str, - bytes: &[u8], - body_style: OutputBodyStyle, + record: &[u8], ) -> std::io::Result<()> where W: AsyncWriteExt + Unpin + Send, { - let colorize = should_colorize_output(); - let mut render_state = OutputRenderState::default(); - - for &byte in bytes { - write_output_byte_to_writer( - writer, - source_label, - byte, - colorize, - body_style, - &mut render_state, - ) - .await?; - } - - flush_rendered_output_to_writer(writer, &mut render_state, false).await + let mut writer = writer.lock().await; + writer.write_all(record).await?; + writer.flush().await } -async fn write_output_byte_to_writer( - writer: &Arc>, +fn prepare_output_byte( source_label: &str, byte: u8, colorize: bool, body_style: OutputBodyStyle, render_state: &mut OutputRenderState, -) -> std::io::Result<()> -where - W: AsyncWriteExt + Unpin + Send, -{ +) -> Option> { if byte == b'\r' { render_state.body_style = body_style; render_state.colorize = colorize; - flush_rendered_output_to_writer(writer, render_state, true).await?; + let record = prepare_flush_rendered_output(render_state, true); render_state.last_was_carriage_return = true; - return Ok(()); + return record; } if byte == b'\n' { @@ -617,10 +1025,9 @@ where render_state.colorize = colorize; if render_state.last_was_carriage_return { render_state.last_was_carriage_return = false; - return Ok(()); + return None; } - flush_rendered_output_to_writer(writer, render_state, true).await?; - return Ok(()); + return prepare_flush_rendered_output(render_state, true); } render_state.last_was_carriage_return = false; @@ -637,24 +1044,20 @@ where let rendered = render_output_byte(byte, colorize, body_style, render_state); render_state.rendered_line.push_str(&rendered); - Ok(()) + None } -async fn flush_rendered_output_to_writer( - writer: &Arc>, +fn prepare_flush_rendered_output( render_state: &mut OutputRenderState, add_newline: bool, -) -> std::io::Result<()> -where - W: AsyncWriteExt + Unpin + Send, -{ +) -> Option> { flush_pending_utf8(render_state); if render_state.rendered_line.is_empty() && !add_newline { - return Ok(()); + return None; } - let mut writer = writer.lock().await; + let mut record = Vec::new(); if !render_state.rendered_line.is_empty() { if render_state.dim_active { render_state @@ -662,16 +1065,73 @@ where .push_str(style_reset(render_state.colorize)); render_state.dim_active = false; } - writer - .write_all(render_state.rendered_line.as_bytes()) - .await?; + record.extend_from_slice(render_state.rendered_line.as_bytes()); render_state.rendered_line.clear(); } if add_newline { - writer.write_all(b"\n").await?; + record.push(b'\n'); } - writer.flush().await?; render_state.at_line_start = true; + Some(record) +} + +async fn flush_rendered_output_to_writer( + writer: &Arc>, + render_state: &mut OutputRenderState, + add_newline: bool, +) -> std::io::Result<()> +where + W: AsyncWriteExt + Unpin + Send, +{ + let Some(record) = prepare_flush_rendered_output(render_state, add_newline) else { + return Ok(()); + }; + write_output_record_to_writer(writer, &record).await +} + +async fn write_captured_output_to_writer( + writer: &Arc>, + source_label: &str, + bytes: &[u8], + body_style: OutputBodyStyle, +) -> std::io::Result<()> +where + W: AsyncWriteExt + Unpin + Send, +{ + let colorize = should_colorize_output(); + let mut render_state = OutputRenderState::default(); + + for &byte in bytes { + write_output_byte_to_writer( + writer, + source_label, + byte, + colorize, + body_style, + &mut render_state, + ) + .await?; + } + + flush_rendered_output_to_writer(writer, &mut render_state, false).await +} + +async fn write_output_byte_to_writer( + writer: &Arc>, + source_label: &str, + byte: u8, + colorize: bool, + body_style: OutputBodyStyle, + render_state: &mut OutputRenderState, +) -> std::io::Result<()> +where + W: AsyncWriteExt + Unpin + Send, +{ + if let Some(record) = + prepare_output_byte(source_label, byte, colorize, body_style, render_state) + { + write_output_record_to_writer(writer, &record).await?; + } Ok(()) } @@ -750,6 +1210,7 @@ fn take_utf8_buffer_lossy(render_state: &mut OutputRenderState) -> String { rendered } +#[cfg(test)] fn process_output_byte_for_rules( process_name: &str, byte: u8, @@ -757,9 +1218,29 @@ fn process_output_byte_for_rules( last_was_carriage_return: &mut bool, rules: &[CompiledOutputRule], state: &SessionState, +) { + process_output_byte_for_rules_guarded( + process_name, + byte, + line_buffer, + last_was_carriage_return, + rules, + state, + None, + ); +} + +fn process_output_byte_for_rules_guarded( + process_name: &str, + byte: u8, + line_buffer: &mut Vec, + last_was_carriage_return: &mut bool, + rules: &[CompiledOutputRule], + state: &SessionState, + output_generation: Option<&OutputStateGeneration>, ) { if byte == b'\r' { - process_output_line(process_name, line_buffer, rules, state); + process_output_line_guarded(process_name, line_buffer, rules, state, output_generation); line_buffer.clear(); *last_was_carriage_return = true; return; @@ -767,7 +1248,7 @@ fn process_output_byte_for_rules( if byte == b'\n' { if !*last_was_carriage_return { - process_output_line(process_name, line_buffer, rules, state); + process_output_line_guarded(process_name, line_buffer, rules, state, output_generation); line_buffer.clear(); } *last_was_carriage_return = false; @@ -778,12 +1259,30 @@ fn process_output_byte_for_rules( line_buffer.push(byte); } -fn process_output_line( +fn process_output_line_guarded( process_name: &str, bytes: &[u8], rules: &[CompiledOutputRule], state: &SessionState, + output_generation: Option<&OutputStateGeneration>, ) { + let _generation_guard = match output_generation { + Some(generation) => { + let guard = match generation.current.lock() { + Ok(guard) => guard, + Err(error) => { + warn!("output generation mutex for {process_name} was poisoned: {error}"); + return; + } + }; + if *guard != generation.value { + return; + } + Some(guard) + } + None => None, + }; + let line = String::from_utf8_lossy(bytes) .trim_end_matches(['\n', '\r']) .to_owned(); @@ -1076,89 +1575,315 @@ fn timeout_error(name: &str, probe: &ProbeSpec) -> anyhow::Error { #[cfg(test)] mod tests { use super::*; - use crate::config::{OutputConfig, OutputExtract, ProbeSpec}; + use crate::config::{OutputConfig, OutputExtract, OutputRule, ProbeSpec}; use crate::test_support::{EnvVarGuard, RustLogGuard}; use rustix::process::test_kill_process; use serde_json::Value; use std::sync::Arc; + use std::sync::atomic::{AtomicU64, Ordering}; use std::time::{SystemTime, UNIX_EPOCH}; use tempfile::tempdir; - use tokio::io::AsyncReadExt; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::sync::Mutex; + static TEST_STATE_SEQUENCE: AtomicU64 = AtomicU64::new(0); + fn unique_state_path() -> PathBuf { let unique = SystemTime::now() .duration_since(UNIX_EPOCH) .expect("system time") .as_nanos(); - std::env::temp_dir().join(format!("devloop-process-state-{unique}.json")) + let sequence = TEST_STATE_SEQUENCE.fetch_add(1, Ordering::Relaxed); + std::env::temp_dir().join(format!("devloop-process-state-{unique}-{sequence}.json")) + } + + fn test_config(root: &Path) -> Config { + Config { + root: root.to_path_buf(), + debounce_ms: 100, + watcher: crate::config::WatcherConfig::default(), + state_file: Some(unique_state_path()), + startup_workflows: vec![], + watch: BTreeMap::new(), + process: BTreeMap::new(), + hook: BTreeMap::new(), + event_server: crate::config::EventServerConfig::default(), + browser_reload_server: crate::config::BrowserReloadServerConfig::default(), + event: BTreeMap::new(), + workflow: BTreeMap::new(), + } + } + + #[cfg(unix)] + async fn wait_for_path(path: &Path) { + let started = Instant::now(); + loop { + if path.exists() { + return; + } + assert!( + started.elapsed() < Duration::from_secs(5), + "timed out waiting for {}", + path.display() + ); + sleep(Duration::from_millis(20)).await; + } + } + + #[cfg(unix)] + async fn assert_process_gone(raw_pid: i32) { + let pid = Pid::from_raw(raw_pid).expect("pid"); + let started = Instant::now(); + loop { + if matches!(test_kill_process(pid), Err(Errno::SRCH)) { + return; + } + assert!( + started.elapsed() < Duration::from_secs(5), + "process {raw_pid} still exists after managed process stop" + ); + sleep(Duration::from_millis(20)).await; + } } - fn test_config(root: &Path) -> Config { - Config { - root: root.to_path_buf(), - debounce_ms: 100, - watcher: crate::config::WatcherConfig::default(), - state_file: Some(unique_state_path()), - startup_workflows: vec![], - watch: BTreeMap::new(), - process: BTreeMap::new(), - hook: BTreeMap::new(), - event_server: crate::config::EventServerConfig::default(), - browser_reload_server: crate::config::BrowserReloadServerConfig::default(), - event: BTreeMap::new(), - workflow: BTreeMap::new(), - } + #[cfg(unix)] + #[tokio::test] + async fn stop_named_terminates_descendant_processes() { + let dir = tempdir().expect("tempdir"); + let pid_path = dir.path().join("descendant.pid"); + let script_path = dir.path().join("spawn-descendant.sh"); + std::fs::write( + &script_path, + format!( + r#"#!/bin/sh +sh -c 'trap "" TERM; while :; do sleep 1; done' & +echo "$!" > "{}" +wait +"#, + pid_path.display() + ), + ) + .expect("write script"); + + let mut config = test_config(dir.path()); + config.process.insert( + "server".into(), + ProcessSpec { + command: vec!["sh".into(), script_path.display().to_string()], + cwd: Some(dir.path().to_path_buf()), + autostart: false, + readiness: None, + liveness: None, + restart: crate::config::RestartPolicy::Never, + env: BTreeMap::new(), + output: OutputConfig::default(), + }, + ); + let state = SessionState::load(unique_state_path()).expect("state"); + let mut manager = ProcessManager::new(&config); + + manager + .start_named("server", &state) + .await + .expect("start process"); + wait_for_path(&pid_path).await; + let descendant_pid = std::fs::read_to_string(&pid_path) + .expect("read pid") + .trim() + .parse::() + .expect("parse pid"); + + manager.stop_named("server").await.expect("stop process"); + + assert_process_gone(descendant_pid).await; + } + + #[cfg(unix)] + #[tokio::test] + async fn stop_named_waits_for_final_process_output_to_reach_the_session_log() { + let dir = tempdir().expect("tempdir"); + let started_path = dir.path().join("started"); + let script_path = dir.path().join("final-output.sh"); + std::fs::write( + &script_path, + format!( + r#"#!/bin/sh +touch "{}" +trap 'printf final-output; exit 0' TERM +while :; do sleep 1; done +"#, + started_path.display() + ), + ) + .expect("write script"); + + let mut config = test_config(dir.path()); + config.process.insert( + "server".into(), + ProcessSpec { + command: vec!["sh".into(), script_path.display().to_string()], + cwd: Some(dir.path().to_path_buf()), + autostart: false, + readiness: None, + liveness: None, + restart: crate::config::RestartPolicy::Never, + env: BTreeMap::new(), + output: OutputConfig { + inherit: false, + ..OutputConfig::default() + }, + }, + ); + let state_file = dir.path().join(".devloop/state.json"); + let state = SessionState::load(state_file.clone()).expect("state"); + let log = SessionLog::create(&state_file).expect("create log"); + let mut manager = ProcessManager::new(&config).with_session_log(log.clone()); + + manager + .start_named("server", &state) + .await + .expect("start process"); + wait_for_path(&started_path).await; + manager.stop_named("server").await.expect("stop process"); + + assert!( + std::fs::read_to_string(log.path()) + .expect("read log") + .contains("[sh server] final-output") + ); } #[cfg(unix)] - async fn wait_for_path(path: &Path) { - let started = Instant::now(); - loop { - if path.exists() { - return; - } - assert!( - started.elapsed() < Duration::from_secs(5), - "timed out waiting for {}", - path.display() - ); - sleep(Duration::from_millis(20)).await; - } + #[tokio::test] + async fn restart_named_continues_after_session_log_flush_failure() { + let dir = tempdir().expect("tempdir"); + let script_path = dir.path().join("restartable.sh"); + std::fs::write( + &script_path, + r#"#!/bin/sh +printf 'started\n' +while :; do sleep 1; done +"#, + ) + .expect("write script"); + + let mut config = test_config(dir.path()); + config.process.insert( + "server".into(), + ProcessSpec { + command: vec!["sh".into(), script_path.display().to_string()], + cwd: Some(dir.path().to_path_buf()), + autostart: false, + readiness: Some(ProbeSpec::StateKey { + key: "server_started".into(), + interval_ms: 10, + timeout_ms: 5000, + }), + liveness: None, + restart: crate::config::RestartPolicy::Never, + env: BTreeMap::new(), + output: OutputConfig { + inherit: false, + rules: vec![OutputRule { + state_key: "server_started".into(), + pattern: Some("(started)".into()), + extract: OutputExtract::Regex, + capture_group: 1, + }], + ..OutputConfig::default() + }, + }, + ); + let state_file = dir.path().join(".devloop/state.json"); + let state = SessionState::load(state_file.clone()).expect("state"); + let log = SessionLog::create(&state_file).expect("create log"); + let mut manager = ProcessManager::new(&config).with_session_log(log.clone()); + + manager + .start_named("server", &state) + .await + .expect("start process"); + manager + .wait_for_named("server", &state) + .await + .expect("wait for first process readiness"); + log.fail_for_test( + std::io::ErrorKind::BrokenPipe, + "simulated session log failure", + ); + + manager + .restart_named("server", &state) + .await + .expect("restart process"); + + manager + .wait_for_named("server", &state) + .await + .expect("wait for restarted process readiness"); + assert!(manager.children.contains_key("server")); + manager.stop_named("server").await.expect("stop process"); } #[cfg(unix)] - async fn assert_process_gone(raw_pid: i32) { - let pid = Pid::from_raw(raw_pid).expect("pid"); - let started = Instant::now(); - loop { - if matches!(test_kill_process(pid), Err(Errno::SRCH)) { - return; + #[tokio::test] + async fn maintain_does_not_block_when_exited_wrapper_leaves_writing_descendant() { + let dir = tempdir().expect("tempdir"); + let script_path = dir.path().join("writing-descendant.sh"); + std::fs::write( + &script_path, + r#"#!/bin/sh +(while :; do printf x; done) & +exit 0 +"#, + ) + .expect("write script"); + + let mut config = test_config(dir.path()); + config.process.insert( + "server".into(), + ProcessSpec { + command: vec!["sh".into(), script_path.display().to_string()], + cwd: Some(dir.path().to_path_buf()), + autostart: false, + readiness: None, + liveness: None, + restart: crate::config::RestartPolicy::Never, + env: BTreeMap::new(), + output: OutputConfig { + inherit: false, + ..OutputConfig::default() + }, + }, + ); + let state_file = dir.path().join(".devloop/state.json"); + let state = SessionState::load(state_file).expect("state"); + let mut manager = ProcessManager::new(&config); + + manager + .start_named("server", &state) + .await + .expect("start process"); + timeout(Duration::from_secs(5), async { + while manager.children.contains_key("server") { + manager.maintain(&state).await.expect("maintain process"); } - assert!( - started.elapsed() < Duration::from_secs(5), - "process {raw_pid} still exists after managed process stop" - ); - sleep(Duration::from_millis(20)).await; - } + }) + .await + .expect("maintain should not block on descendant output"); } #[cfg(unix)] #[tokio::test] - async fn stop_named_terminates_descendant_processes() { + async fn stop_all_waits_for_output_cleanup_scheduled_by_natural_exit() { let dir = tempdir().expect("tempdir"); - let pid_path = dir.path().join("descendant.pid"); - let script_path = dir.path().join("spawn-descendant.sh"); + let script_path = dir.path().join("natural-output.sh"); std::fs::write( &script_path, - format!( - r#"#!/bin/sh -sh -c 'trap "" TERM; while :; do sleep 1; done' & -echo "$!" > "{}" -wait + r#"#!/bin/sh +printf natural-output +exit 0 "#, - pid_path.display() - ), ) .expect("write script"); @@ -1173,26 +1898,160 @@ wait liveness: None, restart: crate::config::RestartPolicy::Never, env: BTreeMap::new(), - output: OutputConfig::default(), + output: OutputConfig { + inherit: false, + ..OutputConfig::default() + }, }, ); - let state = SessionState::load(unique_state_path()).expect("state"); - let mut manager = ProcessManager::new(&config); + let state_file = dir.path().join(".devloop/state.json"); + let state = SessionState::load(state_file.clone()).expect("state"); + let log = SessionLog::create(&state_file).expect("create log"); + let mut manager = ProcessManager::new(&config).with_session_log(log.clone()); manager .start_named("server", &state) .await .expect("start process"); - wait_for_path(&pid_path).await; - let descendant_pid = std::fs::read_to_string(&pid_path) - .expect("read pid") - .trim() - .parse::() - .expect("parse pid"); + while manager.children.contains_key("server") { + manager.maintain(&state).await.expect("maintain process"); + } + manager.stop_all(&state).await.expect("stop all"); - manager.stop_named("server").await.expect("stop process"); + assert!( + std::fs::read_to_string(log.path()) + .expect("read log") + .contains("[sh server] natural-output") + ); + } - assert_process_gone(descendant_pid).await; + #[cfg(unix)] + #[tokio::test] + async fn stop_all_treats_session_log_flush_failure_as_non_fatal() { + let dir = tempdir().expect("tempdir"); + let first_started = dir.path().join("first-started"); + let second_started = dir.path().join("second-started"); + let first_script = dir.path().join("first.sh"); + let second_script = dir.path().join("second.sh"); + std::fs::write( + &first_script, + format!( + r#"#!/bin/sh +touch "{}" +exec sleep 600 +"#, + first_started.display() + ), + ) + .expect("write first script"); + std::fs::write( + &second_script, + format!( + r#"#!/bin/sh +touch "{}" +exec sleep 600 +"#, + second_started.display() + ), + ) + .expect("write second script"); + + let mut config = test_config(dir.path()); + for (name, script) in [("first", first_script), ("second", second_script)] { + config.process.insert( + name.into(), + ProcessSpec { + command: vec!["sh".into(), script.display().to_string()], + cwd: Some(dir.path().to_path_buf()), + autostart: false, + readiness: None, + liveness: None, + restart: crate::config::RestartPolicy::Never, + env: BTreeMap::new(), + output: OutputConfig { + inherit: false, + ..OutputConfig::default() + }, + }, + ); + } + let state_file = dir.path().join(".devloop/state.json"); + let state = SessionState::load(state_file.clone()).expect("state"); + let log = SessionLog::create(&state_file).expect("create log"); + let mut manager = ProcessManager::new(&config).with_session_log(log.clone()); + + manager + .start_named("first", &state) + .await + .expect("start first process"); + manager + .start_named("second", &state) + .await + .expect("start second process"); + wait_for_path(&first_started).await; + wait_for_path(&second_started).await; + log.fail_for_test( + std::io::ErrorKind::BrokenPipe, + "simulated session log failure", + ); + + manager + .stop_all(&state) + .await + .expect("shutdown should continue despite session log failure"); + + assert!(manager.children.is_empty()); + } + + #[tokio::test] + async fn finish_output_cleanup_tasks_drains_all_failures() { + let dir = tempdir().expect("tempdir"); + let config = test_config(dir.path()); + let mut manager = ProcessManager::new(&config); + manager + .output_cleanup_tasks + .spawn(async { Err(anyhow!("first cleanup failure")) }); + manager + .output_cleanup_tasks + .spawn(async { Err(anyhow!("second cleanup failure")) }); + + let error = manager + .finish_output_cleanup_tasks() + .await + .expect_err("cleanup should fail"); + + assert!(manager.output_cleanup_tasks.is_empty()); + let error = format!("{error:#}"); + assert!(error.contains("first cleanup failure")); + assert!(error.contains("second cleanup failure")); + } + + #[tokio::test] + async fn abandoned_output_drain_writes_truncation_marker_to_session_log() { + let dir = tempdir().expect("tempdir"); + let state_file = dir.path().join(".devloop/state.json"); + let log = SessionLog::create(&state_file).expect("create log"); + let output_task = OutputTask { + handle: tokio::spawn(async { + std::future::pending::<()>().await; + }), + }; + + wait_for_output_task_with_deadline( + "server".into(), + output_task, + Some(log.clone()), + Duration::from_millis(10), + ) + .await + .expect("wait for output task"); + log.flush_queued().await.expect("flush session log"); + + assert!( + std::fs::read_to_string(log.path()) + .expect("read log") + .contains("[devloop] process output may be truncated") + ); } #[test] @@ -1254,6 +2113,110 @@ wait assert!(rendered.contains("\u{1b}[2mINF ready\u{1b}[0m")); } + #[tokio::test] + async fn persisted_process_output_is_source_labeled_when_terminal_inheritance_is_disabled() { + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + let mut line = Vec::new(); + let mut last_was_carriage_return = false; + + for byte in b"first\r\nsecond" { + persist_output_byte( + &log, + "server example", + *byte, + &mut line, + &mut last_was_carriage_return, + ) + .expect("persist output byte"); + } + log.write_labeled_line("server example", &line) + .expect("persist final output line"); + log.flush_queued().await.expect("flush session log"); + + assert_eq!( + std::fs::read_to_string(log.path()).expect("read log"), + "[server example] first\n[server example] second\n" + ); + } + + #[tokio::test] + async fn hook_output_is_persisted_when_terminal_inheritance_is_disabled() { + let dir = tempdir().expect("tempdir"); + let mut config = test_config(dir.path()); + config.hook.insert( + "capture".into(), + HookSpec { + command: vec!["sh".into(), "-c".into(), "printf durable-hook".into()], + cwd: None, + env: BTreeMap::new(), + output: HookOutputConfig { + inherit: false, + body_style: OutputBodyStyle::Dim, + }, + capture: None, + state_key: None, + observe: None, + }, + ); + let state_file = dir.path().join(".devloop/state.json"); + let state = SessionState::load(state_file.clone()).expect("state"); + let log = SessionLog::create(&state_file).expect("create log"); + let manager = ProcessManager::new(&config).with_session_log(log.clone()); + + manager + .run_hook("capture", &state, &[], "test") + .await + .expect("run hook"); + + assert_eq!( + std::fs::read_to_string(log.path()).expect("read log"), + "[sh capture] durable-hook\n" + ); + } + + #[tokio::test] + async fn hook_capture_continues_after_session_log_flush_failure() { + let dir = tempdir().expect("tempdir"); + let mut config = test_config(dir.path()); + config.hook.insert( + "capture".into(), + HookSpec { + command: vec!["sh".into(), "-c".into(), "printf captured-value".into()], + cwd: None, + env: BTreeMap::new(), + output: HookOutputConfig { + inherit: false, + body_style: OutputBodyStyle::Dim, + }, + capture: Some(crate::config::CaptureMode::Text), + state_key: Some("captured".into()), + observe: None, + }, + ); + let state_file = dir.path().join(".devloop/state.json"); + let state = SessionState::load(state_file.clone()).expect("state"); + let log = SessionLog::create(&state_file).expect("create log"); + log.fail_for_test( + std::io::ErrorKind::BrokenPipe, + "simulated session log failure", + ); + let manager = ProcessManager::new(&config).with_session_log(log); + + manager + .run_hook("capture", &state, &[], "test") + .await + .expect("run hook"); + + assert_eq!( + state + .get_string("captured") + .expect("read captured state") + .as_deref(), + Some("captured-value") + ); + } + #[test] fn render_output_byte_does_not_dim_newlines() { assert_eq!( @@ -1806,6 +2769,47 @@ wait let _ = std::fs::remove_file(state_path); } + #[tokio::test] + async fn retired_process_output_does_not_update_output_rule_state() { + let state_path = unique_state_path(); + let state = SessionState::load(state_path.clone()).expect("load state"); + let rules = vec![CompiledOutputRule { + regex: Some(Regex::new(r"(https://\S+)").expect("regex")), + state_key: "url".into(), + extract: OutputExtract::Regex, + capture_group: 1, + }]; + let (mut writer, reader) = tokio::io::duplex(128); + writer + .write_all(b"https://stale.example.test\n") + .await + .expect("write stale process output"); + drop(writer); + + forward_output_lines( + reader, + ForwardOutputConfig { + output: OutputSink::Stdout(Arc::new(Mutex::new(tokio::io::stdout()))), + source_label: "server".into(), + inherit_output: false, + body_style: OutputBodyStyle::Plain, + session_log: None, + }, + "server".into(), + rules, + state.clone(), + OutputStateGeneration { + current: Arc::new(StdMutex::new(1)), + value: 0, + }, + ) + .await; + + assert_eq!(state.get_string("url").expect("get url"), None); + + let _ = std::fs::remove_file(state_path); + } + #[tokio::test] async fn state_key_probe_reads_shared_session_state() { let state_path = unique_state_path(); diff --git a/src/session_log.rs b/src/session_log.rs new file mode 100644 index 0000000..d14a43a --- /dev/null +++ b/src/session_log.rs @@ -0,0 +1,417 @@ +use std::fs::{File, OpenOptions}; +use std::io::{self, Write}; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::{SystemTime, UNIX_EPOCH}; + +use anyhow::{Context, Result, anyhow}; +use tokio::sync::{mpsc, oneshot}; + +static SESSION_SEQUENCE: AtomicU64 = AtomicU64::new(0); +const SESSION_LOG_QUEUE_CAPACITY: usize = 1024; +const SESSION_LOG_BATCH_RECORDS: usize = 64; + +/// A durable, run-scoped record of Devloop output. +/// +/// The log lives beside the session state so a client can ignore one owned +/// directory while preserving a separate record for each `devloop run`. +#[derive(Clone)] +pub(crate) struct SessionLog { + path: PathBuf, + error_state: Arc>>, + writer: mpsc::Sender, +} + +enum SessionLogWrite { + Bytes(Vec), + Flush(oneshot::Sender>), +} + +#[derive(Clone)] +struct StoredIoError { + kind: io::ErrorKind, + message: String, +} + +impl StoredIoError { + fn from_error(error: io::Error) -> Self { + Self { + kind: error.kind(), + message: error.to_string(), + } + } + + fn to_error(&self) -> io::Error { + io::Error::new(self.kind, self.message.clone()) + } +} + +impl SessionLog { + pub(crate) fn create(state_file: &Path) -> Result { + let state_dir = state_file.parent().ok_or_else(|| { + anyhow!( + "state file '{}' has no parent directory", + state_file.display() + ) + })?; + let logs_dir = state_dir.join("logs"); + std::fs::create_dir_all(&logs_dir).with_context(|| { + format!( + "failed to create session log directory {}", + logs_dir.display() + ) + })?; + + let now = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis(); + let process_id = std::process::id(); + + for _ in 0..32 { + let sequence = SESSION_SEQUENCE.fetch_add(1, Ordering::Relaxed); + let path = logs_dir.join(format!("session-{now}-{process_id}-{sequence}.log")); + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600); + } + match options.open(&path) { + Ok(file) => { + let error_state = Arc::new(Mutex::new(None)); + let (writer, records) = mpsc::channel(SESSION_LOG_QUEUE_CAPACITY); + spawn_session_log_writer(file, error_state.clone(), records); + return Ok(Self { + path, + error_state, + writer, + }); + } + Err(error) if error.kind() == io::ErrorKind::AlreadyExists => continue, + Err(error) => { + return Err(error).with_context(|| { + format!("failed to create session log {}", path.display()) + }); + } + } + } + + Err(anyhow!( + "failed to allocate a unique session log in {}", + logs_dir.display() + )) + } + + pub(crate) fn path(&self) -> &Path { + &self.path + } + + #[cfg(test)] + pub(crate) fn write_labeled_output(&self, source_label: &str, bytes: &[u8]) -> io::Result<()> { + let mut line = Vec::new(); + let mut last_was_carriage_return = false; + + for &byte in bytes { + if byte == b'\r' { + self.write_labeled_line(source_label, &line)?; + line.clear(); + last_was_carriage_return = true; + continue; + } + if byte == b'\n' { + if !last_was_carriage_return { + self.write_labeled_line(source_label, &line)?; + line.clear(); + } + last_was_carriage_return = false; + continue; + } + last_was_carriage_return = false; + line.push(byte); + } + + if !line.is_empty() { + self.write_labeled_line(source_label, &line)?; + } + Ok(()) + } + + pub(crate) fn write_labeled_line(&self, source_label: &str, bytes: &[u8]) -> io::Result<()> { + let record = labeled_record(source_label, bytes); + self.enqueue_ordered(record) + } + + pub(crate) fn queue_raw(&self, bytes: Vec) -> io::Result<()> { + self.enqueue_ordered(bytes) + } + + fn enqueue_ordered(&self, bytes: Vec) -> io::Result<()> { + if let Some(error) = self.stored_error()? { + return Err(error); + } + match self.writer.try_send(SessionLogWrite::Bytes(bytes)) { + Ok(()) => Ok(()), + Err(mpsc::error::TrySendError::Full(record)) => self.blocking_send(record), + Err(mpsc::error::TrySendError::Closed(_)) => { + Err(io::Error::other("session log writer stopped")) + } + } + } + + pub(crate) async fn queue_labeled_line( + &self, + source_label: &str, + bytes: Vec, + ) -> io::Result<()> { + if let Some(error) = self.stored_error()? { + return Err(error); + } + self.writer + .send(SessionLogWrite::Bytes(labeled_record(source_label, &bytes))) + .await + .map_err(|_| io::Error::other("session log writer stopped")) + } + + pub(crate) async fn flush_queued(&self) -> io::Result<()> { + let (tx, rx) = oneshot::channel(); + self.writer + .send(SessionLogWrite::Flush(tx)) + .await + .map_err(|_| io::Error::other("session log writer stopped"))?; + rx.await + .map_err(|_| io::Error::other("session log writer stopped"))? + } + + fn lock_error_state(&self) -> io::Result>> { + self.error_state + .lock() + .map_err(|_| io::Error::other("session log error mutex was poisoned")) + } + + fn blocking_send(&self, record: SessionLogWrite) -> io::Result<()> { + if tokio::runtime::Handle::try_current().is_ok() { + tokio::task::block_in_place(|| { + self.writer + .blocking_send(record) + .map_err(|_| io::Error::other("session log writer stopped")) + }) + } else { + self.writer + .blocking_send(record) + .map_err(|_| io::Error::other("session log writer stopped")) + } + } + + fn stored_error(&self) -> io::Result> { + Ok(self + .lock_error_state()? + .as_ref() + .map(|error| error.to_error())) + } + + #[cfg(test)] + pub(crate) fn fail_for_test(&self, kind: io::ErrorKind, message: &str) { + *self.lock_error_state().expect("lock log error state") = Some(StoredIoError { + kind, + message: message.into(), + }); + } +} + +fn spawn_session_log_writer( + file: File, + error_state: Arc>>, + mut records: mpsc::Receiver, +) { + std::thread::Builder::new() + .name("devloop-session-log-writer".into()) + .spawn(move || { + let mut file = Some(file); + let mut pending = Vec::new(); + while let Some(record) = records.blocking_recv() { + match record { + SessionLogWrite::Bytes(bytes) => { + pending.extend_from_slice(&bytes); + drain_available_records( + &mut records, + &mut file, + &error_state, + &mut pending, + ); + } + SessionLogWrite::Flush(reply) => { + let result = flush_pending_records(&mut file, &error_state, &mut pending); + let _ = reply.send(result); + } + } + } + }) + .expect("session log writer thread must start"); +} + +fn drain_available_records( + records: &mut mpsc::Receiver, + file: &mut Option, + error_state: &Arc>>, + pending: &mut Vec, +) { + for _ in 1..SESSION_LOG_BATCH_RECORDS { + match records.try_recv() { + Ok(SessionLogWrite::Bytes(bytes)) => pending.extend_from_slice(&bytes), + Ok(SessionLogWrite::Flush(reply)) => { + let result = flush_pending_records(file, error_state, pending); + let _ = reply.send(result); + return; + } + Err(mpsc::error::TryRecvError::Empty) => break, + Err(mpsc::error::TryRecvError::Disconnected) => break, + } + } + let _ = flush_pending_records(file, error_state, pending); +} + +fn flush_pending_records( + file: &mut Option, + error_state: &Arc>>, + pending: &mut Vec, +) -> io::Result<()> { + let had_error = error_state + .lock() + .map_err(|_| io::Error::other("session log error mutex was poisoned"))? + .is_some(); + let result = write_batch_to_file(file, error_state, pending); + if let Err(error) = &result + && !had_error + { + eprintln!("devloop: failed to persist session log output: {error}"); + } + pending.clear(); + result +} + +fn labeled_record(source_label: &str, bytes: &[u8]) -> Vec { + let mut record = Vec::new(); + let prefix = format!("[{source_label}] "); + record.extend_from_slice(prefix.as_bytes()); + record.extend_from_slice(bytes); + record.push(b'\n'); + record +} + +fn write_batch_to_file( + file: &mut Option, + error_state: &Arc>>, + bytes: &[u8], +) -> io::Result<()> { + if let Some(error) = error_state + .lock() + .map_err(|_| io::Error::other("session log error mutex was poisoned"))? + .as_ref() + .map(|error| error.to_error()) + { + return Err(error); + } + if bytes.is_empty() { + return Ok(()); + } + let Some(writer) = file.as_mut() else { + return Err(io::Error::other("session log file is unavailable")); + }; + if let Err(error) = writer.write_all(bytes).and_then(|_| writer.flush()) { + let stored = StoredIoError::from_error(error); + *error_state + .lock() + .map_err(|_| io::Error::other("session log error mutex was poisoned"))? = + Some(stored.clone()); + *file = None; + return Err(stored.to_error()); + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::SessionLog; + use std::fs; + use std::io; + use tempfile::tempdir; + + #[test] + fn creates_a_unique_log_beside_session_state() { + let dir = tempdir().expect("tempdir"); + let state_file = dir.path().join(".devloop").join("state.json"); + let first = SessionLog::create(&state_file).expect("create first log"); + let second = SessionLog::create(&state_file).expect("create second log"); + + assert_ne!(first.path(), second.path()); + assert_eq!( + first.path().parent(), + Some(dir.path().join(".devloop/logs").as_path()) + ); + } + + #[cfg(unix)] + #[test] + fn creates_session_log_with_owner_only_permissions() { + use std::os::unix::fs::PermissionsExt; + + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + + let mode = fs::metadata(log.path()) + .expect("read log metadata") + .permissions() + .mode() + & 0o777; + assert_eq!(mode, 0o600); + } + + #[tokio::test] + async fn labels_each_persisted_output_line() { + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + + log.write_labeled_output("echo server", b"first\r\nsecond\nthird") + .expect("write output"); + log.flush_queued().await.expect("flush log"); + + assert_eq!( + fs::read_to_string(log.path()).expect("read log"), + "[echo server] first\n[echo server] second\n[echo server] third\n" + ); + } + + #[tokio::test] + async fn flush_reports_prior_writer_failure() { + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + { + log.fail_for_test(io::ErrorKind::BrokenPipe, "simulated write failure"); + } + + let error = log.flush_queued().await.expect_err("flush must fail"); + + assert_eq!(error.kind(), io::ErrorKind::BrokenPipe); + assert!(error.to_string().contains("simulated write failure")); + } + + #[tokio::test] + async fn queue_reports_prior_writer_failure() { + let dir = tempdir().expect("tempdir"); + let log = SessionLog::create(&dir.path().join(".devloop/state.json")).expect("create log"); + { + log.fail_for_test(io::ErrorKind::BrokenPipe, "simulated write failure"); + } + + let error = log + .queue_labeled_line("server", b"lost line".to_vec()) + .await + .expect_err("queue must fail"); + + assert_eq!(error.kind(), io::ErrorKind::BrokenPipe); + assert!(error.to_string().contains("simulated write failure")); + } +} diff --git a/tests/session_logs.rs b/tests/session_logs.rs new file mode 100644 index 0000000..9bd0b5f --- /dev/null +++ b/tests/session_logs.rs @@ -0,0 +1,197 @@ +use std::io::{BufRead, BufReader}; +use std::process::{Child, Command, Stdio}; +use std::sync::mpsc::{self, Receiver}; +use std::thread; +use std::time::Duration; + +use tempfile::TempDir; + +#[test] +fn run_persists_engine_and_hidden_hook_output_in_a_session_log() { + let fixture = SessionLogFixture::new(); + let mut child = DevloopChild::spawn(&fixture); + + child.wait_for_stderr("startup durable log", Duration::from_secs(10)); + fixture.wait_for_single_session_log_containing("startup durable log", Duration::from_secs(10)); + + let logs_dir = fixture.path().join(".devloop/logs"); + let entries = std::fs::read_dir(&logs_dir) + .expect("read session logs") + .collect::, _>>() + .expect("read session log entries"); + assert_eq!(entries.len(), 1, "one log per devloop run"); + + let content = std::fs::read_to_string(entries[0].path()).expect("read session log"); + assert!(content.contains("writing session log")); + assert!(content.contains("startup durable log")); + assert!(content.contains("[sh hidden] durable-hook")); +} + +#[test] +fn session_log_path_report_ignores_rust_log_filter() { + let fixture = SessionLogFixture::new(); + let mut child = DevloopChild::spawn_with_env(&fixture, &[("RUST_LOG", "warn")]); + + child.wait_for_stderr("writing session log:", Duration::from_secs(10)); + fixture.wait_for_single_session_log_containing( + "[devloop] writing session log:", + Duration::from_secs(10), + ); + + let content = fixture.read_single_session_log(); + assert!(content.contains("[devloop] writing session log:")); +} + +#[test] +fn runtime_start_failure_is_persisted_in_the_session_log() { + let fixture = SessionLogFixture::new(); + fixture.write_invalid_state_file(); + let mut child = DevloopChild::spawn(&fixture); + + child.wait_for_stderr("devloop run failed", Duration::from_secs(10)); + assert!(!child.wait_for_exit().success()); + + let content = fixture.read_single_session_log(); + assert!(content.contains("devloop run failed")); + assert!(content.contains("failed to parse state file")); +} + +struct SessionLogFixture { + dir: TempDir, +} + +impl SessionLogFixture { + fn new() -> Self { + let dir = tempfile::tempdir().expect("create fixture directory"); + let fixture = Self { dir }; + std::fs::write( + fixture.path().join("devloop.toml"), + r#"root = "." +startup_workflows = ["startup"] + +[watch.config] +paths = ["devloop.toml"] +workflow = "startup" + +[hook.hidden] +command = ["sh", "-c", "printf durable-hook"] +output = { inherit = false } + +[workflow.startup] +steps = [ + { action = "run_hook", hook = "hidden" }, + { action = "log", message = "startup durable log" }, +] +"#, + ) + .expect("write fixture config"); + fixture + } + + fn path(&self) -> &std::path::Path { + self.dir.path() + } + + fn read_single_session_log(&self) -> String { + let logs_dir = self.path().join(".devloop/logs"); + let entries = std::fs::read_dir(&logs_dir) + .expect("read session logs") + .collect::, _>>() + .expect("read session log entries"); + assert_eq!(entries.len(), 1, "one log per devloop run"); + std::fs::read_to_string(entries[0].path()).expect("read session log") + } + + fn wait_for_single_session_log_containing(&self, needle: &str, timeout: Duration) { + let started = std::time::Instant::now(); + loop { + if let Ok(content) = std::panic::catch_unwind(|| self.read_single_session_log()) + && content.contains(needle) + { + return; + } + assert!( + started.elapsed() < timeout, + "timed out waiting for session log containing '{needle}'" + ); + thread::sleep(Duration::from_millis(20)); + } + } + + fn write_invalid_state_file(&self) { + let devloop_dir = self.path().join(".devloop"); + std::fs::create_dir_all(&devloop_dir).expect("create .devloop directory"); + std::fs::write(devloop_dir.join("state.json"), "not-json").expect("write invalid state"); + } +} + +struct DevloopChild { + child: Child, + stderr: Receiver, +} + +impl DevloopChild { + fn spawn(fixture: &SessionLogFixture) -> Self { + Self::spawn_with_env(fixture, &[]) + } + + fn spawn_with_env(fixture: &SessionLogFixture, env: &[(&str, &str)]) -> Self { + let mut command = Command::new(env!("CARGO_BIN_EXE_devloop")); + command + .arg("run") + .arg("--config") + .arg(fixture.path().join("devloop.toml")) + .current_dir(fixture.path()) + .env("RUST_LOG", "info") + .stdout(Stdio::null()) + .stderr(Stdio::piped()); + for (name, value) in env { + command.env(name, value); + } + let mut child = command.spawn().expect("spawn devloop"); + let stderr = child.stderr.take().expect("take devloop stderr"); + let (tx, rx) = mpsc::channel(); + thread::spawn(move || { + for line in BufReader::new(stderr).lines() { + match line { + Ok(line) => { + if tx.send(line).is_err() { + return; + } + } + Err(_) => return, + } + } + }); + Self { child, stderr: rx } + } + + fn wait_for_stderr(&mut self, needle: &str, timeout: Duration) { + let deadline = std::time::Instant::now() + timeout; + loop { + let now = std::time::Instant::now(); + assert!( + now < deadline, + "timed out waiting for stderr containing '{needle}'" + ); + let line = self + .stderr + .recv_timeout(deadline - now) + .unwrap_or_else(|_| panic!("timed out waiting for stderr containing '{needle}'")); + if line.contains(needle) { + return; + } + } + } + + fn wait_for_exit(&mut self) -> std::process::ExitStatus { + self.child.wait().expect("wait for devloop") + } +} + +impl Drop for DevloopChild { + fn drop(&mut self) { + let _ = self.child.kill(); + let _ = self.child.wait(); + } +}