use std::io::{BufReader, BufWriter};
use std::path::Path;
use std::process::{Child, Command, Stdio};
use std::sync::mpsc;
use std::time::{Duration, Instant};
use tau_proto::{
ClientKind, Disconnect, Event, EventName, EventSelector, HarnessInputMessage,
HarnessOutputMessage, Hello, PROTOCOL_VERSION, PeerInputReader, PeerOutputWriter, SessionId,
Subscribe,
};
mod support;
use support::bounded_runtime_tempdir;
struct InitialUiChild {
process: Child,
}
impl InitialUiChild {
fn stop(&mut self) {
if !matches!(self.process.try_wait(), Ok(Some(_))) {
let _ = self.process.kill();
let _ = self.process.wait();
}
}
}
impl Drop for InitialUiChild {
fn drop(&mut self) {
self.stop();
}
}
fn wait_for_clean_exit(child: &mut Child) {
let deadline = Instant::now() + Duration::from_secs(2);
loop {
if let Some(status) = child.try_wait().expect("query child exit") {
assert!(status.success(), "initial-UI daemon failed: {status}");
return;
}
if deadline <= Instant::now() {
let _ = child.kill();
let _ = child.wait();
panic!("initial-UI daemon did not shut down after disconnect");
}
std::thread::sleep(Duration::from_millis(10));
}
}
fn initial_ui_stdio_command(
tau_bin: &str,
temp: &tempfile::TempDir,
config_home: &Path,
state_home: &Path,
runtime_dir: &Path,
) -> Command {
let home = temp.path().join("home");
let cache_home = temp.path().join("cache");
std::fs::create_dir_all(&home).expect("mkdir home");
std::fs::create_dir_all(&cache_home).expect("mkdir cache");
let mut command = Command::new(tau_bin);
command
.env_clear()
.arg("component")
.arg("harness")
.arg("--initial-ui-stdio")
.env("HOME", home)
.env("XDG_CONFIG_HOME", config_home)
.env("XDG_STATE_HOME", state_home)
.env("XDG_CACHE_HOME", cache_home)
.env("XDG_RUNTIME_DIR", runtime_dir)
.env("LANG", "C.UTF-8")
.env("TERM", "xterm-256color")
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::null());
command
}
fn write_initial_ui_handshake(writer: &mut PeerOutputWriter<BufWriter<std::process::ChildStdin>>) {
writer
.write_message(&HarnessInputMessage::Hello(Hello {
declaration_inspection: false,
protocol_version: PROTOCOL_VERSION,
client_name: "tau-cli".parse().expect("valid terminal UI client name"),
client_kind: ClientKind::Ui,
expected_session_id: None,
capabilities: Vec::new(),
}))
.expect("write hello");
writer
.write_message(&HarnessInputMessage::Subscribe(Subscribe {
historical_selectors: vec![EventSelector::Exact(EventName::SESSION_REPLAY_COMPLETE)],
live_selectors: Vec::new(),
}))
.expect("write subscribe");
writer.flush().expect("flush handshake");
}
#[test]
fn initial_ui_introduction_notice_requires_conversational_launch() {
for (eligible, expected) in [(true, 1), (false, 0)] {
let temp = tempfile::tempdir().expect("tempdir");
let config_home = temp.path().join("config");
let state_home = temp.path().join("state");
let runtime_temp = bounded_runtime_tempdir();
let runtime_dir = runtime_temp.path().to_path_buf();
std::fs::create_dir_all(config_home.join("tau")).expect("mkdir config");
std::fs::create_dir_all(&state_home).expect("mkdir state");
std::fs::create_dir_all(&runtime_dir).expect("mkdir runtime");
std::fs::write(
config_home.join("tau/harness.yaml"),
"extensions:\n provider-builtin:\n enable: false\n core-shell:\n enable: false\n",
)
.expect("write minimal harness config");
let tau_bin = std::env::var("CARGO_BIN_EXE_tau").expect("CARGO_BIN_EXE_tau");
let mut command =
initial_ui_stdio_command(&tau_bin, &temp, &config_home, &state_home, &runtime_dir);
command.env("TAU_SESSION_ID", format!("introduction-{eligible}"));
if eligible {
command.env(tau_harness::INITIAL_UI_INTRODUCTION_NOTICE_ENV, "1");
}
let mut child = command.spawn().expect("spawn harness");
let mut writer = PeerOutputWriter::new(BufWriter::new(child.stdin.take().expect("stdin")));
write_initial_ui_handshake(&mut writer);
let stdout = child.stdout.take().expect("stdout");
let (sender, receiver) = mpsc::channel();
let reader_thread = std::thread::spawn(move || {
let mut reader = PeerInputReader::new(BufReader::new(stdout));
while let Ok(Some(message)) = reader.read_message() {
if sender.send(message).is_err() {
break;
}
}
});
let deadline = Instant::now() + Duration::from_secs(10);
let mut introductions = 0;
let mut startup_complete = false;
while Instant::now() < deadline && (!startup_complete || introductions < expected) {
let message = match receiver.recv_timeout(Duration::from_millis(100)) {
Ok(message) => message,
Err(mpsc::RecvTimeoutError::Timeout) => continue,
Err(mpsc::RecvTimeoutError::Disconnected) => {
panic!("initial UI reader exited before startup completed")
}
};
match message {
HarnessOutputMessage::Deliver(delivery) => match delivery.event() {
Event::HarnessNotice(notice)
if notice.kind == tau_proto::notice_kind::HARNESS_INTRODUCTION =>
{
introductions += 1;
}
Event::SessionReplayComplete(_) => startup_complete = true,
_ => {}
},
HarnessOutputMessage::Disconnect(disconnect) => {
panic!(
"initial UI disconnected before startup completed: {:?}",
disconnect.reason
);
}
_ => {}
}
}
assert!(startup_complete, "spawned harness did not finish startup");
std::thread::sleep(Duration::from_millis(100));
while let Ok(message) = receiver.try_recv() {
if matches!(
message,
HarnessOutputMessage::Deliver(ref delivery)
if matches!(
delivery.event(),
Event::HarnessNotice(notice)
if notice.kind == tau_proto::notice_kind::HARNESS_INTRODUCTION
)
) {
introductions += 1;
}
}
assert_eq!(introductions, expected);
writer
.write_message(&HarnessInputMessage::Disconnect(Disconnect {
reason: Some("test complete".to_owned()),
}))
.expect("write disconnect");
writer.flush().expect("flush disconnect");
drop(writer);
wait_for_clean_exit(&mut child);
reader_thread.join().expect("reader thread");
}
}
#[test]
fn initial_ui_stdio_startup_error_reaches_child_stdout() {
let temp = tempfile::tempdir().expect("tempdir");
let config_home = temp.path().join("config");
let state_home = temp.path().join("state");
let runtime_temp = bounded_runtime_tempdir();
let runtime_dir = runtime_temp.path().to_path_buf();
let tau_config_dir = config_home.join("tau");
std::fs::create_dir_all(&tau_config_dir).expect("mkdir config");
std::fs::create_dir_all(&state_home).expect("mkdir state");
std::fs::create_dir_all(&runtime_dir).expect("mkdir runtime");
std::fs::write(
tau_config_dir.join("harness.yaml"),
r#"
extensions:
startup-secret-test:
command: [tau]
secrets:
missing_token: {}
"#,
)
.expect("write harness config");
let tau_bin = std::env::var("CARGO_BIN_EXE_tau").expect("CARGO_BIN_EXE_tau");
let mut child =
initial_ui_stdio_command(&tau_bin, &temp, &config_home, &state_home, &runtime_dir)
.spawn()
.expect("spawn tau harness");
let _stdin = child.stdin.take();
let stdout = child.stdout.take().expect("stdout");
let (sender, receiver) = mpsc::channel();
let reader_thread = std::thread::spawn(move || {
let mut reader = PeerInputReader::new(BufReader::new(stdout));
let _ = sender.send(reader.read_message());
});
let message = match receiver.recv_timeout(Duration::from_secs(10)) {
Ok(message) => message
.expect("read startup disconnect")
.expect("startup disconnect"),
Err(mpsc::RecvTimeoutError::Timeout) => {
let _ = child.kill();
let _ = child.wait();
let _ = reader_thread.join();
panic!("timed out waiting for startup disconnect on child stdout");
}
Err(mpsc::RecvTimeoutError::Disconnected) => {
panic!("startup disconnect reader thread exited without reporting a result");
}
};
reader_thread.join().expect("reader thread");
let HarnessOutputMessage::Disconnect(disconnect) = message else {
panic!("expected disconnect frame");
};
let reason = disconnect.reason.expect("disconnect reason");
assert!(reason.contains("harness startup failed"));
assert!(reason.contains("missing_token"));
let status = child.wait().expect("wait child");
assert!(!status.success());
}
#[test]
fn resumed_harness_process_does_not_recreate_deleted_session() {
let temp = tempfile::tempdir().expect("tempdir");
let config_home = temp.path().join("config");
let state_home = temp.path().join("state");
let runtime_temp = bounded_runtime_tempdir();
let runtime_dir = runtime_temp.path().to_path_buf();
let tau_config_dir = config_home.join("tau");
let session_dir = state_home.join("tau/sessions/deleted-session");
std::fs::create_dir_all(&tau_config_dir).expect("mkdir config");
std::fs::create_dir_all(&session_dir).expect("mkdir selected session");
std::fs::create_dir_all(&runtime_dir).expect("mkdir runtime");
std::fs::write(session_dir.join("lock"), b"").expect("write session lock");
std::fs::write(
session_dir.join("meta.json"),
br#"{"created_at":1,"last_touched":1}"#,
)
.expect("write session metadata");
std::fs::write(session_dir.join("events.cbor"), b"").expect("write ordinary journal");
std::fs::write(session_dir.join("restore-events.cbor"), b"").expect("write restore journal");
std::fs::remove_dir_all(&session_dir).expect("delete selected session before startup lock");
let tau_bin = std::env::var("CARGO_BIN_EXE_tau").expect("CARGO_BIN_EXE_tau");
let mut child =
initial_ui_stdio_command(&tau_bin, &temp, &config_home, &state_home, &runtime_dir)
.env("TAU_SESSION_ID", "deleted-session")
.env("TAU_SESSION_STATUS", "resumed")
.spawn()
.expect("spawn resumed tau harness");
let _stdin = child.stdin.take();
let stdout = child.stdout.take().expect("stdout");
let (sender, receiver) = mpsc::channel();
let reader_thread = std::thread::spawn(move || {
let mut reader = PeerInputReader::new(BufReader::new(stdout));
let _ = sender.send(reader.read_message());
});
let message = match receiver.recv_timeout(Duration::from_secs(10)) {
Ok(message) => message
.expect("read startup disconnect")
.expect("startup disconnect"),
Err(error) => {
let _ = child.kill();
let _ = child.wait();
let _ = reader_thread.join();
panic!("failed waiting for resume deletion result: {error}");
}
};
reader_thread.join().expect("reader thread");
let HarnessOutputMessage::Disconnect(disconnect) = message else {
panic!("expected disconnect frame");
};
assert!(disconnect.reason.as_deref().is_some_and(
|reason| reason.contains("deleted-session") && reason.contains("no longer exists")
));
assert!(!session_dir.exists());
let status = child.wait().expect("wait child");
assert!(!status.success());
}
#[test]
fn resumed_configured_harness_process_creates_relay_log() {
let temp = tempfile::tempdir().expect("tempdir");
let config_home = temp.path().join("config");
let state_home = temp.path().join("state");
let runtime_temp = bounded_runtime_tempdir();
let runtime_dir = runtime_temp.path().to_path_buf();
let tau_config_dir = config_home.join("tau");
let session_dir = state_home.join("tau/sessions/resumed-session");
std::fs::create_dir_all(&tau_config_dir).expect("mkdir config");
std::fs::create_dir_all(&session_dir).expect("mkdir selected session");
std::fs::create_dir_all(&runtime_dir).expect("mkdir runtime");
std::fs::write(session_dir.join("lock"), b"").expect("write session lock");
std::fs::write(
session_dir.join("meta.json"),
br#"{"created_at":1,"last_touched":1}"#,
)
.expect("write session metadata");
std::fs::write(session_dir.join("events.cbor"), b"").expect("write ordinary journal");
std::fs::write(session_dir.join("restore-events.cbor"), b"").expect("write restore journal");
let tau_bin = std::env::var("CARGO_BIN_EXE_tau").expect("CARGO_BIN_EXE_tau");
let mut child = InitialUiChild {
process: initial_ui_stdio_command(&tau_bin, &temp, &config_home, &state_home, &runtime_dir)
.env("TAU_SESSION_ID", "resumed-session")
.env("TAU_SESSION_STATUS", "resumed")
.spawn()
.expect("spawn resumed tau harness"),
};
let expected_session = SessionId::parse("resumed-session").expect("valid expected session id");
let mut writer =
PeerOutputWriter::new(BufWriter::new(child.process.stdin.take().expect("stdin")));
write_initial_ui_handshake(&mut writer);
let stdout = child.process.stdout.take().expect("stdout");
let (sender, receiver) = mpsc::channel();
let reader_thread = std::thread::spawn(move || {
let mut reader = PeerInputReader::new(BufReader::new(stdout));
while let Ok(Some(message)) = reader.read_message() {
if sender.send(message).is_err() {
break;
}
}
});
let deadline = Instant::now() + Duration::from_secs(10);
let startup = loop {
if Instant::now() >= deadline {
break Err("resumed harness did not complete startup".to_owned());
}
let message = match receiver.recv_timeout(Duration::from_millis(100)) {
Ok(message) => message,
Err(mpsc::RecvTimeoutError::Timeout) => continue,
Err(mpsc::RecvTimeoutError::Disconnected) => {
break Err("resumed harness disconnected before completing startup".to_owned());
}
};
match message {
HarnessOutputMessage::Deliver(delivery) => {
if let Event::SessionReplayComplete(replay) = delivery.event() {
if replay.session_id == expected_session && replay.error.is_none() {
break Ok(());
}
break Err(format!(
"resumed harness completed the wrong or failed replay: {replay:?}"
));
}
}
HarnessOutputMessage::Disconnect(disconnect) => {
break Err(format!(
"resumed harness disconnected before completing startup: {:?}",
disconnect.reason
));
}
_ => {}
}
};
if let Err(error) = startup {
drop(writer);
child.stop();
let _ = reader_thread.join();
panic!("{error}");
}
let harness_log = session_dir.join("logs/tau-harness.log");
assert!(
harness_log.is_file(),
"configured resume must create the parent relay target"
);
let mut sessions = std::fs::read_dir(state_home.join("tau/sessions"))
.expect("read sessions")
.map(|entry| entry.expect("read session").file_name())
.collect::<Vec<_>>();
sessions.sort();
assert_eq!(sessions, [std::ffi::OsString::from("resumed-session")]);
let shutdown = writer
.write_message(&HarnessInputMessage::Disconnect(Disconnect {
reason: Some("test complete".to_owned()),
}))
.map_err(|error| format!("write clean shutdown: {error}"))
.and_then(|()| {
writer
.flush()
.map_err(|error| format!("flush clean shutdown: {error}"))
});
drop(writer);
shutdown.unwrap_or_else(|error| panic!("{error}"));
wait_for_clean_exit(&mut child.process);
reader_thread.join().expect("reader thread");
}
#[test]
fn harness_process_rejects_attach_session_mismatch() {
let temp = tempfile::tempdir().expect("tempdir");
let config_home = temp.path().join("config");
let state_home = temp.path().join("state");
let runtime_temp = bounded_runtime_tempdir();
let runtime_dir = runtime_temp.path().to_path_buf();
std::fs::create_dir_all(config_home.join("tau")).expect("mkdir config");
std::fs::create_dir_all(&state_home).expect("mkdir state");
std::fs::create_dir_all(&runtime_dir).expect("mkdir runtime");
let tau_bin = std::env::var("CARGO_BIN_EXE_tau").expect("CARGO_BIN_EXE_tau");
let mut child =
initial_ui_stdio_command(&tau_bin, &temp, &config_home, &state_home, &runtime_dir)
.env("TAU_SESSION_ID", "actual-session")
.spawn()
.expect("spawn tau harness");
let mut writer = PeerOutputWriter::new(BufWriter::new(child.stdin.take().expect("stdin")));
writer
.write_message(&HarnessInputMessage::Hello(Hello {
declaration_inspection: false,
protocol_version: PROTOCOL_VERSION,
client_name: "attach-test".parse().expect("valid client name"),
client_kind: ClientKind::Ui,
expected_session_id: Some(
SessionId::parse("requested-session").expect("valid session id"),
),
capabilities: Vec::new(),
}))
.expect("write attach hello");
writer.flush().expect("flush attach hello");
let stdout = child.stdout.take().expect("stdout");
let (sender, receiver) = mpsc::channel();
let reader_thread = std::thread::spawn(move || {
let mut reader = PeerInputReader::new(BufReader::new(stdout));
let _ = sender.send(reader.read_message());
});
let message = match receiver.recv_timeout(Duration::from_secs(10)) {
Ok(message) => message
.expect("read handshake response")
.expect("handshake response"),
Err(error) => {
let _ = child.kill();
let _ = child.wait();
let _ = reader_thread.join();
panic!("failed waiting for attach mismatch result: {error}");
}
};
reader_thread.join().expect("reader thread");
let HarnessOutputMessage::Disconnect(disconnect) = message else {
panic!("expected disconnect frame");
};
let reason = disconnect.reason.expect("disconnect reason");
assert!(reason.contains("requested-session"));
assert!(reason.contains("actual-session"));
assert!(reason.contains("tau attach requested-session"));
drop(writer);
let status = child.wait().expect("wait child");
assert!(!status.success());
}
#[test]
fn late_startup_failure_does_not_emit_introduction_notice() {
let temp = tempfile::tempdir().expect("tempdir");
let config_home = temp.path().join("config");
let state_home = temp.path().join("state");
let runtime_dir = temp.path().join("runtime");
std::fs::create_dir_all(config_home.join("tau")).expect("mkdir config");
std::fs::create_dir_all(&state_home).expect("mkdir state");
std::fs::create_dir_all(&runtime_dir).expect("mkdir runtime");
std::fs::write(
config_home.join("tau/harness.yaml"),
"extensions:\n provider-builtin:\n enable: false\n core-shell:\n enable: false\n",
)
.expect("write minimal harness config");
let tau_bin = std::env::var("CARGO_BIN_EXE_tau").expect("CARGO_BIN_EXE_tau");
let instance = "0123456789abcdef";
let mut child =
initial_ui_stdio_command(&tau_bin, &temp, &config_home, &state_home, &runtime_dir)
.env("TAU_SESSION_ID", "late-startup-failure")
.env("TAU_HARNESS_INSTANCE_ID", instance)
.env(tau_harness::INITIAL_UI_INTRODUCTION_NOTICE_ENV, "1")
.spawn()
.expect("spawn harness");
let metadata = runtime_dir
.join("tau/harnesses")
.join(format!("{}-{instance}.json", child.id()));
let deadline = Instant::now() + Duration::from_secs(10);
while Instant::now() < deadline
&& !metadata
.parent()
.expect("metadata parent")
.try_exists()
.expect("check runtime directory")
{
std::thread::sleep(Duration::from_millis(10));
}
std::fs::create_dir(&metadata).expect("block metadata file replacement");
let mut writer = PeerOutputWriter::new(BufWriter::new(child.stdin.take().expect("stdin")));
let _ = writer.write_message(&HarnessInputMessage::Hello(Hello {
declaration_inspection: false,
protocol_version: PROTOCOL_VERSION,
client_name: "tau-cli".parse().expect("valid terminal UI client name"),
client_kind: ClientKind::Ui,
expected_session_id: None,
capabilities: Vec::new(),
}));
let _ = writer.flush();
let mut reader = PeerInputReader::new(BufReader::new(child.stdout.take().expect("stdout")));
let mut introductions = 0;
let disconnect = loop {
let Some(message) = reader.read_message().expect("read startup result") else {
break None;
};
match message {
HarnessOutputMessage::Deliver(delivery) => {
if matches!(
delivery.event(),
Event::HarnessNotice(notice)
if notice.kind == tau_proto::notice_kind::HARNESS_INTRODUCTION
) {
introductions += 1;
}
}
HarnessOutputMessage::Disconnect(disconnect) => break Some(disconnect),
_ => {}
}
};
assert_eq!(introductions, 0);
assert!(disconnect.is_none_or(|disconnect| {
disconnect
.reason
.as_deref()
.is_some_and(|reason| reason.contains("harness startup failed"))
}));
drop(writer);
assert!(!child.wait().expect("wait child").success());
}