Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

198 changes: 197 additions & 1 deletion crates/cli/tests/client_host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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.
Expand Down
11 changes: 11 additions & 0 deletions crates/client/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,8 @@ struct Inner {
tasks: Vec<tokio::task::JoinHandle<()>>,
/// 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 {
Expand Down Expand Up @@ -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 {
Expand All @@ -334,6 +337,7 @@ impl Client {
next_request_id: AtomicU32::new(3),
tasks: vec![writer, reader],
started: std::time::Instant::now(),
receive_budget,
}),
}
}
Expand All @@ -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()
Expand Down
Loading
Loading