diff --git a/crates/cli/tests/client_host.rs b/crates/cli/tests/client_host.rs index dd1895ef..19890f60 100644 --- a/crates/cli/tests/client_host.rs +++ b/crates/cli/tests/client_host.rs @@ -17,7 +17,11 @@ use yas_client::{Client, Error, HelloOptions}; const TIMEOUT: Duration = Duration::from_secs(30); fn options() -> HostOptions { - HostOptions::new(env!("CARGO_BIN_EXE_yas")) + options_for(env!("CARGO_BIN_EXE_yas")) +} + +fn options_for(binary: impl Into) -> HostOptions { + HostOptions::new(binary) .arg("--no-persistent-extensions") .env("YAS_EXT", "0") .env("YAS_CHANNEL", "0") @@ -524,6 +528,104 @@ async fn a_dropped_output_stream_leaves_its_process_and_session_working() { still_runs_commands(&client).await; } +/// A spawned process whose attachment goes before its exit (it was detached, or its streams +/// were dropped) is asked for its exit with WAIT, which its attachment would have reported. +/// Nothing holds it to the session meanwhile: a detached one gives back its slot, and either +/// can be attached to again. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_process_whose_attachment_went_is_waited_for_and_held_by_nothing() { + use yas_client::wire::schema::process as schema; + let server = start().await; + let client = server.connect().await.unwrap(); + assert_ne!( + client.launcher_flags() & schema::SPAWN_REPORT_EXIT as u32, + 0 + ); + let per_session = client.process_limits().unwrap().max_processes_per_session as usize; + // Detached, each gives back its slot: one more than the session holds still spawns. + let mut detached = Vec::new(); + for _ in 0..per_session { + let process = client + .spawn(Command::new("sleep").arg("30").detachable(true)) + .await + .unwrap(); + within("a detach", process.detach()).await.unwrap(); + detached.push(process); + } + let one_more = within( + "a spawn beside the detached", + client.spawn(Command::new("sh").args(["-c", "exit 6"])), + ) + .await + .unwrap(); + assert_eq!( + within("its exit", one_more.wait()).await.unwrap().code(), + Some(6) + ); + // A detached one can be attached to again, and its exit comes (by WAIT). + let again = within("an attach", client.attach(detached[0].handle(), false)) + .await + .unwrap(); + detached[0].kill().await.unwrap(); + let status = within("a detached one's exit", detached[0].wait()) + .await + .unwrap(); + // (Each answer stamps the time it is given.) + let attached = within("its exit, attached", again.wait()).await.unwrap(); + assert_eq!( + (&status.kind, &status.reason, status.raw_code), + (&attached.kind, &attached.reason, attached.raw_code) + ); + for process in &detached[1..] { + process.kill().await.unwrap(); + within("a detached one's exit", process.wait()) + .await + .unwrap(); + } + // An ordinary one whose streams were dropped: the same. + let mut sleeper = client.spawn(Command::new("sleep").arg("30")).await.unwrap(); + drop(sleeper.take_stdout()); + drop(sleeper.take_stderr()); + let again = within("an attach", client.attach(sleeper.handle(), false)) + .await + .unwrap_or_else(|error| panic!("{error}\n{}", server_log(&server))); + assert_eq!( + within( + "a short wait", + sleeper.wait_timeout(Duration::from_millis(200)) + ) + .await + .unwrap(), + None + ); + sleeper.kill().await.unwrap(); + let status = within("its exit", sleeper.wait()).await.unwrap(); + // (Each answer stamps the time it is given.) + let attached = within("its exit, attached", again.wait()).await.unwrap(); + assert_eq!( + (&status.kind, &status.reason, status.raw_code), + (&attached.kind, &attached.reason, attached.raw_code) + ); + // One whose stdin was aborted: its attachment goes with it. + let mut reader = client + .spawn( + Command::new("sh") + .args(["-c", "sleep 0.3; exit 3"]) + .stdin(Stdin::Piped), + ) + .await + .unwrap(); + reader.take_stdin().unwrap().abort(); + assert_eq!( + within("its exit, stdin aborted", reader.wait()) + .await + .unwrap() + .code(), + Some(3) + ); + still_runs_commands(&client).await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn a_watcher_that_falls_behind_is_dropped_alone() { let server = start().await; @@ -568,6 +670,347 @@ async fn a_watcher_that_falls_behind_is_dropped_alone() { still_runs_commands(&owner).await; } +/// A server started with `--process-max-waits 1`, with that one WAIT held on a +/// `sleep` and shown to be held: a second WAIT is refused. +async fn hold_the_only_wait( + client: &Client, +) -> (yas_client::process::Process, tokio::task::JoinHandle<()>) { + let sleeper = client.spawn(Command::new("sleep").arg("30")).await.unwrap(); + let handle = sleeper.handle(); + let holder = { + let client = client.clone(); + tokio::spawn(async move { + let _ = client.wait_process(handle, None).await; + }) + }; + // The holder's WAIT goes first on the one connection. + tokio::time::sleep(Duration::from_millis(200)).await; + let refused = client + .wait_process(handle, Some(Duration::from_millis(1))) + .await + .unwrap_err(); + assert_eq!( + refused.status(), + Some(yas_client::wire::core::Status::ResourceExhausted), + "{refused}" + ); + (sleeper, holder) +} + +/// Commands whose exits must come whole: a non-zero exit with stderr, output +/// beyond the stream buffer read while the exit is awaited, a signal. +async fn exits_come_with_all_their_output(server: &HostedServer, client: &Client) { + let output = client + .spawn(Command::new("sh").args(["-c", "echo out; echo 'went wrong' >&2; exit 42"])) + .await + .unwrap() + .output() + .await + .unwrap(); + assert_eq!(output.status.code(), Some(42), "{}", output.status); + assert_eq!(output.stdout, b"out\n"); + assert_eq!(output.stderr, b"went wrong\n"); + + // Three times the largest stream buffer a default server keeps: the + // command blocks on its pipe until this reads, and its exit arrives while + // the tail is still on its way, as Ultimator's bash reads it. + const BYTES: u64 = 3 * 8 * 1024 * 1024; + let mut process = client + .spawn(Command::new("sh").args([ + "-c".to_owned(), + format!("head -c {BYTES} /dev/zero; printf end; echo tail >&2; exit 9"), + ])) + .await + .unwrap(); + let stdout = process.take_stdout().unwrap(); + let stderr = process.take_stderr().unwrap(); + let reader = tokio::spawn(async move { + ( + stdout.read_to_end(BYTES + 1024).await.unwrap(), + stderr.read_to_end(1024).await, + ) + }); + let status = tokio::time::timeout(TIMEOUT, process.wait()) + .await + .unwrap_or_else(|_| panic!("no exit\n{}", server_log(server))) + .unwrap(); + assert_eq!(status.code(), Some(9), "{status}"); + let (stdout, stderr) = tokio::time::timeout(TIMEOUT, reader) + .await + .expect("output after the exit") + .unwrap(); + let stderr = stderr.unwrap_or_else(|error| panic!("{error}\n{}", server_log(server))); + assert_eq!(stdout.len() as u64, BYTES + 3); + assert!(stdout.ends_with(b"end")); + assert!(stdout[..BYTES as usize].iter().all(|byte| *byte == 0)); + assert_eq!(stderr, b"tail\n"); + + let sleeper = client.spawn(Command::new("sleep").arg("30")).await.unwrap(); + assert_eq!( + sleeper + .wait_timeout(Duration::from_millis(100)) + .await + .unwrap(), + None + ); + sleeper.signal(Signal::Terminate).await.unwrap(); + let status = sleeper.wait().await.unwrap(); + assert_eq!(status.signal(), Some(libc::SIGTERM), "{status}"); + // Asked again, the same exit (an older server's WAITs may each give it another time). + let same = |other: &yas_client::process::ExitStatus| { + (&other.kind, &other.reason, other.raw_code) + == (&status.kind, &status.reason, status.raw_code) + }; + let again = sleeper.wait().await.unwrap(); + assert!(same(&again), "{again:?} after {status:?}"); + let again = sleeper + .wait_timeout(Duration::from_millis(1)) + .await + .unwrap() + .expect("exited"); + assert!(same(&again), "{again:?} after {status:?}"); + + for _ in 0..20 { + let output = client + .spawn(&Command::new("true")) + .await + .unwrap() + .output() + .await + .unwrap(); + assert!(output.status.success()); + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_spawned_process_reports_its_exit_without_a_wait() { + use yas_client::wire::schema::process as schema; + let server = tokio::time::timeout( + TIMEOUT, + HostedServer::start( + options() + .arg("--verbose") + .args(["--process-max-waits", "1"]), + ), + ) + .await + .expect("hosted server start timed out") + .expect("hosted server starts"); + let client = server.connect().await.unwrap(); + assert_ne!( + client.launcher_flags() & schema::SPAWN_REPORT_EXIT as u32, + 0 + ); + // With the only WAIT held, every exit below comes unasked. + let (sleeper, holder) = hold_the_only_wait(&client).await; + exits_come_with_all_their_output(&server, &client).await; + // A process attached to (not spawned) is still waited for: another session, its only + // WAIT held too, attaches to this one's sleeper and is refused a second WAIT. + let other = server.connect().await.unwrap(); + let (other_sleeper, other_holder) = hold_the_only_wait(&other).await; + let attached = other.attach(sleeper.handle(), false).await.unwrap(); + let refused = attached + .wait_timeout(Duration::from_millis(1)) + .await + .unwrap_err(); + assert_eq!( + refused.status(), + Some(yas_client::wire::core::Status::ResourceExhausted), + "{refused}" + ); + for (which, (sleeper, holder)) in [(sleeper, holder), (other_sleeper, other_holder)] + .into_iter() + .enumerate() + { + sleeper.kill().await.unwrap(); + // Killed by YAS at this client's request, as a WAIT would answer it. + let status = sleeper.wait().await.unwrap(); + assert_eq!( + (&status.kind, &status.reason), + ( + &yas_client::process::ExitKind::Killed, + &yas_client::process::ExitReason::Client + ), + "sleeper {which}: {status:?}" + ); + tokio::time::timeout(TIMEOUT, holder) + .await + .unwrap() + .unwrap(); + } +} + +/// A process spawned with `keep_output`: its streams' heads and tails, and what was dropped. +async fn kept(client: &Client, command: &Command) -> yas_client::process::Output { + let process = client.spawn(command).await.unwrap(); + within("the kept output", process.output_limited(4 << 20)) + .await + .unwrap() +} + +/// What a WHATWG UTF-8 decoder makes of `text`, as an elision counts it. +fn elision_of(offset: usize, text: &str) -> yas_client::process::OutputElision { + yas_client::process::OutputElision { + offset: offset as u64, + bytes: text.len() as u64, + lines: text.matches('\n').count() as u64, + code_points: text.chars().count() as u64, + utf16_units: text.encode_utf16().count() as u64, + } +} + +/// With `keep_output`, each stream is sent as its head and its tail, cut between characters, +/// and the exit says exactly what was dropped between them: far more output than the session +/// takes runs at the speed of its pipe. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn kept_output_is_its_head_and_tail_and_the_exit_counts_the_rest() { + use yas_client::wire::schema::process as schema; + let server = start().await; + let client = server.connect().await.unwrap(); + assert_ne!( + client.launcher_flags() & schema::SPAWN_KEEP_OUTPUT as u32, + 0 + ); + + // 64 MiB of `yes` on stdout, 1 MiB on stderr: 1,000 + 2,000 bytes of each come. + const BYTES: usize = 64 << 20; + let yes = kept( + &client, + Command::new("sh") + .args([ + "-c".to_owned(), + format!("yes | head -c {BYTES}; yes e | head -c 1048576 >&2; exit 4"), + ]) + .keep_output(1_000, 2_000), + ) + .await; + assert_eq!(yes.status.code(), Some(4), "{}", server_log(&server)); + assert_eq!(yes.stdout, "y\n".repeat(1_500).into_bytes()); + assert_eq!(yes.stderr, "e\n".repeat(1_500).into_bytes()); + assert_eq!( + yes.elided, + [ + Some(elision_of(1_000, &"y\n".repeat((BYTES - 3_000) / 2))), + Some(elision_of(1_000, &"e\n".repeat((1_048_576 - 3_000) / 2))), + ] + ); + // Mixed widths: the cuts fall inside emoji and CJK, and move to their ends. + let directory = tempfile::tempdir().unwrap(); + let file = directory.path().join("mixed"); + let middle = "δΈ­ζ–‡ πŸ™‚ ascii\n".repeat(700); + let whole = format!("abπŸ™‚δΈ­{middle}πŸ™‚δΈ­cd"); + std::fs::write(&file, &whole).unwrap(); + // 3 bytes of head end within πŸ™‚ (2 + 4), 8 of tail start within πŸ™‚ (β€¦πŸ™‚ δΈ­ c d). + let (head, tail) = (3, 8); + let head_len = 6; + let tail_start = whole.len() - "δΈ­cd".len(); + assert!(whole.is_char_boundary(head_len) && whole.is_char_boundary(tail_start)); + let mixed = kept( + &client, + Command::new("sh") + .args([ + "-c".to_owned(), + format!("cat '{0}'; cat '{0}' >&2", file.display()), + ]) + .keep_output(head, tail), + ) + .await; + assert!(mixed.status.success(), "{}", mixed.status); + let expected = [&whole[..head_len], &whole[tail_start..]] + .concat() + .into_bytes(); + assert_eq!( + String::from_utf8_lossy(&mixed.stdout), + String::from_utf8_lossy(&expected) + ); + assert_eq!(mixed.stderr, expected); + let dropped = Some(elision_of(head_len, &whole[head_len..tail_start])); + assert_eq!(mixed.elided, [dropped, dropped]); + + // What fits is all sent, and nothing said dropped; merged stderr is stdout's. + let fits = kept( + &client, + Command::new("sh") + .args(["-c", "printf 'out πŸ™‚'; printf ' err' >&2"]) + .merge_stderr(true) + .keep_output(100, 100), + ) + .await; + assert_eq!(fits.stdout, "out πŸ™‚ err".as_bytes()); + assert_eq!(fits.elided, [None, None]); + let merged = kept( + &client, + Command::new("sh") + .args(["-c", "yes o | head -c 100000; yes e | head -c 100000 >&2"]) + .merge_stderr(true) + .keep_output(10, 10), + ) + .await; + assert_eq!(merged.stdout, b"o\no\no\no\no\ne\ne\ne\ne\ne\n"); + assert_eq!(merged.elided[0].unwrap().bytes, 200_000 - 20); + assert_eq!(merged.elided[1], None); + + // A residue holding the streams past the grace: the tail kept so far goes out before + // the exit, which counts what was dropped until then. + let residue = kept( + &client, + Command::new("sh") + .args(["-c", "yes | head -c 3000000; (sleep 2; echo late) &"]) + .leave_residue(Some(Duration::from_millis(300))) + .keep_output(10, 10), + ) + .await; + assert_eq!(residue.status.detail, "residual process group left running"); + assert_eq!(residue.stdout, "y\n".repeat(10).into_bytes()); + assert_eq!( + residue.elided[0], + Some(elision_of(10, &"y\n".repeat((3_000_000 - 20) / 2))) + ); + still_runs_commands(&client).await; +} + +/// Against a server from before SPAWN_REPORT_EXIT, the same commands are +/// waited for with WAIT. Run with +/// `YAS_OLD_SERVER=/path/to/yas cargo test -p yas-cli --test client_host -- --ignored`, with a +/// binary that paces process output (its long output must not reset the stream). +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +#[ignore = "needs YAS_OLD_SERVER: a yas binary from before SPAWN_REPORT_EXIT"] +async fn a_server_from_before_report_exit_is_waited_for() { + use yas_client::wire::schema::process as schema; + let binary = std::env::var_os("YAS_OLD_SERVER").expect("YAS_OLD_SERVER names a yas binary"); + let server = tokio::time::timeout( + TIMEOUT, + HostedServer::start(options_for(binary).args(["--process-max-waits", "1"])), + ) + .await + .expect("hosted server start timed out") + .expect("hosted server starts"); + let client = server.connect().await.unwrap(); + assert_eq!( + client.launcher_flags() & schema::SPAWN_REPORT_EXIT as u32, + 0 + ); + // Its exits take a WAIT: with the only one held, waiting is refused. + let (sleeper, holder) = hold_the_only_wait(&client).await; + let process = client + .spawn(Command::new("sh").args(["-c", "exit 5"])) + .await + .unwrap(); + let refused = process.wait().await.unwrap_err(); + assert_eq!( + refused.status(), + Some(yas_client::wire::core::Status::ResourceExhausted), + "{refused}" + ); + sleeper.kill().await.unwrap(); + tokio::time::timeout(TIMEOUT, holder) + .await + .unwrap() + .unwrap(); + assert_eq!(process.wait().await.unwrap().code(), Some(5)); + exits_come_with_all_their_output(&server, &client).await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn a_commands_background_children_die_with_it() { let server = start().await; @@ -920,7 +1363,10 @@ async fn a_command_can_leave_its_background_running_with_a_null_stdin() { let client = server.connect().await.unwrap(); assert_eq!( client.launcher_flags(), - (schema::SPAWN_LEAVE_RESIDUE | schema::SPAWN_STDIN_NULL) as u32 + (schema::SPAWN_LEAVE_RESIDUE + | schema::SPAWN_STDIN_NULL + | schema::SPAWN_REPORT_EXIT + | schema::SPAWN_KEEP_OUTPUT) as u32 ); // The background `sleep` holds stdout: the exit comes after the grace, and // the sleep keeps running. diff --git a/crates/client/src/client.rs b/crates/client/src/client.rs index 76f1ff28..4f9be8f8 100644 --- a/crates/client/src/client.rs +++ b/crates/client/src/client.rs @@ -20,7 +20,7 @@ //! terminates every ordinary process it spawned; see [`crate::process`]. use std::collections::{HashMap, HashSet, VecDeque}; -use std::sync::atomic::{AtomicU32, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use std::sync::{Arc, Mutex, RwLock}; use std::time::Duration; @@ -51,6 +51,8 @@ pub(crate) enum Route { Transfer(u32), /// State events for one `(family, subscription_id)`. State(u16, u32), + /// Process EXIT events for one process handle (a SPAWN with REPORT_EXIT). + ProcessExit(u64), } /// Registered by a call: run by the reader on an OK `Result`, before any later @@ -63,6 +65,8 @@ pub(crate) type FrameReceiver = mpsc::UnboundedReceiver; pub(crate) struct Reply { pub(crate) prefix: ResultPrefix, routes: Vec<(Route, FrameReceiver)>, + /// With a [`Route::ProcessExit`]: told if that process's attachment goes first. + report_lost: Option>, } impl Reply { @@ -70,6 +74,41 @@ impl Reply { let index = self.routes.iter().position(|(key, _)| *key == route)?; Some(self.routes.swap_remove(index).1) } + + pub(crate) fn take_report_lost(&mut self) -> Option> { + self.report_lost.take() + } +} + +/// Marked when the attachment that would report a process's exit (SPAWN_REPORT_EXIT) goes +/// before it: a Transfer RESET, sent or received, on any of its streams, or a DETACH. The +/// server then sends no EXIT, and the exit is asked for with WAIT. +#[derive(Debug, Default)] +pub(crate) struct ReportLost { + lost: AtomicBool, + notify: tokio::sync::Notify, +} + +impl ReportLost { + pub(crate) fn mark(&self) { + self.lost.store(true, Ordering::Release); + self.notify.notify_waiters(); + } + + pub(crate) fn is_lost(&self) -> bool { + self.lost.load(Ordering::Acquire) + } + + /// Once marked. + pub(crate) async fn lost(&self) { + loop { + let notified = self.notify.notified(); + if self.is_lost() { + return; + } + notified.await; + } + } } struct Pending { @@ -86,6 +125,8 @@ struct Router { orphan_bytes: usize, released: HashSet, released_order: VecDeque, + /// The open transfers of processes that report their exit: a RESET on one marks it. + report_watch: HashMap>, closed: Option, } @@ -106,6 +147,9 @@ impl Router { fn release(&mut self, route: Route) { self.routes.remove(&route); + if let Route::Transfer(transfer_id) = route { + self.report_watch.remove(&transfer_id); + } if let Some(frames) = self.orphans.remove(&route) { for frame in frames { self.orphan_bytes = self.orphan_bytes.saturating_sub(frame.payload.len()); @@ -122,6 +166,16 @@ impl Router { } } + /// A transfer closed, or reset (either way): a reset one's process, if it reports its + /// exit, will not. + fn transfer_ended(&mut self, transfer_id: u32, reset: bool) { + if let Some(lost) = self.report_watch.remove(&transfer_id) + && reset + { + lost.mark(); + } + } + fn deliver(&mut self, route: Route, frame: Frame) { if let Some(sender) = self.routes.get(&route) { if sender.send(frame).is_err() { @@ -161,6 +215,7 @@ impl Router { // Dropping the senders ends every stream and subscription; they then // report the session error. self.routes.clear(); + self.report_watch.clear(); self.orphans.clear(); self.orphan_order.clear(); self.orphan_bytes = 0; @@ -412,6 +467,16 @@ impl Client { if let Some(error) = self.closed_reason() { return Err(error); } + if frame.header.kind == yas_wire::transfer::kind::RESET + && let Some(transfer_id) = transfer_event_id(&frame) + { + self.inner + .shared + .router + .lock() + .unwrap() + .transfer_ended(transfer_id, true); + } self.inner .outbound .send(frame) @@ -617,7 +682,26 @@ fn dispatch(shared: &Shared, frame: Frame) -> Result<()> { routes.push((route, receiver)); } } - if pending.reply.send(Ok(Reply { prefix, routes })).is_err() { + // A process that reports its exit: any of its transfers reset (its attachment + // went) means no EXIT comes. + let report_lost = routes + .iter() + .any(|(route, _)| matches!(route, Route::ProcessExit(_))) + .then(|| { + let lost = Arc::new(ReportLost::default()); + for (route, _) in &routes { + if let Route::Transfer(transfer_id) = route { + router.report_watch.insert(*transfer_id, lost.clone()); + } + } + lost + }); + let reply = Reply { + prefix, + routes, + report_lost, + }; + if pending.reply.send(Ok(reply)).is_err() { // The caller went away between sending and now. } Ok(()) @@ -625,11 +709,27 @@ fn dispatch(shared: &Shared, frame: Frame) -> Result<()> { Class::Event => { if frame.header.family == family::TRANSFER { if let Some(transfer_id) = transfer_event_id(&frame) { + let mut router = shared.router.lock().unwrap(); + let kind = frame.header.kind; + if matches!( + kind, + yas_wire::transfer::kind::CLOSE | yas_wire::transfer::kind::RESET + ) { + router.transfer_ended(transfer_id, kind == yas_wire::transfer::kind::RESET); + } + router.deliver(Route::Transfer(transfer_id), frame); + } + return Ok(()); + } + if frame.header.family == family::PROCESS + && frame.header.kind == yas_wire::process::event_kind::EXIT + { + if let Some(handle) = yas_wire::process::ExitReport::handle_of(&frame.payload) { shared .router .lock() .unwrap() - .deliver(Route::Transfer(transfer_id), frame); + .deliver(Route::ProcessExit(handle), frame); } return Ok(()); } diff --git a/crates/client/src/process.rs b/crates/client/src/process.rs index 855b2c67..3bbee215 100644 --- a/crates/client/src/process.rs +++ b/crates/client/src/process.rs @@ -86,6 +86,7 @@ //! never more than 1 MiB nor less than 16 KiB. use std::ffi::OsStr; +use std::sync::Arc; use std::time::Duration; use yas_wire::{ @@ -94,8 +95,8 @@ use yas_wire::{ family, process::{ self as wire, Attach, Control, ControlAction, ControlResult, Cwd, EnvEntry, - EnvironmentKind, ExitRecord, ProcessRecord, RemovedProcess, Spawn, StreamBundle, Wait, - request_kind, + EnvironmentKind, ExitRecord, ExitReport, ProcessRecord, RemovedProcess, Spawn, + StreamBundle, Wait, request_kind, }, schema::process as schema, state::{Phase, RecordKind, Watch as StateWatch}, @@ -109,7 +110,7 @@ use crate::transfer::{ByteSink, ByteStream, DEFAULT_WINDOW}; /// The smallest output window [`Client::default_process_window`] picks. const MIN_AUTO_WINDOW: u64 = 16 * 1024; -pub use yas_wire::process::ExitKind; +pub use yas_wire::process::{ExitKind, OutputElision}; /// What the child's stdin is connected to. #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] @@ -147,6 +148,7 @@ pub struct Command { window: Option, operation_id: [u8; 16], leave_residue: Option>, + keep_output: Option<(u64, u64)>, } impl Command { @@ -163,6 +165,7 @@ impl Command { window: None, operation_id: nonzero_id(), leave_residue: None, + keep_output: None, } } @@ -248,6 +251,18 @@ impl Command { self } + /// Send only the first `head` bytes of each output stream (a few more, to end between + /// UTF-8 characters) and its last `tail` bytes (at most 1 MiB, + /// `MAX_KEEP_OUTPUT_TAIL_BYTES`; a few fewer, to start between characters), where the + /// server offers `SPAWN_KEEP_OUTPUT` with `SPAWN_REPORT_EXIT` ([`Client::launcher_flags`]): + /// what comes between is dropped as the server reads it, so the command runs at the speed + /// of its pipe rather than of this session, and counted ([`Process::elided`]). Other + /// servers send it all. + pub fn keep_output(&mut self, head: u64, tail: u64) -> &mut Self { + self.keep_output = Some((head, tail.min(schema::MAX_KEEP_OUTPUT_TAIL_BYTES))); + self + } + /// Connect stdin. pub fn stdin(&mut self, stdin: Stdin) -> &mut Self { self.stdin = stdin; @@ -468,10 +483,14 @@ impl std::fmt::Display for ExitStatus { pub struct Output { /// How it ended. pub status: ExitStatus, - /// Everything it wrote to stdout (and stderr, when merged). + /// Everything it wrote to stdout (and stderr, when merged); with + /// [`Command::keep_output`], its head then its tail, as `elided` says. pub stdout: Vec, - /// Everything it wrote to stderr (empty when merged). + /// Everything it wrote to stderr (empty when merged), or its head then its tail. pub stderr: Vec, + /// What `KEEP_OUTPUT` dropped of stdout and of stderr ([`Process::elided`]): None when + /// nothing was, the stream then whole. + pub elided: [Option; 2], } /// A running (or finished) process and the streams this session holds. @@ -485,6 +504,93 @@ pub struct Process { stdout_offset: u64, stderr_offset: u64, merged_stderr: bool, + /// Where the server sends its exit, for a process spawned with REPORT_EXIT. + reported: Option, +} + +/// The exit a server reports unasked (an EXIT event), received once and kept. Its attachment +/// reports it: once that goes before the exit (a stream dropped or reset, the process +/// detached), the exit is asked for with WAIT. +struct ReportedExit { + frames: tokio::sync::Mutex, + /// The exit, and what KEEP_OUTPUT dropped of stdout and of stderr. + status: std::sync::OnceLock<(ExitStatus, [Option; 2])>, + lost: Arc, +} + +impl std::fmt::Debug for ReportedExit { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ReportedExit") + .field("status", &self.status.get()) + .finish_non_exhaustive() + } +} + +impl ReportedExit { + /// The exit, waiting at most `timeout`: None if it is still running then. + async fn wait( + &self, + client: &Client, + handle: u64, + timeout: Option, + ) -> Result> { + use yas_wire::Decode; + let deadline = timeout.map(|timeout| tokio::time::Instant::now() + timeout); + let until_deadline = async { + match deadline { + Some(deadline) => tokio::time::sleep_until(deadline).await, + None => std::future::pending().await, + } + }; + tokio::pin!(until_deadline); + if let Some((status, _)) = self.status.get() { + return Ok(Some(status.clone())); + } + let mut frames = tokio::select! { + frames = self.frames.lock() => frames, + () = &mut until_deadline => return Ok(None), + }; + if let Some((status, _)) = self.status.get() { + return Ok(Some(status.clone())); + } + let frame = tokio::select! { + biased; + frame = frames.recv() => frame, + // Its attachment went: an EXIT it sent first may still be here. + () = self.lost.lost() => frames.try_recv().ok(), + () = &mut until_deadline => return Ok(None), + }; + let status = match frame { + Some(frame) => { + let report = ExitReport::decode(&frame.payload)?; + let elided = [report.elided(false)?, report.elided(true)?]; + (ExitStatus::from_wire(report.exit), elided) + } + None if self.lost.is_lost() => { + let left = deadline.map(|deadline| { + deadline.saturating_duration_since(tokio::time::Instant::now()) + }); + match client.wait_process(handle, left).await? { + Some(status) => (status, [None; 2]), + None => return Ok(None), + } + } + None => { + return Err(client + .closed_reason() + .unwrap_or_else(|| Error::protocol("Process EXIT route closed"))); + } + }; + Ok(Some(self.status.get_or_init(|| status).0.clone())) + } +} + +impl Drop for Process { + fn drop(&mut self) { + if self.reported.is_some() { + self.client.release(Route::ProcessExit(self.handle)); + } + } } impl Process { @@ -536,16 +642,26 @@ impl Process { self.stderr.take() } - /// Wait for the process to exit. + /// Wait for the process to exit. A process this session spawned from a server that + /// reports exits (SPAWN_REPORT_EXIT, in [`Client::launcher_flags`]) takes no request: + /// the server sends the exit as it happens. Otherwise, or once the attachment that would + /// report it went first (a stream dropped or reset before its end, a detach), this asks + /// with WAIT. pub async fn wait(&self) -> Result { - self.client - .wait_process(self.handle, None) - .await? - .ok_or_else(|| Error::protocol("Process WAIT without timeout timed out")) + let status = match &self.reported { + Some(reported) => reported.wait(&self.client, self.handle, None).await?, + None => self.client.wait_process(self.handle, None).await?, + }; + status.ok_or_else(|| Error::protocol("Process WAIT without timeout timed out")) } /// Wait at most `timeout`; `None` if it is still running. pub async fn wait_timeout(&self, timeout: Duration) -> Result> { + if let Some(reported) = &self.reported { + return reported + .wait(&self.client, self.handle, Some(timeout)) + .await; + } self.client.wait_process(self.handle, Some(timeout)).await } @@ -571,8 +687,20 @@ impl Process { .await } + /// What `KEEP_OUTPUT` ([`Command::keep_output`]) dropped of stdout (or, with `stderr`, of + /// stderr), once [`Process::wait`] answered: None when nothing was. The stream's bytes + /// are its head, up to `offset`, then its tail; the elision's counts say what was between. + pub fn elided(&self, stderr: bool) -> Option { + let (_, elided) = self.reported.as_ref()?.status.get()?; + elided[usize::from(stderr)] + } + /// Make a detachable process independent of this session's attachment. pub async fn detach(&self) -> Result<()> { + // Its attachment goes, and would have reported the exit. + if let Some(reported) = &self.reported { + reported.lost.mark(); + } self.client .control_process(self.handle, ControlAction::Detach, 0) .await @@ -605,6 +733,7 @@ impl Process { status, stdout, stderr, + elided: [self.elided(false), self.elided(true)], }) } @@ -765,13 +894,18 @@ fn change_from_record(record: &yas_wire::state::Record) -> Result Hook { - Box::new(|prefix: &ResultPrefix| { +/// The routes a SPAWN or ATTACH Result announces: its streams, and with `report_exit` the +/// process's EXIT events. +fn bundle_hook(report_exit: bool) -> Hook { + Box::new(move |prefix: &ResultPrefix| { use yas_wire::Decode; let Ok(bundle) = StreamBundle::decode(&prefix.body) else { return Vec::new(); }; let mut routes = vec![Route::Transfer(bundle.stdout.transfer_id)]; + if report_exit { + routes.push(Route::ProcessExit(bundle.process_handle)); + } if let Some(stdin) = &bundle.stdin { routes.push(Route::Transfer(stdin.transfer_id)); } @@ -787,10 +921,26 @@ impl Client { /// happens to it and its children. pub async fn spawn(&self, command: &Command) -> Result { let stdin_null = self.launcher_flags() & schema::SPAWN_STDIN_NULL as u32 != 0; + // The server sends the exit unasked: waiting for it takes no round trip. + let report_exit = self.launcher_flags() & schema::SPAWN_REPORT_EXIT as u32 != 0; let window = command .window .unwrap_or_else(|| self.default_process_window()); - let spawn = command.to_wire(stdin_null, window)?; + let mut spawn = command.to_wire(stdin_null, window)?; + if report_exit { + spawn.flags |= schema::SPAWN_REPORT_EXIT as u16; + } + // Its last extension (tag 4): the others' tags are lower. + if let Some((head, tail)) = command.keep_output + && report_exit + && self.launcher_flags() & schema::SPAWN_KEEP_OUTPUT as u32 != 0 + { + spawn.flags |= schema::SPAWN_KEEP_OUTPUT as u16; + spawn + .extensions + .0 + .push(Spawn::keep_output_extension(head, tail)); + } // Servers from before the extended limits refuse more than 256 // entries as undecodable; say why instead. if let Some(limits) = self.process_limits() @@ -808,7 +958,7 @@ impl Client { request_kind::SPAWN, spawn.encode()?, Some(DEFAULT_REQUEST_TIMEOUT), - Some(bundle_hook()), + Some(bundle_hook(report_exit)), ) .await?; let mut process = self.process_from_reply(reply, window)?; @@ -847,7 +997,7 @@ impl Client { request_kind::ATTACH, attach.encode()?, Some(DEFAULT_REQUEST_TIMEOUT), - Some(bundle_hook()), + Some(bundle_hook(false)), ) .await?; self.process_from_reply(reply, window) @@ -893,6 +1043,14 @@ impl Client { } None => None, }; + let lost = reply.take_report_lost().unwrap_or_default(); + let reported = reply + .take(Route::ProcessExit(bundle.process_handle)) + .map(|frames| ReportedExit { + frames: tokio::sync::Mutex::new(frames), + status: std::sync::OnceLock::new(), + lost, + }); Ok(Process { client: self.clone(), handle: bundle.process_handle, @@ -902,6 +1060,7 @@ impl Client { stdout_offset: bundle.stdout_lifetime_offset, stderr_offset: bundle.stderr_lifetime_offset, merged_stderr: bundle.merged_stderr, + reported, }) } @@ -1035,7 +1194,9 @@ impl Client { .and_then(|limits| wire::Limits::from_extensions(&limits).ok()) } - /// The opt-in SPAWN flags this server honours (`SPAWN_LEAVE_RESIDUE`, `SPAWN_STDIN_NULL`); + /// The opt-in SPAWN flags this server honours (`SPAWN_LEAVE_RESIDUE`, `SPAWN_STDIN_NULL`, + /// `SPAWN_REPORT_EXIT`, which [`Client::spawn`] sets itself so that [`Process::wait`] takes + /// no round trip); /// 0 for servers that predate them. [`Stdin::Null`] uses the null device where offered. pub fn launcher_flags(&self) -> u32 { self.process_limits() diff --git a/crates/server/src/lib.rs b/crates/server/src/lib.rs index c8d9d2c9..411b1e4c 100644 --- a/crates/server/src/lib.rs +++ b/crates/server/src/lib.rs @@ -62,6 +62,8 @@ mod net; mod nvdec_decode; mod nvenc_encode; #[cfg(any(unix, windows))] +mod output_keep; +#[cfg(any(unix, windows))] mod process; mod pty; mod read_only_stream; diff --git a/crates/server/src/output_keep.rs b/crates/server/src/output_keep.rs new file mode 100644 index 00000000..a08ae1b8 --- /dev/null +++ b/crates/server/src/output_keep.rs @@ -0,0 +1,616 @@ +//! A process output stream kept as its head and tail (SPAWN `KEEP_OUTPUT`): what comes between +//! is dropped as it is read, and counted, so a command writing far more than its client keeps +//! runs at the speed of its pipe rather than of the client's window. +//! +//! The cuts fall between characters as a WHATWG UTF-8 decoder with replacement (a JavaScript +//! `TextDecoder`, Rust's `from_utf8_lossy`) reads the whole stream: the head ends where such a +//! decoder is between characters, or where the next byte cannot continue the character it is in; +//! the tail starts between characters, with no continuation byte. Decoding the head and the tail, apart or +//! one after the other, then gives exactly the characters it gives them within the whole, and +//! the counts of what was dropped (bytes, lines, code points, UTF-16 units) complete it: a client +//! can say how much it did not get as it would have counted it. + +use std::collections::VecDeque; + +/// What was dropped from a stream: from `offset` on (the head's length), `bytes` bytes that +/// decode to `code_points` characters, `utf16_units` UTF-16 code units, `lines` of them `\n`. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub(crate) struct Elision { + pub(crate) offset: u64, + pub(crate) bytes: u64, + pub(crate) lines: u64, + pub(crate) code_points: u64, + pub(crate) utf16_units: u64, +} + +/// A WHATWG UTF-8 decoder's state, reading one byte at a time. +#[derive(Clone, Copy, Debug)] +struct Utf8Scan { + needed: u8, + seen: u8, + lower: u8, + upper: u8, +} + +impl Default for Utf8Scan { + fn default() -> Self { + Self { + needed: 0, + seen: 0, + lower: 0x80, + upper: 0xBF, + } + } +} + +/// What a byte did to a [`Utf8Scan`]. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum Step { + /// Part of a character still incomplete. + Pending, + /// It ended a character of this many UTF-16 units (a replacement for an invalid byte is 1). + Char { units: u8, newline: bool }, + /// It cannot continue the pending character: that one ends as a replacement (1 unit) before + /// it, and the byte must be read again, from between characters. + Again, +} + +impl Utf8Scan { + fn between_characters(&self) -> bool { + self.needed == 0 + } + + fn step(&mut self, byte: u8) -> Step { + if self.needed == 0 { + return match byte { + 0x00..=0x7F => Step::Char { + units: 1, + newline: byte == b'\n', + }, + 0xC2..=0xDF => { + self.needed = 1; + Step::Pending + } + 0xE0..=0xEF => { + if byte == 0xE0 { + self.lower = 0xA0; + } else if byte == 0xED { + self.upper = 0x9F; + } + self.needed = 2; + Step::Pending + } + 0xF0..=0xF4 => { + if byte == 0xF0 { + self.lower = 0x90; + } else if byte == 0xF4 { + self.upper = 0x8F; + } + self.needed = 3; + Step::Pending + } + _ => Step::Char { + units: 1, + newline: false, + }, + }; + } + if !(self.lower..=self.upper).contains(&byte) { + *self = Self::default(); + return Step::Again; + } + self.lower = 0x80; + self.upper = 0xBF; + self.seen += 1; + if self.seen < self.needed { + return Step::Pending; + } + // Four bytes make a code point beyond the BMP: two UTF-16 units. + let units = if self.needed == 3 { 2 } else { 1 }; + *self = Self::default(); + Step::Char { + units, + newline: false, + } + } +} + +/// Counts what goes through a [`Utf8Scan`], byte by byte. +#[derive(Debug, Default)] +struct Counter { + scan: Utf8Scan, + counts: Elision, +} + +impl Counter { + fn count(&mut self, units: u8, newline: bool) { + self.counts.code_points += 1; + self.counts.utf16_units += u64::from(units); + self.counts.lines += u64::from(newline); + } + + fn byte(&mut self, byte: u8) { + self.counts.bytes += 1; + loop { + match self.scan.step(byte) { + Step::Pending => return, + Step::Char { units, newline } => { + self.count(units, newline); + return; + } + Step::Again => self.count(1, false), + } + } + } + + /// Bytes in order: runs of ASCII at once. + fn bytes(&mut self, data: &[u8]) { + let mut at = 0; + while at < data.len() { + if self.scan.between_characters() { + let run = data[at..] + .iter() + .position(|byte| !byte.is_ascii()) + .unwrap_or(data.len() - at); + if run > 0 { + let ascii = &data[at..at + run]; + let run = run as u64; + self.counts.bytes += run; + self.counts.code_points += run; + self.counts.utf16_units += run; + self.counts.lines += ascii.iter().filter(|byte| **byte == b'\n').count() as u64; + at += ascii.len(); + continue; + } + } + self.byte(data[at]); + at += 1; + } + } + + /// A byte that may continue the pending character: false when it cannot (the character + /// ends as a replacement before it, and the byte is not taken). + fn continue_with(&mut self, byte: u8) -> bool { + match self.scan.step(byte) { + Step::Again => { + self.count(1, false); + false + } + Step::Pending => { + self.counts.bytes += 1; + true + } + Step::Char { units, newline } => { + self.counts.bytes += 1; + self.count(units, newline); + true + } + } + } + + /// The end of the stream: a character left incomplete is a replacement. + fn end(&mut self) { + if !self.scan.between_characters() { + self.scan = Utf8Scan::default(); + self.count(1, false); + } + } +} + +/// One output stream kept as its head and tail. +#[derive(Debug)] +pub(crate) struct KeptOutput { + /// Bytes to keep at the start, at least: the head ends between characters. + head: u64, + /// Bytes to keep at the end, at most. + tail: usize, + /// The head so far, and how a decoder stands at its end. + head_len: u64, + head_scan: Utf8Scan, + head_done: bool, + /// The last bytes read since the head, at most `tail` of them. + ring: VecDeque, + /// What left the ring at its front: dropped, in order. + dropped: Counter, + finished: bool, + /// Once finished: the tail, and how much of it went out. + tail_out: Vec, + tail_sent: usize, +} + +impl KeptOutput { + pub(crate) fn new(head: u64, tail: usize) -> Self { + Self { + head, + tail, + head_len: 0, + head_scan: Utf8Scan::default(), + head_done: false, + ring: VecDeque::new(), + dropped: Counter::default(), + finished: false, + tail_out: Vec::new(), + tail_sent: 0, + } + } + + /// Whether what is fed now goes out (the head is still being taken): only then need the + /// stream wait for its reader to take it. + pub(crate) fn sends_now(&self) -> bool { + !self.head_done && !self.finished + } + + /// Whether the stream is finished and its tail all went out. + pub(crate) fn tail_out(&self) -> bool { + self.finished && self.tail_sent == self.tail_out.len() + } + + /// Once finished: the next at most `max` bytes of the tail, taken as sent. + pub(crate) fn next_tail_chunk(&mut self, max: usize) -> Option> { + let rest = &self.tail_out[self.tail_sent..]; + if !self.finished || rest.is_empty() { + return None; + } + let chunk = rest[..rest.len().min(max)].to_vec(); + self.tail_sent += chunk.len(); + Some(chunk) + } + + /// The stream stops before its tail was taken (its reader was stopped, or the exit is + /// reported first): what it kept counts as dropped, to the stream's end. A tail already + /// being sent is left as it is. + pub(crate) fn drop_rest(&mut self) { + if self.finished { + return; + } + self.finish(); + let rest = std::mem::take(&mut self.tail_out); + self.dropped.bytes(&rest); + self.dropped.end(); + } + + /// Take `data` as read from the stream, and answer what of it goes out now: the start of + /// it that belongs to the head (nothing once that is complete, or once finished). + pub(crate) fn feed<'a>(&mut self, data: &'a [u8]) -> &'a [u8] { + if self.finished { + return &[]; + } + let mut taken = 0; + while !self.head_done && taken < data.len() { + if self.head_len >= self.head && self.head_scan.between_characters() { + self.head_done = true; + break; + } + match self.head_scan.step(data[taken]) { + // The pending character ends (a replacement) before this byte: past the head's + // size, so does the head; else the byte is read again, between characters. + Step::Again if self.head_len >= self.head => self.head_done = true, + Step::Again => {} + Step::Pending | Step::Char { .. } => { + taken += 1; + self.head_len += 1; + } + } + } + self.keep(&data[taken..]); + &data[..taken] + } + + /// Push bytes past the head into the ring; what that pushes out of it is dropped. + fn keep(&mut self, data: &[u8]) { + if data.is_empty() { + return; + } + let excess = (self.ring.len() + data.len()).saturating_sub(self.tail); + let from_ring = excess.min(self.ring.len()); + if from_ring > 0 { + let (first, second) = self.ring.as_slices(); + let in_first = from_ring.min(first.len()); + self.dropped.bytes(&first[..in_first]); + self.dropped.bytes(&second[..from_ring - in_first]); + self.ring.drain(..from_ring); + } + let from_data = excess - from_ring; + self.dropped.bytes(&data[..from_data]); + self.ring.extend(&data[from_data..]); + } + + /// The stream ended (or stops being forwarded): its tail, which goes out last + /// ([`next_tail_chunk`](Self::next_tail_chunk)), is what the ring holds from the first + /// character boundary on. Then nothing more is taken. + pub(crate) fn finish(&mut self) { + if self.finished { + return; + } + self.finished = true; + // The dropped bytes may end within a character: the ring's first bytes complete it, or + // it ends as a replacement before them. + while !self.dropped.scan.between_characters() { + let Some(&byte) = self.ring.front() else { + self.dropped.end(); + break; + }; + if !self.dropped.continue_with(byte) { + break; + } + self.ring.pop_front(); + } + // The head may end within a character (one that the next byte broke): a tail starting + // with a byte that could continue it would decode differently after it than apart. + // Once anything was dropped, the tail starts with no continuation byte; those it would + // start with are lone ones there, dropped as a replacement each. + if self.dropped.counts.bytes > 0 { + while let Some(&byte) = self.ring.front() + && (0x80..=0xBF).contains(&byte) + { + self.dropped.byte(byte); + self.ring.pop_front(); + } + } + self.tail_out = self.ring.drain(..).collect(); + self.ring = VecDeque::new(); + } + + /// What was dropped so far, where; None while nothing was. + pub(crate) fn elision(&self) -> Option { + (self.dropped.counts.bytes > 0).then_some(Elision { + offset: self.head_len, + ..self.dropped.counts + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + /// xorshift64: reproducible without a dependency. + struct Rng(u64); + + impl Rng { + fn below(&mut self, bound: usize) -> usize { + let mut x = self.0; + x ^= x << 13; + x ^= x >> 7; + x ^= x << 17; + self.0 = x; + (x % bound as u64) as usize + } + } + + /// Text of every width with newlines, invalid bytes, truncated and overlong sequences, a + /// surrogate's encoding and a BOM among it. + fn mixed(rng: &mut Rng, len: usize) -> Vec { + const PIECES: &[&[u8]] = &[ + b"a", + b"\n", + "Γ©".as_bytes(), + "δΈ­".as_bytes(), + "ζ–‡ε­—".as_bytes(), + "πŸ˜€".as_bytes(), + "π„ž".as_bytes(), + b"\xFF", + b"\x80", + b"\xC3", + b"\xC0\xAF", + b"\xE4\xB8", + b"\xF0\x9F\x98", + b"\xED\xA0\x80", + b"\xF4\x90\x80\x80", + b"\xEF\xBB\xBF", + ]; + let mut out = Vec::with_capacity(len + 4); + while out.len() < len { + out.extend_from_slice(PIECES[rng.below(PIECES.len())]); + } + out + } + + fn utf16(bytes: &[u8]) -> Vec { + String::from_utf8_lossy(bytes).encode_utf16().collect() + } + + fn chars(bytes: &[u8]) -> u64 { + String::from_utf8_lossy(bytes).chars().count() as u64 + } + + fn lines(bytes: &[u8]) -> u64 { + bytes.iter().filter(|byte| **byte == b'\n').count() as u64 + } + + /// Keep `whole` fed in pieces of at most `piece` bytes; answer the head and tail sent. + fn keep( + whole: &[u8], + head: u64, + tail: usize, + pieces: &mut dyn FnMut() -> usize, + ) -> (Vec, Vec, KeptOutput) { + let mut kept = KeptOutput::new(head, tail); + let mut sent = Vec::new(); + let mut rest = whole; + while !rest.is_empty() { + let n = pieces().clamp(1, rest.len()); + sent.extend_from_slice(kept.feed(&rest[..n])); + rest = &rest[n..]; + } + kept.finish(); + let mut tail_sent = Vec::new(); + while let Some(chunk) = kept.next_tail_chunk(7) { + tail_sent.extend_from_slice(&chunk); + } + assert!(kept.tail_out()); + (sent, tail_sent, kept) + } + + /// The head and tail decode, apart or together, to the whole's first and last characters, + /// and what was dropped counts the rest in every unit. + fn check(whole: &[u8], head: u64, tail: usize, pieces: &mut dyn FnMut() -> usize) { + let (head_sent, tail_sent, kept) = keep(whole, head, tail, pieces); + let what = format!("{} bytes, head {head}, tail {tail}", whole.len()); + let all = utf16(whole); + let (first, last) = (utf16(&head_sent), utf16(&tail_sent)); + assert_eq!( + utf16(&[head_sent.as_slice(), &tail_sent].concat()), + [first.as_slice(), &last].concat(), + "{what}" + ); + assert_eq!(&all[..first.len()], first.as_slice(), "{what}"); + assert_eq!(&all[all.len() - last.len()..], last.as_slice(), "{what}"); + let dropped = kept.elision().unwrap_or_default(); + assert_eq!( + head_sent.len() as u64 + dropped.bytes + tail_sent.len() as u64, + whole.len() as u64, + "{what}" + ); + assert_eq!( + first.len() as u64 + dropped.utf16_units + last.len() as u64, + all.len() as u64, + "{what}" + ); + assert_eq!( + chars(&head_sent) + dropped.code_points + chars(&tail_sent), + chars(whole), + "{what}" + ); + assert_eq!( + lines(&head_sent) + dropped.lines + lines(&tail_sent), + lines(whole), + "{what}" + ); + // The head is its size or all there is, and at most one character more; the tail at + // most its size, and at most one character less once anything was dropped. + assert!( + head_sent.len() as u64 >= head.min(whole.len() as u64), + "{what}" + ); + assert!(head_sent.len() as u64 <= head + 3, "{what}"); + assert!(tail_sent.len() <= tail, "{what}"); + assert!(whole.ends_with(&tail_sent), "{what}"); + if dropped.bytes > 0 { + assert_eq!(dropped.offset, head_sent.len() as u64, "{what}"); + let first = tail_sent.first(); + assert!( + first.is_none_or(|byte| !(0x80..=0xBF).contains(byte)), + "{what}" + ); + } else { + assert_eq!([head_sent, tail_sent].concat(), whole, "{what}"); + } + } + + #[test] + fn head_and_tail_of_any_text_decode_as_within_the_whole_and_the_counts_make_up_the_rest() { + let mut rng = Rng(0x9E37_79B9_7F4A_7C15); + for _ in 0..4000 { + let len = rng.below(600); + let whole = mixed(&mut rng, len); + let head = rng.below(whole.len() + 8) as u64; + let tail = rng.below(whole.len() + 8); + let mut seed = Rng(rng.below(usize::MAX) as u64 | 1); + check(&whole, head, tail, &mut || 1 + seed.below(40)); + } + } + + #[test] + fn cuts_within_wide_characters_move_to_their_ends() { + // An emoji across the head's end, CJK across the tail's start. + let whole = [ + b"ab".as_slice(), + "πŸ˜€".as_bytes(), + "δΈ­".repeat(10).as_bytes(), + "ζ–‡".as_bytes(), + b"yz", + ] + .concat(); + for head in 2..8 { + for tail in 0..12 { + for piece in [1, 2, 3, 5, 64] { + check(&whole, head, tail, &mut || piece); + } + } + } + let (head_sent, tail_sent, kept) = keep(&whole, 3, 7, &mut || 1); + assert_eq!(head_sent, [b"ab".as_slice(), "πŸ˜€".as_bytes()].concat()); + assert_eq!(tail_sent, ["ζ–‡".as_bytes(), b"yz"].concat()); + assert_eq!( + kept.elision(), + Some(Elision { + offset: 6, + bytes: 30, + lines: 0, + code_points: 10, + utf16_units: 10 + }) + ); + } + + #[test] + fn a_head_ending_within_a_character_never_meets_a_tail_that_continues_it() { + // The head ends with ED (the byte after, A0, cannot continue it); the tail would start + // with 90 80, which can (ED 90 80 is U+D400): those go. + let whole = [b"ab\xED\xA0".as_slice(), &[b'x'; 20], b"\x90\x80z"].concat(); + let (head_sent, tail_sent, kept) = keep(&whole, 3, 3, &mut || 1); + assert_eq!(head_sent, b"ab\xED"); + assert_eq!(tail_sent, b"z"); + let joined = [head_sent.as_slice(), &tail_sent].concat(); + assert_eq!(String::from_utf8_lossy(&joined), "ab\u{FFFD}z"); + // A0, the x's, then 90 and 80: a replacement each. + assert_eq!(kept.elision().unwrap().utf16_units, 1 + 20 + 2); + check(&whole, 3, 3, &mut || 1); + } + + #[test] + fn nothing_is_dropped_from_what_fits() { + let whole = "line\n".repeat(100).into_bytes(); + let (head_sent, tail_sent, kept) = keep(&whole, 300, 300, &mut || 7); + assert_eq!([head_sent, tail_sent].concat(), whole); + assert_eq!(kept.elision(), None); + } + + #[test] + fn a_long_ascii_stream_is_counted_in_bulk() { + let whole = "y\n".repeat(1 << 20).into_bytes(); + let (head_sent, tail_sent, kept) = keep(&whole, 1000, 500, &mut || 65536); + assert_eq!(head_sent.len(), 1000); + assert_eq!(tail_sent.len(), 500); + let dropped = kept.elision().unwrap(); + assert_eq!(dropped.bytes, (2 << 20) - 1500); + assert_eq!(dropped.utf16_units, dropped.bytes); + assert_eq!(dropped.lines, dropped.bytes / 2); + } + + #[test] + fn a_stream_stopped_before_its_tail_drops_all_it_kept() { + let mut whole = "head βœ“ then πŸ™‚ and δΈ­ζ–‡\n".repeat(50).into_bytes(); + // It stops within a character: that ends as a replacement. + whole.extend_from_slice(&[0xF0, 0x9F]); + let mut kept = KeptOutput::new(10, 64); + let sent = kept.feed(&whole).to_vec(); + assert_eq!(sent, b"head \xE2\x9C\x93 t"); + assert!(!kept.sends_now()); + kept.drop_rest(); + assert!(kept.tail_out()); + assert_eq!(kept.next_tail_chunk(1024), None); + assert!(kept.feed(b"more").is_empty()); + let rest = &whole[sent.len()..]; + assert_eq!( + kept.elision(), + Some(Elision { + offset: sent.len() as u64, + bytes: rest.len() as u64, + lines: lines(rest), + code_points: chars(rest), + utf16_units: utf16(rest).len() as u64, + }) + ); + // Once its tail is taken, nothing more is dropped. + let mut finished = KeptOutput::new(4, 8); + assert_eq!(finished.feed(b"abcdefghijklmnop"), b"abcd"); + finished.finish(); + assert_eq!(finished.next_tail_chunk(5).unwrap(), b"ijklm"); + finished.drop_rest(); + assert_eq!(finished.next_tail_chunk(5).unwrap(), b"nop"); + assert!(finished.tail_out()); + assert_eq!(finished.elision().unwrap().bytes, 4); + } +} diff --git a/crates/server/src/process.rs b/crates/server/src/process.rs index daeff672..a0b1b455 100644 --- a/crates/server/src/process.rs +++ b/crates/server/src/process.rs @@ -26,6 +26,7 @@ use tokio::task::AbortHandle; use yas_wire::process as wire; use yas_wire::schema::process as process_schema; +use crate::output_keep::{Elision, KeptOutput}; #[cfg(unix)] use crate::pty; @@ -501,6 +502,22 @@ struct FinalRecord { kill_cause: u8, code: u32, detail: &'static str, + /// KEEP_OUTPUT: what was dropped of stdout and of stderr. + elided: [Option; 2], +} + +impl FinalRecord { + fn exit(&self) -> NativeExit { + NativeExit { + elided: self.elided, + ..native_exit( + self.reason, + self.kill_cause, + self.code, + self.detail.as_bytes(), + ) + } + } } /// Transport-neutral process catalogue snapshot used by the YAS adapter. @@ -531,6 +548,8 @@ pub(crate) struct NativeExit { pub(crate) reason: u8, pub(crate) code: i32, pub(crate) detail: Vec, + /// KEEP_OUTPUT: what was dropped of stdout and of stderr. + pub(crate) elided: [Option; 2], } #[derive(Clone, Debug)] @@ -544,6 +563,9 @@ pub(crate) struct NativeSpawnRequest { /// LEAVE_RESIDUE: how long the streams are forwarded after the direct child exits (None: /// until they close). Only read with the flag. pub(crate) residue_grace: Option, + /// KEEP_OUTPUT: the bytes of each output stream's head and tail that are sent (the middle + /// is dropped, and counted). + pub(crate) keep_output: Option<(u64, usize)>, pub(crate) cwd: Option>, pub(crate) argv: Vec>, pub(crate) env: Vec<(Vec, Vec)>, @@ -838,12 +860,7 @@ impl Server { stdin_received: record.stdin_received, stdout_produced: record.stdout_next, stderr_produced: record.stderr_next, - exit: Some(native_exit( - record.reason, - record.kill_cause, - record.code, - record.detail.as_bytes(), - )), + exit: Some(record.exit()), }); } records.sort_unstable_by_key(|record| record.process_handle); @@ -1104,6 +1121,7 @@ struct Pending { preserve_residual: bool, leave_residue: bool, residue_grace: Option, + keep_output: Option<(u64, usize)>, stdin_null: bool, request_bytes: usize, endpoint: Weak, @@ -1149,6 +1167,17 @@ struct Binding { struct StreamState { next: u64, + /// KEEP_OUTPUT: the head and tail of this stream that are sent, the middle dropped. + kept: Option, +} + +impl StreamState { + fn new(keep_output: Option<(u64, usize)>) -> Self { + Self { + next: 0, + kept: keep_output.map(|(head, tail)| KeptOutput::new(head, tail)), + } + } } #[derive(Clone, Copy)] @@ -1195,6 +1224,8 @@ struct RecordInner { stderr_readers: u8, /// Output readers waiting for the owner to take its window (`owner_with_room`). paced_readers: u8, + /// KEEP_OUTPUT streams are to send their tails now (`flush_kept`). + flush_kept: bool, child_outcome: Option, tree_cleanup_done: bool, exit_override: Option, @@ -1346,6 +1377,7 @@ impl Manager { preserve_residual: request.preserve_residual, leave_residue: owned.flags & PROCESS_SPAWN_LEAVE_RESIDUE != 0, residue_grace: request.residue_grace, + keep_output: request.keep_output, stdin_null: owned.flags & PROCESS_SPAWN_STDIN_NULL != 0, request_bytes, endpoint: Arc::downgrade(&self.endpoint), @@ -1525,8 +1557,9 @@ impl Manager { }, stdin_closed_by_child: false, stdin_writer_done: stdin.is_none(), - stdout: StreamState { next: 0 }, - stderr: (!merged).then_some(StreamState { next: 0 }), + stdout: StreamState::new(pending.keep_output), + stderr: (!merged).then(|| StreamState::new(pending.keep_output)), + flush_kept: false, stdout_readers: 1, stderr_readers: if merged { 0 } else { 1 }, paced_readers: 0, @@ -1950,12 +1983,7 @@ impl Manager { stdout_next: record.stdout_next, stderr_next: record.stderr_next, stdin_window: 0, - exit: Some(native_exit( - record.reason, - record.kill_cause, - record.code, - record.detail.as_bytes(), - )), + exit: Some(record.exit()), }) } else { Err(NativeError::NotFound) @@ -2068,6 +2096,8 @@ impl Manager { ) .await; for record in &ordinary { + // The owner is gone: kept tails go to the watchers at once. + flush_kept(record).await; finish_pipes(record); } // Pipe abortion makes terminal publication eligible. Keep shutdown @@ -2898,121 +2928,274 @@ async fn stdin_writer( async fn output_reader(record: Arc, stream: u8, mut reader: impl AsyncRead + Unpin) { let mut buffer = vec![0u8; OUTPUT_FRAME_PAYLOAD]; + let kept = output_state(&mut record.inner.lock().unwrap(), stream) + .kept + .is_some(); + // Whether the stream stopped for a flush (`flush_kept`) rather than at its end. + let mut flushed = false; loop { // The process's owner takes all of its output: the pipe is read no further ahead of // what the owner has taken than its window, so a writer faster than the owner's // Transfer blocks on its pipe. (It used to be evicted, which closed the owner's whole // Process endpoint.) Other watchers are still dropped when they fall a window behind. - let owner = owner_with_room(&record, stream).await; - match reader.read(&mut buffer).await { + // KEEP_OUTPUT: once the head is out, what the pipe gives is kept or dropped here, so it + // is read at the writer's speed; only what goes out waits for the owner. + let paced = !kept + || output_state(&mut record.inner.lock().unwrap(), stream) + .kept + .as_ref() + .is_some_and(KeptOutput::sends_now); + let owner = match (paced, kept) { + (false, _) => None, + (true, false) => owner_with_room(&record, stream).await, + (true, true) => tokio::select! { + owner = owner_with_room(&record, stream) => owner, + () = flush_requested(&record) => { + flushed = true; + break; + } + }, + }; + let read = if kept { + tokio::select! { + read = reader.read(&mut buffer) => read, + () = flush_requested(&record) => { + flushed = true; + break; + } + } + } else { + reader.read(&mut buffer).await + }; + match read { Ok(0) => break, Err(_) => { host_failure(&record, "process output pipe read failed"); break; } Ok(n) => { - // Room in the owner's event queue too, taken before the lock: its frame never - // finds the queue full. - let owner = match owner { - Some((endpoint_id, process_id, events)) => events - .reserve_owned() - .await - .ok() - .map(|permit| (endpoint_id, process_id, permit)), - None => None, - }; - let mut owner = owner; - let mut inner = record.inner.lock().unwrap(); - let state = if stream == PROCESS_STREAM_STDOUT { - &mut inner.stdout - } else { - inner.stderr.as_mut().expect("separate stderr") - }; - let offset = state.next; - let Some(next) = offset.checked_add(n as u64) else { - drop(inner); + let owner = reserve(owner).await; + if deliver(&record, stream, Output::Read(&buffer[..n]), owner).is_err() { protocol_violation(&record); return; - }; - state.next = next; - let mut evicted = Vec::new(); - let mut index = 0; - while index < inner.bindings.len() { - if let Some((endpoint_id, process_id, _)) = owner - && inner.bindings[index].endpoint_id == endpoint_id - && inner.bindings[index].process_id == process_id - { - let (_, _, permit) = owner.take().expect("owner permit"); - permit.send(NativeEventEnvelope { - event: NativeEvent::Output { - process_id, - stream, - offset, - data: buffer[..n].to_vec(), - }, - _guard: None, - }); - let binding = &mut inner.bindings[index]; - let credit = if stream == PROCESS_STREAM_STDOUT { - &mut binding.stdout - } else { - binding.stderr.as_mut().expect("separate stderr binding") - }; - credit.frames.push_back(next); - index += 1; - continue; - } - let has_credit = { - let binding = &inner.bindings[index]; - let credit = if stream == PROCESS_STREAM_STDOUT { - &binding.stdout - } else { - binding.stderr.as_ref().expect("separate stderr binding") - }; - let available = offset - .checked_sub(credit.acked) - .and_then(|debt| PROCESS_DEFAULT_STREAM_WINDOW.checked_sub(debt)); - available.is_some_and(|bytes| bytes >= n as u64) - && credit.frames.len() < PROCESS_MAX_UNACKED_PACKETS - }; - if !has_credit { - evicted.push(remove_binding_at(&mut inner, index)); - continue; - } - let process_id = inner.bindings[index].process_id; - let sent = inner.bindings[index].out.send_output( - process_id, - stream, - offset, - &buffer[..n], - ); - if sent { - let binding = &mut inner.bindings[index]; - let credit = if stream == PROCESS_STREAM_STDOUT { - &mut binding.stdout - } else { - binding.stderr.as_mut().expect("separate stderr binding") - }; - credit.frames.push_back(next); - index += 1; - } else { - evicted.push(remove_binding_at(&mut inner, index)); - } - } - drop(inner); - for binding in evicted { - // Its endpoint goes on: free the slot, or the process holds one for good. - if let Some(endpoint) = binding.endpoint.upgrade() { - remove_bound_slot(&endpoint, binding.process_id, &record); - } - binding.out.evict(binding.process_id); } } } } + if kept { + if send_kept_tail(&record, stream).await.is_err() { + protocol_violation(&record); + return; + } + // Flushed, the stream may go on (a residue holds it, or it is about to be aborted): + // what it gives now goes to nobody. + while flushed && matches!(reader.read(&mut buffer).await, Ok(1..)) {} + } stream_closed(&record, stream); } +/// KEEP_OUTPUT: send what the stream kept of its end, a frame at a time as the owner takes them. +/// The stream is finished: nothing more of it goes out. +async fn send_kept_tail(record: &Arc, stream: u8) -> Result<(), ()> { + if let Some(kept) = output_state(&mut record.inner.lock().unwrap(), stream) + .kept + .as_mut() + { + kept.finish(); + } + loop { + let owner = reserve(owner_with_room(record, stream).await).await; + if !deliver(record, stream, Output::Tail, owner)? { + break; + } + } + record.changed.notify_waiters(); + Ok(()) +} + +/// Once the KEEP_OUTPUT streams are to send their tails now ([`flush_kept`]). +async fn flush_requested(record: &Record) { + loop { + let changed = record.changed.notified(); + if record.inner.lock().unwrap().flush_kept { + return; + } + changed.await; + } +} + +/// KEEP_OUTPUT: have the streams send their tails now, as the exit is about to be reported or +/// the streams stopped, and wait for them as [`drain_paced`] waits: while a reader waits for +/// the owner to take its window, or until none has for `DRAIN_IDLE`. What a stream gives after +/// its tail goes to nobody. +async fn flush_kept(record: &Record) { + let tails_out = |inner: &RecordInner| { + let out = |state: Option<&StreamState>, readers: u8| { + state + .and_then(|state| state.kept.as_ref()) + .is_none_or(|kept| kept.tail_out() || readers == 0) + }; + out(Some(&inner.stdout), inner.stdout_readers) + && out(inner.stderr.as_ref(), inner.stderr_readers) + }; + { + let mut inner = record.inner.lock().unwrap(); + if tails_out(&inner) { + return; + } + inner.flush_kept = true; + } + record.changed.notify_waiters(); + loop { + let changed = record.changed.notified(); + let paced = { + let inner = record.inner.lock().unwrap(); + if tails_out(&inner) { + return; + } + inner.paced_readers > 0 + }; + if paced { + changed.await; + } else if tokio::time::timeout(DRAIN_IDLE, changed).await.is_err() { + return; + } + } +} + +/// What an output reader sends: bytes it read (only their share of a KEEP_OUTPUT stream's head, +/// when it keeps one), or the next frame of a kept tail. +enum Output<'a> { + Read(&'a [u8]), + Tail, +} + +/// The owner's Process events, with room for one frame of output reserved. +type OwnerPermit = (u64, u32, mpsc::OwnedPermit); + +/// Reserve room in the owner's event queue ([`owner_with_room`]'s answer), before the record's +/// lock is taken: its frame never finds the queue full. +async fn reserve( + owner: Option<(u64, u32, mpsc::Sender)>, +) -> Option { + let (endpoint_id, process_id, events) = owner?; + let permit = events.reserve_owned().await.ok()?; + Some((endpoint_id, process_id, permit)) +} + +fn output_state(inner: &mut RecordInner, stream: u8) -> &mut StreamState { + if stream == PROCESS_STREAM_STDOUT { + &mut inner.stdout + } else { + inner.stderr.as_mut().expect("separate stderr") + } +} + +/// Send output of `stream` to the process's bindings: to the owner through `owner`, reserved +/// when it had room, and to each watcher that keeps up (the others are dropped). False when +/// there was no tail left to send; Err past a u64 of offset (a protocol violation). +fn deliver( + record: &Arc, + stream: u8, + output: Output<'_>, + owner: Option, +) -> Result { + let mut owner = owner; + let mut inner = record.inner.lock().unwrap(); + let tail; + let state = output_state(&mut inner, stream); + let data: &[u8] = match (output, state.kept.as_mut()) { + (Output::Read(data), None) => data, + (Output::Read(data), Some(kept)) => kept.feed(data), + (Output::Tail, Some(kept)) => match kept.next_tail_chunk(OUTPUT_FRAME_PAYLOAD) { + Some(chunk) => { + tail = chunk; + &tail + } + None => return Ok(false), + }, + (Output::Tail, None) => return Ok(false), + }; + if data.is_empty() { + return Ok(true); + } + let offset = state.next; + let Some(next) = offset.checked_add(data.len() as u64) else { + return Err(()); + }; + state.next = next; + let mut evicted = Vec::new(); + let mut index = 0; + while index < inner.bindings.len() { + if let Some((endpoint_id, process_id, _)) = owner + && inner.bindings[index].endpoint_id == endpoint_id + && inner.bindings[index].process_id == process_id + { + let (_, _, permit) = owner.take().expect("owner permit"); + permit.send(NativeEventEnvelope { + event: NativeEvent::Output { + process_id, + stream, + offset, + data: data.to_vec(), + }, + _guard: None, + }); + let binding = &mut inner.bindings[index]; + let credit = if stream == PROCESS_STREAM_STDOUT { + &mut binding.stdout + } else { + binding.stderr.as_mut().expect("separate stderr binding") + }; + credit.frames.push_back(next); + index += 1; + continue; + } + let has_credit = { + let binding = &inner.bindings[index]; + let credit = if stream == PROCESS_STREAM_STDOUT { + &binding.stdout + } else { + binding.stderr.as_ref().expect("separate stderr binding") + }; + let available = offset + .checked_sub(credit.acked) + .and_then(|debt| PROCESS_DEFAULT_STREAM_WINDOW.checked_sub(debt)); + available.is_some_and(|bytes| bytes >= data.len() as u64) + && credit.frames.len() < PROCESS_MAX_UNACKED_PACKETS + }; + if !has_credit { + evicted.push(remove_binding_at(&mut inner, index)); + continue; + } + let process_id = inner.bindings[index].process_id; + let sent = inner.bindings[index] + .out + .send_output(process_id, stream, offset, data); + if sent { + let binding = &mut inner.bindings[index]; + let credit = if stream == PROCESS_STREAM_STDOUT { + &mut binding.stdout + } else { + binding.stderr.as_mut().expect("separate stderr binding") + }; + credit.frames.push_back(next); + index += 1; + } else { + evicted.push(remove_binding_at(&mut inner, index)); + } + } + drop(inner); + for binding in evicted { + // Its endpoint goes on: free the slot, or the process holds one for good. + if let Some(endpoint) = binding.endpoint.upgrade() { + remove_bound_slot(&endpoint, binding.process_id, record); + } + binding.out.evict(binding.process_id); + } + Ok(true) +} + /// Waits until the process's owner, while it is bound, can take another whole frame of /// `stream` (its window and unacknowledged frames); answers where to send it. None when the /// owner is not bound (it left, or detached): the output then goes to the watchers alone. @@ -3224,6 +3407,7 @@ fn schedule_residual_cleanup(record: Arc) { record.changed.notify_waiters(); try_queue_terminal(&record); } else { + flush_kept(&record).await; abandon_residue(&record); } return; @@ -3284,6 +3468,8 @@ fn schedule_residual_cleanup(record: Arc) { if !cleanup_failed { drain_paced(&record).await; } + // What the streams kept of their ends goes out before they are stopped. + flush_kept(&record).await; let (stdin_abort, output_aborts) = { let mut inner = record.inner.lock().unwrap(); if inner.tree_cleanup_done { @@ -3524,6 +3710,16 @@ fn try_queue_terminal(record: &Arc) { } inner.terminal_queued = true; let (reason, kill_cause, code) = outcome_fields(outcome, inner.exit_override); + // A kept tail that never went out (its reader was stopped first) counts as dropped. + let elided = |state: Option<&mut StreamState>| { + let kept = state?.kept.as_mut()?; + kept.drop_rest(); + kept.elision() + }; + let elided = [ + elided(Some(&mut inner.stdout)), + elided(inner.stderr.as_mut()), + ]; let final_record = Arc::new(FinalRecord { generation: record.generation, pid: record.pid, @@ -3541,6 +3737,7 @@ fn try_queue_terminal(record: &Arc) { kill_cause, code, detail: inner.cleanup_detail, + elided, }); inner.stdin_controller = None; (std::mem::take(&mut inner.bindings), final_record) @@ -3566,16 +3763,9 @@ fn try_queue_terminal(record: &Arc) { finish_terminal(record_for_guard, final_for_guard); } }); - let _ = binding.out.send_exit( - binding.process_id, - native_exit( - final_record.reason, - final_record.kill_cause, - final_record.code, - final_record.detail.as_bytes(), - ), - guard, - ); + let _ = binding + .out + .send_exit(binding.process_id, final_record.exit(), guard); } } @@ -3590,12 +3780,14 @@ fn native_exit(reason: u8, kill_cause: u8, code: u32, detail: &[u8]) -> NativeEx reason: process_schema::EXIT_REASON_UNKNOWN as u8, code: code as i32, detail: detail.to_vec(), + elided: [None; 2], }, PROCESS_EXIT_SIGNALLED => NativeExit { kind: wire::ExitKind::Signal, reason: portable_signal_reason(code), code: code as i32, detail: detail.to_vec(), + elided: [None; 2], }, PROCESS_EXIT_KILLED => NativeExit { kind: wire::ExitKind::Killed, @@ -3608,6 +3800,7 @@ fn native_exit(reason: u8, kill_cause: u8, code: u32, detail: &[u8]) -> NativeEx }, code: 0, detail: detail.to_vec(), + elided: [None; 2], }, _ => NativeExit { kind: wire::ExitKind::Other, @@ -3618,6 +3811,7 @@ fn native_exit(reason: u8, kill_cause: u8, code: u32, detail: &[u8]) -> NativeEx } else { detail.to_vec() }, + elided: [None; 2], }, } } @@ -3724,6 +3918,7 @@ async fn terminate_record(record: &Arc, cause: u8, grace: Duration) { }; let _ = tokio::time::timeout(grace.max(Duration::from_millis(100)), forced).await; } + flush_kept(record).await; finish_pipes(record); let _ = tokio::time::timeout( grace.max(Duration::from_millis(100)), diff --git a/crates/server/src/yas.rs b/crates/server/src/yas.rs index 178306f9..17caa6eb 100644 --- a/crates/server/src/yas.rs +++ b/crates/server/src/yas.rs @@ -3047,6 +3047,8 @@ enum ProcessOperationOutcome { fingerprint: [u8; 32], stdout_receive_credit: u64, stderr_receive_credit: u64, + /// SPAWN_REPORT_EXIT: send the exit to this session as an EXIT event. + report_exit: bool, outcome: Result, }, Attach { @@ -18211,6 +18213,7 @@ impl Session { let operation_id = request.operation_id; let stdout_receive_credit = request.stdout_receive_credit; let stderr_receive_credit = request.stderr_receive_credit; + let report_exit = request.flags & yas_wire::schema::process::SPAWN_REPORT_EXIT as u16 != 0; let session = self.process.as_ref().ok_or(())?.session.clone(); let internal = self.internal.clone(); let cancellation = self.cancellation.clone(); @@ -18228,6 +18231,7 @@ impl Session { fingerprint, stdout_receive_credit, stderr_receive_credit, + report_exit, outcome, }, _slot: None, @@ -25385,6 +25389,7 @@ impl Session { fingerprint, stdout_receive_credit, stderr_receive_credit, + report_exit, outcome, } => { if kind != yas_process_wire::request_kind::SPAWN { @@ -25442,6 +25447,7 @@ impl Session { } }; let body = bundle.encode().map_err(|_| ())?; + let report_exit = report_exit.then_some(bundle.process_handle); self.record_process_replay( operation_id, kind, @@ -25465,7 +25471,7 @@ impl Session { self.remove_process_attachment(attachment_id, true); return Err(()); } - self.activate_process_attachment(attachment_id, events)?; + self.activate_process_attachment(attachment_id, events, report_exit)?; Ok(()) } ProcessOperationOutcome::Attach { @@ -25526,7 +25532,7 @@ impl Session { self.remove_process_attachment(attachment_id, true); return Err(()); } - self.activate_process_attachment(attachment_id, events)?; + self.activate_process_attachment(attachment_id, events, None)?; Ok(()) } ProcessOperationOutcome::Control { @@ -25801,10 +25807,13 @@ impl Session { )) } + /// Start forwarding an installed attachment's streams; with `report_exit` (the handle of a + /// process spawned with SPAWN_REPORT_EXIT), its exit is sent as an EXIT event too. fn activate_process_attachment( &mut self, attachment_id: u32, events: super::yas_process::AttachmentEvents, + report_exit: Option, ) -> Result<(), ()> { let attachment = self .process @@ -25828,6 +25837,7 @@ impl Session { attachment.stdout_transfer, stdout_flow, attachment.stderr_transfer.zip(stderr_flow), + report_exit, self.out.clone(), self.internal.clone(), self.cancellation.clone(), @@ -26083,6 +26093,8 @@ impl Session { self.outbound_sensitive.remove(&transfer_id); } } + // A SPAWN_REPORT_EXIT process whose attachment goes before its exit reports none: the + // client knows (it dropped a stream, detached, or saw its streams reset) and WAITs. if detach { tokio::spawn(async move { let _ = attachment.control.detach().await; @@ -31003,6 +31015,36 @@ struct ProcessOutputChunk { data: Vec, } +/// A SPAWN_REPORT_EXIT process's exit, as the EXIT event its session gets unasked. +async fn send_exit_report( + out: &FrameSender, + process_handle: u64, + exit: super::yas_process::ExitInfo, + connection: &ConnectionCancellation, +) { + // KEEP_OUTPUT: what was dropped of each stream. + let elided = exit + .elided + .iter() + .zip([false, true]) + .filter_map(|(elided, stderr)| elided.map(|elided| elided.extension(stderr))) + .collect(); + let report = yas_process_wire::ExitReport { + process_handle, + exit: exit.into_record(monotonic_ns()), + extensions: Extensions(elided), + }; + let _ = send_event_with_sensitivity( + out, + family::PROCESS, + yas_process_wire::event_kind::EXIT, + &report, + connection, + true, + ) + .await; +} + #[allow(clippy::too_many_arguments)] fn spawn_process_attachment( attachment_id: u32, @@ -31012,12 +31054,14 @@ fn spawn_process_attachment( stdout_transfer: u32, stdout_flow: Arc, stderr: Option<(u32, Arc)>, + report_exit: Option, out: FrameSender, internal: mpsc::Sender, connection: ConnectionCancellation, cancellation: ConnectionCancellation, ) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { + let exit_out = out.clone(); let (stdout_tx, stdout_rx) = mpsc::channel(PROCESS_STREAM_QUEUE); let stdout_task = tokio::spawn(run_process_output( attachment_id, @@ -31113,8 +31157,18 @@ fn spawn_process_attachment( return; } } - Some(super::yas_process::Event::Exit(_)) => { + Some(super::yas_process::Event::Exit(exit)) => { exited = true; + // As a WAIT would answer it, at once: the streams go on with what + // the process wrote before, at the pace of their credit. + // In a task of its own: removing the attachment meanwhile aborts this + // one, and the exit is already this attachment's to report. + if let Some(process_handle) = report_exit { + let (out, connection) = (exit_out.clone(), connection.clone()); + tokio::spawn(async move { + send_exit_report(&out, process_handle, exit, &connection).await; + }); + } let _ = internal .send(Internal::ProcessExited { attachment_id }) .await; @@ -51805,7 +51859,7 @@ mod tests { yas_process_wire::Limits { max_mutation_replays: MAX_PROCESS_OPERATION_REPLAYS as u32, launcher_flags: if cfg!(unix) { - yas_wire::schema::process::SPAWN_LAUNCHER_FLAGS as u32 + yas_wire::schema::process::SPAWN_LAUNCHER_FLAGS_EXTENDED as u32 } else { yas_wire::schema::process::SPAWN_STDIN_NULL as u32 }, @@ -52170,6 +52224,142 @@ mod tests { timeout(TEST_TIMEOUT, server_task).await.unwrap().unwrap(); } + #[cfg(unix)] + #[tokio::test(flavor = "multi_thread")] + async fn process_spawn_with_report_exit_is_sent_its_exit_without_a_wait() { + use std::os::unix::ffi::OsStrExt; + + let sh = std::env::split_paths(&std::env::var_os("PATH").expect("test PATH is set")) + .map(|directory| directory.join("sh")) + .find(|path| path.is_file()) + .expect("sh is on PATH") + .as_os_str() + .as_bytes() + .to_vec(); + let state = super::super::tests::process_transport::test_state( + super::super::process::Server::new(false, true), + ); + let (mut client, codec, hello, server_task) = + start_registered_session(state, &[family::TRANSFER, family::PROCESS]).await; + let descriptor = hello + .families + .iter() + .find(|descriptor| descriptor.family_id == family::PROCESS) + .expect("Process negotiated"); + let limits = yas_process_wire::Limits::from_extensions(&descriptor.limits).unwrap(); + assert_ne!( + limits.launcher_flags & yas_wire::schema::process::SPAWN_REPORT_EXIT as u32, + 0 + ); + + let spawn = yas_process_wire::Spawn { + operation_id: [0x72; 16], + flags: (yas_wire::schema::process::SPAWN_REPORT_EXIT + | yas_wire::schema::process::SPAWN_STDIN_NULL) as u16, + environment_kind: yas_process_wire::EnvironmentKind::Empty, + cwd: yas_process_wire::Cwd::ServerDefault, + argv: vec![ + sh, + b"-c".to_vec(), + b"printf out; printf err >&2; exit 3".to_vec(), + ], + env: Vec::new(), + stdout_receive_credit: 1024 * 1024, + stderr_receive_credit: 1024 * 1024, + extensions: Extensions::default(), + }; + write_request( + &mut client, + &codec, + family::PROCESS, + yas_process_wire::request_kind::SPAWN, + 11, + &spawn, + ) + .await; + let (spawned, _) = next_process_result_collecting_state( + &mut client, + &codec, + yas_process_wire::request_kind::SPAWN, + 11, + ) + .await; + assert_eq!(spawned.status, Status::Ok); + let streams = yas_process_wire::StreamBundle::decode(&spawned.body).unwrap(); + assert!(streams.stdin.is_none()); + let stderr = streams.stderr.as_ref().expect("separate stderr Transfer"); + + // No WAIT is sent: the exit comes on its own, next to the streams. + let mut stdout = Vec::new(); + let mut stderr_bytes = Vec::new(); + let mut stdout_closed = false; + let mut stderr_closed = false; + let mut exit = None; + while !stdout_closed || !stderr_closed || exit.is_none() { + let frame = next_frame(&mut client, &codec).await; + match (frame.header.family, frame.header.kind) { + (family::PROCESS, yas_process_wire::event_kind::EXIT) => { + assert_eq!( + frame.header, + FrameHeader { + sensitive: true, + ..FrameHeader::event( + family::PROCESS, + yas_process_wire::event_kind::EXIT, + ) + } + ); + assert!(exit.is_none(), "a second EXIT"); + let report = yas_process_wire::ExitReport::decode(&frame.payload).unwrap(); + assert_eq!(report.process_handle, streams.process_handle); + assert_eq!( + yas_process_wire::ExitReport::handle_of(&frame.payload), + Some(streams.process_handle) + ); + exit = Some(report.exit); + } + (family::TRANSFER, yas_wire::schema::transfer::event::BYTE_DATA) => { + let data = ByteData::decode(&frame.payload).unwrap(); + let target = if data.transfer_id == streams.stdout.transfer_id { + &mut stdout + } else if data.transfer_id == stderr.transfer_id { + &mut stderr_bytes + } else { + panic!("bytes for unknown Process Transfer") + }; + assert_eq!(data.offset, target.len() as u64); + target.extend_from_slice(&data.data); + } + (family::TRANSFER, yas_wire::schema::transfer::event::CLOSE) => { + let close = Close::decode(&frame.payload).unwrap(); + assert_eq!(close.status, Status::Ok.code()); + if close.transfer_id == streams.stdout.transfer_id { + assert_eq!(close.final_data_bytes, stdout.len() as u64); + stdout_closed = true; + } else if close.transfer_id == stderr.transfer_id { + assert_eq!(close.final_data_bytes, stderr_bytes.len() as u64); + stderr_closed = true; + } else { + panic!("CLOSE for unknown Process Transfer") + } + } + _ => panic!( + "unexpected native frame from a reported Process: {:?}", + frame.header + ), + } + } + let exit = exit.unwrap(); + assert_eq!(exit.kind, yas_process_wire::ExitKind::Code); + assert_eq!(exit.code, 3); + assert_ne!(exit.exited_server_ns, 0); + assert_eq!(stdout, b"out"); + assert_eq!(stderr_bytes, b"err"); + + drop(client); + timeout(TEST_TIMEOUT, server_task).await.unwrap().unwrap(); + } + #[cfg(unix)] #[tokio::test(flavor = "multi_thread")] async fn terminal_query_transfer_and_cancel_are_native_and_correlated() { diff --git a/crates/server/src/yas_process.rs b/crates/server/src/yas_process.rs index 9c7637cc..42cbbfa7 100644 --- a/crates/server/src/yas_process.rs +++ b/crates/server/src/yas_process.rs @@ -139,6 +139,8 @@ pub(crate) struct ExitInfo { pub(crate) reason: u8, pub(crate) code: i32, pub(crate) detail: Vec, + /// KEEP_OUTPUT: what was dropped of stdout and of stderr (in the EXIT event only). + pub(crate) elided: [Option; 2], } impl ExitInfo { @@ -293,9 +295,10 @@ impl Runtime { } pub(crate) fn limits(&self) -> wire::Limits { - // LEAVE_RESIDUE works with Unix process groups and Windows jobs alike. + // LEAVE_RESIDUE works with Unix process groups and Windows jobs alike; REPORT_EXIT is + // the YAS connection's own. wire::Limits { - launcher_flags: schema::process::SPAWN_LAUNCHER_FLAGS as u32, + launcher_flags: schema::process::SPAWN_LAUNCHER_FLAGS_EXTENDED as u32, ..self.server.maxima().limits() } } @@ -381,8 +384,18 @@ impl Session { resolved_cwd: Option>, ) -> Result { let cwd = resolve_cwd(&request.cwd, resolved_cwd)?; - let flags = u8::try_from(request.flags) - .map_err(|_| Error::Invalid("Process SPAWN flags do not fit v1".to_owned()))?; + // REPORT_EXIT asks the YAS connection for an EXIT event, and KEEP_OUTPUT (with it) the + // output's head and tail alone: the process is the same. + let flags = u8::try_from( + request.flags + & !((schema::process::SPAWN_REPORT_EXIT | schema::process::SPAWN_KEEP_OUTPUT) + as u16), + ) + .map_err(|_| Error::Invalid("Process SPAWN flags do not fit v1".to_owned()))?; + let keep_output = request + .keep_output() + .map_err(|error| Error::Invalid(error.to_string()))? + .map(|(head, tail)| (head, tail as usize)); let process_id = self.allocate_process_id()?; let (route, events) = self.install_route(process_id, false)?; let session_env = (request.environment_kind == wire::EnvironmentKind::Session) @@ -406,6 +419,7 @@ impl Session { flags, preserve_residual, residue_grace, + keep_output, cwd, argv: request.argv.clone(), env: request @@ -1110,6 +1124,15 @@ fn native_exit_info(exit: process::NativeExit) -> ExitInfo { reason: exit.reason, code: exit.code, detail: exit.detail, + elided: exit.elided.map(|elided| { + elided.map(|elided| wire::OutputElision { + offset: elided.offset, + bytes: elided.bytes, + lines: elided.lines, + code_points: elided.code_points, + utf16_units: elided.utf16_units, + }) + }), } } @@ -1146,6 +1169,7 @@ mod tests { reason: 0, code: process_handle as i32, detail: Vec::new(), + elided: [None; 2], }, ); } diff --git a/crates/yas/src/generated.rs b/crates/yas/src/generated.rs index a3b9caf4..24df68c2 100644 --- a/crates/yas/src/generated.rs +++ b/crates/yas/src/generated.rs @@ -3791,6 +3791,7 @@ pub const WAIT: u16 = 0x0005; pub mod event { pub const STATE: u16 = 0x0000; pub const STATE_ACK: u16 = 0x0001; +pub const EXIT: u16 = 0x0002; } pub const SPAWN_MERGE_STDERR: u64 = 1; pub const SPAWN_DETACHABLE: u64 = 2; @@ -3798,6 +3799,9 @@ pub const SPAWN_LEAVE_RESIDUE: u64 = 4; pub const SPAWN_STDIN_NULL: u64 = 8; pub const SPAWN_FLAGS: u64 = 3; pub const SPAWN_LAUNCHER_FLAGS: u64 = 12; +pub const SPAWN_REPORT_EXIT: u64 = 16; +pub const SPAWN_KEEP_OUTPUT: u64 = 32; +pub const SPAWN_LAUNCHER_FLAGS_EXTENDED: u64 = 60; pub const ENV_EMPTY: u64 = 0; pub const ENV_SESSION: u64 = 1; pub const CWD_SERVER_DEFAULT: u64 = 0; @@ -3844,6 +3848,10 @@ pub const STREAM_STDERR_CONTENT_KIND: u64 = 2; pub const SPAWN_SURFACE_APP_EXTENSION: u64 = 1; pub const SPAWN_RESOURCE_TAG_EXTENSION: u64 = 2; pub const SPAWN_RESIDUE_GRACE_EXTENSION: u64 = 3; +pub const SPAWN_KEEP_OUTPUT_EXTENSION: u64 = 4; +pub const MAX_KEEP_OUTPUT_TAIL_BYTES: u64 = 1048576; +pub const EXIT_STDOUT_ELIDED_EXTENSION: u64 = 1; +pub const EXIT_STDERR_ELIDED_EXTENSION: u64 = 2; pub const MAX_ARGC: u64 = 1024; pub const MAX_ARG_BYTES: u64 = 1048576; pub const MAX_ARG_LEN: u64 = 65536; @@ -3886,6 +3894,7 @@ pub const LIMIT_MAX_DETACHED_RETENTION_NS: u64 = 9; pub const MAX_MUTATION_REPLAYS: u64 = 65536; pub const LIMIT_MAX_MUTATION_REPLAYS: u64 = 10; pub const LIMIT_LAUNCHER_FLAGS: u64 = 11; +pub const LIMIT_LAUNCHER_FLAGS_EXTENDED: u64 = 19; pub static OPERATIONS: &[super::OperationMetadata] = &[ super::OperationMetadata { name: "WATCH", class: 1, kind: 0, direction: 0, sensitive: 1, compression: 0, datagram: 0, layout: "StateWatch; ResultPrefix + StateWatchResult" }, super::OperationMetadata { name: "UNWATCH", class: 1, kind: 1, direction: 0, sensitive: 0, compression: 0, datagram: 0, layout: "subscription_id:u32; ResultPrefix" }, @@ -3895,11 +3904,15 @@ super::OperationMetadata { name: "CONTROL", class: 1, kind: 4, direction: 0, sen super::OperationMetadata { name: "WAIT", class: 1, kind: 5, direction: 0, sensitive: 1, compression: 0, datagram: 0, layout: "process_handle:u64,timeout_ns:u64,Extensions; ResultPrefix + ExitRecord" }, super::OperationMetadata { name: "STATE", class: 0, kind: 0, direction: 1, sensitive: 1, compression: 0, datagram: 0, layout: "StateEvent" }, super::OperationMetadata { name: "STATE_ACK", class: 0, kind: 1, direction: 0, sensitive: 0, compression: 0, datagram: 0, layout: "StateAck" }, +super::OperationMetadata { name: "EXIT", class: 0, kind: 2, direction: 1, sensitive: 1, compression: 0, datagram: 0, layout: "ExitReport" }, ]; pub static TYPES: &[super::TypeMetadata] = &[ super::TypeMetadata { name: "cwd", layout: "kind:u8,reserved:[u8;3]=0; SERVER_DEFAULT empty, PATH path:bytes_u32, TERMINAL terminal_handle:u64, FS root_handle:u64,component_count:u16,repeated component:bytes_u16" }, super::TypeMetadata { name: "process_record", layout: "process_handle:u64,lifecycle:u8,stream_state:u8,flags:u16,native_pid:u64,owner_session:[u8;16],argv0:bytes_u32,stdin_received:u64,stdout_produced:u64,stderr_produced:u64,retention_deadline_server_ns:u64,exit_present:u8,reserved:[u8;7]=0,optional exit:bytes_u32 containing ExitRecord,Extensions" }, super::TypeMetadata { name: "remove_record", layout: "process_handle:u64" }, +super::TypeMetadata { name: "exit_report", layout: "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions; extension tag 1 stdout OutputElision, tag 2 stderr OutputElision, each present iff KEEP_OUTPUT dropped bytes of that stream" }, +super::TypeMetadata { name: "keep_output_extension", layout: "SPAWN extension tag 4 exact value head_bytes:u64,tail_bytes:u64; only with SPAWN_KEEP_OUTPUT and SPAWN_REPORT_EXIT; tail_bytes at most MAX_KEEP_OUTPUT_TAIL_BYTES" }, +super::TypeMetadata { name: "output_elision", layout: "offset:u64,bytes:u64,lines:u64,code_points:u64,utf16_units:u64; offset is the stream offset where the dropped bytes were (the head's length); lines, code points and UTF-16 units count them as a WHATWG UTF-8 decoder with replacement reads them within the whole stream" }, super::TypeMetadata { name: "exit_record", layout: "kind:u8,reason:u8,reserved:u16=0,code:i32,exited_server_ns:u64,detail:bytes_u32" }, super::TypeMetadata { name: "stream_bundle", layout: "process_handle:u64,flags:u16,reserved:u16=0,stdout_lifetime_offset:u64,stderr_lifetime_offset:u64,optional stdin/stdout/stderr descriptor:bytes_u32 containing sensitive BYTE TransferDescriptor,Extensions" }, super::TypeMetadata { name: "state_entity_body", layout: "ADD/REPLACE complete ProcessRecord; REMOVE process_handle:u64" }, @@ -3924,6 +3937,7 @@ super::LimitMetadata { name: "MAX_STREAM_BUFFER_BYTES_EXTENDED", tag: 15, value_ super::LimitMetadata { name: "MAX_ENVC_EXTENDED", tag: 16, value_type: super::LimitValueType::U32, required: false, hard_min: 1, hard_max: 16384 }, super::LimitMetadata { name: "MAX_PENDING_WAITS", tag: 17, value_type: super::LimitValueType::U32, required: false, hard_min: 1, hard_max: 65536 }, super::LimitMetadata { name: "MAX_PENDING_OPERATIONS", tag: 18, value_type: super::LimitValueType::U32, required: false, hard_min: 1, hard_max: 16384 }, +super::LimitMetadata { name: "LAUNCHER_FLAGS_EXTENDED", tag: 19, value_type: super::LimitValueType::U32, required: false, hard_min: 0, hard_max: 65535 }, ]; pub static CONSTANTS: &[super::ConstantMetadata] = &[ super::ConstantMetadata { name: "SPAWN_MERGE_STDERR", value: 1 }, @@ -3932,6 +3946,9 @@ super::ConstantMetadata { name: "SPAWN_LEAVE_RESIDUE", value: 4 }, super::ConstantMetadata { name: "SPAWN_STDIN_NULL", value: 8 }, super::ConstantMetadata { name: "SPAWN_FLAGS", value: 3 }, super::ConstantMetadata { name: "SPAWN_LAUNCHER_FLAGS", value: 12 }, +super::ConstantMetadata { name: "SPAWN_REPORT_EXIT", value: 16 }, +super::ConstantMetadata { name: "SPAWN_KEEP_OUTPUT", value: 32 }, +super::ConstantMetadata { name: "SPAWN_LAUNCHER_FLAGS_EXTENDED", value: 60 }, super::ConstantMetadata { name: "ENV_EMPTY", value: 0 }, super::ConstantMetadata { name: "ENV_SESSION", value: 1 }, super::ConstantMetadata { name: "CWD_SERVER_DEFAULT", value: 0 }, @@ -3978,6 +3995,10 @@ super::ConstantMetadata { name: "STREAM_STDERR_CONTENT_KIND", value: 2 }, super::ConstantMetadata { name: "SPAWN_SURFACE_APP_EXTENSION", value: 1 }, super::ConstantMetadata { name: "SPAWN_RESOURCE_TAG_EXTENSION", value: 2 }, super::ConstantMetadata { name: "SPAWN_RESIDUE_GRACE_EXTENSION", value: 3 }, +super::ConstantMetadata { name: "SPAWN_KEEP_OUTPUT_EXTENSION", value: 4 }, +super::ConstantMetadata { name: "MAX_KEEP_OUTPUT_TAIL_BYTES", value: 1048576 }, +super::ConstantMetadata { name: "EXIT_STDOUT_ELIDED_EXTENSION", value: 1 }, +super::ConstantMetadata { name: "EXIT_STDERR_ELIDED_EXTENSION", value: 2 }, super::ConstantMetadata { name: "MAX_ARGC", value: 1024 }, super::ConstantMetadata { name: "MAX_ARG_BYTES", value: 1048576 }, super::ConstantMetadata { name: "MAX_ARG_LEN", value: 65536 }, @@ -4020,6 +4041,7 @@ super::ConstantMetadata { name: "LIMIT_MAX_DETACHED_RETENTION_NS", value: 9 }, super::ConstantMetadata { name: "MAX_MUTATION_REPLAYS", value: 65536 }, super::ConstantMetadata { name: "LIMIT_MAX_MUTATION_REPLAYS", value: 10 }, super::ConstantMetadata { name: "LIMIT_LAUNCHER_FLAGS", value: 11 }, +super::ConstantMetadata { name: "LIMIT_LAUNCHER_FLAGS_EXTENDED", value: 19 }, ]; } pub mod net { @@ -5033,6 +5055,7 @@ GoldenVector { name: "yas.process.request.control.header", hex: "400004000901000 GoldenVector { name: "yas.process.request.wait.header", hex: "400005000901000000" }, GoldenVector { name: "yas.process.event.state.header", hex: "4000000008" }, GoldenVector { name: "yas.process.event.state_ack.header", hex: "4000010000" }, +GoldenVector { name: "yas.process.event.exit.header", hex: "4000020008" }, GoldenVector { name: "yas.net.request.open.header", hex: "410000000901000000" }, GoldenVector { name: "yas.net.request.close.header", hex: "410001000901000000" }, GoldenVector { name: "yas.net.event.datagram.header", hex: "4100000008" }, diff --git a/crates/yas/src/process.rs b/crates/yas/src/process.rs index 5a70d2cf..a3f9df9a 100644 --- a/crates/yas/src/process.rs +++ b/crates/yas/src/process.rs @@ -198,8 +198,8 @@ pub struct Spawn { impl Spawn { fn validate(&self) -> Result<()> { validate_operation_id(&self.operation_id)?; - let known = - crate::schema::process::SPAWN_FLAGS | crate::schema::process::SPAWN_LAUNCHER_FLAGS; + let known = crate::schema::process::SPAWN_FLAGS + | crate::schema::process::SPAWN_LAUNCHER_FLAGS_EXTENDED; if self.flags & !(known as u16) != 0 { return Err(Error::Invalid("Process spawn flags")); } @@ -220,9 +220,56 @@ impl Spawn { "Process residue grace without LEAVE_RESIDUE", )); } + let keep_output = self.flags & crate::schema::process::SPAWN_KEEP_OUTPUT as u16 != 0; + let report_exit = self.flags & crate::schema::process::SPAWN_REPORT_EXIT as u16 != 0; + if keep_output != self.keep_output()?.is_some() { + return Err(Error::Invalid( + "Process KEEP_OUTPUT flag and extension go together", + )); + } + if keep_output && !report_exit { + return Err(Error::Invalid("Process KEEP_OUTPUT without REPORT_EXIT")); + } Ok(()) } + /// `KEEP_OUTPUT`: how many bytes of each output stream's head and tail are sent (the + /// middle is dropped, and counted in the EXIT event), or None: all of it. + pub fn keep_output(&self) -> Result> { + let Some(extension) = self.extensions.0.iter().find(|extension| { + extension.tag == crate::schema::process::SPAWN_KEEP_OUTPUT_EXTENSION as u16 + }) else { + return Ok(None); + }; + let value: [u8; 16] = extension + .value + .as_slice() + .try_into() + .map_err(|_| Error::Invalid("Process keep-output extension"))?; + let head = u64::from_le_bytes(value[..8].try_into().expect("8 bytes")); + let tail = u64::from_le_bytes(value[8..].try_into().expect("8 bytes")); + if tail > crate::schema::process::MAX_KEEP_OUTPUT_TAIL_BYTES { + return Err(limit( + "Process keep-output tail bytes", + tail, + crate::schema::process::MAX_KEEP_OUTPUT_TAIL_BYTES, + )); + } + Ok(Some((head, tail))) + } + + /// The `KEEP_OUTPUT` extension keeping `head` and `tail` bytes of each output stream. + pub fn keep_output_extension(head: u64, tail: u64) -> Extension { + let mut value = Vec::with_capacity(16); + value.extend_from_slice(&head.to_le_bytes()); + value.extend_from_slice(&tail.to_le_bytes()); + Extension { + tag: crate::schema::process::SPAWN_KEEP_OUTPUT_EXTENSION as u16, + required: true, + value, + } + } + /// How long a `LEAVE_RESIDUE` process's streams are forwarded after its direct child /// exits, or None: until they close. pub fn residue_grace_ns(&self) -> Result> { @@ -914,6 +961,124 @@ impl Decode for RemovedProcess { } } +/// EXIT: a process's final exit, sent to the session that spawned it with +/// `SPAWN_REPORT_EXIT`, once the exit is final (as WAIT would answer it). +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct ExitReport { + pub process_handle: u64, + pub exit: ExitRecord, + pub extensions: Extensions, +} + +impl Encode for ExitReport { + fn encode_to(&self, out: &mut Vec) -> Result<()> { + validate_handle(self.process_handle, "Process handle")?; + put_u64(out, self.process_handle); + put_bytes_u32(out, &self.exit.encode()?)?; + self.extensions.encode_tail(out) + } +} + +impl Decode for ExitReport { + fn decode(input: &[u8]) -> Result { + let mut decoder = Decoder::new(input); + let process_handle = decoder.u64()?; + let exit = ExitRecord::decode(decoder.len_bytes_u32()?)?; + let value = Self { + process_handle, + exit, + extensions: decoder.extensions()?, + }; + decoder.finish()?; + validate_handle(value.process_handle, "Process handle")?; + Ok(value) + } +} + +impl ExitReport { + /// The process handle of an EXIT event's payload, without decoding the rest. + pub fn handle_of(payload: &[u8]) -> Option { + Some(u64::from_le_bytes(payload.get(..8)?.try_into().ok()?)) + } + + /// What `KEEP_OUTPUT` dropped of stdout (`stderr` false) or stderr, if anything. + pub fn elided(&self, stderr: bool) -> Result> { + let tag = if stderr { + crate::schema::process::EXIT_STDERR_ELIDED_EXTENSION + } else { + crate::schema::process::EXIT_STDOUT_ELIDED_EXTENSION + }; + self.extensions + .0 + .iter() + .find(|extension| extension.tag == tag as u16) + .map(|extension| OutputElision::decode(&extension.value)) + .transpose() + } +} + +/// What `KEEP_OUTPUT` dropped of an output stream: from `offset` on (the head's length), +/// `bytes` bytes that decode, as a WHATWG UTF-8 decoder with replacement reads them within the +/// whole stream, to `code_points` characters and `utf16_units` UTF-16 code units, `lines` of +/// them `\n`. The tail follows the head in the stream. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct OutputElision { + pub offset: u64, + pub bytes: u64, + pub lines: u64, + pub code_points: u64, + pub utf16_units: u64, +} + +impl OutputElision { + /// As the EXIT event's extension for stdout (`stderr` false) or stderr. + pub fn extension(&self, stderr: bool) -> Extension { + let tag = if stderr { + crate::schema::process::EXIT_STDERR_ELIDED_EXTENSION + } else { + crate::schema::process::EXIT_STDOUT_ELIDED_EXTENSION + }; + let mut value = Vec::with_capacity(40); + for field in [ + self.offset, + self.bytes, + self.lines, + self.code_points, + self.utf16_units, + ] { + put_u64(&mut value, field); + } + Extension { + tag: tag as u16, + required: false, + value, + } + } +} + +impl Decode for OutputElision { + fn decode(input: &[u8]) -> Result { + let mut decoder = Decoder::new(input); + let value = Self { + offset: decoder.u64()?, + bytes: decoder.u64()?, + lines: decoder.u64()?, + code_points: decoder.u64()?, + utf16_units: decoder.u64()?, + }; + decoder.finish()?; + if value.bytes == 0 + || value.lines > value.code_points + || value.code_points > value.bytes + || value.utf16_units < value.code_points + || value.utf16_units > value.code_points.saturating_mul(2) + { + return Err(Error::Invalid("Process output elision")); + } + Ok(value) + } +} + /// Process family maxima, as a server selects them in HELLO. /// /// The first ten fields are the family's original limits (tags 1–10). A @@ -937,8 +1102,10 @@ pub struct Limits { pub max_stream_buffer_bytes: u64, pub max_detached_retention_ns: u64, pub max_mutation_replays: u32, - /// SPAWN flags of `SPAWN_LAUNCHER_FLAGS` the server honours (LEAVE_RESIDUE, STDIN_NULL); - /// 0 from servers that predate them. + /// SPAWN flags of `SPAWN_LAUNCHER_FLAGS_EXTENDED` the server honours (LEAVE_RESIDUE, + /// STDIN_NULL, REPORT_EXIT); 0 from servers that predate them. LAUNCHER_FLAGS (tag 11) + /// carries those of `SPAWN_LAUNCHER_FLAGS`, all older clients accept; + /// LAUNCHER_FLAGS_EXTENDED (tag 19) carries all of them, when there are more. pub launcher_flags: u32, /// Pending `WAIT`s one session may hold. pub max_pending_waits: u32, @@ -962,7 +1129,7 @@ impl Limits { max_mutation_replays: crate::schema::process::MAX_MUTATION_REPLAYS as u32, max_pending_waits: crate::schema::process::MAX_PENDING_WAITS as u32, max_pending_operations: crate::schema::process::MAX_PENDING_OPERATIONS as u32, - launcher_flags: crate::schema::process::SPAWN_LAUNCHER_FLAGS as u32, + launcher_flags: crate::schema::process::SPAWN_LAUNCHER_FLAGS_EXTENDED as u32, }; /// The original hard maxima: what an unconfigured server enforces, and @@ -1096,9 +1263,17 @@ impl Limits { self.max_pending_operations, )); } - if self.launcher_flags != 0 { + let legacy_launcher_flags = + self.launcher_flags & crate::schema::process::SPAWN_LAUNCHER_FLAGS as u32; + if legacy_launcher_flags != 0 { extensions.push(limit_u32( crate::schema::process::LIMIT_LAUNCHER_FLAGS, + legacy_launcher_flags, + )); + } + if self.launcher_flags != legacy_launcher_flags { + extensions.push(limit_u32( + crate::schema::process::LIMIT_LAUNCHER_FLAGS_EXTENDED, self.launcher_flags, )); } @@ -1126,6 +1301,7 @@ impl Limits { crate::schema::process::LIMIT_MAX_ENVC_EXTENDED as u16, crate::schema::process::LIMIT_MAX_PENDING_WAITS as u16, crate::schema::process::LIMIT_MAX_PENDING_OPERATIONS as u16, + crate::schema::process::LIMIT_LAUNCHER_FLAGS_EXTENDED as u16, ]; reject_unknown_required(extensions, &known)?; let legacy = Self::DEFAULT; @@ -1145,6 +1321,14 @@ impl Limits { extensions, crate::schema::process::LIMIT_MAX_STREAM_BUFFER_BYTES, )?; + let legacy_launcher_flags = + if extensions.0.iter().any(|extension| { + extension.tag == crate::schema::process::LIMIT_LAUNCHER_FLAGS as u16 + }) { + read_limit_u32(extensions, crate::schema::process::LIMIT_LAUNCHER_FLAGS)? + } else { + 0 + }; let value = Self { max_argc: read_limit_u32(extensions, crate::schema::process::LIMIT_MAX_ARGC)?, max_arg_bytes: read_limit_u32(extensions, crate::schema::process::LIMIT_MAX_ARG_BYTES)?, @@ -1175,13 +1359,14 @@ impl Limits { extensions, crate::schema::process::LIMIT_MAX_MUTATION_REPLAYS, )?, - launcher_flags: if extensions.0.iter().any(|extension| { - extension.tag == crate::schema::process::LIMIT_LAUNCHER_FLAGS as u16 - }) { - read_limit_u32(extensions, crate::schema::process::LIMIT_LAUNCHER_FLAGS)? - } else { - 0 - }, + // A flag this side does not know is one it never sets: ignored, so a later flag + // needs no new tag. + launcher_flags: read_optional_limit_u32( + extensions, + crate::schema::process::LIMIT_LAUNCHER_FLAGS_EXTENDED, + )? + .map(|flags| flags & crate::schema::process::SPAWN_LAUNCHER_FLAGS_EXTENDED as u32) + .unwrap_or(legacy_launcher_flags), max_pending_waits: u32_or( crate::schema::process::LIMIT_MAX_PENDING_WAITS, legacy.max_pending_waits, @@ -1204,6 +1389,11 @@ impl Limits { || value.max_processes < max_processes || value.max_pending_spawns < max_pending_spawns || value.max_stream_buffer_bytes < max_stream_buffer_bytes + // Tag 11 carries the v1 launcher flags and tag 19 all of them: the + // extended set keeps those the legacy tag promised, and no others of theirs. + || legacy_launcher_flags & !(crate::schema::process::SPAWN_LAUNCHER_FLAGS as u32) != 0 + || value.launcher_flags & crate::schema::process::SPAWN_LAUNCHER_FLAGS as u32 + != legacy_launcher_flags { return Err(Error::Invalid("Process family limit")); } @@ -1332,6 +1522,7 @@ fn validate_spawn_extensions(extensions: &Extensions) -> Result<()> { crate::schema::process::SPAWN_SURFACE_APP_EXTENSION as u16, crate::schema::process::SPAWN_RESOURCE_TAG_EXTENSION as u16, crate::schema::process::SPAWN_RESIDUE_GRACE_EXTENSION as u16, + crate::schema::process::SPAWN_KEEP_OUTPUT_EXTENSION as u16, ]; reject_unknown_required(extensions, &known)?; extension_u64( @@ -1496,6 +1687,50 @@ fn read_limit_u64(extensions: &Extensions, tag: u64) -> Result { mod tests { use super::*; + fn elision_bytes(fields: [u64; 5]) -> Vec { + fields + .iter() + .flat_map(|field| field.to_le_bytes()) + .collect() + } + + #[test] + fn an_output_elision_is_its_five_counts_and_they_must_agree() { + let elision = OutputElision { + offset: 7, + bytes: 10, + lines: 2, + code_points: 6, + utf16_units: 8, + }; + let encoded = elision.extension(false).value; + assert_eq!(encoded, elision_bytes([7, 10, 2, 6, 8])); + assert_eq!(OutputElision::decode(&encoded).unwrap(), elision); + for end in 0..encoded.len() { + assert!( + OutputElision::decode(&encoded[..end]).is_err(), + "prefix {end}" + ); + } + for fields in [ + [0, 0, 0, 0, 0], + [0, 10, 7, 6, 6], + [0, 10, 0, 11, 11], + [0, 10, 0, 6, 5], + [0, 10, 0, 6, 13], + ] { + assert!( + OutputElision::decode(&elision_bytes(fields)).is_err(), + "{fields:?}" + ); + } + // Twice these code points is past u64: checked without overflowing. + let widest = [u64::MAX; 5]; + assert!(OutputElision::decode(&elision_bytes(widest)).is_ok()); + let halves = [0, u64::MAX, 0, u64::MAX / 2 + 1, u64::MAX]; + assert!(OutputElision::decode(&elision_bytes(halves)).is_ok()); + } + fn every_truncation(value: &T) where T: Encode + Decode + PartialEq + std::fmt::Debug, @@ -1754,6 +1989,152 @@ mod tests { ); } + #[test] + fn report_exit_is_advertised_in_a_tag_older_clients_ignore() { + use crate::schema::process as p; + let v1 = p::SPAWN_LAUNCHER_FLAGS as u32; + let all = p::SPAWN_LAUNCHER_FLAGS_EXTENDED as u32; + assert_eq!( + all & !v1, + (p::SPAWN_REPORT_EXIT | p::SPAWN_KEEP_OUTPUT) as u32 + ); + + // Only the v1 flags: tag 11 alone, as before REPORT_EXIT. + let before = Limits { + launcher_flags: v1, + ..Limits::DEFAULT + }; + let extensions = before.to_extensions().unwrap(); + assert_eq!(limit_value(&extensions, p::LIMIT_LAUNCHER_FLAGS), Some(12)); + assert_eq!( + limit_value(&extensions, p::LIMIT_LAUNCHER_FLAGS_EXTENDED), + None + ); + assert_eq!(Limits::from_extensions(&extensions).unwrap(), before); + + // With REPORT_EXIT: tag 11 keeps within the maximum clients from before it + // validate, and the optional tag 19 carries every flag. + let extended = Limits { + launcher_flags: all, + ..Limits::DEFAULT + }; + let extensions = extended.to_extensions().unwrap(); + assert_eq!( + limit_value(&extensions, p::LIMIT_LAUNCHER_FLAGS), + Some(p::SPAWN_LAUNCHER_FLAGS) + ); + assert_eq!( + limit_value(&extensions, p::LIMIT_LAUNCHER_FLAGS_EXTENDED), + Some(p::SPAWN_LAUNCHER_FLAGS_EXTENDED) + ); + assert!(extensions.0.iter().all(|extension| !extension.required)); + assert_eq!(Limits::from_extensions(&extensions).unwrap(), extended); + // A client that does not know tag 19 reads the v1 flags: no REPORT_EXIT. + let without_19 = Extensions( + extensions + .0 + .iter() + .filter(|e| u64::from(e.tag) != p::LIMIT_LAUNCHER_FLAGS_EXTENDED) + .cloned() + .collect(), + ); + assert_eq!(Limits::from_extensions(&without_19).unwrap(), before); + + let with = |tag_11: Option, tag_19: Option| { + let mut extensions = Limits::DEFAULT.to_extensions().unwrap(); + for (tag, value) in [ + (p::LIMIT_LAUNCHER_FLAGS, tag_11), + (p::LIMIT_LAUNCHER_FLAGS_EXTENDED, tag_19), + ] { + if let Some(value) = value { + extensions.0.push(limit_u32(tag, value)); + } + } + Limits::from_extensions(&extensions) + }; + assert_eq!( + with(None, Some(p::SPAWN_REPORT_EXIT as u32)) + .unwrap() + .launcher_flags, + p::SPAWN_REPORT_EXIT as u32 + ); + // REPORT_EXIT has no place in tag 11, which older clients bound to 12. + assert!(with(Some(all), None).is_err()); + assert!(with(Some(all), Some(all)).is_err()); + // Tag 19 keeps the v1 flags tag 11 promised, and adds none of theirs. + assert!(with(Some(v1), Some(p::SPAWN_REPORT_EXIT as u32)).is_err()); + assert!(with(None, Some(all)).is_err()); + assert!(with(Some(v1), Some(all << 1)).is_err()); + // Tag 19 is a set of SPAWN flags: those this side does not know (a later server's) are + // ignored, and any u16 of them passes the family's limit bounds. + assert_eq!( + with(Some(v1), Some(all | 1 << 15)).unwrap().launcher_flags, + all + ); + let tag_19 = crate::schema::family_metadata(crate::family::PROCESS, 1) + .unwrap() + .limits + .iter() + .find(|limit| u64::from(limit.tag) == p::LIMIT_LAUNCHER_FLAGS_EXTENDED) + .unwrap(); + assert_eq!(tag_19.hard_max, u64::from(u16::MAX)); + } + + #[test] + fn exit_report_round_trips_and_names_its_process_cheaply() { + let report = ExitReport { + process_handle: 0x0102_0304_0506_0708, + exit: ExitRecord { + kind: ExitKind::Code, + reason: crate::schema::process::EXIT_REASON_UNKNOWN as u8, + code: 3, + exited_server_ns: 9, + detail: b"exited".to_vec(), + }, + extensions: Extensions::default(), + }; + every_truncation(&report); + let bytes = report.encode().unwrap(); + assert_eq!(ExitReport::handle_of(&bytes), Some(report.process_handle)); + assert_eq!(ExitReport::handle_of(&bytes[..7]), None); + let mut trailing = bytes.clone(); + trailing.push(0); + assert!(ExitReport::decode(&trailing).is_err()); + assert!( + ExitReport { + process_handle: 0, + ..report.clone() + } + .encode() + .is_err() + ); + let mut zero = bytes; + zero[..8].fill(0); + assert!(ExitReport::decode(&zero).is_err()); + } + + #[test] + fn spawn_accepts_report_exit_among_its_flags() { + use crate::schema::process as p; + let spawn = |flags: u64| Spawn { + operation_id: [1; 16], + flags: flags as u16, + environment_kind: EnvironmentKind::Session, + cwd: Cwd::ServerDefault, + argv: vec![b"true".to_vec()], + env: Vec::new(), + stdout_receive_credit: 65_536, + stderr_receive_credit: 65_536, + extensions: Extensions::default(), + }; + every_truncation(&spawn(p::SPAWN_REPORT_EXIT | p::SPAWN_STDIN_NULL)); + assert!( + spawn(p::SPAWN_LAUNCHER_FLAGS_EXTENDED << 1) + .encode() + .is_err() + ); + } + #[test] fn invalid_environment_and_transfer_policy_fail() { let mut spawn = Spawn { diff --git a/docs/design/processes.md b/docs/design/processes.md index c0ad004c..362583d0 100644 --- a/docs/design/processes.md +++ b/docs/design/processes.md @@ -91,7 +91,21 @@ observe output, while at most one attachment owns stdin. `CONTROL` provides typed signal, terminate, kill, and detach actions under nonzero operation IDs. `WAIT` returns the final portable exit record or -`TIMEOUT`. Closing stdin half-closes the child stream. An ordinary child belongs +`TIMEOUT`. A client that sets `SPAWN_REPORT_EXIT`, where the server offers it +(family-limit tag 19), is sent that record as an `EXIT` event once the exit is +final, so a command's exit costs no round trip of its own: SPAWN's Result, its +output and its exit all travel from one request. yas-client sets it whenever +offered, and `Process::wait` then takes no request and no pending-`WAIT` slot +unless the attachment that reports the exit went first (a stream reset or +dropped before its end, a detach), when it sends `WAIT`, as it does against +older servers. With it, `SPAWN_KEEP_OUTPUT` sends only a head and a tail of +each output stream, as the SPAWN asks: the server drops the middle as it reads +it, so a command writing far more than its client keeps runs at the speed of +its pipe, and the `EXIT` event says what was dropped, counted as a UTF-8 +decoder would count it (bytes, lines, code points, UTF-16 units). +yas-client's `Command::keep_output` asks for it where offered, and +`Process::elided` reads the counts. Closing stdin half-closes the +child stream. An ordinary child belongs to its spawning session and is terminated when it disappears; a detachable child remains discoverable until its retained final record expires. diff --git a/docs/design/yas.md b/docs/design/yas.md index 7f7cbd1e..4d0738bf 100644 --- a/docs/design/yas.md +++ b/docs/design/yas.md @@ -3242,7 +3242,7 @@ of defining separate data and ACK messages. | Class | Kinds | | ------- | -------------------------------------------- | | Request | WATCH, UNWATCH, SPAWN, ATTACH, CONTROL, WAIT | -| Event | STATE, STATE_ACK | +| Event | STATE, STATE_ACK, EXIT | SPAWN executes exact argv and environment bytes without an implicit shell. It accepts explicit cwd, inherited terminal cwd, FS root/path, session environment, @@ -3268,6 +3268,52 @@ TERMINATE's escalation SIGKILLs members left after the kill grace (terminates the job on Windows) and stops waiting for their streams. The exit's detail says `residual process group left running` when members held the streams. +`REPORT_EXIT` (16) sends the spawning session the exit as it becomes final, +without a WAIT: one EXIT Event (`0x0002`, sensitive), `[process_handle: u64, +exit: bytes_u32 containing ExitRecord, Extensions]`, the record a WAIT would +return at that moment. The streams go on with what the process wrote before its +exit, at the pace of their credit, so the event can arrive before their last +bytes and their CLOSE. The spawning session's attachment sends it: when that +attachment goes before the exit (a Transfer RESET on any of its streams, stdin +included, from either side, or a DETACH), no EXIT is sent and the client WAITs +for the exit instead, as it would without the flag. yas-client does so on its +own. The report takes none of the session's pending WAITs, and +nothing changes for other sessions: they, and sessions that ATTACH, still WAIT. +A SPAWN retried under its operation ID shares the original attachment, whose +exit is reported once. + +`KEEP_OUTPUT` (32), with `REPORT_EXIT` and SPAWN extension tag 4 +`head_bytes: u64, tail_bytes: u64` (the tail at most +`MAX_KEEP_OUTPUT_TAIL_BYTES`, 1 MiB), sends only the head and the tail of each +output stream: what comes between is dropped as the server reads it, never +held for the client's credit, so the pipe drains at the writer's speed. The +head is at least `head_bytes` and ends between characters, as a WHATWG UTF-8 +decoder with replacement reads the whole stream: where it is between +characters, or where the next byte cannot continue the character it is in. The +tail is at most `tail_bytes`, starts between characters, and once anything was +dropped never starts with a continuation byte. Decoding the head and the tail, +apart or one after the other, gives exactly the characters they have within the +whole. The Transfer carries the head then the tail, contiguous; the EXIT event +says what was dropped of each stream in extension tag 1 (stdout) and 2 +(stderr), present only when something was: `OutputElision` +`[offset: u64, bytes: u64, lines: u64, code_points: u64, utf16_units: u64]`, +where `offset` is the head's length and the counts are the dropped bytes as +that decoder reads them within the whole stream. A client can then say how +much it did not get in the units it counts. The tail goes out when the stream +ends, or before the exit is reported when the stream outlives it (a residue +past its grace, a TERMINATE, a lost owner, a forced cleanup): what is written +after that goes to nobody. A tail cut short by an aborted stream is not counted +in the elision; that stream does not end cleanly. + +Servers advertise the opt-in flags they honour in two optional family limits: +tag 11 `LAUNCHER_FLAGS` carries those of v1 (`LEAVE_RESIDUE`, `STDIN_NULL`), +at most 12, which is all clients from before `REPORT_EXIT` accept; tag 19 +`LAUNCHER_FLAGS_EXTENDED` carries every flag the server honours, when that is +more. Tag 19 names each flag tag 11 does and adds none of v1's. It is a set of +SPAWN flags, any u16: a client ignores the flags it does not know, so later +flags need no new tag. A client sets `REPORT_EXIT` only when tag 19 offers it, +and otherwise WAITs, so either side may be older. + Catalog records contain argv0, native PID for diagnostics, lifecycle, owner session, detachable flag, stream offsets, exit record, and retention deadline. An ordinary process is owned by its spawning session and terminated when that @@ -3282,7 +3328,8 @@ pipe is read at most a window ahead); an attachment of another session that fall a window behind has its Transfers reset `RESOURCE_EXHAUSTED`, and nothing else changes. CONTROL provides portable signal, terminate, kill, and detach actions with operation IDs. Closing the stdin Transfer half-closes -stdin. WAIT returns the final exit record or TIMEOUT. +stdin. WAIT returns the final exit record or TIMEOUT; a process spawned with +`REPORT_EXIT` needs none, its EXIT Event carries the same record. The canonical v1 payloads are generated from `protocol/yas/families/process.toml`. SPAWN carries `[operation_id, flags, diff --git a/js/core/src/yas/generated.ts b/js/core/src/yas/generated.ts index 1fe93767..6d36d946 100644 --- a/js/core/src/yas/generated.ts +++ b/js/core/src/yas/generated.ts @@ -1649,12 +1649,16 @@ export const YAS_PROCESS_CONTROL = 4 as const; export const YAS_PROCESS_WAIT = 5 as const; export const YAS_PROCESS_STATE = 0 as const; export const YAS_PROCESS_STATE_ACK = 1 as const; +export const YAS_PROCESS_EXIT = 2 as const; export const YAS_PROCESS_SPAWN_MERGE_STDERR = 1 as const; export const YAS_PROCESS_SPAWN_DETACHABLE = 2 as const; export const YAS_PROCESS_SPAWN_LEAVE_RESIDUE = 4 as const; export const YAS_PROCESS_SPAWN_STDIN_NULL = 8 as const; export const YAS_PROCESS_SPAWN_FLAGS = 3 as const; export const YAS_PROCESS_SPAWN_LAUNCHER_FLAGS = 12 as const; +export const YAS_PROCESS_SPAWN_REPORT_EXIT = 16 as const; +export const YAS_PROCESS_SPAWN_KEEP_OUTPUT = 32 as const; +export const YAS_PROCESS_SPAWN_LAUNCHER_FLAGS_EXTENDED = 60 as const; export const YAS_PROCESS_ENV_EMPTY = 0 as const; export const YAS_PROCESS_ENV_SESSION = 1 as const; export const YAS_PROCESS_CWD_SERVER_DEFAULT = 0 as const; @@ -1701,6 +1705,10 @@ export const YAS_PROCESS_STREAM_STDERR_CONTENT_KIND = 2 as const; export const YAS_PROCESS_SPAWN_SURFACE_APP_EXTENSION = 1 as const; export const YAS_PROCESS_SPAWN_RESOURCE_TAG_EXTENSION = 2 as const; export const YAS_PROCESS_SPAWN_RESIDUE_GRACE_EXTENSION = 3 as const; +export const YAS_PROCESS_SPAWN_KEEP_OUTPUT_EXTENSION = 4 as const; +export const YAS_PROCESS_MAX_KEEP_OUTPUT_TAIL_BYTES = 1048576 as const; +export const YAS_PROCESS_EXIT_STDOUT_ELIDED_EXTENSION = 1 as const; +export const YAS_PROCESS_EXIT_STDERR_ELIDED_EXTENSION = 2 as const; export const YAS_PROCESS_MAX_ARGC = 1024 as const; export const YAS_PROCESS_MAX_ARG_BYTES = 1048576 as const; export const YAS_PROCESS_MAX_ARG_LEN = 65536 as const; @@ -1743,6 +1751,7 @@ export const YAS_PROCESS_LIMIT_MAX_DETACHED_RETENTION_NS = 9 as const; export const YAS_PROCESS_MAX_MUTATION_REPLAYS = 65536 as const; export const YAS_PROCESS_LIMIT_MAX_MUTATION_REPLAYS = 10 as const; export const YAS_PROCESS_LIMIT_LAUNCHER_FLAGS = 11 as const; +export const YAS_PROCESS_LIMIT_LAUNCHER_FLAGS_EXTENDED = 19 as const; export const YAS_FAMILY_NET = 65 as const; export const YAS_NET_VERSION = 1 as const; export const YAS_NET_OPEN = 0 as const; @@ -2251,6 +2260,7 @@ export const YAS_FAMILY_LIMIT_POLICIES: Readonly "64/2/5": [1, 0, 0], "64/0/0": [1, 0, 0], "64/0/1": [0, 0, 0], + "64/0/2": [1, 0, 0], "65/1/0": [1, 0, 0], "65/2/0": [1, 0, 0], "65/1/1": [1, 0, 0], @@ -2876,6 +2887,7 @@ export const YAS_OPERATION_DIRECTION_MASKS: Readonly> = { "64/1/5": 1, "64/0/0": 2, "64/0/1": 1, + "64/0/2": 2, "65/1/0": 1, "65/1/1": 1, "65/0/0": 3, @@ -11733,6 +11745,14 @@ export const YAS_SCHEMA = { "required": false, "hard_min": 1, "hard_max": 16384 + }, + { + "name": "LAUNCHER_FLAGS_EXTENDED", + "tag": 19, + "type": "u32", + "required": false, + "hard_min": 0, + "hard_max": 65535 } ], "requests": [ @@ -11809,6 +11829,15 @@ export const YAS_SCHEMA = { "compression": "allowed", "datagram": "forbidden", "layout": "StateAck" + }, + { + "name": "EXIT", + "kind": 2, + "direction": "server_to_client", + "sensitive": "required", + "compression": "allowed", + "datagram": "forbidden", + "layout": "ExitReport" } ], "types": [ @@ -11824,6 +11853,18 @@ export const YAS_SCHEMA = { "name": "remove_record", "layout": "process_handle:u64" }, + { + "name": "exit_report", + "layout": "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions; extension tag 1 stdout OutputElision, tag 2 stderr OutputElision, each present iff KEEP_OUTPUT dropped bytes of that stream" + }, + { + "name": "keep_output_extension", + "layout": "SPAWN extension tag 4 exact value head_bytes:u64,tail_bytes:u64; only with SPAWN_KEEP_OUTPUT and SPAWN_REPORT_EXIT; tail_bytes at most MAX_KEEP_OUTPUT_TAIL_BYTES" + }, + { + "name": "output_elision", + "layout": "offset:u64,bytes:u64,lines:u64,code_points:u64,utf16_units:u64; offset is the stream offset where the dropped bytes were (the head's length); lines, code points and UTF-16 units count them as a WHATWG UTF-8 decoder with replacement reads them within the whole stream" + }, { "name": "exit_record", "layout": "kind:u8,reason:u8,reserved:u16=0,code:i32,exited_server_ns:u64,detail:bytes_u32" @@ -11866,6 +11907,18 @@ export const YAS_SCHEMA = { "name": "SPAWN_LAUNCHER_FLAGS", "value": 12 }, + { + "name": "SPAWN_REPORT_EXIT", + "value": 16 + }, + { + "name": "SPAWN_KEEP_OUTPUT", + "value": 32 + }, + { + "name": "SPAWN_LAUNCHER_FLAGS_EXTENDED", + "value": 60 + }, { "name": "ENV_EMPTY", "value": 0 @@ -12050,6 +12103,22 @@ export const YAS_SCHEMA = { "name": "SPAWN_RESIDUE_GRACE_EXTENSION", "value": 3 }, + { + "name": "SPAWN_KEEP_OUTPUT_EXTENSION", + "value": 4 + }, + { + "name": "MAX_KEEP_OUTPUT_TAIL_BYTES", + "value": 1048576 + }, + { + "name": "EXIT_STDOUT_ELIDED_EXTENSION", + "value": 1 + }, + { + "name": "EXIT_STDERR_ELIDED_EXTENSION", + "value": 2 + }, { "name": "MAX_ARGC", "value": 1024 @@ -12217,6 +12286,10 @@ export const YAS_SCHEMA = { { "name": "LIMIT_LAUNCHER_FLAGS", "value": 11 + }, + { + "name": "LIMIT_LAUNCHER_FLAGS_EXTENDED", + "value": 19 } ] }, @@ -15393,6 +15466,10 @@ export const YAS_GOLDEN_VECTORS = { "name": "yas.process.event.state_ack.header", "hex": "4000010000" }, + { + "name": "yas.process.event.exit.header", + "hex": "4000020008" + }, { "name": "yas.net.request.open.header", "hex": "410000000901000000" diff --git a/protocol/yas/families/process.toml b/protocol/yas/families/process.toml index 76f0f40d..055040d7 100644 --- a/protocol/yas/families/process.toml +++ b/protocol/yas/families/process.toml @@ -20,6 +20,8 @@ limits = [ { name = "MAX_ENVC_EXTENDED", tag = 16, type = "u32", required = false, hard_min = 1, hard_max = 16384 }, { name = "MAX_PENDING_WAITS", tag = 17, type = "u32", required = false, hard_min = 1, hard_max = 65536 }, { name = "MAX_PENDING_OPERATIONS", tag = 18, type = "u32", required = false, hard_min = 1, hard_max = 16384 }, + # A set of SPAWN flags (u16): bits a client does not know are ignored, so later flags need no new tag. + { name = "LAUNCHER_FLAGS_EXTENDED", tag = 19, type = "u32", required = false, hard_min = 0, hard_max = 65535 }, ] [[constant]] @@ -40,6 +42,22 @@ value = 3 [[constant]] name = "SPAWN_LAUNCHER_FLAGS" value = 12 +# The spawning session is sent an EXIT event once the exit is final, so it +# learns it without a WAIT. Advertised by LAUNCHER_FLAGS_EXTENDED (tag 19), +# which old clients ignore: LAUNCHER_FLAGS (tag 11) keeps its v1 maximum. +[[constant]] +name = "SPAWN_REPORT_EXIT" +value = 16 +# Only the head and the tail of each output stream are sent, as the +# KEEP_OUTPUT extension sizes them: the middle is dropped as it is read, cut +# between UTF-8 characters, and counted in the EXIT event. Needs REPORT_EXIT. +# Advertised by LAUNCHER_FLAGS_EXTENDED (tag 19). +[[constant]] +name = "SPAWN_KEEP_OUTPUT" +value = 32 +[[constant]] +name = "SPAWN_LAUNCHER_FLAGS_EXTENDED" +value = 60 [[constant]] name = "ENV_EMPTY" @@ -190,6 +208,18 @@ value = 2 [[constant]] name = "SPAWN_RESIDUE_GRACE_EXTENSION" value = 3 +[[constant]] +name = "SPAWN_KEEP_OUTPUT_EXTENSION" +value = 4 +[[constant]] +name = "MAX_KEEP_OUTPUT_TAIL_BYTES" +value = 1048576 +[[constant]] +name = "EXIT_STDOUT_ELIDED_EXTENSION" +value = 1 +[[constant]] +name = "EXIT_STDERR_ELIDED_EXTENSION" +value = 2 [[constant]] name = "MAX_ARGC" @@ -325,6 +355,9 @@ value = 10 [[constant]] name = "LIMIT_LAUNCHER_FLAGS" value = 11 +[[constant]] +name = "LIMIT_LAUNCHER_FLAGS_EXTENDED" +value = 19 [[request]] name = "WATCH" @@ -398,6 +431,15 @@ compression = "allowed" datagram = "forbidden" layout = "StateAck" +[[event]] +name = "EXIT" +kind = 0x0002 +direction = "server_to_client" +sensitive = "required" +compression = "allowed" +datagram = "forbidden" +layout = "ExitReport" + [[type]] name = "cwd" layout = "kind:u8,reserved:[u8;3]=0; SERVER_DEFAULT empty, PATH path:bytes_u32, TERMINAL terminal_handle:u64, FS root_handle:u64,component_count:u16,repeated component:bytes_u16" @@ -410,6 +452,18 @@ layout = "process_handle:u64,lifecycle:u8,stream_state:u8,flags:u16,native_pid:u name = "remove_record" layout = "process_handle:u64" +[[type]] +name = "exit_report" +layout = "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions; extension tag 1 stdout OutputElision, tag 2 stderr OutputElision, each present iff KEEP_OUTPUT dropped bytes of that stream" + +[[type]] +name = "keep_output_extension" +layout = "SPAWN extension tag 4 exact value head_bytes:u64,tail_bytes:u64; only with SPAWN_KEEP_OUTPUT and SPAWN_REPORT_EXIT; tail_bytes at most MAX_KEEP_OUTPUT_TAIL_BYTES" + +[[type]] +name = "output_elision" +layout = "offset:u64,bytes:u64,lines:u64,code_points:u64,utf16_units:u64; offset is the stream offset where the dropped bytes were (the head's length); lines, code points and UTF-16 units count them as a WHATWG UTF-8 decoder with replacement reads them within the whole stream" + [[type]] name = "exit_record" layout = "kind:u8,reason:u8,reserved:u16=0,code:i32,exited_server_ns:u64,detail:bytes_u32" diff --git a/protocol/yas/inspection.json b/protocol/yas/inspection.json index f4abc74d..dd2d5565 100644 --- a/protocol/yas/inspection.json +++ b/protocol/yas/inspection.json @@ -5599,6 +5599,23 @@ "name": "STATE_ACK", "sensitive": "allowed" }, + { + "class": "event", + "class_id": 0, + "compression": "allowed", + "correlated": false, + "datagram": "forbidden", + "direction": "server_to_client", + "family": 64, + "family_name": "yas.process", + "family_version": 1, + "header_bytes": 5, + "key": "64/0/2", + "kind": 2, + "layout": "ExitReport", + "name": "EXIT", + "sensitive": "required" + }, { "class": "request", "class_id": 1, diff --git a/protocol/yas/schema.json b/protocol/yas/schema.json index 1cea9b7b..57bc83a7 100644 --- a/protocol/yas/schema.json +++ b/protocol/yas/schema.json @@ -8817,6 +8817,14 @@ "required": false, "hard_min": 1, "hard_max": 16384 + }, + { + "name": "LAUNCHER_FLAGS_EXTENDED", + "tag": 19, + "type": "u32", + "required": false, + "hard_min": 0, + "hard_max": 65535 } ], "requests": [ @@ -8893,6 +8901,15 @@ "compression": "allowed", "datagram": "forbidden", "layout": "StateAck" + }, + { + "name": "EXIT", + "kind": 2, + "direction": "server_to_client", + "sensitive": "required", + "compression": "allowed", + "datagram": "forbidden", + "layout": "ExitReport" } ], "types": [ @@ -8908,6 +8925,18 @@ "name": "remove_record", "layout": "process_handle:u64" }, + { + "name": "exit_report", + "layout": "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions; extension tag 1 stdout OutputElision, tag 2 stderr OutputElision, each present iff KEEP_OUTPUT dropped bytes of that stream" + }, + { + "name": "keep_output_extension", + "layout": "SPAWN extension tag 4 exact value head_bytes:u64,tail_bytes:u64; only with SPAWN_KEEP_OUTPUT and SPAWN_REPORT_EXIT; tail_bytes at most MAX_KEEP_OUTPUT_TAIL_BYTES" + }, + { + "name": "output_elision", + "layout": "offset:u64,bytes:u64,lines:u64,code_points:u64,utf16_units:u64; offset is the stream offset where the dropped bytes were (the head's length); lines, code points and UTF-16 units count them as a WHATWG UTF-8 decoder with replacement reads them within the whole stream" + }, { "name": "exit_record", "layout": "kind:u8,reason:u8,reserved:u16=0,code:i32,exited_server_ns:u64,detail:bytes_u32" @@ -8950,6 +8979,18 @@ "name": "SPAWN_LAUNCHER_FLAGS", "value": 12 }, + { + "name": "SPAWN_REPORT_EXIT", + "value": 16 + }, + { + "name": "SPAWN_KEEP_OUTPUT", + "value": 32 + }, + { + "name": "SPAWN_LAUNCHER_FLAGS_EXTENDED", + "value": 60 + }, { "name": "ENV_EMPTY", "value": 0 @@ -9134,6 +9175,22 @@ "name": "SPAWN_RESIDUE_GRACE_EXTENSION", "value": 3 }, + { + "name": "SPAWN_KEEP_OUTPUT_EXTENSION", + "value": 4 + }, + { + "name": "MAX_KEEP_OUTPUT_TAIL_BYTES", + "value": 1048576 + }, + { + "name": "EXIT_STDOUT_ELIDED_EXTENSION", + "value": 1 + }, + { + "name": "EXIT_STDERR_ELIDED_EXTENSION", + "value": 2 + }, { "name": "MAX_ARGC", "value": 1024 @@ -9301,6 +9358,10 @@ { "name": "LIMIT_LAUNCHER_FLAGS", "value": 11 + }, + { + "name": "LIMIT_LAUNCHER_FLAGS_EXTENDED", + "value": 19 } ] }, diff --git a/protocol/yas/vectors.json b/protocol/yas/vectors.json index 27f1e0fb..a6ebb842 100644 --- a/protocol/yas/vectors.json +++ b/protocol/yas/vectors.json @@ -753,6 +753,10 @@ "name": "yas.process.event.state_ack.header", "hex": "4000010000" }, + { + "name": "yas.process.event.exit.header", + "hex": "4000020008" + }, { "name": "yas.net.request.open.header", "hex": "410000000901000000" diff --git a/protocol/yas/wire.md b/protocol/yas/wire.md index 1d2f075f..8dadf360 100644 --- a/protocol/yas/wire.md +++ b/protocol/yas/wire.md @@ -893,6 +893,7 @@ Every Request kind has a correlated Result with the same family and kind. | ---: | --- | --- | --- | --- | --- | --- | | `0x0000` | `STATE` | `server_to_client` | `required` | `allowed` | `forbidden` | StateEvent | | `0x0001` | `STATE_ACK` | `client_to_server` | `allowed` | `allowed` | `forbidden` | StateAck | +| `0x0002` | `EXIT` | `server_to_client` | `required` | `allowed` | `forbidden` | ExitReport | ### Limits @@ -916,6 +917,7 @@ Every Request kind has a correlated Result with the same family and kind. | 16 | `MAX_ENVC_EXTENDED` | 4 | false | 1 | 16384 | | 17 | `MAX_PENDING_WAITS` | 4 | false | 1 | 65536 | | 18 | `MAX_PENDING_OPERATIONS` | 4 | false | 1 | 16384 | +| 19 | `LAUNCHER_FLAGS_EXTENDED` | 4 | false | 0 | 65535 | ### Shared types @@ -924,6 +926,9 @@ Every Request kind has a correlated Result with the same family and kind. | `cwd` | kind:u8,reserved:[u8;3]=0; SERVER_DEFAULT empty, PATH path:bytes_u32, TERMINAL terminal_handle:u64, FS root_handle:u64,component_count:u16,repeated component:bytes_u16 | | `process_record` | process_handle:u64,lifecycle:u8,stream_state:u8,flags:u16,native_pid:u64,owner_session:[u8;16],argv0:bytes_u32,stdin_received:u64,stdout_produced:u64,stderr_produced:u64,retention_deadline_server_ns:u64,exit_present:u8,reserved:[u8;7]=0,optional exit:bytes_u32 containing ExitRecord,Extensions | | `remove_record` | process_handle:u64 | +| `exit_report` | process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions; extension tag 1 stdout OutputElision, tag 2 stderr OutputElision, each present iff KEEP_OUTPUT dropped bytes of that stream | +| `keep_output_extension` | SPAWN extension tag 4 exact value head_bytes:u64,tail_bytes:u64; only with SPAWN_KEEP_OUTPUT and SPAWN_REPORT_EXIT; tail_bytes at most MAX_KEEP_OUTPUT_TAIL_BYTES | +| `output_elision` | offset:u64,bytes:u64,lines:u64,code_points:u64,utf16_units:u64; offset is the stream offset where the dropped bytes were (the head's length); lines, code points and UTF-16 units count them as a WHATWG UTF-8 decoder with replacement reads them within the whole stream | | `exit_record` | kind:u8,reason:u8,reserved:u16=0,code:i32,exited_server_ns:u64,detail:bytes_u32 | | `stream_bundle` | process_handle:u64,flags:u16,reserved:u16=0,stdout_lifetime_offset:u64,stderr_lifetime_offset:u64,optional stdin/stdout/stderr descriptor:bytes_u32 containing sensitive BYTE TransferDescriptor,Extensions | | `state_entity_body` | ADD/REPLACE complete ProcessRecord; REMOVE process_handle:u64 |