use std::process::{Child, Command};
use std::sync::Mutex;
use std::time::Duration;
use futures::StreamExt as _;
use theway_transport::client::{GrpcClient, port_file_path, probe, wait_ready};
use theway_transport::proto::theway_grpc::stream_frame;
static DAEMON_E2E_LOCK: Mutex<()> = Mutex::new(());
struct DaemonGuard {
child: Child,
}
impl Drop for DaemonGuard {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
async fn spawn_daemon(dir: &std::path::Path) -> (DaemonGuard, String) {
spawn_daemon_with(dir, |cmd| cmd).await
}
async fn spawn_daemon_with(
dir: &std::path::Path,
customize: impl FnOnce(&mut Command) -> &mut Command,
) -> (DaemonGuard, String) {
let binary = env!("CARGO_BIN_EXE_thewayd");
let mut command = Command::new(binary);
customize(
command
.arg("--port")
.arg("0")
.arg("--cwd")
.arg(dir)
.env("THEWAY_DIR", dir)
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null()),
);
let child = command.spawn().expect("spawn thewayd");
let child_pid = child.id();
let addr = tokio::time::timeout(
Duration::from_secs(20),
wait_ready(Duration::from_secs(20), dir, child_pid),
)
.await
.expect("daemon never became ready")
.expect("wait_ready failed");
(DaemonGuard { child }, addr)
}
#[tokio::test]
async fn wait_ready_ignores_stale_entry_from_a_dead_daemon() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
std::fs::write(port_file_path(dir.path()), format!("1 {}", u32::MAX)).unwrap();
let (_daemon, addr) = spawn_daemon(dir.path()).await;
assert_ne!(
addr, "127.0.0.1:1",
"wait_ready accepted the stale port instead of the spawned daemon's"
);
let mut client = GrpcClient::connect(&addr).await.unwrap();
let state = client.get_snapshot().await.unwrap();
let expected_cwd =
std::fs::canonicalize(dir.path()).unwrap_or_else(|_| dir.path().to_path_buf());
assert_eq!(
state.info.as_ref().unwrap().cwd,
expected_cwd.display().to_string()
);
unsafe { std::env::remove_var("THEWAY_DIR") };
}
#[tokio::test]
async fn spawned_daemon_serves_get_snapshot_immediately() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
let (_daemon, addr) = spawn_daemon(dir.path()).await;
let mut client = GrpcClient::connect(&addr).await.unwrap();
let state = tokio::time::timeout(Duration::from_secs(5), client.get_snapshot())
.await
.expect("GetSnapshot hung after readiness")
.unwrap();
assert!(!state.session_id.is_empty(), "daemon created a session");
let expected_cwd =
std::fs::canonicalize(dir.path()).unwrap_or_else(|_| dir.path().to_path_buf());
assert_eq!(
state.info.as_ref().unwrap().cwd,
expected_cwd.display().to_string()
);
unsafe { std::env::remove_var("THEWAY_DIR") };
}
#[tokio::test]
async fn client_round_trip_against_spawned_daemon() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
let (_daemon, addr) = spawn_daemon(dir.path()).await;
let mut client = GrpcClient::connect(&addr).await.unwrap();
let state = client.get_snapshot().await.unwrap();
let session_id = state.session_id.clone();
let accepted = client
.send_message("hello from the client".into(), vec![], false)
.await
.unwrap();
assert!(accepted, "send_message accepted");
let mut stream = client.stream_events().await.unwrap();
let mut saw_snapshot = false;
for _ in 0..8 {
let frame = tokio::time::timeout(Duration::from_secs(15), stream.next())
.await
.expect("timed out waiting for stream frame")
.expect("stream ended")
.expect("stream frame error");
if let Some(stream_frame::Payload::Snapshot(snapshot)) = frame.payload {
saw_snapshot = true;
assert_eq!(snapshot.session_id, session_id);
break;
}
}
assert!(saw_snapshot, "stream carried a snapshot frame");
unsafe { std::env::remove_var("THEWAY_DIR") };
}
#[tokio::test]
async fn two_clients_both_receive_frames_from_spawned_daemon() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
let (_daemon, addr) = spawn_daemon(dir.path()).await;
let mut client_a = GrpcClient::connect(&addr).await.unwrap();
let mut client_b = GrpcClient::connect(&addr).await.unwrap();
let mut stream_a = client_a.stream_events().await.unwrap();
let mut stream_b = client_b.stream_events().await.unwrap();
client_b
.send_message("broadcast test".into(), vec![], false)
.await
.unwrap();
for (label, stream) in [("a", &mut stream_a), ("b", &mut stream_b)] {
let mut saw_snapshot = false;
for _ in 0..8 {
let frame = tokio::time::timeout(Duration::from_secs(15), stream.next())
.await
.expect("timed out")
.expect("stream ended")
.expect("stream frame error");
if let Some(stream_frame::Payload::Snapshot(_)) = frame.payload {
saw_snapshot = true;
break;
}
}
assert!(saw_snapshot, "client {label} received a snapshot frame");
}
unsafe { std::env::remove_var("THEWAY_DIR") };
}
#[tokio::test]
async fn dead_daemon_probe_fails_promptly() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
let (_daemon, addr) = spawn_daemon(dir.path()).await;
drop(_daemon); tokio::time::sleep(Duration::from_millis(200)).await;
let err = probe(&addr, Duration::from_millis(500)).await;
assert!(
err.is_err(),
"probe against a dead daemon must fail: {err:?}"
);
unsafe { std::env::remove_var("THEWAY_DIR") };
}
#[tokio::test]
async fn grpc_surface_has_health_check() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
let (_daemon, addr) = spawn_daemon(dir.path()).await;
let mut health = theway_transport::proto::health::health_client::HealthClient::connect(
format!("http://{addr}"),
)
.await
.unwrap();
let response = health
.check(theway_transport::proto::health::HealthCheckRequest {
service: String::new(),
})
.await
.unwrap()
.into_inner();
assert_eq!(
response.status,
theway_transport::proto::health::health_check_response::ServingStatus::Serving as i32
);
unsafe { std::env::remove_var("THEWAY_DIR") };
}
#[tokio::test]
async fn list_sessions_marks_current_after_spawn() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
let (_daemon, addr) = spawn_daemon(dir.path()).await;
let mut client = GrpcClient::connect(&addr).await.unwrap();
let (sessions, current) = client.list_sessions().await.unwrap();
assert!(!current.is_empty(), "daemon has a current session");
assert!(
!sessions.iter().any(|s| s.session_id == current),
"unmaterialized fresh session must not be listed yet: {sessions:?}"
);
let accepted = client
.send_message("hello".into(), vec![], false)
.await
.unwrap();
assert!(accepted, "send_message accepted");
let mut materialized = false;
for _ in 0..100 {
tokio::time::sleep(Duration::from_millis(100)).await;
let (sessions, current_after) = client.list_sessions().await.unwrap();
assert_eq!(
current_after, current,
"session id stable across materialization"
);
if sessions.iter().any(|s| s.session_id == current) {
materialized = true;
break;
}
}
assert!(
materialized,
"materialized session listed after first message"
);
unsafe { std::env::remove_var("THEWAY_DIR") };
}
#[tokio::test]
async fn path_context_round_trip_against_spawned_daemon() {
let _guard = DAEMON_E2E_LOCK.lock().unwrap();
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("THEWAY_DIR", dir.path()) };
let startup_extra = dir.path().join("extra-skills");
std::fs::create_dir_all(&startup_extra).unwrap();
let startup_extra_arg = startup_extra.clone();
let (_daemon, addr) = spawn_daemon_with(dir.path(), move |cmd| {
cmd.arg("--skills-dir").arg(&startup_extra_arg)
})
.await;
let mut client = GrpcClient::connect(&addr).await.unwrap();
let ctx = client.get_path_context().await.unwrap();
let expected_work_dir =
std::fs::canonicalize(dir.path()).unwrap_or_else(|_| dir.path().to_path_buf());
assert_eq!(ctx.work_dir, expected_work_dir.display().to_string());
assert_eq!(ctx.base, dir.path().display().to_string());
assert!(!ctx.home.is_empty(), "home resolved at startup");
assert_eq!(
ctx.skills_dirs,
vec![startup_extra.display().to_string()],
"startup --skills-dir extras served by GetPathContext"
);
let replacement = dir.path().join("replacement-skills");
std::fs::create_dir_all(&replacement).unwrap();
let accepted = client
.set_skill_dirs(&[replacement.display().to_string()])
.await
.unwrap();
assert!(accepted, "SetSkillDirs command queued");
let ctx = client.get_path_context().await.unwrap();
assert_eq!(
ctx.skills_dirs,
vec![replacement.display().to_string()],
"SetSkillDirs update visible to GetPathContext"
);
assert_eq!(ctx.work_dir, expected_work_dir.display().to_string());
assert_eq!(ctx.base, dir.path().display().to_string());
unsafe { std::env::remove_var("THEWAY_DIR") };
}