diff --git a/Cargo.lock b/Cargo.lock index 597a01ce..4ab1fab9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8332,6 +8332,7 @@ dependencies = [ "rustls", "rustls-native-certs", "serde_json", + "socket2", "tempfile", "tokio", "tokio-rustls", diff --git a/crates/cli/tests/client_host.rs b/crates/cli/tests/client_host.rs index 0fbf0f5d..e00529da 100644 --- a/crates/cli/tests/client_host.rs +++ b/crates/cli/tests/client_host.rs @@ -840,6 +840,135 @@ async fn a_spawned_process_reports_its_exit_without_a_wait() { } } +/// 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 @@ -882,6 +1011,70 @@ async fn a_server_from_before_report_exit_is_waited_for() { exits_come_with_all_their_output(&server, &client).await; } +/// Every stdout/stderr stream holds its window of the session's receive budget +/// while it is open, writing or not. A client that asks for a wider budget +/// gets windows that still all fit: at 256 processes a session in 256 MiB, +/// 384 KiB each, where 16 MiB gives 24 KiB. With 255 quiet processes holding +/// theirs (510 streams, 191 MiB: more than the default budget), one more still +/// gets credit, and all of its output at once. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn a_wider_receive_budget_fits_every_process_window() { + const BUDGET: u64 = 256 << 20; + const QUIET: usize = 255; + let server = within( + "hosted server start", + HostedServer::start( + options() + .args(["--process-max-per-session", "256"]) + .args(["--process-max", "1024"]) + .hello(HelloOptions::named("yas-test").receive_budget(BUDGET)), + ), + ) + .await + .expect("hosted server starts"); + let narrow = server + .connect_with(&HelloOptions::named("yas-test")) + .await + .unwrap(); + assert_eq!(narrow.default_process_window(), 24 * 1024); + drop(narrow); + let client = server.connect().await.unwrap(); + assert_eq!(client.receive_budget(), BUDGET); + assert_eq!(client.default_process_window(), 384 * 1024); + let mut quiet = Vec::with_capacity(QUIET); + for _ in 0..QUIET { + quiet.push(client.spawn(Command::new("sleep").arg("60")).await.unwrap()); + } + // Three windows of output: it comes whole only if credit keeps coming. + const BYTES: usize = 3 * 384 * 1024; + let started = Instant::now(); + let output = within( + "a command beside the quiet ones", + client + .spawn(Command::new("sh").args(["-c", &format!("head -c {BYTES} /dev/zero; echo end")])) + .await + .unwrap() + .output(), + ) + .await + .unwrap(); + let took = started.elapsed(); + assert!(output.status.success(), "{}", output.status); + assert_eq!(output.stdout.len(), BYTES + 4); + assert!(output.stdout.ends_with(b"end\n")); + assert!(took < Duration::from_secs(5), "{took:?}"); + for process in &quiet { + process.kill().await.unwrap(); + } + for process in &quiet { + within("a quiet process's exit", process.wait()) + .await + .unwrap(); + } + drop(quiet); + still_runs_commands(&client).await; +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn a_commands_background_children_die_with_it() { let server = start().await; @@ -1234,7 +1427,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 | schema::SPAWN_REPORT_EXIT) 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 4f9be8f8..0ec61fc1 100644 --- a/crates/client/src/client.rs +++ b/crates/client/src/client.rs @@ -249,6 +249,8 @@ struct Inner { tasks: Vec>, /// The session clock that input events' `client_monotonic_ns` count on. started: std::time::Instant, + /// What this client offered to receive at once ([`HelloOptions::receive_budget`]). + receive_budget: u64, } impl Drop for Inner { @@ -317,6 +319,7 @@ impl Client { /// runtime: it spawns the reader and writer tasks. pub fn from_native(native: NativeClient) -> Self { let hello = native.hello().clone(); + let receive_budget = native.receive_budget(); let (reader, sender) = native.into_framed(); let (closed, _) = watch::channel(None); let shared = Arc::new(Shared { @@ -334,6 +337,7 @@ impl Client { next_request_id: AtomicU32::new(3), tasks: vec![writer, reader], started: std::time::Instant::now(), + receive_budget, }), } } @@ -343,6 +347,13 @@ impl Client { self.inner.shared.hello.read().unwrap().clone() } + /// How many bytes the server may have on their way to this session at + /// once, as this client offered in HELLO + /// ([`HelloOptions::receive_budget`]). + pub fn receive_budget(&self) -> u64 { + self.inner.receive_budget + } + /// The server instance name (`default`, or the `--name` it runs under). pub fn server_name(&self) -> String { self.inner.shared.hello.read().unwrap().server_name.clone() diff --git a/crates/client/src/native.rs b/crates/client/src/native.rs index 63d69364..2c5b468f 100644 --- a/crates/client/src/native.rs +++ b/crates/client/src/native.rs @@ -37,8 +37,11 @@ use crate::{ConnectOptions, HelloOptions}; const HELLO_REQUEST_ID: u32 = 1; const WATCH_CREDIT: u64 = yas_wire::schema::transport::RECOMMENDED_BUFFERED; const REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); -const MAX_PENDING_FRAMES: usize = 1_024; -const MAX_PENDING_BYTES: usize = yas_wire::schema::transport::RECOMMENDED_BUFFERED as usize; +/// Frames parked while a caller waits for another, at least, and per this many bytes of the +/// declared receive budget (1,024 in the default 16 MiB): a peer within its credit parks no +/// more bytes than the budget. +const MIN_PENDING_FRAMES: usize = 1_024; +const PENDING_BYTES_PER_FRAME: u64 = 16 * 1024; pub const MAX_COLLECTED_TRANSFER_BYTES: u64 = 256 * 1024 * 1024; type Reader = Box; @@ -113,7 +116,10 @@ impl NativeClient { let hello_request = ClientHello { min_minor: 1, max_minor: 1, - receive: ReceiveLimits::recommended(max_datagram), + receive: ReceiveLimits { + max_buffered: options.receive_budget, + ..ReceiveLimits::recommended(max_datagram) + }, client_instance: rand::random(), client_name: options.client_name.clone(), client_release: options.client_release.clone(), @@ -214,6 +220,12 @@ impl NativeClient { &self.hello } + /// The receive budget this client offered in HELLO + /// ([`crate::HelloOptions::receive_budget`]). + pub fn receive_budget(&self) -> u64 { + self.local_receive.max_buffered + } + pub fn supports_datagrams(&self) -> bool { self.datagram.is_some() && self.local_receive.max_datagram > 0 @@ -973,7 +985,13 @@ impl NativeClient { .pending_bytes .checked_add(frame.payload.len()) .ok_or_else(|| Error::protocol("pending YAS frame accounting overflow"))?; - if self.pending.len() >= MAX_PENDING_FRAMES || next_bytes > MAX_PENDING_BYTES { + let budget = self.local_receive.max_buffered; + let max_frames = usize::try_from(budget / PENDING_BYTES_PER_FRAME) + .unwrap_or(usize::MAX) + .max(MIN_PENDING_FRAMES); + if self.pending.len() >= max_frames + || u64::try_from(next_bytes).unwrap_or(u64::MAX) > budget + { return Err(Error::protocol( "native YAS peer exceeded the bounded pending-frame queue", )); @@ -1430,6 +1448,124 @@ mod tests { } } + /// Answer a client's HELLO on `server_stream` with [`test_server_hello`]: the codec + /// for what follows, and the client's HELLO. + async fn answer_hello( + server_stream: &mut tokio::io::DuplexStream, + ) -> (FrameCodec, ClientHello) { + let mut preface = [0; yas_wire::PREFACE.len()]; + server_stream.read_exact(&mut preface).await.unwrap(); + let pre_hello = FrameCodec::pre_hello(); + let hello_frame = read_frame(server_stream, &pre_hello).await.unwrap(); + let client_hello = ClientHello::decode(&hello_frame.payload).unwrap(); + let result_frame = Frame { + header: FrameHeader::result( + family::CORE, + yas_wire::core::request_kind::HELLO, + HELLO_REQUEST_ID, + ), + payload: ResultPrefix { + status: Status::Ok, + detail: Extensions::default(), + body: test_server_hello().encode().unwrap(), + } + .encode() + .unwrap(), + }; + server_stream + .write_all(&pre_hello.encode_stream(&result_frame).unwrap()) + .await + .unwrap(); + let codec = FrameCodec::new( + FrameLimits { + max_wire_frame: client_hello.receive.max_frame, + max_decoded_frame: client_hello.receive.max_decoded, + }, + [], + ) + .unwrap(); + (codec, client_hello) + } + + /// A client that declares a wider receive budget parks as much as it declared while it + /// waits for another frame: 17 MiB of Transfer data in 1,372 frames within 64 MiB, past + /// both the default's caps (16 MiB, 1,024 frames). + #[tokio::test] + async fn frames_parked_within_a_wider_receive_budget_keep_the_session() { + const BUDGET: u64 = 64 << 20; + const CHUNK: usize = 64 * 1024; + const CHUNKS: usize = 272; + const SMALL: usize = 1_100; + let (client_stream, mut server_stream) = tokio::io::duplex(64 * 1024); + let server = tokio::spawn(async move { + let (codec, client_hello) = answer_hello(&mut server_stream).await; + assert_eq!(client_hello.receive.max_buffered, BUDGET); + for index in 0..CHUNKS + SMALL { + let (offset, size) = if index < CHUNKS { + (index * CHUNK, CHUNK) + } else { + (CHUNKS * CHUNK + index - CHUNKS, 1) + }; + let data = Frame { + header: FrameHeader::event( + family::TRANSFER, + yas_wire::transfer::kind::BYTE_DATA, + ), + payload: ByteData { + transfer_id: 7, + offset: offset as u64, + data: vec![0; size], + } + .encode() + .unwrap(), + }; + server_stream + .write_all(&codec.encode_stream(&data).unwrap()) + .await + .unwrap(); + } + // What the client waits for: another Transfer's end. + let close = Frame { + header: FrameHeader::event(family::TRANSFER, yas_wire::transfer::kind::CLOSE), + payload: TransferClose { + transfer_id: 8, + final_data_bytes: 0, + status: Status::Ok.code(), + detail: Vec::new(), + } + .encode() + .unwrap(), + }; + server_stream + .write_all(&codec.encode_stream(&close).unwrap()) + .await + .unwrap(); + server_stream + }); + let mut client = NativeClient::connect_transport( + transport::Transport::Duplex(client_stream), + &HelloOptions::named("yas-test").receive_budget(BUDGET), + ) + .await + .unwrap(); + assert_eq!(client.receive_budget(), BUDGET); + let close = client + .next_matching_event(family::TRANSFER, yas_wire::transfer::kind::CLOSE) + .await + .unwrap(); + assert_eq!( + TransferClose::decode(&close.payload).unwrap().transfer_id, + 8 + ); + assert_eq!(client.pending.len(), CHUNKS + SMALL); + assert!( + client.pending_bytes > CHUNKS * CHUNK, + "{}", + client.pending_bytes + ); + drop(server.await.unwrap()); + } + #[tokio::test] async fn native_session_answers_peer_ping_without_legacy_fallback() { let (client_stream, mut server_stream) = tokio::io::duplex(64 * 1024); diff --git a/crates/client/src/options.rs b/crates/client/src/options.rs index ae617bc4..4f9de35d 100644 --- a/crates/client/src/options.rs +++ b/crates/client/src/options.rs @@ -24,6 +24,14 @@ pub struct HelloOptions { /// HELLO extension, so a server that does not understand it refuses the /// session instead of silently granting full control. pub read_only: bool, + /// How many bytes the server may have on their way to this client at once + /// (HELLO's `max_buffered`): every Transfer and State window the session + /// grants comes out of it, and [`crate::Client::default_process_window`] + /// sizes process output windows so that all of them fit in it. Each open + /// stream holds its window whether or not it is writing, so a client that + /// runs many processes at once and wants wide windows raises it. It bounds + /// what may wait here unread, not what is allocated. 16 MiB by default. + pub receive_budget: u64, } impl Default for HelloOptions { @@ -34,6 +42,7 @@ impl Default for HelloOptions { families: None, required: Vec::new(), read_only: false, + receive_budget: yas_wire::schema::transport::RECOMMENDED_BUFFERED, } } } @@ -65,6 +74,13 @@ impl HelloOptions { self } + /// Set [`HelloOptions::receive_budget`], at least 1 byte and at most + /// 1 GiB (the protocol's hard maximum). + pub fn receive_budget(mut self, bytes: u64) -> Self { + self.receive_budget = bytes.clamp(1, yas_wire::schema::transport::HARD_MAX_BUFFERED); + self + } + pub(crate) fn family_offers(&self) -> Vec { yas_wire::schema::FAMILIES .iter() diff --git a/crates/client/src/process.rs b/crates/client/src/process.rs index bca92bee..2b8e6ab3 100644 --- a/crates/client/src/process.rs +++ b/crates/client/src/process.rs @@ -80,10 +80,12 @@ //! flags to a hosted server. //! //! Every stdout/stderr stream holds its [`Command::window`] of the -//! session's receive budget (16 MiB) while it is open. Unless a command sets +//! session's receive budget ([`crate::HelloOptions::receive_budget`], 16 MiB +//! unless the client asks for more) while it is open. Unless a command sets //! one, [`Client::default_process_window`] sizes it so that the server's -//! per-session maximum of processes fits: 384 KiB at the default of 16, -//! never more than 1 MiB nor less than 16 KiB. +//! per-session maximum of processes fits: 384 KiB at the default of 16 in +//! 16 MiB (or 256 in 256 MiB), never more than 1 MiB nor less than 16 KiB. +//! Over a network a stream carries at most about a window a round trip. use std::ffi::OsStr; use std::sync::Arc; @@ -110,7 +112,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)] @@ -148,6 +150,7 @@ pub struct Command { window: Option, operation_id: [u8; 16], leave_residue: Option>, + keep_output: Option<(u64, u64)>, } impl Command { @@ -164,6 +167,7 @@ impl Command { window: None, operation_id: nonzero_id(), leave_residue: None, + keep_output: None, } } @@ -249,6 +253,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; @@ -469,10 +485,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. @@ -495,7 +515,8 @@ pub struct Process { /// detached), the exit is asked for with WAIT. struct ReportedExit { frames: tokio::sync::Mutex, - status: std::sync::OnceLock, + /// The exit, and what KEEP_OUTPUT dropped of stdout and of stderr. + status: std::sync::OnceLock<(ExitStatus, [Option; 2])>, lost: Arc, } @@ -524,14 +545,14 @@ impl ReportedExit { } }; tokio::pin!(until_deadline); - if let Some(status) = self.status.get() { + 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() { + if let Some((status, _)) = self.status.get() { return Ok(Some(status.clone())); } let frame = tokio::select! { @@ -542,13 +563,17 @@ impl ReportedExit { () = &mut until_deadline => return Ok(None), }; let status = match frame { - Some(frame) => ExitStatus::from_wire(ExitReport::decode(&frame.payload)?.exit), + 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, + Some(status) => (status, [None; 2]), None => return Ok(None), } } @@ -558,7 +583,7 @@ impl ReportedExit { .unwrap_or_else(|| Error::protocol("Process EXIT route closed"))); } }; - Ok(Some(self.status.get_or_init(|| status).clone())) + Ok(Some(self.status.get_or_init(|| status).0.clone())) } } @@ -664,6 +689,14 @@ 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. @@ -702,6 +735,7 @@ impl Process { status, stdout, stderr, + elided: [self.elided(false), self.elided(true)], }) } @@ -898,6 +932,17 @@ impl Client { 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() @@ -1135,13 +1180,13 @@ impl Client { /// The output window a [`Command`] gets unless it sets one: 1 MiB, or /// less when the server admits so many processes per session that their /// stdout and stderr windows would not fit in three quarters of the - /// session's receive budget (16 MiB), leaving the rest for everything - /// else the session receives. Never below 16 KiB. + /// session's receive budget ([`Client::receive_budget`]), leaving the + /// rest for everything else the session receives. Never below 16 KiB. pub fn default_process_window(&self) -> u64 { let per_session = self.process_limits().map_or(1, |limits| { u64::from(limits.max_processes_per_session).max(1) }); - let budget = yas_wire::schema::transport::RECOMMENDED_BUFFERED / 4 * 3; + let budget = self.receive_budget() / 4 * 3; (budget / (2 * per_session)).clamp(MIN_AUTO_WINDOW, DEFAULT_WINDOW) } diff --git a/crates/proxy/Cargo.toml b/crates/proxy/Cargo.toml index 6eb4f038..0651af5e 100644 --- a/crates/proxy/Cargo.toml +++ b/crates/proxy/Cargo.toml @@ -41,6 +41,7 @@ rand = "0.10" url = "2" reqwest = { version = "0.13", default-features = false, features = ["rustls"] } serde_json = "1" +socket2 = "0.6" yas-webrtc-forwarder = { workspace = true } yas-composite-transport = { workspace = true } yas-uplink = { workspace = true } diff --git a/crates/proxy/src/uplink_producer.rs b/crates/proxy/src/uplink_producer.rs index 773db6a4..b9ae6f8d 100644 --- a/crates/proxy/src/uplink_producer.rs +++ b/crates/proxy/src/uplink_producer.rs @@ -61,6 +61,38 @@ const MAX_PENDING: usize = 64; /// Liveness of a session: a keepalive this often, dead after this long silent. const KEEPALIVE: Duration = Duration::from_secs(10); const IDLE_TIMEOUT: Duration = Duration::from_secs(30); +/// The congestion window a WebTransport session starts with: CUBIC's, but +/// large enough that an answer goes out at once. quinn paces a window over +/// the round trip (1.25 windows per RTT), and a burst leaves the connection +/// app-limited, which keeps the window from growing past it: from quinn's +/// 12,000 bytes (ten 1,200-byte datagrams), a 1 MiB answer settles at two +/// round trips. Loss shrinks the window as ever, and a stream's receive window +/// (1.25 MB) still bounds what one consumer has in flight. +const INITIAL_WINDOW: u64 = 16 << 20; +/// The UDP receive buffer asked for (the system may allow less): a paced +/// burst is up to 256 datagrams, well over Linux's default of 208 KiB. +const RECEIVE_BUFFER: usize = 8 << 20; +/// Whether the system reports double the receive buffer it allows (for its +/// bookkeeping), as Linux does. +const REPORTS_DOUBLE: bool = cfg!(any(target_os = "linux", target_os = "android")); +/// What the system reports for all of [`RECEIVE_BUFFER`], uncapped. +const RECEIVE_BUFFER_REPORTED: usize = if REPORTS_DOUBLE { + 2 * RECEIVE_BUFFER +} else { + RECEIVE_BUFFER +}; +/// What caps a UDP receive buffer, for whoever would raise it. +const RECEIVE_BUFFER_CAP: &str = if REPORTS_DOUBLE { + "net.core.rmem_max" +} else if cfg!(any( + target_vendor = "apple", + target_os = "freebsd", + target_os = "dragonfly" +)) { + "kern.ipc.maxsockbuf" +} else { + "the system" +}; /// How long a WebSocket (session or stream) may take to open. const WEBSOCKET_CONNECT: Duration = Duration::from_secs(10); /// How long [`Transport::Auto`] waits for a WebTransport session before it @@ -270,6 +302,11 @@ pub enum Event { /// A consumer stream couldn't be served (a WebSocket stream the relay /// asked for that didn't open, a datagram lane this session can't carry). StreamFailed { error: String }, + /// The UDP receive buffer a WebTransport session got, in bytes, as the + /// system reports it (Linux doubles what it allows, for its bookkeeping). + /// A system allowing less than the uplink asks for only makes bursts lose + /// more packets; the session goes on. + ReceiveBuffer { bytes: usize }, } impl fmt::Display for Event { @@ -300,6 +337,17 @@ impl fmt::Display for Event { write!(f, "local yas server unavailable at {local}: {error}") } Self::StreamFailed { error } => write!(f, "consumer stream failed: {error}"), + Self::ReceiveBuffer { bytes } if *bytes < RECEIVE_BUFFER_REPORTED => { + write!(f, "UDP receive buffer: {bytes} bytes, under the ")?; + if REPORTS_DOUBLE { + write!(f, "{RECEIVE_BUFFER_REPORTED} Linux reports for the ")?; + } + write!( + f, + "{RECEIVE_BUFFER} asked for ({RECEIVE_BUFFER_CAP} caps it)" + ) + } + Self::ReceiveBuffer { bytes } => write!(f, "UDP receive buffer: {bytes} bytes"), } } } @@ -535,8 +583,8 @@ impl Producer { limit: Option, active: &Active, ) -> SessionEnd { - let client = match webtransport_client(relay.cert_hash.as_deref()) { - Ok(client) => client, + let (client, receive_buffer) = match webtransport_client(relay.cert_hash.as_deref()) { + Ok(made) => made, Err(error) => return SessionEnd::NeverConnected(error), }; // Careful: the URL is the credential β€” show `relay.label` only. @@ -562,6 +610,9 @@ impl Producer { relay: relay.label.clone(), carrier: Carrier::WebTransport, }); + if let Some(bytes) = receive_buffer { + self.emit(Event::ReceiveBuffer { bytes }); + } active.set(Session::WebTransport(Box::new(session.clone()))); let routes = DatagramRoutes::new(MAX_DATAGRAM_ROUTES); @@ -946,9 +997,11 @@ enum SessionEnd { } /// A WebTransport client with the liveness settings of docs/uplink.md (10 s -/// keepalive, 30 s idle timeout), on YAS's rustls provider, verifying the -/// relay with the platform's roots or a certificate pin. -fn webtransport_client(cert_hash: Option<&[u8]>) -> Result { +/// keepalive, 30 s idle timeout) and a window for bursts ([`INITIAL_WINDOW`], +/// [`RECEIVE_BUFFER`]), on YAS's rustls provider, verifying the relay with +/// the platform's roots or a certificate pin; with the UDP receive buffer it +/// got, where the system says. +fn webtransport_client(cert_hash: Option<&[u8]>) -> Result<(wt::Client, Option), String> { let provider = yas_webrtc_forwarder::tls::provider(); let builder = rustls::ClientConfig::builder_with_provider(provider.clone()) .with_protocol_versions(&[&rustls::version::TLS13]) @@ -974,11 +1027,73 @@ fn webtransport_client(cert_hash: Option<&[u8]>) -> Result { transport.max_idle_timeout(Some( wt::quinn::IdleTimeout::try_from(IDLE_TIMEOUT).expect("30s fits in an idle timeout"), )); + let mut cubic = wt::quinn::congestion::CubicConfig::default(); + cubic.initial_window(INITIAL_WINDOW); + transport.congestion_controller_factory(Arc::new(cubic)); config.transport_config(Arc::new(transport)); - let endpoint = wt::quinn::Endpoint::client((std::net::Ipv6Addr::UNSPECIFIED, 0).into()) - .or_else(|_| wt::quinn::Endpoint::client((std::net::Ipv4Addr::UNSPECIFIED, 0).into())) + let socket = udp_socket((std::net::Ipv6Addr::UNSPECIFIED, 0).into()) + .or_else(|_| udp_socket((std::net::Ipv4Addr::UNSPECIFIED, 0).into())) .map_err(|error| format!("UDP socket: {error}"))?; - Ok(wt::Client::new(endpoint, config)) + let receive_buffer = socket2::SockRef::from(&socket).recv_buffer_size().ok(); + let runtime = wt::quinn::default_runtime().ok_or("UDP socket: no async runtime")?; + let endpoint = + wt::quinn::Endpoint::new(wt::quinn::EndpointConfig::default(), None, socket, runtime) + .map_err(|error| format!("UDP socket: {error}"))?; + Ok((wt::Client::new(endpoint, config), receive_buffer)) +} + +/// A UDP socket bound to `address` (dual-stack for IPv6's unspecified one, as +/// quinn's own client endpoint is), with as much of [`RECEIVE_BUFFER`] as the +/// system allows. +fn udp_socket(address: std::net::SocketAddr) -> std::io::Result { + let socket = socket2::Socket::new( + socket2::Domain::for_address(address), + socket2::Type::DGRAM, + Some(socket2::Protocol::UDP), + )?; + if address.is_ipv6() { + let _ = socket.set_only_v6(false); + } + grow_receive_buffer(&socket); + socket.bind(&address.into())?; + Ok(socket.into()) +} + +/// Gives `socket` as much of [`RECEIVE_BUFFER`] as the system allows, which is +/// still better than its default; none is no reason to fail. Linux caps a size +/// over net.core.rmem_max, but macOS and the BSDs refuse one over their cap +/// (kern.ipc.maxsockbuf less mbuf overhead: 7,456,540 bytes of macOS's usual +/// 8 MiB) with ENOBUFS and keep their default, so there the uplink looks for +/// the largest size they take, between the two. +fn grow_receive_buffer(socket: &socket2::Socket) { + if socket.set_recv_buffer_size(RECEIVE_BUFFER).is_ok() { + return; + } + if let Ok(default) = socket.recv_buffer_size() { + largest_taken(default, RECEIVE_BUFFER, |size| { + socket.set_recv_buffer_size(size).is_ok() + }); + } +} + +/// The largest size from `taken` (one the system has) up to `refused` (one it +/// refused, excluded) that `take` takes, halving the gap each try: a system +/// takes every size up to its cap. A refused size leaves the buffer as it was, +/// so what `take` took last is the size returned. +fn largest_taken( + mut taken: usize, + mut refused: usize, + mut take: impl FnMut(usize) -> bool, +) -> usize { + while refused.saturating_sub(taken) > 1 { + let size = taken + (refused - taken) / 2; + if take(size) { + taken = size; + } else { + refused = size; + } + } + taken } /// A WebSocket to `url` (`wss`, a session's or a stream's) speaking @@ -1015,6 +1130,8 @@ async fn connect_websocket(url: &url::Url, cert_hash: Option<&[u8]>) -> Result bool { .is_err() } +/// A 1 MiB answer goes out in one round trip, not two: the session starts +/// with a window for it (quinn paces a window over the round trip, and an +/// app-limited connection never grows its window past what it sent). +#[tokio::test] +async fn webtransport_sessions_start_with_a_window_for_bursts() { + yas_webrtc_forwarder::tls::install_default_provider(); + tokio::time::timeout(Duration::from_secs(10), async { + let cert = rcgen::generate_simple_self_signed(vec!["localhost".into()]).unwrap(); + let hash = wt::crypto::sha256(&yas_webrtc_forwarder::tls::provider(), cert.cert.der()); + let mut relay = wt::ServerBuilder::new() + .with_addr("127.0.0.1:0".parse().unwrap()) + .with_certificate( + vec![cert.cert.der().clone()], + rustls::pki_types::PrivatePkcs8KeyDer::from(cert.signing_key.serialize_der()) + .into(), + ) + .unwrap(); + let (client, receive_buffer) = webtransport_client(Some(hash.as_ref())).unwrap(); + assert!(receive_buffer.is_some_and(|bytes| bytes > 0)); + let url: url::Url = format!("https://127.0.0.1:{}/", relay.local_addr().unwrap().port()) + .parse() + .unwrap(); + let (connected, _accepted) = tokio::join!(client.connect(url), async { + relay.accept().await.unwrap().ok().await.unwrap() + }); + let session = connected.unwrap(); + let window = (*session).stats().path.cwnd; + assert!( + window >= INITIAL_WINDOW, + "the session starts with a {window}-byte window" + ); + }) + .await + .expect("window test stalled"); +} + #[cfg(unix)] #[tokio::test] async fn producer_authenticates_before_ipc_and_encrypts_datagrams() { @@ -88,7 +124,7 @@ async fn producer_authenticates_before_ipc_and_encrypts_datagrams() { .into(), ) .unwrap(); - let outer = webtransport_client(Some(hash.as_ref())).unwrap(); + let (outer, _) = webtransport_client(Some(hash.as_ref())).unwrap(); let url: url::Url = format!("https://127.0.0.1:{}/", worker.local_addr().unwrap().port()) .parse() .unwrap(); @@ -302,6 +338,8 @@ impl WebSocketRelay { let mut session_taken = false; loop { let (tcp, _) = listener.accept().await.unwrap(); + // As a relay does: its answers go out at once. + let _ = tcp.set_nodelay(true); let Ok(tls) = tls.accept(tcp).await else { continue; }; @@ -649,6 +687,78 @@ async fn websocket_session_serves_consumers_and_takes_allowlist_changes() { .expect("WebSocket session test stalled"); } +/// Small writes a moment apart, the shape of a terminal's echo, go out at +/// once: were Nagle's algorithm on for the relay's TCP connection, the second +/// would wait for the relay to acknowledge the first, which it delays by +/// 40 ms or more while it has nothing to send back. +#[cfg(unix)] +#[tokio::test] +async fn websocket_streams_answer_small_writes_without_nagle_stalls() { + tokio::time::timeout(Duration::from_secs(20), async { + let dir = tempfile::tempdir().unwrap(); + let keys = keys(); + let (local, listener) = local_socket(dir.path()); + let mut relay = WebSocketRelay::start().await; + let producer = Producer::new( + "https://127.0.0.1:1/control", + "token", + keys.server.clone(), + local, + ) + .unwrap(); + let active = Active::default(); + let _session = { + let (producer, active, relay) = (producer.clone(), active.clone(), relay.relay()); + tokio::spawn(async move { producer.websocket_session(&relay, &active).await }) + }; + relay.ask("echo"); + let (_, stream) = relay.stream().await; + let mut consumer = yas_uplink::connect(stream, keys.client.clone()) + .await + .unwrap(); + consumer + .write_all(yas_wire::PREFACE.as_slice()) + .await + .unwrap(); + consumer.flush().await.unwrap(); + let (mut ipc, _) = listener.accept().await.unwrap(); + let mut preface = vec![0; yas_wire::PREFACE.len()]; + ipc.read_exact(&mut preface).await.unwrap(); + // The local server answers each byte with two, a millisecond apart. + tokio::spawn(async move { + let mut byte = [0; 1]; + while ipc.read_exact(&mut byte).await.is_ok() { + if ipc.write_all(b"a").await.is_err() { + break; + } + tokio::time::sleep(Duration::from_millis(1)).await; + if ipc.write_all(b"b").await.is_err() { + break; + } + } + }); + let mut rounds = Vec::new(); + for _ in 0..30 { + let start = std::time::Instant::now(); + consumer.write_all(b"k").await.unwrap(); + consumer.flush().await.unwrap(); + let mut answer = [0; 2]; + consumer.read_exact(&mut answer).await.unwrap(); + assert_eq!(&answer, b"ab"); + rounds.push(start.elapsed()); + } + rounds.sort(); + let median = rounds[rounds.len() / 2]; + assert!( + median < Duration::from_millis(20), + "a round took {median:?} (median of {rounds:?}): is Nagle's algorithm on?" + ); + active.close().await; + }) + .await + .expect("Nagle test stalled"); +} + #[tokio::test] async fn websocket_relay_must_speak_the_subprotocol() { // A WebSocket server that selects no subprotocol is not an uplink relay. @@ -903,6 +1013,105 @@ fn events_read_as_yas_uplink_prints_them() { .to_string(), "relay pool exhausted; re-querying in 4s" ); + assert_eq!( + Event::ReceiveBuffer { bytes: 16 << 20 }.to_string(), + "UDP receive buffer: 16777216 bytes" + ); + // Linux reports double what it allows: all 8 MiB reads as 16 MiB, and + // 8 MiB is what a 4 MiB net.core.rmem_max allows. + #[cfg(any(target_os = "linux", target_os = "android"))] + for (bytes, note) in [ + (8 << 20, "8388608 bytes"), + (15_000_000, "15000000 bytes"), + (425_984, "425984 bytes"), + ] { + assert_eq!( + Event::ReceiveBuffer { bytes }.to_string(), + format!( + "UDP receive buffer: {note}, under the 16777216 Linux reports for the \ + 8388608 asked for (net.core.rmem_max caps it)" + ) + ); + } + // macOS reports what it allows, and takes at most 2048/2304 of + // kern.ipc.maxsockbuf (8 MiB unless raised). + #[cfg(target_vendor = "apple")] + { + assert_eq!( + Event::ReceiveBuffer { bytes: 8 << 20 }.to_string(), + "UDP receive buffer: 8388608 bytes" + ); + assert_eq!( + Event::ReceiveBuffer { bytes: 7_456_540 }.to_string(), + "UDP receive buffer: 7456540 bytes, under the 8388608 asked for \ + (kern.ipc.maxsockbuf caps it)" + ); + } + #[cfg(windows)] + { + assert_eq!( + Event::ReceiveBuffer { bytes: 8 << 20 }.to_string(), + "UDP receive buffer: 8388608 bytes" + ); + assert_eq!( + Event::ReceiveBuffer { bytes: 65_536 }.to_string(), + "UDP receive buffer: 65536 bytes, under the 8388608 asked for \ + (the system caps it)" + ); + } +} + +/// macOS and the BSDs refuse a receive buffer over their cap rather than cap +/// it, keeping their default: the uplink finds the most they take. +#[test] +fn receive_buffers_get_the_most_a_refusing_system_takes() { + // macOS's sbreserve takes up to kern.ipc.maxsockbuf Γ— MCLBYTES / (MSIZE + + // MCLBYTES), from a default of net.inet.udp.recvspace. + let cap = (8 << 20) * 2048 / 2304; + assert_eq!(cap, 7_456_540); + let mut buffer = 786_896; + let mut tries = 0; + let got = largest_taken(buffer, RECEIVE_BUFFER, |size| { + tries += 1; + let taken = size <= cap; + if taken { + buffer = size; + } + taken + }); + assert_eq!((got, buffer), (cap, cap)); + assert!(tries <= 23, "{tries} tries"); + // A system taking nothing over its default keeps it, and one whose + // default is already as large isn't asked again. + assert_eq!(largest_taken(786_896, RECEIVE_BUFFER, |_| false), 786_896); + assert_eq!( + largest_taken(16 << 20, RECEIVE_BUFFER, |_| unreachable!()), + 16 << 20 + ); +} + +/// On Linux, the uplink's socket gets as much of its ask as +/// net.core.rmem_max allows (reported doubled), and the event notes a cap +/// exactly when there is one. +#[cfg(target_os = "linux")] +#[test] +fn receive_buffers_note_linux_caps() { + // Only the initial network namespace shows it. + let Some(rmem_max) = std::fs::read_to_string("/proc/sys/net/core/rmem_max") + .ok() + .and_then(|max| max.trim().parse::().ok()) + else { + return; + }; + let socket = udp_socket((std::net::Ipv4Addr::LOCALHOST, 0).into()).unwrap(); + let bytes = socket2::SockRef::from(&socket).recv_buffer_size().unwrap(); + assert_eq!(bytes, 2 * rmem_max.min(RECEIVE_BUFFER)); + let event = Event::ReceiveBuffer { bytes }.to_string(); + assert_eq!( + event.contains("net.core.rmem_max caps it"), + rmem_max < RECEIVE_BUFFER, + "rmem_max {rmem_max}: {event}" + ); } #[test] 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 c88984cd..17caa6eb 100644 --- a/crates/server/src/yas.rs +++ b/crates/server/src/yas.rs @@ -31022,10 +31022,17 @@ async fn send_exit_report( 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::default(), + extensions: Extensions(elided), }; let _ = send_event_with_sensitivity( out, diff --git a/crates/server/src/yas_process.rs b/crates/server/src/yas_process.rs index 68d60e88..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 { @@ -382,9 +384,18 @@ impl Session { resolved_cwd: Option>, ) -> Result { let cwd = resolve_cwd(&request.cwd, resolved_cwd)?; - // 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()))?; + // 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) @@ -408,6 +419,7 @@ impl Session { flags, preserve_residual, residue_grace, + keep_output, cwd, argv: request.argv.clone(), env: request @@ -1112,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, + }) + }), } } @@ -1148,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 3b7ecee5..24df68c2 100644 --- a/crates/yas/src/generated.rs +++ b/crates/yas/src/generated.rs @@ -3800,7 +3800,8 @@ 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 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; @@ -3847,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; @@ -3905,7 +3910,9 @@ 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_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" }, @@ -3940,7 +3947,8 @@ 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: "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 }, @@ -3987,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 }, diff --git a/crates/yas/src/process.rs b/crates/yas/src/process.rs index e0dd3587..a3f9df9a 100644 --- a/crates/yas/src/process.rs +++ b/crates/yas/src/process.rs @@ -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> { @@ -953,6 +1000,83 @@ impl ExitReport { 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. @@ -1398,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( @@ -1562,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, @@ -1825,7 +1994,10 @@ mod tests { 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); + 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 { diff --git a/docs/design/processes.md b/docs/design/processes.md index 133f4c4d..f77f9d5e 100644 --- a/docs/design/processes.md +++ b/docs/design/processes.md @@ -98,7 +98,13 @@ 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 +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. @@ -153,12 +159,17 @@ So do the server-wide budgets: `YAS_PROCESS_MAX_SPAWNING` (concurrent native spawn calls server-wide) defaults to the pending-spawn maximum. -A session's streams share its receive budgets, 16 MiB each way: +A session's streams share its receive budgets, each side's HELLO +`receive_max_buffered` (16 MiB unless a client declares more, up to 1 GiB): - **stdout and stderr**: every open stream holds the receive credit its client - granted, and that credit comes out of the client's declared buffer. A client - that runs many processes on one session should use smaller windows. yas-client - divides three quarters of the budget between the server's per-session maximum. + granted, whether or not it writes, and that credit comes out of the client's + declared buffer: once it is all held, a new stream gets none until another + lets go. A client that runs many processes on one session should use smaller + windows, or declare a wider buffer. yas-client divides three quarters of its + budget (`HelloOptions::receive_budget`) between the stdout and stderr of the + server's per-session maximum: 24 KiB each at 256 processes in 16 MiB, 384 KiB + in 256 MiB. Over a network a stream carries about a window a round trip. The server sends output as far as the credit reaches, so a window smaller than one 64 KiB chunk still makes progress (servers that predate configurable maxima waited for credit through a whole chunk). diff --git a/docs/design/yas.md b/docs/design/yas.md index a1f84ede..4d0738bf 100644 --- a/docs/design/yas.md +++ b/docs/design/yas.md @@ -3282,6 +3282,29 @@ 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 diff --git a/docs/uplink.md b/docs/uplink.md index fa1003f5..1c3800a4 100644 --- a/docs/uplink.md +++ b/docs/uplink.md @@ -270,6 +270,22 @@ A WebTransport (HTTP/3 CONNECT) session to the relay URL. Liveness settings are a **10s keepalive** and a **30s idle timeout**, so a dead relay is noticed within 30 seconds without any application-level pings. +Its congestion control is CUBIC with a **16 MiB initial window**, and the +uplink asks for an **8 MiB UDP receive buffer**. quinn paces a window over the +round trip, and answers leave the connection app-limited, so the window +never grows past what they need. From quinn's usual 12,000 bytes (ten +1,200-byte datagrams, as RFC 9002 recommends), a 1 MiB answer settles at two +round trips; from 16 MiB, it takes one. Loss still shrinks the window, and +each stream's 1.25 MB receive window still bounds what one consumer has in +flight. A system may allow a smaller buffer: Linux caps it at +`net.core.rmem_max` (then doubles it for bookkeeping, so all 8 MiB reads as +16 MiB), while macOS and the BSDs refuse a size over their cap +(`kern.ipc.maxsockbuf`, less their overhead: about 7.1 MiB of macOS's usual +8 MiB), so there the uplink asks for the largest size they take. The session +goes on either way, and the uplink says what it got as it connects (`UDP +receive buffer: N bytes`, and what caps it when something does). A relay +should do the same for what it sends. + The uplink never opens streams. The relay opens **one bidirectional stream per consumer**. After Noise authentication, the uplink bridges decrypted bytes to a fresh local YAS socket. Direct streams carry the normal YAS preface diff --git a/js/core/src/yas/generated.ts b/js/core/src/yas/generated.ts index 181bd360..6d36d946 100644 --- a/js/core/src/yas/generated.ts +++ b/js/core/src/yas/generated.ts @@ -1657,7 +1657,8 @@ 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_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; @@ -1704,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; @@ -11850,7 +11855,15 @@ export const YAS_SCHEMA = { }, { "name": "exit_report", - "layout": "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions" + "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", @@ -11898,9 +11911,13 @@ export const YAS_SCHEMA = { "name": "SPAWN_REPORT_EXIT", "value": 16 }, + { + "name": "SPAWN_KEEP_OUTPUT", + "value": 32 + }, { "name": "SPAWN_LAUNCHER_FLAGS_EXTENDED", - "value": 28 + "value": 60 }, { "name": "ENV_EMPTY", @@ -12086,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 diff --git a/protocol/yas/families/process.toml b/protocol/yas/families/process.toml index 11f19714..055040d7 100644 --- a/protocol/yas/families/process.toml +++ b/protocol/yas/families/process.toml @@ -48,9 +48,16 @@ value = 12 [[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 = 28 +value = 60 [[constant]] name = "ENV_EMPTY" @@ -201,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" @@ -435,7 +454,15 @@ layout = "process_handle:u64" [[type]] name = "exit_report" -layout = "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions" +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" diff --git a/protocol/yas/schema.json b/protocol/yas/schema.json index f30f6530..57bc83a7 100644 --- a/protocol/yas/schema.json +++ b/protocol/yas/schema.json @@ -8927,7 +8927,15 @@ }, { "name": "exit_report", - "layout": "process_handle:u64,exit:bytes_u32 containing ExitRecord,Extensions" + "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", @@ -8975,9 +8983,13 @@ "name": "SPAWN_REPORT_EXIT", "value": 16 }, + { + "name": "SPAWN_KEEP_OUTPUT", + "value": 32 + }, { "name": "SPAWN_LAUNCHER_FLAGS_EXTENDED", - "value": 28 + "value": 60 }, { "name": "ENV_EMPTY", @@ -9163,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 diff --git a/protocol/yas/wire.md b/protocol/yas/wire.md index 60191502..8dadf360 100644 --- a/protocol/yas/wire.md +++ b/protocol/yas/wire.md @@ -926,7 +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 | +| `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 |