Skip to content
Merged
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
318 changes: 316 additions & 2 deletions crates/cli/tests/client_host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,11 @@ use yas_client::{Client, Error, HelloOptions};
const TIMEOUT: Duration = Duration::from_secs(30);

fn options() -> HostOptions {
HostOptions::new(env!("CARGO_BIN_EXE_yas"))
options_for(env!("CARGO_BIN_EXE_yas"))
}

fn options_for(binary: impl Into<std::path::PathBuf>) -> HostOptions {
HostOptions::new(binary)
.arg("--no-persistent-extensions")
.env("YAS_EXT", "0")
.env("YAS_CHANNEL", "0")
Expand Down Expand Up @@ -524,6 +528,104 @@ async fn a_dropped_output_stream_leaves_its_process_and_session_working() {
still_runs_commands(&client).await;
}

/// A spawned process whose attachment goes before its exit (it was detached, or its streams
/// were dropped) is asked for its exit with WAIT, which its attachment would have reported.
/// Nothing holds it to the session meanwhile: a detached one gives back its slot, and either
/// can be attached to again.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_process_whose_attachment_went_is_waited_for_and_held_by_nothing() {
use yas_client::wire::schema::process as schema;
let server = start().await;
let client = server.connect().await.unwrap();
assert_ne!(
client.launcher_flags() & schema::SPAWN_REPORT_EXIT as u32,
0
);
let per_session = client.process_limits().unwrap().max_processes_per_session as usize;
// Detached, each gives back its slot: one more than the session holds still spawns.
let mut detached = Vec::new();
for _ in 0..per_session {
let process = client
.spawn(Command::new("sleep").arg("30").detachable(true))
.await
.unwrap();
within("a detach", process.detach()).await.unwrap();
detached.push(process);
}
let one_more = within(
"a spawn beside the detached",
client.spawn(Command::new("sh").args(["-c", "exit 6"])),
)
.await
.unwrap();
assert_eq!(
within("its exit", one_more.wait()).await.unwrap().code(),
Some(6)
);
// A detached one can be attached to again, and its exit comes (by WAIT).
let again = within("an attach", client.attach(detached[0].handle(), false))
.await
.unwrap();
detached[0].kill().await.unwrap();
let status = within("a detached one's exit", detached[0].wait())
.await
.unwrap();
// (Each answer stamps the time it is given.)
let attached = within("its exit, attached", again.wait()).await.unwrap();
assert_eq!(
(&status.kind, &status.reason, status.raw_code),
(&attached.kind, &attached.reason, attached.raw_code)
);
for process in &detached[1..] {
process.kill().await.unwrap();
within("a detached one's exit", process.wait())
.await
.unwrap();
}
// An ordinary one whose streams were dropped: the same.
let mut sleeper = client.spawn(Command::new("sleep").arg("30")).await.unwrap();
drop(sleeper.take_stdout());
drop(sleeper.take_stderr());
let again = within("an attach", client.attach(sleeper.handle(), false))
.await
.unwrap_or_else(|error| panic!("{error}\n{}", server_log(&server)));
assert_eq!(
within(
"a short wait",
sleeper.wait_timeout(Duration::from_millis(200))
)
.await
.unwrap(),
None
);
sleeper.kill().await.unwrap();
let status = within("its exit", sleeper.wait()).await.unwrap();
// (Each answer stamps the time it is given.)
let attached = within("its exit, attached", again.wait()).await.unwrap();
assert_eq!(
(&status.kind, &status.reason, status.raw_code),
(&attached.kind, &attached.reason, attached.raw_code)
);
// One whose stdin was aborted: its attachment goes with it.
let mut reader = client
.spawn(
Command::new("sh")
.args(["-c", "sleep 0.3; exit 3"])
.stdin(Stdin::Piped),
)
.await
.unwrap();
reader.take_stdin().unwrap().abort();
assert_eq!(
within("its exit, stdin aborted", reader.wait())
.await
.unwrap()
.code(),
Some(3)
);
still_runs_commands(&client).await;
}

#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_watcher_that_falls_behind_is_dropped_alone() {
let server = start().await;
Expand Down Expand Up @@ -568,6 +670,218 @@ async fn a_watcher_that_falls_behind_is_dropped_alone() {
still_runs_commands(&owner).await;
}

/// A server started with `--process-max-waits 1`, with that one WAIT held on a
/// `sleep` and shown to be held: a second WAIT is refused.
async fn hold_the_only_wait(
client: &Client,
) -> (yas_client::process::Process, tokio::task::JoinHandle<()>) {
let sleeper = client.spawn(Command::new("sleep").arg("30")).await.unwrap();
let handle = sleeper.handle();
let holder = {
let client = client.clone();
tokio::spawn(async move {
let _ = client.wait_process(handle, None).await;
})
};
// The holder's WAIT goes first on the one connection.
tokio::time::sleep(Duration::from_millis(200)).await;
let refused = client
.wait_process(handle, Some(Duration::from_millis(1)))
.await
.unwrap_err();
assert_eq!(
refused.status(),
Some(yas_client::wire::core::Status::ResourceExhausted),
"{refused}"
);
(sleeper, holder)
}

/// Commands whose exits must come whole: a non-zero exit with stderr, output
/// beyond the stream buffer read while the exit is awaited, a signal.
async fn exits_come_with_all_their_output(server: &HostedServer, client: &Client) {
let output = client
.spawn(Command::new("sh").args(["-c", "echo out; echo 'went wrong' >&2; exit 42"]))
.await
.unwrap()
.output()
.await
.unwrap();
assert_eq!(output.status.code(), Some(42), "{}", output.status);
assert_eq!(output.stdout, b"out\n");
assert_eq!(output.stderr, b"went wrong\n");

// Three times the largest stream buffer a default server keeps: the
// command blocks on its pipe until this reads, and its exit arrives while
// the tail is still on its way, as Ultimator's bash reads it.
const BYTES: u64 = 3 * 8 * 1024 * 1024;
let mut process = client
.spawn(Command::new("sh").args([
"-c".to_owned(),
format!("head -c {BYTES} /dev/zero; printf end; echo tail >&2; exit 9"),
]))
.await
.unwrap();
let stdout = process.take_stdout().unwrap();
let stderr = process.take_stderr().unwrap();
let reader = tokio::spawn(async move {
(
stdout.read_to_end(BYTES + 1024).await.unwrap(),
stderr.read_to_end(1024).await,
)
});
let status = tokio::time::timeout(TIMEOUT, process.wait())
.await
.unwrap_or_else(|_| panic!("no exit\n{}", server_log(server)))
.unwrap();
assert_eq!(status.code(), Some(9), "{status}");
let (stdout, stderr) = tokio::time::timeout(TIMEOUT, reader)
.await
.expect("output after the exit")
.unwrap();
let stderr = stderr.unwrap_or_else(|error| panic!("{error}\n{}", server_log(server)));
assert_eq!(stdout.len() as u64, BYTES + 3);
assert!(stdout.ends_with(b"end"));
assert!(stdout[..BYTES as usize].iter().all(|byte| *byte == 0));
assert_eq!(stderr, b"tail\n");

let sleeper = client.spawn(Command::new("sleep").arg("30")).await.unwrap();
assert_eq!(
sleeper
.wait_timeout(Duration::from_millis(100))
.await
.unwrap(),
None
);
sleeper.signal(Signal::Terminate).await.unwrap();
let status = sleeper.wait().await.unwrap();
assert_eq!(status.signal(), Some(libc::SIGTERM), "{status}");
// Asked again, the same exit (an older server's WAITs may each give it another time).
let same = |other: &yas_client::process::ExitStatus| {
(&other.kind, &other.reason, other.raw_code)
== (&status.kind, &status.reason, status.raw_code)
};
let again = sleeper.wait().await.unwrap();
assert!(same(&again), "{again:?} after {status:?}");
let again = sleeper
.wait_timeout(Duration::from_millis(1))
.await
.unwrap()
.expect("exited");
assert!(same(&again), "{again:?} after {status:?}");

for _ in 0..20 {
let output = client
.spawn(&Command::new("true"))
.await
.unwrap()
.output()
.await
.unwrap();
assert!(output.status.success());
}
}

#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_spawned_process_reports_its_exit_without_a_wait() {
use yas_client::wire::schema::process as schema;
let server = tokio::time::timeout(
TIMEOUT,
HostedServer::start(
options()
.arg("--verbose")
.args(["--process-max-waits", "1"]),
),
)
.await
.expect("hosted server start timed out")
.expect("hosted server starts");
let client = server.connect().await.unwrap();
assert_ne!(
client.launcher_flags() & schema::SPAWN_REPORT_EXIT as u32,
0
);
// With the only WAIT held, every exit below comes unasked.
let (sleeper, holder) = hold_the_only_wait(&client).await;
exits_come_with_all_their_output(&server, &client).await;
// A process attached to (not spawned) is still waited for: another session, its only
// WAIT held too, attaches to this one's sleeper and is refused a second WAIT.
let other = server.connect().await.unwrap();
let (other_sleeper, other_holder) = hold_the_only_wait(&other).await;
let attached = other.attach(sleeper.handle(), false).await.unwrap();
let refused = attached
.wait_timeout(Duration::from_millis(1))
.await
.unwrap_err();
assert_eq!(
refused.status(),
Some(yas_client::wire::core::Status::ResourceExhausted),
"{refused}"
);
for (which, (sleeper, holder)) in [(sleeper, holder), (other_sleeper, other_holder)]
.into_iter()
.enumerate()
{
sleeper.kill().await.unwrap();
// Killed by YAS at this client's request, as a WAIT would answer it.
let status = sleeper.wait().await.unwrap();
assert_eq!(
(&status.kind, &status.reason),
(
&yas_client::process::ExitKind::Killed,
&yas_client::process::ExitReason::Client
),
"sleeper {which}: {status:?}"
);
tokio::time::timeout(TIMEOUT, holder)
.await
.unwrap()
.unwrap();
}
}

/// Against a server from before SPAWN_REPORT_EXIT, the same commands are
/// waited for with WAIT. Run with
/// `YAS_OLD_SERVER=/path/to/yas cargo test -p yas-cli --test client_host -- --ignored`, with a
/// binary that paces process output (its long output must not reset the stream).
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[ignore = "needs YAS_OLD_SERVER: a yas binary from before SPAWN_REPORT_EXIT"]
async fn a_server_from_before_report_exit_is_waited_for() {
use yas_client::wire::schema::process as schema;
let binary = std::env::var_os("YAS_OLD_SERVER").expect("YAS_OLD_SERVER names a yas binary");
let server = tokio::time::timeout(
TIMEOUT,
HostedServer::start(options_for(binary).args(["--process-max-waits", "1"])),
)
.await
.expect("hosted server start timed out")
.expect("hosted server starts");
let client = server.connect().await.unwrap();
assert_eq!(
client.launcher_flags() & schema::SPAWN_REPORT_EXIT as u32,
0
);
// Its exits take a WAIT: with the only one held, waiting is refused.
let (sleeper, holder) = hold_the_only_wait(&client).await;
let process = client
.spawn(Command::new("sh").args(["-c", "exit 5"]))
.await
.unwrap();
let refused = process.wait().await.unwrap_err();
assert_eq!(
refused.status(),
Some(yas_client::wire::core::Status::ResourceExhausted),
"{refused}"
);
sleeper.kill().await.unwrap();
tokio::time::timeout(TIMEOUT, holder)
.await
.unwrap()
.unwrap();
assert_eq!(process.wait().await.unwrap().code(), Some(5));
exits_come_with_all_their_output(&server, &client).await;
}

#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn a_commands_background_children_die_with_it() {
let server = start().await;
Expand Down Expand Up @@ -920,7 +1234,7 @@ async fn a_command_can_leave_its_background_running_with_a_null_stdin() {
let client = server.connect().await.unwrap();
assert_eq!(
client.launcher_flags(),
(schema::SPAWN_LEAVE_RESIDUE | schema::SPAWN_STDIN_NULL) as u32
(schema::SPAWN_LEAVE_RESIDUE | schema::SPAWN_STDIN_NULL | schema::SPAWN_REPORT_EXIT) as u32
);
// The background `sleep` holds stdout: the exit comes after the grace, and
// the sleep keeps running.
Expand Down
Loading
Loading