diff --git a/crates/cli/tests/client_host.rs b/crates/cli/tests/client_host.rs index dd1895ef..0fbf0f5d 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,218 @@ 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(); + } +} + +/// 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 +1234,7 @@ 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) 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..bca92bee 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}, @@ -485,6 +486,88 @@ 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, + status: std::sync::OnceLock, + 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) => ExitStatus::from_wire(ExitReport::decode(&frame.payload)?.exit), + 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 => 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).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 +619,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 } @@ -573,6 +666,10 @@ impl Process { /// 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 @@ -765,13 +862,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 +889,15 @@ 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; + } // 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 +915,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 +954,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 +1000,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 +1017,7 @@ impl Client { stdout_offset: bundle.stdout_lifetime_offset, stderr_offset: bundle.stderr_lifetime_offset, merged_stderr: bundle.merged_stderr, + reported, }) } @@ -1035,7 +1151,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/yas.rs b/crates/server/src/yas.rs index 178306f9..c88984cd 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,29 @@ 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, +) { + let report = yas_process_wire::ExitReport { + process_handle, + exit: exit.into_record(monotonic_ns()), + extensions: Extensions::default(), + }; + 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 +31047,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 +31150,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 +51852,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 +52217,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..68d60e88 100644 --- a/crates/server/src/yas_process.rs +++ b/crates/server/src/yas_process.rs @@ -293,9 +293,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,7 +382,8 @@ impl Session { resolved_cwd: Option>, ) -> Result { let cwd = resolve_cwd(&request.cwd, resolved_cwd)?; - let flags = u8::try_from(request.flags) + // REPORT_EXIT asks the YAS connection for an EXIT event; the process is the same. + let flags = u8::try_from(request.flags & !(schema::process::SPAWN_REPORT_EXIT as u16)) .map_err(|_| Error::Invalid("Process SPAWN flags do not fit v1".to_owned()))?; let process_id = self.allocate_process_id()?; let (route, events) = self.install_route(process_id, false)?; diff --git a/crates/yas/src/generated.rs b/crates/yas/src/generated.rs index a3b9caf4..3b7ecee5 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,8 @@ 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_LAUNCHER_FLAGS_EXTENDED: u64 = 28; pub const ENV_EMPTY: u64 = 0; pub const ENV_SESSION: u64 = 1; pub const CWD_SERVER_DEFAULT: u64 = 0; @@ -3886,6 +3889,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 +3899,13 @@ 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" }, 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 +3930,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 +3939,8 @@ 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_LAUNCHER_FLAGS_EXTENDED", value: 28 }, super::ConstantMetadata { name: "ENV_EMPTY", value: 0 }, super::ConstantMetadata { name: "ENV_SESSION", value: 1 }, super::ConstantMetadata { name: "CWD_SERVER_DEFAULT", value: 0 }, @@ -4020,6 +4029,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 +5043,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..e0dd3587 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")); } @@ -914,6 +914,47 @@ 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()?)) + } +} + /// 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 +978,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 +1005,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 +1139,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 +1177,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 +1197,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 +1235,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 +1265,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")); } @@ -1754,6 +1820,149 @@ 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 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..133f4c4d 100644 --- a/docs/design/processes.md +++ b/docs/design/processes.md @@ -91,7 +91,15 @@ 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. 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..a1f84ede 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,29 @@ 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. + +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 +3305,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..181bd360 100644 --- a/js/core/src/yas/generated.ts +++ b/js/core/src/yas/generated.ts @@ -1649,12 +1649,15 @@ 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_LAUNCHER_FLAGS_EXTENDED = 28 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; @@ -1743,6 +1746,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 +2255,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 +2882,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 +11740,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 +11824,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 +11848,10 @@ export const YAS_SCHEMA = { "name": "remove_record", "layout": "process_handle:u64" }, + { + "name": "exit_report", + "layout": "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions" + }, { "name": "exit_record", "layout": "kind:u8,reason:u8,reserved:u16=0,code:i32,exited_server_ns:u64,detail:bytes_u32" @@ -11866,6 +11894,14 @@ export const YAS_SCHEMA = { "name": "SPAWN_LAUNCHER_FLAGS", "value": 12 }, + { + "name": "SPAWN_REPORT_EXIT", + "value": 16 + }, + { + "name": "SPAWN_LAUNCHER_FLAGS_EXTENDED", + "value": 28 + }, { "name": "ENV_EMPTY", "value": 0 @@ -12217,6 +12253,10 @@ export const YAS_SCHEMA = { { "name": "LIMIT_LAUNCHER_FLAGS", "value": 11 + }, + { + "name": "LIMIT_LAUNCHER_FLAGS_EXTENDED", + "value": 19 } ] }, @@ -15393,6 +15433,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..11f19714 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,15 @@ 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 +[[constant]] +name = "SPAWN_LAUNCHER_FLAGS_EXTENDED" +value = 28 [[constant]] name = "ENV_EMPTY" @@ -325,6 +336,9 @@ value = 10 [[constant]] name = "LIMIT_LAUNCHER_FLAGS" value = 11 +[[constant]] +name = "LIMIT_LAUNCHER_FLAGS_EXTENDED" +value = 19 [[request]] name = "WATCH" @@ -398,6 +412,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 +433,10 @@ 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" + [[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..f30f6530 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,10 @@ "name": "remove_record", "layout": "process_handle:u64" }, + { + "name": "exit_report", + "layout": "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions" + }, { "name": "exit_record", "layout": "kind:u8,reason:u8,reserved:u16=0,code:i32,exited_server_ns:u64,detail:bytes_u32" @@ -8950,6 +8971,14 @@ "name": "SPAWN_LAUNCHER_FLAGS", "value": 12 }, + { + "name": "SPAWN_REPORT_EXIT", + "value": 16 + }, + { + "name": "SPAWN_LAUNCHER_FLAGS_EXTENDED", + "value": 28 + }, { "name": "ENV_EMPTY", "value": 0 @@ -9301,6 +9330,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..60191502 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,7 @@ 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 | | `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 |