use std::ffi::OsString;
use std::sync::Arc;
use std::time::{Duration, Instant};
use arc_swap::ArcSwapOption;
#[cfg(unix)]
use std::ffi::OsStr;
#[cfg(unix)]
use std::path::{Path, PathBuf};
const READ_INTERVAL: Duration = Duration::from_mins(10);
const READ_RETRY_MAX: Duration = Duration::from_mins(80);
#[cfg(unix)]
const READ_TIMEOUT: Duration = Duration::from_secs(15);
#[cfg(unix)]
const DUMP_STDOUT_CAP: usize = 1024 * 1024;
#[cfg(unix)]
const READ_CHUNK: usize = 4096;
#[cfg(unix)]
const KNOWN_SHELLS: &[&str] = &[
"sh", "bash", "dash", "ash", "zsh", "fish", "ksh", "ksh93", "mksh", "pdksh",
];
#[cfg(unix)]
#[must_use]
pub(crate) fn knows_shell(shell: &str) -> bool {
KNOWN_SHELLS.contains(&shell)
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum ReadFailure {
#[cfg(unix)]
NoShell,
#[cfg(unix)]
UnknownShell(String),
#[cfg(unix)]
Spawn(String),
#[cfg(unix)]
Timeout,
#[cfg(unix)]
NoDump(String),
#[cfg(windows)]
AssemblyFailed(String),
}
impl std::fmt::Display for ReadFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
#[cfg(unix)]
Self::NoShell => write!(f, "no shell recorded for the account"),
#[cfg(unix)]
Self::UnknownShell(name) => write!(f, "unrecognized shell {name}"),
#[cfg(unix)]
Self::Spawn(detail) => write!(f, "could not start the shell ({detail})"),
#[cfg(unix)]
Self::Timeout => write!(f, "the shell did not finish in time"),
#[cfg(unix)]
Self::NoDump(detail) => write!(f, "the shell produced no environment dump ({detail})"),
#[cfg(windows)]
Self::AssemblyFailed(detail) => {
write!(f, "the OS could not assemble the environment ({detail})")
}
}
}
}
pub(crate) struct OwnerEnv {
vars: Vec<(OsString, OsString)>,
}
impl OwnerEnv {
#[must_use]
pub(crate) fn new(vars: Vec<(OsString, OsString)>) -> Self {
Self { vars }
}
#[must_use]
pub(crate) fn vars(&self) -> &[(OsString, OsString)] {
&self.vars
}
}
static SNAPSHOT: ArcSwapOption<OwnerEnv> = ArcSwapOption::const_empty();
#[must_use]
pub(crate) fn snapshot() -> Option<Arc<OwnerEnv>> {
SNAPSHOT.load_full()
}
fn publish(env: Arc<OwnerEnv>) {
SNAPSHOT.store(Some(env));
}
#[cfg(all(test, unix))]
pub(crate) fn set_snapshot(env: Option<OwnerEnv>) {
SNAPSHOT.store(env.map(Arc::new));
}
pub async fn run_reader_loop() {
let mut failures: u32 = 0;
loop {
if crate::shutdown::aborting() {
return;
}
let started = Instant::now();
#[cfg(unix)]
let result = read_unix().await;
#[cfg(windows)]
let result = read_windows();
match result {
Ok(env) => {
let names = env.vars().len();
let path_entries = path_entry_count(env.vars());
let elapsed_ms = started.elapsed().as_millis();
let pairs = crate::tools::shell::agent_env_pairs_from(&env);
crate::tools::shell::grep_engine::invalidate_verdict_for(&pairs);
let env = Arc::new(env);
publish(Arc::clone(&env));
failures = 0;
crate::util::owner_path::sync(env).await;
#[cfg(unix)]
let search = if refresh_grep_parity(pairs).await {
"verdict measured for the environment in effect"
} else {
"no verdict yet (the real search runs)"
};
#[cfg(windows)]
let search = "not applicable (the served path is the only search)";
tracing::debug!(
names,
path_entries,
elapsed_ms,
search,
"owner environment read"
);
}
Err(reason) => {
failures = failures.saturating_add(1);
let retry_secs = retry_delay(failures).as_secs();
tracing::warn!(
consecutive_failures = failures,
retry_secs,
reason = %reason,
"owner environment read failed"
);
refresh_grep_parity(crate::tools::shell::agent_env_pairs()).await;
}
}
if !crate::shutdown::sleep_or_shutdown_or_drain(retry_delay(failures)).await {
return;
}
}
}
async fn refresh_grep_parity(env: Vec<(OsString, OsString)>) -> bool {
if cfg!(windows) {
return false;
}
tokio::task::spawn_blocking(move || {
crate::tools::shell::grep_engine::refresh_if_environment_changed(&env)
})
.await
.unwrap_or(false)
}
#[must_use]
fn retry_delay(failures: u32) -> Duration {
let doublings = failures.saturating_sub(1).min(3);
let secs = READ_INTERVAL.as_secs() << doublings;
Duration::from_secs(secs.min(READ_RETRY_MAX.as_secs()))
}
#[must_use]
fn path_entry_count(vars: &[(OsString, OsString)]) -> usize {
vars.iter()
.find(|(name, _)| name.eq_ignore_ascii_case("PATH"))
.map_or(0, |(_, value)| {
std::env::split_paths(value)
.filter(|entry| !entry.as_os_str().is_empty())
.count()
})
}
#[must_use]
fn start_marker(nonce: &str) -> String {
format!("\u{1e}{nonce}\u{1f}")
}
#[must_use]
fn end_marker(nonce: &str) -> String {
format!("\u{1e}{nonce}\u{1d}")
}
pub fn dump_environment(args: &[String]) -> i32 {
use std::io::Write as _;
let nonce = args.first().map_or("", String::as_str);
let mut out = std::io::BufWriter::new(std::io::stdout().lock());
let _ = out.write_all(start_marker(nonce).as_bytes());
#[cfg(unix)]
{
use std::os::unix::ffi::OsStrExt as _;
for (name, value) in std::env::vars_os() {
let _ = out.write_all(name.as_bytes());
let _ = out.write_all(b"=");
let _ = out.write_all(value.as_bytes());
let _ = out.write_all(b"\0");
}
}
#[cfg(not(unix))]
{
for (name, value) in std::env::vars_os() {
let _ = out.write_all(name.to_string_lossy().as_bytes());
let _ = out.write_all(b"=");
let _ = out.write_all(value.to_string_lossy().as_bytes());
let _ = out.write_all(b"\0");
}
}
let _ = out.write_all(end_marker(nonce).as_bytes());
let _ = out.flush();
0
}
#[cfg(unix)]
#[must_use]
fn shell_recipe(shell_basename: &str, macos: bool) -> Option<Vec<&'static str>> {
if !knows_shell(shell_basename) {
return None;
}
Some(if macos { vec!["-l", "-i"] } else { vec!["-i"] })
}
#[cfg(unix)]
fn parse_dump(bytes: &[u8], nonce: &str) -> Result<Vec<(OsString, OsString)>, ReadFailure> {
use std::os::unix::ffi::OsStrExt as _;
let start = start_marker(nonce);
let end = end_marker(nonce);
let at = bytes
.windows(start.len())
.rposition(|window| window == start.as_bytes())
.ok_or_else(|| ReadFailure::NoDump("no start marker".into()))?;
let body = &bytes[at + start.len()..];
let dump_end = body
.windows(end.len())
.position(|window| window == end.as_bytes())
.ok_or_else(|| ReadFailure::NoDump("no end marker".into()))?;
let mut vars = Vec::new();
for chunk in body[..dump_end].split(|&byte| byte == 0) {
let Some(eq) = chunk.iter().position(|&byte| byte == b'=') else {
continue;
};
let name = &chunk[..eq];
if name.is_empty() {
continue;
}
if name.iter().any(|&byte| byte < b' ') {
return Err(ReadFailure::NoDump("a mangled entry in the dump".into()));
}
vars.push((
OsStr::from_bytes(name).to_os_string(),
OsStr::from_bytes(&chunk[eq + 1..]).to_os_string(),
));
}
if vars.is_empty() {
return Err(ReadFailure::NoDump("no variables".into()));
}
Ok(vars)
}
#[cfg(unix)]
async fn read_capped(
mut reader: impl tokio::io::AsyncRead + Unpin,
cap: usize,
stop_at: Option<&[u8]>,
) -> Vec<u8> {
use tokio::io::AsyncReadExt as _;
let mut buf = Vec::new();
let mut chunk = [0u8; READ_CHUNK];
while buf.len() < cap {
match reader.read(&mut chunk).await {
Ok(0) | Err(_) => break,
Ok(n) => {
let room = cap - buf.len();
buf.extend_from_slice(&chunk[..n.min(room)]);
if let Some(marker) = stop_at
&& !marker.is_empty()
&& buf.windows(marker.len()).any(|window| window == marker)
{
break;
}
}
}
}
buf
}
#[cfg(unix)]
async fn drain(mut reader: impl tokio::io::AsyncRead + Unpin) {
use tokio::io::AsyncReadExt as _;
let mut chunk = [0u8; READ_CHUNK];
while let Ok(n) = reader.read(&mut chunk).await {
if n == 0 {
break;
}
}
}
#[cfg(unix)]
#[derive(Clone)]
struct CaptureDone(Arc<tokio::sync::Notify>);
#[cfg(unix)]
impl CaptureDone {
fn new() -> Self {
CaptureDone(Arc::new(tokio::sync::Notify::new()))
}
async fn wait(&self) {
self.0.notified().await;
}
}
#[cfg(unix)]
impl Drop for CaptureDone {
fn drop(&mut self) {
self.0.notify_one();
}
}
#[cfg(unix)]
fn kill_group(pgid: Option<u32>) {
let Some(pgid) = pgid else {
return;
};
let Ok(pgid) = libc::pid_t::try_from(pgid) else {
return;
};
let ret = unsafe { libc::kill(-pgid, libc::SIGKILL) };
if ret != 0 {
let err = std::io::Error::last_os_error();
if err.raw_os_error() != Some(libc::ESRCH) {
tracing::warn!(pgid, err = %err, "kill(-pgid) failed — a read stray may survive");
}
}
}
#[cfg(unix)]
async fn read_with_shell(
shell: &Path,
dump_cmd: &str,
nonce: &str,
cwd: Option<&Path>,
timeout: Duration,
) -> Result<OwnerEnv, ReadFailure> {
let end = end_marker(nonce);
let bytes = run_shell(shell, dump_cmd, cwd, timeout, Some(end.as_bytes())).await?;
Ok(OwnerEnv::new(parse_dump(&bytes, nonce)?))
}
#[cfg(unix)]
async fn run_shell(
shell: &Path,
cmd: &str,
cwd: Option<&Path>,
timeout: Duration,
stop_at: Option<&[u8]>,
) -> Result<Vec<u8>, ReadFailure> {
let basename = shell_basename(shell);
let recipe = shell_recipe(&basename, cfg!(target_os = "macos"))
.ok_or(ReadFailure::UnknownShell(basename))?;
let mut child = {
let mut cmd_builder = tokio::process::Command::new(shell);
cmd_builder.args(recipe);
cmd_builder.arg("-c").arg(cmd);
cmd_builder
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
if let Some(dir) = cwd {
cmd_builder.current_dir(dir);
}
unsafe {
cmd_builder.pre_exec(|| {
if libc::setsid() == -1 {
return Err(std::io::Error::last_os_error());
}
Ok(())
});
}
cmd_builder
.spawn()
.map_err(|e| ReadFailure::Spawn(e.to_string()))?
};
let pgid = child.id();
let Some(stdout) = child.stdout.take() else {
return Err(ReadFailure::Spawn("no stdout pipe".into()));
};
let Some(stderr) = child.stderr.take() else {
return Err(ReadFailure::Spawn("no stderr pipe".into()));
};
let capture_done = CaptureDone::new();
tokio::spawn({
let capture_done = capture_done.clone();
async move {
let draining = tokio::spawn(drain(stderr));
capture_done.wait().await;
if tokio::time::timeout(timeout, child.wait()).await.is_err() {
kill_group(pgid);
let _ = tokio::time::timeout(timeout, child.wait()).await;
}
draining.abort();
}
});
let captured =
tokio::time::timeout(timeout, read_capped(stdout, DUMP_STDOUT_CAP, stop_at)).await;
if captured.is_err() {
kill_group(pgid);
}
drop(capture_done);
captured.map_err(|_elapsed| ReadFailure::Timeout)
}
#[cfg(unix)]
#[must_use]
pub(crate) fn shell_basename(shell: &Path) -> String {
shell
.file_name()
.map_or_else(String::new, |name| name.to_string_lossy().into_owned())
}
#[cfg(unix)]
#[must_use]
fn new_nonce() -> String {
format!("{:016x}", rand::random::<u64>())
}
#[cfg(unix)]
#[must_use]
fn dump_command(exe: &Path, nonce: &str) -> String {
format!(
"{} __env-dump {nonce}",
crate::tools::path::shell_quote(&exe.to_string_lossy())
)
}
#[cfg(unix)]
async fn read_unix() -> Result<OwnerEnv, ReadFailure> {
let (shell, home) = owner_shell_and_home();
let shell = shell.ok_or(ReadFailure::NoShell)?;
let exe = std::env::current_exe().map_err(|e| ReadFailure::Spawn(e.to_string()))?;
let nonce = new_nonce();
let dump_cmd = dump_command(&exe, &nonce);
read_with_shell(&shell, &dump_cmd, &nonce, home.as_deref(), READ_TIMEOUT).await
}
#[cfg(unix)]
#[must_use]
pub(crate) fn owner_shell_and_home() -> (Option<PathBuf>, Option<PathBuf>) {
let (recorded_shell, recorded_home) = passwd_record().unwrap_or((None, None));
let shell = recorded_shell.or_else(|| {
std::env::var_os("SHELL")
.filter(|shell| !shell.is_empty())
.map(PathBuf::from)
});
let home = std::env::var_os("HOME")
.filter(|home| !home.is_empty())
.map(PathBuf::from)
.filter(|home| home.is_dir())
.or_else(|| recorded_home.filter(|home| home.is_dir()));
(shell, home)
}
#[cfg(unix)]
#[must_use]
fn passwd_record() -> Option<(Option<PathBuf>, Option<PathBuf>)> {
let mut buf = vec![0; 16 * 1024];
let mut entry: libc::passwd = unsafe { std::mem::zeroed() };
let mut found: *mut libc::passwd = std::ptr::null_mut();
let ret = unsafe {
libc::getpwuid_r(
libc::geteuid(),
&raw mut entry,
buf.as_mut_ptr(),
buf.len(),
&raw mut found,
)
};
if ret != 0 || found.is_null() {
return None;
}
Some(unsafe { (c_path(entry.pw_shell), c_path(entry.pw_dir)) })
}
#[cfg(unix)]
#[must_use]
unsafe fn c_path(ptr: *const libc::c_char) -> Option<PathBuf> {
use std::os::unix::ffi::OsStrExt as _;
if ptr.is_null() {
return None;
}
let bytes = unsafe { std::ffi::CStr::from_ptr(ptr) }.to_bytes();
(!bytes.is_empty()).then(|| PathBuf::from(OsStr::from_bytes(bytes)))
}
#[cfg(windows)]
fn read_windows() -> Result<OwnerEnv, ReadFailure> {
use windows_sys::Win32::Foundation::{CloseHandle, HANDLE};
use windows_sys::Win32::Security::{TOKEN_DUPLICATE, TOKEN_QUERY};
use windows_sys::Win32::System::Environment::{
CreateEnvironmentBlock, DestroyEnvironmentBlock,
};
use windows_sys::Win32::System::Threading::{GetCurrentProcess, OpenProcessToken};
let mut token: HANDLE = 0;
let process = unsafe { GetCurrentProcess() };
if unsafe { OpenProcessToken(process, TOKEN_QUERY | TOKEN_DUPLICATE, &raw mut token) } == 0 {
return Err(ReadFailure::AssemblyFailed("OpenProcessToken".into()));
}
let mut block: *mut std::ffi::c_void = std::ptr::null_mut();
let created = unsafe { CreateEnvironmentBlock(&raw mut block, token, 0) };
unsafe { CloseHandle(token) };
if created == 0 || block.is_null() {
return Err(ReadFailure::AssemblyFailed("CreateEnvironmentBlock".into()));
}
let vars = unsafe { walk_block(block.cast()) };
unsafe { DestroyEnvironmentBlock(block) };
if vars.is_empty() {
return Err(ReadFailure::AssemblyFailed("empty block".into()));
}
Ok(OwnerEnv::new(vars))
}
#[cfg(windows)]
unsafe fn walk_block(block: *const u16) -> Vec<(OsString, OsString)> {
use std::os::windows::ffi::OsStringExt as _;
let mut vars = Vec::new();
let mut cursor = block;
while unsafe { *cursor } != 0 {
let mut len = 0usize;
while unsafe { *cursor.add(len) } != 0 {
len += 1;
}
let entry = unsafe { std::slice::from_raw_parts(cursor, len) };
if let Some(eq) = entry.iter().position(|&unit| unit == u16::from(b'='))
&& eq > 0
{
vars.push((
OsString::from_wide(&entry[..eq]),
OsString::from_wide(&entry[eq + 1..]),
));
}
cursor = unsafe { cursor.add(len + 1) };
}
vars
}
#[cfg(all(test, unix))]
mod tests {
use super::*;
use std::os::unix::ffi::OsStrExt as _;
#[test]
fn parse_dump_round_trips_exotic_bytes() {
let nonce = "0123456789abcdef";
let invalid = OsStr::from_bytes(b"b\xffd");
let mut bytes = Vec::new();
bytes.extend_from_slice(start_marker(nonce).as_bytes());
for (name, value) in [
("PLAIN", OsStr::new("value")),
("EQUALS", OsStr::new("a=b=c")),
("NEWLINE", OsStr::new("line one\nline two")),
("SPACE", OsStr::new("a b")),
("INVALID", invalid),
] {
bytes.extend_from_slice(name.as_bytes());
bytes.push(b'=');
bytes.extend_from_slice(value.as_bytes());
bytes.push(0);
}
bytes.extend_from_slice(end_marker(nonce).as_bytes());
let vars = parse_dump(&bytes, nonce).expect("the round trip parses");
assert_eq!(
vars,
vec![
(OsString::from("PLAIN"), OsString::from("value")),
(OsString::from("EQUALS"), OsString::from("a=b=c")),
(
OsString::from("NEWLINE"),
OsString::from("line one\nline two")
),
(OsString::from("SPACE"), OsString::from("a b")),
(OsString::from("INVALID"), invalid.to_os_string()),
]
);
}
#[test]
fn parse_dump_ignores_output_around_the_dump() {
let nonce = "feedfacefeedface";
let mut bytes = Vec::new();
bytes.extend_from_slice(b"nvm: warning\n");
bytes.extend_from_slice(start_marker(nonce).as_bytes());
bytes.extend_from_slice(b"NAME=value");
bytes.push(0);
bytes.extend_from_slice(b"OTHER=x");
bytes.push(0);
bytes.extend_from_slice(end_marker(nonce).as_bytes());
bytes.extend_from_slice(b"stray noise\n");
let vars = parse_dump(&bytes, nonce).expect("stray output around the dump is tolerated");
assert_eq!(
vars,
vec![
(OsString::from("NAME"), OsString::from("value")),
(OsString::from("OTHER"), OsString::from("x")),
]
);
}
#[test]
fn parse_dump_refuses_an_interleaved_write() {
let nonce = "deadbeefdeadbeef";
let mut bytes = Vec::new();
bytes.extend_from_slice(start_marker(nonce).as_bytes());
bytes.extend_from_slice(b"NAME=value");
bytes.push(0);
bytes.extend_from_slice(b"stray noise\n");
bytes.extend_from_slice(b"OTHER=x");
bytes.push(0);
bytes.extend_from_slice(end_marker(nonce).as_bytes());
assert!(matches!(
parse_dump(&bytes, nonce),
Err(ReadFailure::NoDump(_))
));
}
#[test]
fn parse_dump_requires_markers_and_entries() {
let nonce = "cafebabecafebabe";
assert!(matches!(
parse_dump(b"just some output\n", nonce),
Err(ReadFailure::NoDump(_))
));
let mut unterminated = Vec::new();
unterminated.extend_from_slice(start_marker(nonce).as_bytes());
unterminated.extend_from_slice(b"NAME=value\0");
assert!(matches!(
parse_dump(&unterminated, nonce),
Err(ReadFailure::NoDump(_))
));
let empty = format!("{}{}", start_marker(nonce), end_marker(nonce));
assert!(matches!(
parse_dump(empty.as_bytes(), nonce),
Err(ReadFailure::NoDump(_))
));
}
#[tokio::test]
async fn read_with_shell_parses_a_synthetic_dump() {
const NONCE: &str = "0123456789abcdef";
let dump = format!("printf '\\036{NONCE}\\037NAME=VALUE\\000\\036{NONCE}\\035'");
let env = read_with_shell(
Path::new("/bin/sh"),
&dump,
NONCE,
None,
Duration::from_secs(10),
)
.await
.expect("the synthetic dump parses");
assert_eq!(
env.vars(),
&[(OsString::from("NAME"), OsString::from("VALUE"))]
);
}
#[tokio::test]
async fn read_with_shell_reports_no_dump_when_the_shell_is_silent() {
let failure = read_with_shell(
Path::new("/bin/sh"),
":",
"0123456789abcdef",
None,
Duration::from_secs(10),
)
.await
.err()
.expect("a silent shell has no dump");
assert!(matches!(failure, ReadFailure::NoDump(_)), "{failure}");
}
#[tokio::test]
async fn read_with_shell_times_out_on_a_hanging_shell() {
let failure = read_with_shell(
Path::new("/bin/sh"),
"sleep 30",
"0123456789abcdef",
None,
Duration::from_millis(250),
)
.await
.err()
.expect("a hanging shell times out");
assert_eq!(failure, ReadFailure::Timeout);
}
#[tokio::test]
async fn read_with_shell_refuses_an_unknown_shell() {
let failure = read_with_shell(
Path::new("/bin/pwsh"),
":",
"0123456789abcdef",
None,
Duration::from_secs(10),
)
.await
.err()
.expect("an unknown shell is refused");
assert_eq!(failure, ReadFailure::UnknownShell("pwsh".to_string()));
}
}
#[cfg(all(test, unix))]
mod evidence {
use super::*;
use std::collections::BTreeSet;
fn listed_names(bytes: &[u8]) -> BTreeSet<String> {
String::from_utf8_lossy(bytes)
.lines()
.filter_map(|line| line.split_once('=').map(|(name, _)| name.to_string()))
.collect()
}
fn env_value(bytes: &[u8], name: &str) -> Option<Vec<u8>> {
let prefix = format!("{name}=");
bytes
.split(|&b| b == b'\n')
.find_map(|line| line.strip_prefix(prefix.as_bytes()))
.map(<[u8]>::to_vec)
}
fn shell_bookkeeping(name: &str) -> bool {
matches!(name, "HOME" | "LOGNAME" | "OLDPWD" | "PWD" | "SHLVL" | "_")
}
fn recovered_names(env: &OwnerEnv) -> BTreeSet<String> {
env.vars()
.iter()
.map(|(name, _)| name.to_string_lossy().into_owned())
.filter(|name| !name.is_empty())
.collect()
}
fn built_binary() -> Option<PathBuf> {
let harness = std::env::current_exe().ok()?;
let candidate = harness.parent()?.parent()?.join("mahbot");
candidate.is_file().then_some(candidate)
}
async fn marked_env(shell: &Path, home: Option<&Path>) -> Vec<u8> {
let nonce = new_nonce();
let start = start_marker(&nonce);
let end = end_marker(&nonce);
let cmd = format!(
"printf '%s\\n' {}; env; printf '%s\\n' {}",
crate::tools::path::shell_quote(&start),
crate::tools::path::shell_quote(&end)
);
let bytes = run_shell(shell, &cmd, home, READ_TIMEOUT, Some(end.as_bytes()))
.await
.expect("the same shell lists its environment");
let text = String::from_utf8_lossy(&bytes).into_owned();
match (text.find(&start), text.find(&end)) {
(Some(from), Some(to)) if from < to => text.as_bytes()[from + start.len()..to].to_vec(),
_ => Vec::new(),
}
}
#[serial_test::serial(shell_env)]
#[tokio::test]
#[ignore = "manual evidence run: needs `cargo build --bin mahbot` and the owner's shell"]
async fn owner_environment_end_to_end() {
let Some(binary) = built_binary() else {
println!("EVIDENCE SKIPPED: build the binary first (cargo build --bin mahbot)");
return;
};
let (shell, home) = owner_shell_and_home();
let shell = shell.expect("the account records a shell");
println!(
"shell: {} | home resolved: {}",
shell_basename(&shell),
home.is_some()
);
let nonce = new_nonce();
let started = Instant::now();
let recovered = read_with_shell(
&shell,
&dump_command(&binary, &nonce),
&nonce,
home.as_deref(),
READ_TIMEOUT,
)
.await
.expect("the read succeeds on this host");
let elapsed_ms = started.elapsed().as_millis();
let read_names = recovered_names(&recovered);
let listed = marked_env(&shell, home.as_deref()).await;
let shell_names = listed_names(&listed);
let missing: Vec<&String> = shell_names.difference(&read_names).collect();
let extra: Vec<&String> = read_names.difference(&shell_names).collect();
println!(
"read: {} names, {} PATH entries, {elapsed_ms} ms",
read_names.len(),
path_entry_count(recovered.vars())
);
println!(
"login shell: {} names | names the read missed: {} | names the read added: {}",
shell_names.len(),
missing.len(),
extra.len()
);
assert!(
!read_names.is_empty(),
"a read that recovered nothing is not a read"
);
let marker = "/nonexistent-owner-value-evidence";
let temp_vars = crate::temp::shell_temp_vars();
let mut replaced: Vec<OsString> = vec![OsString::from("HOME")];
replaced.extend(temp_vars.iter().map(|(name, _)| OsString::from(name)));
let mut vars = recovered.vars().to_vec();
vars.retain(|(name, _)| !replaced.contains(name));
for name in &replaced {
vars.push((name.clone(), OsString::from(marker)));
}
set_snapshot(Some(OwnerEnv::new(vars)));
let mut command = {
let mut cmd = tokio::process::Command::new("sh");
cmd.arg("-c").arg("env");
crate::tools::shell::apply_agent_env(&mut cmd);
cmd
};
let command_started = Instant::now();
let command_out = command.output().await.expect("the command runs");
let command_ms = command_started.elapsed().as_millis();
let command_names = listed_names(&command_out.stdout);
let command_home = env_value(&command_out.stdout, "HOME");
let imposed_home = crate::tools::shell::agent_env_pairs()
.into_iter()
.find_map(|(name, value)| (name == "HOME").then_some(value));
let imposed_home = imposed_home
.as_deref()
.map(std::ffi::OsStr::as_encoded_bytes);
let temp_names: Vec<(String, bool)> = temp_vars
.iter()
.map(|(name, pinned)| {
let value = env_value(&command_out.stdout, name);
let wins = value.as_deref() == Some(pinned.as_bytes())
&& value.as_deref() != Some(marker.as_bytes());
(name.clone(), wins)
})
.collect();
set_snapshot(None);
let added: Vec<&String> = command_names.difference(&read_names).collect();
let lost: Vec<&String> = read_names.difference(&command_names).collect();
let (bookkeeping, lost): (Vec<&String>, Vec<&String>) = lost
.into_iter()
.partition(|name| shell_bookkeeping(name.as_str()));
let (bookkeeping_added, added): (Vec<&String>, Vec<&String>) = added
.into_iter()
.partition(|name| shell_bookkeeping(name.as_str()));
let rewritten: Vec<&String> = bookkeeping.into_iter().chain(bookkeeping_added).collect();
println!(
"agent command: {} names, started in {command_ms} ms",
command_names.len()
);
println!("names the command's own shell rewrote: {rewritten:?}");
println!("owner names it lost: {lost:?} | names it gained: {added:?}");
println!("temp variables: {temp_names:?}");
println!(
"HOME: the command's value is the product's, not the read's own: {}",
command_home.as_deref() == imposed_home
&& command_home.as_deref() != Some(marker.as_bytes())
);
assert!(
lost.is_empty(),
"the owner's own names must reach the command: {lost:?}"
);
assert!(
added.is_empty(),
"nothing but the pinned values may be added: {added:?}"
);
assert!(
command_home.as_deref() == imposed_home,
"the data-location home must win over the read's own"
);
assert!(
temp_names.iter().all(|(_, pinned)| *pinned),
"the pinned temp names must win: {temp_names:?}"
);
let fallback: BTreeSet<String> = crate::tools::shell::agent_env_pairs()
.iter()
.map(|(name, _)| name.to_string_lossy().into_owned())
.collect();
let fallback_path_entries = path_entry_count(&crate::tools::shell::agent_env_pairs());
println!(
"fallback (no read yet): {} names, {fallback_path_entries} PATH entries, HOME pinned: {}, TERM pinned: {}",
fallback.len(),
fallback.contains("HOME"),
fallback.contains("TERM")
);
assert!(
fallback.contains("HOME") && fallback.contains("PATH"),
"the fallback is the reduced environment: {fallback:?}"
);
assert!(
!fallback.contains("OLDPWD"),
"the fallback must not carry the owner's own environment: {fallback:?}"
);
let failure_of = |result: Result<OwnerEnv, ReadFailure>| match result {
Ok(_) => panic!("this read was expected to fail"),
Err(failure) => failure,
};
let nonce = "0123456789abcdef";
let failures = [
(
"shell the product does not know",
failure_of(
read_with_shell(
Path::new("/bin/pwsh"),
":",
nonce,
None,
Duration::from_secs(1),
)
.await,
),
),
(
"shell that produces no dump",
failure_of(
read_with_shell(
Path::new("/bin/sh"),
":",
nonce,
None,
Duration::from_secs(1),
)
.await,
),
),
(
"shell that hangs",
failure_of(
read_with_shell(
Path::new("/bin/sh"),
"sleep 30",
nonce,
None,
Duration::from_millis(250),
)
.await,
),
),
];
for (case, failure) in failures {
let recorded = failure.to_string();
println!("{case}: recorded as \"{recorded}\"");
assert!(
!recorded.is_empty() && !recorded.contains('='),
"a recorded reason names no value: {recorded}"
);
}
let delays: Vec<u64> = (1..=5).map(|n| retry_delay(n).as_secs()).collect();
println!("retry delay after 1..5 consecutive failures: {delays:?} seconds");
assert!(
delays[0] == READ_INTERVAL.as_secs()
&& delays.windows(2).all(|pair| pair[0] <= pair[1]),
"a repeated failure is retried no more often than the interval: {delays:?}"
);
}
}