From aed848030eae4dd7c394b359c2ab175751814ac5 Mon Sep 17 00:00:00 2001 From: Daniel Vianna <1708810+pasunboneleve@users.noreply.github.com> Date: Wed, 22 Jul 2026 21:44:32 +1000 Subject: [PATCH 1/3] Persist devloop session logs Context: Devloop had supervised output and runtime tracing, but no durable per-run record. When terminal inheritance was disabled or an agent shell lost scrollback, there was no stable artifact for diagnosing engine, hook, or process behavior. Decision: Add a per-session log under the state-file parent, defaulting to .devloop/logs/. Devloop creates and verifies the log before runtime startup, tees tracing output into it, persists managed process and hook output even when terminal inheritance is disabled, ignores only the active log file in watcher classification, and flushes output at lifecycle boundaries. Retired process output now carries a generation token so cleanup can drain logs without publishing stale output-rule state into a replacement process. Alternatives considered: Writing logs directly from producers was simpler, but it reordered records and let filesystem stalls leak into runtime tasks. Ignoring the whole log directory was also simpler, but it hid user-owned files under broad watch patterns. Tradeoffs: Session-log producers use bounded backpressure, so a badly stalled log writer can slow producers instead of dropping durable records. Inherited terminal output remains backpressured by terminal behavior, while the outer output-drain deadline remains the truncation boundary. Architectural impact: The runtime now has an explicit session-log adapter at the edge of the engine. Process output forwarding separates durable persistence, terminal rendering, and output-rule state mutation, with generation-scoped state writes to keep restart behavior deterministic. --- .gitignore | 1 + CHANGELOG.md | 6 + README.md | 4 + docs/behavior.md | 14 + docs/configuration.md | 6 + src/engine.rs | 148 ++++- src/main.rs | 233 ++++++- src/processes.rs | 1412 +++++++++++++++++++++++++++++++++++------ src/session_log.rs | 417 ++++++++++++ tests/session_logs.rs | 197 ++++++ 10 files changed, 2208 insertions(+), 230 deletions(-) create mode 100644 src/session_log.rs create mode 100644 tests/session_logs.rs 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..41b8692 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,12 @@ All notable changes to `devloop` will be recorded in this file. ## [Unreleased] +### 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/README.md b/README.md index 30ae8f6..ecc7595 100644 --- a/README.md +++ b/README.md @@ -130,6 +130,10 @@ 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`. + ## Example use case Used as the primary local development workflow for diff --git a/docs/behavior.md b/docs/behavior.md index c995f39..41aad5d 100644 --- a/docs/behavior.md +++ b/docs/behavior.md @@ -202,6 +202,20 @@ server for browser listeners. first. - When output color is enabled, labels are colorized per source. +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/`. The persistent log contains `tracing` output plus labeled +managed-process and hook output, even when an `output.inherit` setting hides +that child output from the terminal. Devloop reports the selected log path at +startup. It never deletes or rotates session logs; the client owns retention. + +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 creates the log before it starts the runtime. If it cannot create the +directory or file, startup fails rather than running without durable evidence. + ### 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..40b556b 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -26,6 +26,12 @@ startup_workflows = ["startup"] - `startup_workflows`: workflows to run after autostart processes have been started. +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. Logs do not need +configuration and are not created by `devloop validate` or `devloop docs`. + 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..46fce95 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!( diff --git a/src/processes.rs b/src/processes.rs index 503f074..0a85830 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,14 +1575,14 @@ 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::time::{SystemTime, UNIX_EPOCH}; use tempfile::tempdir; - use tokio::io::AsyncReadExt; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::sync::Mutex; fn unique_state_path() -> PathBuf { @@ -1094,71 +1593,293 @@ mod tests { std::env::temp_dir().join(format!("devloop-process-state-{unique}.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(), - } + 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; + } + } + + #[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 +1894,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 "{}" +while :; do sleep 1; done +"#, + first_started.display() + ), + ) + .expect("write first script"); + std::fs::write( + &second_script, + format!( + r#"#!/bin/sh +touch "{}" +while :; do sleep 1; done +"#, + 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 +2109,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 +2765,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(); + } +} From afb41ab141d54019d24b67360226392c5f548d97 Mon Sep 17 00:00:00 2001 From: Daniel Vianna <1708810+pasunboneleve@users.noreply.github.com> Date: Thu, 23 Jul 2026 07:59:39 +1000 Subject: [PATCH 2/3] release: prepare devloop 0.10.0 Context: Persistent session logs add a backwards-compatible user-visible capability, but their location and retention rules were not yet prominent in the durable references or embedded CLI documentation. Decision: Release version 0.10.0, move the session-log entry into its dated changelog section, and add explicit state-and-log and behavior sections that `devloop docs` renders and tests. Alternatives considered: Leaving the release notes under Unreleased or documenting logs only in the README would hide the shipped contract from CLI users, so both were rejected. Tradeoffs: Logs still have no rotation policy; clients retain ownership of retention and `.devloop/` ignore rules. Architectural impact: No runtime boundary changes. The embedded documentation boundary now has regression coverage for configuration and behavior session-log sections. --- CHANGELOG.md | 2 ++ Cargo.lock | 2 +- Cargo.toml | 2 +- README.md | 2 ++ docs/README.md | 4 +++- docs/behavior.md | 18 ++++++++++++------ docs/configuration.md | 8 ++++++-- src/main.rs | 11 +++++++++++ 8 files changed, 38 insertions(+), 11 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 41b8692..a42493f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,8 @@ 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 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 ecc7595..27defb3 100644 --- a/README.md +++ b/README.md @@ -133,6 +133,8 @@ 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 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 41aad5d..8b192c3 100644 --- a/docs/behavior.md +++ b/docs/behavior.md @@ -202,19 +202,25 @@ 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/`. The persistent log contains `tracing` output plus labeled -managed-process and hook output, even when an `output.inherit` setting hides -that child output from the terminal. Devloop reports the selected log path at -startup. It never deletes or rotates session logs; the client owns retention. +`.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 creates the log before it starts the runtime. If it cannot create the -directory or file, startup fails rather than running without durable evidence. +`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 diff --git a/docs/configuration.md b/docs/configuration.md index 40b556b..cf994ef 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -26,11 +26,15 @@ 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. Logs do not need -configuration and are not created by `devloop validate` or `devloop docs`. +`.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: diff --git a/src/main.rs b/src/main.rs index 46fce95..9d8dbcf 100644 --- a/src/main.rs +++ b/src/main.rs @@ -618,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] From c9a7a5b2085c7c8213a540820f3e0979ceff6147 Mon Sep 17 00:00:00 2001 From: Daniel Vianna <1708810+pasunboneleve@users.noreply.github.com> Date: Thu, 23 Jul 2026 08:14:58 +1000 Subject: [PATCH 3/3] test: stabilize shutdown CI coverage --- src/processes.rs | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/src/processes.rs b/src/processes.rs index 0a85830..30817a8 100644 --- a/src/processes.rs +++ b/src/processes.rs @@ -1580,17 +1580,21 @@ mod tests { 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, 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 { @@ -1934,7 +1938,7 @@ exit 0 format!( r#"#!/bin/sh touch "{}" -while :; do sleep 1; done +exec sleep 600 "#, first_started.display() ), @@ -1945,7 +1949,7 @@ while :; do sleep 1; done format!( r#"#!/bin/sh touch "{}" -while :; do sleep 1; done +exec sleep 600 "#, second_started.display() ),