mirror of
https://github.com/NVIDIA/OpenShell.git
synced 2026-10-02 07:34:45 +08:00
fix(sandbox): preserve local sessions across host sleep (#3573)
* fix(cli): recover sandbox connect transport Signed-off-by: Drew Newberry <anewberry@nvidia.com> * fix(auth): propagate non-expiring local sandbox sessions Signed-off-by: Drew Newberry <anewberry@nvidia.com> * fix(cli): distinguish main exit from transport loss Signed-off-by: Drew Newberry <anewberry@nvidia.com> * fix(cli): bound sandbox connect recovery Signed-off-by: Drew Newberry <anewberry@nvidia.com> --------- Signed-off-by: Drew Newberry <anewberry@nvidia.com>
This commit is contained in:
+1
-1
@@ -232,7 +232,7 @@ tonic = { version = "0.14", features = ["transport"] }
|
||||
tonic-prost = "0.14"
|
||||
tower = "0.5"
|
||||
url = "2"
|
||||
nix = { version = "0.29", features = ["user"] }
|
||||
nix = { version = "0.29", features = ["process", "signal", "term", "user"] }
|
||||
|
||||
[dev-dependencies]
|
||||
serial_test = "3"
|
||||
|
||||
@@ -3,6 +3,8 @@
|
||||
|
||||
#![cfg(feature = "e2e")]
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
use std::fs;
|
||||
use std::process::Stdio;
|
||||
use std::time::Duration;
|
||||
|
||||
@@ -211,6 +213,64 @@ async fn reconnect_with_input_ownership(
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
fn find_process_with_args(expected_args: &[&str]) -> Option<u32> {
|
||||
for entry in fs::read_dir("/proc").ok()?.filter_map(Result::ok) {
|
||||
let Ok(pid) = entry.file_name().to_string_lossy().parse::<u32>() else {
|
||||
continue;
|
||||
};
|
||||
let cmdline = fs::read(entry.path().join("cmdline")).unwrap_or_default();
|
||||
let args = cmdline.split(|byte| *byte == 0).collect::<Vec<_>>();
|
||||
if expected_args
|
||||
.iter()
|
||||
.all(|expected| args.contains(&expected.as_bytes()))
|
||||
{
|
||||
return Some(pid);
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
fn find_child_process_with_args(parent_pid: u32, expected_args: &[&str]) -> Option<u32> {
|
||||
for entry in fs::read_dir("/proc").ok()?.filter_map(Result::ok) {
|
||||
let Ok(pid) = entry.file_name().to_string_lossy().parse::<u32>() else {
|
||||
continue;
|
||||
};
|
||||
let status = fs::read_to_string(entry.path().join("status")).unwrap_or_default();
|
||||
let process_parent = status.lines().find_map(|line| {
|
||||
line.strip_prefix("PPid:")
|
||||
.and_then(|value| value.trim().parse::<u32>().ok())
|
||||
});
|
||||
if process_parent != Some(parent_pid) {
|
||||
continue;
|
||||
}
|
||||
let cmdline = fs::read(entry.path().join("cmdline")).unwrap_or_default();
|
||||
let args = cmdline.split(|byte| *byte == 0).collect::<Vec<_>>();
|
||||
if expected_args
|
||||
.iter()
|
||||
.all(|expected| args.contains(&expected.as_bytes()))
|
||||
{
|
||||
return Some(pid);
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
async fn wait_for_process_with_args(expected_args: &[&str]) -> u32 {
|
||||
tokio::time::timeout(Duration::from_secs(10), async {
|
||||
loop {
|
||||
if let Some(pid) = find_process_with_args(expected_args) {
|
||||
return pid;
|
||||
}
|
||||
sleep(Duration::from_millis(50)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.unwrap_or_else(|_| panic!("process with arguments {expected_args:?} did not start"))
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial(sandbox_lifecycle)]
|
||||
async fn sandbox_stop_start_preserves_workspace() {
|
||||
@@ -722,6 +782,237 @@ async fn canonical_main_disconnect_reconnect_replays_history_for_same_process()
|
||||
sandbox.cleanup().await;
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
#[tokio::test]
|
||||
#[serial(sandbox_lifecycle)]
|
||||
async fn canonical_main_connect_recovers_its_ssh_transport() {
|
||||
let script = r#"trap 'kill "$writer" 2>/dev/null || true' EXIT; (n=1; while true; do printf 'transport_pid=%s sequence=%04d\n' "$$" "$n"; n=$((n + 1)); sleep 0.2; done) & writer=$!; while IFS= read -r line; do printf 'transport_pid=%s input=%s\n' "$$" "$line"; done"#;
|
||||
let mut sandbox = SandboxGuard::create_detached_main(&["sh", "-lc", script])
|
||||
.await
|
||||
.expect("create retained canonical main process");
|
||||
|
||||
let mut connect_cmd = openshell_cmd();
|
||||
connect_cmd
|
||||
.args(["sandbox", "connect", &sandbox.name])
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
let mut connect = connect_cmd.spawn().expect("spawn supervised attachment");
|
||||
let mut connect_stdin = connect.stdin.take().expect("connect stdin");
|
||||
let connect_stdout = connect.stdout.take().expect("connect stdout");
|
||||
let mut connect_lines = BufReader::new(connect_stdout).lines();
|
||||
let connect_stderr = connect.stderr.take().expect("connect stderr");
|
||||
let mut connect_errors = BufReader::new(connect_stderr).lines();
|
||||
|
||||
let initial_line = tokio::time::timeout(Duration::from_secs(30), connect_lines.next_line())
|
||||
.await
|
||||
.expect("initial attachment output timeout")
|
||||
.expect("read initial attachment output")
|
||||
.expect("initial attachment output should remain open");
|
||||
assert!(
|
||||
initial_line.contains("transport_pid="),
|
||||
"unexpected initial attachment output: {initial_line}"
|
||||
);
|
||||
|
||||
// Recovery intentionally starts only for an established attachment, so
|
||||
// let the initial SSH process live beyond that setup guard before killing
|
||||
// only its ProxyCommand transport.
|
||||
sleep(Duration::from_secs(3)).await;
|
||||
let proxy_pid = wait_for_process_with_args(&["ssh-proxy", "--sandbox", &sandbox.name]).await;
|
||||
let kill_status = tokio::process::Command::new("kill")
|
||||
.args(["-TERM", &proxy_pid.to_string()])
|
||||
.status()
|
||||
.await
|
||||
.expect("terminate SSH proxy transport");
|
||||
assert!(kill_status.success(), "terminate SSH proxy transport");
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(30), async {
|
||||
loop {
|
||||
let line = connect_errors
|
||||
.next_line()
|
||||
.await
|
||||
.expect("read supervised attachment diagnostics")
|
||||
.expect("supervised attachment diagnostics should remain open");
|
||||
if normalize_output(&line).contains("Connection to sandbox lost; reconnecting") {
|
||||
break;
|
||||
}
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("supervised attachment did not enter recovery");
|
||||
|
||||
let input_token = format!("after-transport-recovery-{:x}", rand::random::<u64>());
|
||||
connect_stdin
|
||||
.write_all(format!("{input_token}\n").as_bytes())
|
||||
.await
|
||||
.expect("write input after transport recovery");
|
||||
connect_stdin
|
||||
.flush()
|
||||
.await
|
||||
.expect("flush input after transport recovery");
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(30), async {
|
||||
loop {
|
||||
let line = connect_lines
|
||||
.next_line()
|
||||
.await
|
||||
.expect("read recovered attachment output")
|
||||
.expect("recovered attachment output should remain open");
|
||||
if line.contains(&format!("input={input_token}")) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("replacement SSH session did not reattach to canonical main");
|
||||
assert!(
|
||||
connect
|
||||
.try_wait()
|
||||
.expect("inspect supervised attachment")
|
||||
.is_none(),
|
||||
"sandbox connect parent should remain alive after recovery"
|
||||
);
|
||||
|
||||
connect
|
||||
.kill()
|
||||
.await
|
||||
.expect("disconnect recovered attachment");
|
||||
connect
|
||||
.wait()
|
||||
.await
|
||||
.expect("wait for recovered attachment disconnect");
|
||||
sandbox.cleanup().await;
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
#[tokio::test]
|
||||
#[serial(sandbox_lifecycle)]
|
||||
async fn canonical_main_connect_forwards_pid_targeted_termination_and_reaps_ssh() {
|
||||
use std::os::fd::OwnedFd;
|
||||
|
||||
let mut sandbox = SandboxGuard::create_detached_main(&["sh", "-lc", "exec sleep infinity"])
|
||||
.await
|
||||
.expect("create retained canonical main process");
|
||||
let pty = nix::pty::openpty(None, None).expect("open pseudo-terminal");
|
||||
let controller: OwnedFd = pty.master;
|
||||
let follower: OwnedFd = pty.slave;
|
||||
|
||||
let mut connect_cmd = openshell_cmd();
|
||||
connect_cmd
|
||||
.args(["sandbox", "connect", &sandbox.name])
|
||||
.stdin(
|
||||
follower
|
||||
.try_clone()
|
||||
.expect("duplicate PTY follower for stdin"),
|
||||
)
|
||||
.stdout(
|
||||
follower
|
||||
.try_clone()
|
||||
.expect("duplicate PTY follower for stdout"),
|
||||
)
|
||||
.stderr(
|
||||
follower
|
||||
.try_clone()
|
||||
.expect("duplicate PTY follower for stderr"),
|
||||
);
|
||||
let mut connect = connect_cmd
|
||||
.spawn()
|
||||
.expect("spawn supervised PTY attachment");
|
||||
let connect_pid = connect.id().expect("connect process ID");
|
||||
drop(follower);
|
||||
|
||||
let ssh_pid = tokio::time::timeout(Duration::from_secs(30), async {
|
||||
loop {
|
||||
if let Some(pid) =
|
||||
find_child_process_with_args(connect_pid, &["-s", "sandbox", "openshell-main"])
|
||||
{
|
||||
return pid;
|
||||
}
|
||||
sleep(Duration::from_millis(50)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("supervised SSH child did not start");
|
||||
|
||||
nix::sys::signal::kill(
|
||||
nix::unistd::Pid::from_raw(i32::try_from(connect_pid).expect("connect PID fits i32")),
|
||||
nix::sys::signal::Signal::SIGTERM,
|
||||
)
|
||||
.expect("send SIGTERM to only the OpenShell parent");
|
||||
|
||||
let status = tokio::time::timeout(Duration::from_secs(10), connect.wait())
|
||||
.await
|
||||
.expect("OpenShell parent did not terminate")
|
||||
.expect("wait for OpenShell parent");
|
||||
assert_eq!(status.code(), Some(143));
|
||||
tokio::time::timeout(Duration::from_secs(5), async {
|
||||
while fs::metadata(format!("/proc/{ssh_pid}")).is_ok() {
|
||||
sleep(Duration::from_millis(25)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("SSH child was not reaped after parent termination");
|
||||
|
||||
drop(controller);
|
||||
sandbox.cleanup().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial(sandbox_lifecycle)]
|
||||
async fn canonical_main_exit_255_is_not_retried_as_transport_failure() {
|
||||
const READY_MARKER: &str = "exit-255-ready";
|
||||
const RELEASE_PATH: &str = "/sandbox/.openshell-exit-255-release";
|
||||
let script = format!(
|
||||
"echo {READY_MARKER}; while [ ! -e '{RELEASE_PATH}' ]; do sleep 0.05; done; exit 255"
|
||||
);
|
||||
let mut sandbox = SandboxGuard::create_detached_main(&["sh", "-c", &script])
|
||||
.await
|
||||
.expect("create retained canonical main process");
|
||||
|
||||
let mut connect_cmd = openshell_cmd();
|
||||
connect_cmd
|
||||
.args(["sandbox", "connect", &sandbox.name])
|
||||
.stdin(Stdio::null())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped());
|
||||
let mut connect = connect_cmd.spawn().expect("spawn supervised attachment");
|
||||
let connect_stdout = connect.stdout.take().expect("connect stdout");
|
||||
let mut connect_lines = BufReader::new(connect_stdout).lines();
|
||||
|
||||
tokio::time::timeout(Duration::from_secs(30), async {
|
||||
loop {
|
||||
let line = connect_lines
|
||||
.next_line()
|
||||
.await
|
||||
.expect("read attachment output")
|
||||
.expect("attachment output should remain open");
|
||||
if line.contains(READY_MARKER) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("attachment did not observe the canonical main process");
|
||||
|
||||
// Keep the attachment alive beyond the setup guard used to distinguish
|
||||
// initial SSH failures from established transport failures.
|
||||
sleep(Duration::from_secs(3)).await;
|
||||
sandbox
|
||||
.exec(&["touch", RELEASE_PATH])
|
||||
.await
|
||||
.expect("release canonical main process");
|
||||
|
||||
let status = tokio::time::timeout(Duration::from_secs(10), connect.wait()).await;
|
||||
if status.is_err() {
|
||||
connect.kill().await.expect("stop stuck attachment");
|
||||
}
|
||||
sandbox.cleanup().await;
|
||||
let status = status
|
||||
.expect("exit status 255 must not enter the transport recovery loop")
|
||||
.expect("wait for canonical main attachment");
|
||||
assert_eq!(status.code(), Some(255));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial(sandbox_lifecycle)]
|
||||
async fn sandbox_create_with_no_keep_cleans_up_after_tty_command() {
|
||||
|
||||
Reference in New Issue
Block a user