diff --git a/lib/crates/fabro-cli/tests/it/cmd/asset_cp.rs b/lib/crates/fabro-cli/tests/it/cmd/asset_cp.rs index 27f62b9dd..329d29d77 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/asset_cp.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/asset_cp.rs @@ -1,6 +1,6 @@ use fabro_test::{fabro_snapshot, test_context}; -use super::support::{read_text, setup_asset_sandbox_run, setup_completed_dry_run, text_tree}; +use super::support::{read_text, setup_asset_run, setup_completed_dry_run, text_tree}; #[test] fn help() { @@ -54,7 +54,7 @@ fn asset_cp_empty_run_reports_no_assets() { #[test] fn asset_cp_specific_path_copies_single_asset() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_asset_run(&context); let dest = context.temp_dir.join("asset-one"); let mut cmd = context.command(); cmd.args([ @@ -79,7 +79,7 @@ fn asset_cp_specific_path_copies_single_asset() { #[test] fn asset_cp_ambiguous_path_requires_node_or_retry() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_asset_run(&context); let dest = context.temp_dir.join("asset-one"); let mut cmd = context.command(); cmd.args([ @@ -101,7 +101,7 @@ fn asset_cp_ambiguous_path_requires_node_or_retry() { #[test] fn asset_cp_tree_preserves_structure() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_asset_run(&context); let dest = context.temp_dir.join("asset-tree"); let mut cmd = context.command(); cmd.args([ @@ -135,7 +135,7 @@ fn asset_cp_tree_preserves_structure() { #[test] fn asset_cp_flat_mode_rejects_filename_collisions() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_asset_run(&context); let dest = context.temp_dir.join("asset-flat"); let mut cmd = context.command(); cmd.args(["asset", "cp", &setup.run.run_id, dest.to_str().unwrap()]); diff --git a/lib/crates/fabro-cli/tests/it/cmd/asset_list.rs b/lib/crates/fabro-cli/tests/it/cmd/asset_list.rs index 32d21d600..bc5681295 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/asset_list.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/asset_list.rs @@ -1,6 +1,6 @@ use fabro_test::{fabro_snapshot, test_context}; -use super::support::{setup_asset_sandbox_run, setup_completed_dry_run}; +use super::support::{setup_asset_run, setup_completed_dry_run}; #[test] fn help() { @@ -51,7 +51,7 @@ fn asset_list_empty_run_reports_no_assets() { #[test] fn asset_list_json_outputs_entries() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_asset_run(&context); let mut filters = context.filters(); filters.push(( r"\[STORAGE_DIR\]/runs/\d{8}-\[ULID\]".to_string(), @@ -115,7 +115,7 @@ fn asset_list_json_outputs_entries() { #[test] fn asset_list_filters_by_node_and_retry() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_asset_run(&context); let mut filters = context.filters(); filters.push(( r"\[STORAGE_DIR\]/runs/\d{8}-\[ULID\]".to_string(), diff --git a/lib/crates/fabro-cli/tests/it/cmd/attach.rs b/lib/crates/fabro-cli/tests/it/cmd/attach.rs index b76056e8f..4c2911ad5 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/attach.rs @@ -1,11 +1,13 @@ -use fabro_test::{fabro_snapshot, test_context}; +use std::time::Duration; + +use fabro_test::{fabro_snapshot, run_and_format, test_context}; use serde_json::Value; use crate::support::{ compact_progress_event, example_fixture, fabro_json_snapshot, run_output_filters, }; -use super::support::{output_stdout, write_sleep_workflow}; +use super::support::{output_stdout, resolve_run, wait_for_status, write_gated_workflow}; #[test] fn help() { @@ -100,12 +102,7 @@ fn attach_replays_completed_detached_run() { #[test] fn attach_before_completion_streams_to_finished_state() { let context = test_context!(); - write_sleep_workflow( - &context.temp_dir.join("slow.fabro"), - "slow", - "Run slowly", - 2, - ); + let gate = write_gated_workflow(&context.temp_dir.join("slow.fabro"), "slow", "Run slowly"); let mut run_cmd = context.command(); run_cmd.current_dir(&context.temp_dir); @@ -128,25 +125,32 @@ fn attach_before_completion_streams_to_finished_state() { String::from_utf8_lossy(&run_output.stderr) ); let run_id = output_stdout(&run_output).trim().to_string(); + let run = resolve_run(&context, &run_id); + wait_for_status(&run.run_dir, &["running"]); let mut filters = context.filters(); filters.push(( r"\b\d+(\.\d+)?(ms|s)\b".to_string(), "[DURATION]".to_string(), )); + let release_gate = std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(300)); + gate.release(); + }); let mut attach_cmd = context.command(); attach_cmd.current_dir(&context.temp_dir); attach_cmd.args(["attach", &run_id]); + let (snapshot, _output) = run_and_format(&mut attach_cmd, &filters); + release_gate.join().expect("gate releaser should join"); + wait_for_status(&run.run_dir, &["succeeded"]); - fabro_snapshot!(filters, attach_cmd, @" + insta::assert_snapshot!(snapshot, @" success: true exit_code: 0 ----- stdout ----- ----- stderr ----- Sandbox: local (ready in [TIME]) ✓ start [DURATION] - ✓ wait [DURATION] - ✓ exit [DURATION] "); } diff --git a/lib/crates/fabro-cli/tests/it/cmd/sandbox_cp.rs b/lib/crates/fabro-cli/tests/it/cmd/sandbox_cp.rs index a06e03cd5..4581604d6 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/sandbox_cp.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/sandbox_cp.rs @@ -1,6 +1,6 @@ use fabro_test::{fabro_snapshot, test_context}; -use super::support::{read_text, setup_asset_sandbox_run, setup_created_dry_run, text_tree}; +use super::support::{read_text, setup_created_dry_run, setup_local_sandbox_run, text_tree}; #[test] fn help() { @@ -53,7 +53,7 @@ fn sandbox_cp_run_without_sandbox_json_errors_cleanly() { #[test] fn sandbox_cp_downloads_file_from_run() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_local_sandbox_run(&context); let dest = context.temp_dir.join("downloaded-root.txt"); let mut cmd = context.cp(); cmd.args([ @@ -73,7 +73,7 @@ fn sandbox_cp_downloads_file_from_run() { #[test] fn sandbox_cp_uploads_file_to_run() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_local_sandbox_run(&context); let local = context.temp_dir.join("upload.txt"); std::fs::write(&local, "uploaded-root") .unwrap_or_else(|err| panic!("failed to write {}: {err}", local.display())); @@ -98,7 +98,7 @@ fn sandbox_cp_uploads_file_to_run() { #[test] fn sandbox_cp_recursive_downloads_directory() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_local_sandbox_run(&context); let dest = context.temp_dir.join("download-dir"); let mut cmd = context.cp(); cmd.args([ diff --git a/lib/crates/fabro-cli/tests/it/cmd/sandbox_preview.rs b/lib/crates/fabro-cli/tests/it/cmd/sandbox_preview.rs index d2778a86c..0adaf8d8a 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/sandbox_preview.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/sandbox_preview.rs @@ -1,6 +1,6 @@ use fabro_test::{fabro_snapshot, test_context}; -use super::support::setup_asset_sandbox_run; +use super::support::setup_local_sandbox_run; #[test] fn help() { @@ -37,7 +37,7 @@ fn help() { #[test] fn sandbox_preview_rejects_non_daytona_run() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_local_sandbox_run(&context); let mut cmd = context.preview(); cmd.args([&setup.run.run_id, "3000"]); diff --git a/lib/crates/fabro-cli/tests/it/cmd/sandbox_ssh.rs b/lib/crates/fabro-cli/tests/it/cmd/sandbox_ssh.rs index 016b37dec..40118a89e 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/sandbox_ssh.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/sandbox_ssh.rs @@ -1,6 +1,6 @@ use fabro_test::{fabro_snapshot, test_context}; -use super::support::setup_asset_sandbox_run; +use super::support::setup_local_sandbox_run; #[test] fn help() { @@ -35,7 +35,7 @@ fn help() { #[test] fn sandbox_ssh_rejects_non_daytona_run() { let context = test_context!(); - let setup = setup_asset_sandbox_run(&context); + let setup = setup_local_sandbox_run(&context); let mut cmd = context.ssh(); cmd.args([&setup.run.run_id, "--print"]); diff --git a/lib/crates/fabro-cli/tests/it/cmd/start.rs b/lib/crates/fabro-cli/tests/it/cmd/start.rs index 41d53ceac..7f7a85503 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/start.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/start.rs @@ -2,7 +2,7 @@ use fabro_test::{fabro_snapshot, test_context}; use crate::support::{example_fixture, fabro_json_snapshot, read_json}; -use super::support::{output_stdout, resolve_run, wait_for_status, write_sleep_workflow}; +use super::support::{output_stdout, resolve_run, wait_for_status, write_gated_workflow}; #[test] fn help() { @@ -160,12 +160,7 @@ digraph Smoke { #[test] fn start_rejects_already_active_or_completed_run() { let context = test_context!(); - write_sleep_workflow( - &context.temp_dir.join("slow.fabro"), - "slow", - "Run slowly", - 3, - ); + let gate = write_gated_workflow(&context.temp_dir.join("slow.fabro"), "slow", "Run slowly"); let mut create_cmd = context.command(); create_cmd.current_dir(&context.temp_dir); @@ -195,7 +190,7 @@ fn start_rejects_already_active_or_completed_run() { start_cmd.args(["start", &run_id]); start_cmd.assert().success(); - wait_for_status(&run.run_dir, &["starting", "running"]); + wait_for_status(&run.run_dir, &["running"]); let mut active_cmd = context.command(); active_cmd.current_dir(&context.temp_dir); @@ -208,6 +203,7 @@ fn start_rejects_already_active_or_completed_run() { error: an engine process is still running for this run — cannot start "); + gate.release(); wait_for_status(&run.run_dir, &["succeeded"]); let mut completed_cmd = context.command(); diff --git a/lib/crates/fabro-cli/tests/it/cmd/support.rs b/lib/crates/fabro-cli/tests/it/cmd/support.rs index df23d785d..24ab50a93 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/support.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/support.rs @@ -5,6 +5,7 @@ use std::time::{Duration, Instant}; use fabro_test::TestContext; use serde_json::Value; +use shlex::try_quote; const COMMAND_TIMEOUT: Duration = Duration::from_secs(30); @@ -24,11 +25,15 @@ pub(crate) struct ProjectFixture { pub(crate) fabro_root: PathBuf, } -pub(crate) struct AssetSandboxSetup { +pub(crate) struct WorkspaceRunSetup { pub(crate) run: RunSetup, pub(crate) workspace_dir: PathBuf, } +pub(crate) struct WorkflowGate { + gate_path: PathBuf, +} + #[derive(Clone, Copy)] enum GitWorkflowKind { Changed, @@ -187,20 +192,26 @@ pub(crate) fn setup_project_fixture(context: &TestContext) -> ProjectFixture { } } -pub(crate) fn setup_asset_sandbox_run(context: &TestContext) -> AssetSandboxSetup { - let workspace_dir = context.temp_dir.join("asset-sandbox"); +impl WorkflowGate { + pub(crate) fn release(&self) { + write_text_file(&self.gate_path, "open\n"); + } +} + +pub(crate) fn setup_asset_run(context: &TestContext) -> WorkspaceRunSetup { + let workspace_dir = context.temp_dir.join("asset-run"); std::fs::create_dir_all(&workspace_dir) .unwrap_or_else(|err| panic!("failed to create {}: {err}", workspace_dir.display())); write_text_file( - &workspace_dir.join("asset_sandbox.fabro"), - r#"digraph AssetSandbox { - graph [goal="Exercise asset and sandbox commands", default_max_retries=0] + &workspace_dir.join("asset_run.fabro"), + r#"digraph AssetRun { + graph [goal="Exercise asset commands", default_max_retries=0] start [shape=Mdiamond] exit [shape=Msquare] - create_assets [shape=parallelogram, script="mkdir -p assets/shared assets/node_a sandbox_dir/download_me/nested && printf one > assets/shared/report.txt && printf alpha > assets/node_a/summary.txt && printf keep > sandbox_dir/download_me/root.txt && printf nested > sandbox_dir/download_me/nested/child.txt && sleep 1", max_retries=0] - retry_assets [shape=parallelogram, script="mkdir -p assets/retry && if [ ! -f .retry-sentinel ]; then printf first > assets/retry/report.txt && touch .retry-sentinel && sleep 1; else printf second > assets/retry/report.txt; fi", retry_policy="linear", timeout="50ms"] - create_colliding [shape=parallelogram, script="mkdir -p assets/other assets/retry && printf beta > assets/other/summary.txt && printf second > assets/retry/report.txt", max_retries=0] + create_assets [shape=parallelogram, script="mkdir -p assets/shared assets/node_a && printf one > assets/shared/report.txt && printf alpha > assets/node_a/summary.txt", max_retries=0] + retry_assets [shape=parallelogram, script="mkdir -p assets/retry && touch -c -t 200001010000 assets/shared/report.txt assets/node_a/summary.txt && if [ ! -f .retry-sentinel ]; then printf first > assets/retry/report.txt && touch .retry-sentinel && sleep 0.2; else printf second > assets/retry/report.txt; fi", retry_policy="linear", timeout="50ms"] + create_colliding [shape=parallelogram, script="mkdir -p assets/other assets/retry && touch -c -t 200001010000 assets/shared/report.txt assets/node_a/summary.txt assets/retry/report.txt && printf beta > assets/other/summary.txt && printf second > assets/retry/report.txt", max_retries=0] start -> create_assets -> retry_assets -> create_colliding -> exit } "#, @@ -208,8 +219,8 @@ pub(crate) fn setup_asset_sandbox_run(context: &TestContext) -> AssetSandboxSetu write_text_file( &workspace_dir.join("run.toml"), r#"version = 1 -graph = "asset_sandbox.fabro" -goal = "Exercise asset and sandbox commands" +graph = "asset_run.fabro" +goal = "Exercise asset commands" [sandbox] provider = "local" @@ -223,8 +234,60 @@ include = ["assets/**"] "#, ); + let run = run_local_workflow(context, &workspace_dir, "run.toml"); + assert!( + run.run_dir + .join("cache/artifacts/assets/retry_assets/retry_2/manifest.json") + .exists(), + "setup_asset_run should materialize retry_2 assets" + ); + + WorkspaceRunSetup { run, workspace_dir } +} + +pub(crate) fn setup_local_sandbox_run(context: &TestContext) -> WorkspaceRunSetup { + let workspace_dir = context.temp_dir.join("local-sandbox"); + std::fs::create_dir_all(&workspace_dir) + .unwrap_or_else(|err| panic!("failed to create {}: {err}", workspace_dir.display())); + + write_text_file( + &workspace_dir.join("sandbox_run.fabro"), + r#"digraph SandboxRun { + graph [goal="Exercise sandbox commands", default_max_retries=0] + start [shape=Mdiamond] + exit [shape=Msquare] + populate_sandbox [shape=parallelogram, script="mkdir -p sandbox_dir/download_me/nested && printf keep > sandbox_dir/download_me/root.txt && printf nested > sandbox_dir/download_me/nested/child.txt", max_retries=0] + start -> populate_sandbox -> exit +} +"#, + ); + write_text_file( + &workspace_dir.join("run.toml"), + r#"version = 1 +graph = "sandbox_run.fabro" +goal = "Exercise sandbox commands" + +[sandbox] +provider = "local" +preserve = true + +[sandbox.local] +worktree_mode = "never" +"#, + ); + + let run = run_local_workflow(context, &workspace_dir, "run.toml"); + assert!( + run.run_dir.join("sandbox.json").exists(), + "setup_local_sandbox_run should persist sandbox.json" + ); + + WorkspaceRunSetup { run, workspace_dir } +} + +fn run_local_workflow(context: &TestContext, workspace_dir: &Path, workflow: &str) -> RunSetup { let mut cmd = context.command(); - cmd.current_dir(&workspace_dir); + cmd.current_dir(workspace_dir); cmd.timeout(COMMAND_TIMEOUT); cmd.env("OPENAI_API_KEY", "test"); cmd.args([ @@ -235,30 +298,18 @@ include = ["assets/**"] "local", "--provider", "openai", - "run.toml", + workflow, ]); let output = cmd.output().expect("command should execute"); if !output.status.success() { panic!( - "command failed: fabro run --auto-approve --no-retro --sandbox local --provider openai run.toml\nstdout:\n{}\nstderr:\n{}", + "command failed: fabro run --auto-approve --no-retro --sandbox local --provider openai {workflow}\nstdout:\n{}\nstderr:\n{}", stdout(&output), stderr(&output) ); } - let run = only_run(context); - assert!( - run.run_dir - .join("cache/artifacts/assets/retry_assets/retry_2/manifest.json") - .exists(), - "setup F should materialize retry_2 assets" - ); - assert!( - run.run_dir.join("sandbox.json").exists(), - "setup F should persist sandbox.json" - ); - - AssetSandboxSetup { run, workspace_dir } + only_run(context) } pub(crate) fn add_project_workflow( @@ -296,14 +347,20 @@ pub(crate) fn add_user_workflow(context: &TestContext, name: &str, goal: &str) - workflow_dir } -pub(crate) fn write_sleep_workflow(path: &Path, name: &str, goal: &str, sleep_seconds: u64) { +pub(crate) fn write_gated_workflow(path: &Path, name: &str, goal: &str) -> WorkflowGate { + let gate_path = path.with_extension("gate"); + let _ = std::fs::remove_file(&gate_path); + let gate_path_str = gate_path.to_string_lossy().into_owned(); + let quoted_gate_path = try_quote(&gate_path_str) + .unwrap_or_else(|_| panic!("failed to quote {}", gate_path.display())); write_text_file( path, &format!( - "digraph {} {{\n graph [goal={goal:?}]\n start [shape=Mdiamond]\n exit [shape=Msquare]\n wait [shape=parallelogram, script=\"sleep {sleep_seconds}\"]\n start -> wait -> exit\n}}\n", + "digraph {} {{\n graph [goal={goal:?}]\n start [shape=Mdiamond]\n exit [shape=Msquare]\n wait [shape=parallelogram, script=\"while [ ! -f {quoted_gate_path} ]; do sleep 0.01; done; sleep 0.2\"]\n start -> wait -> exit\n}}\n", to_pascal_case(name), ), ); + WorkflowGate { gate_path } } pub(crate) fn wait_for_status(run_dir: &Path, expected: &[&str]) -> String { diff --git a/lib/crates/fabro-interview/src/file.rs b/lib/crates/fabro-interview/src/file.rs index b425c5a35..edb87f164 100644 --- a/lib/crates/fabro-interview/src/file.rs +++ b/lib/crates/fabro-interview/src/file.rs @@ -7,11 +7,6 @@ use tokio::time; use crate::{Answer, Interviewer, Question}; -#[cfg(test)] -const REATTACH_WINDOW: Duration = Duration::from_millis(300); -#[cfg(not(test))] -const REATTACH_WINDOW: Duration = Duration::from_secs(30); - #[cfg(test)] use std::path::Path; @@ -25,6 +20,28 @@ pub struct FileInterviewer { request_path: PathBuf, response_path: PathBuf, claim_path: PathBuf, + poll_interval: Duration, + reattach_window: Duration, +} + +#[cfg(test)] +const DEFAULT_REATTACH_WINDOW: Duration = Duration::from_millis(300); +#[cfg(not(test))] +const DEFAULT_REATTACH_WINDOW: Duration = Duration::from_secs(30); + +const DEFAULT_POLL_INTERVAL: Duration = Duration::from_millis(100); + +#[cfg(test)] +const TEST_POLL_INTERVAL: Duration = Duration::from_millis(1); +#[cfg(test)] +const TEST_REATTACH_WINDOW: Duration = Duration::from_millis(5); + +fn default_reattach_window() -> Duration { + DEFAULT_REATTACH_WINDOW +} + +fn default_poll_interval() -> Duration { + DEFAULT_POLL_INTERVAL } impl FileInterviewer { @@ -33,6 +50,25 @@ impl FileInterviewer { request_path, response_path, claim_path, + poll_interval: default_poll_interval(), + reattach_window: default_reattach_window(), + } + } + + #[cfg(test)] + fn with_timing( + request_path: PathBuf, + response_path: PathBuf, + claim_path: PathBuf, + poll_interval: Duration, + reattach_window: Duration, + ) -> Self { + Self { + request_path, + response_path, + claim_path, + poll_interval, + reattach_window, } } @@ -110,7 +146,7 @@ impl Interviewer for FileInterviewer { claim_was_seen = true; reattach_deadline = None; } else if claim_was_seen && reattach_deadline.is_none() { - reattach_deadline = Some(time::Instant::now() + REATTACH_WINDOW); + reattach_deadline = Some(time::Instant::now() + self.reattach_window); } if let Some(deadline) = reattach_deadline { @@ -120,7 +156,7 @@ impl Interviewer for FileInterviewer { } } - time::sleep(Duration::from_millis(100)).await; + time::sleep(self.poll_interval).await; } }; @@ -147,6 +183,20 @@ mod tests { use super::*; use crate::{AnswerValue, QuestionType}; + fn test_interviewer( + request_path: PathBuf, + response_path: PathBuf, + claim_path: PathBuf, + ) -> FileInterviewer { + FileInterviewer::with_timing( + request_path, + response_path, + claim_path, + TEST_POLL_INTERVAL, + TEST_REATTACH_WINDOW, + ) + } + fn interviewer_paths(run_dir: &Path) -> (PathBuf, PathBuf, PathBuf) { let runtime_dir = run_dir.join("runtime"); ( @@ -156,12 +206,22 @@ mod tests { ) } + async fn wait_for_exists(path: &Path) { + for _ in 0..200 { + if path.exists() { + return; + } + time::sleep(TEST_POLL_INTERVAL).await; + } + panic!("{} should exist", path.display()); + } + #[tokio::test] async fn write_request_poll_response() { let dir = tempfile::tempdir().unwrap(); let run_dir = dir.path().to_path_buf(); let (request_path, response_path, claim_path) = interviewer_paths(&run_dir); - let interviewer = FileInterviewer::new( + let interviewer = test_interviewer( request_path.clone(), response_path.clone(), claim_path.clone(), @@ -173,13 +233,7 @@ mod tests { let ask_handle = tokio::spawn(async move { interviewer.ask(question).await }); // Wait for the request file to appear - for _ in 0..50 { - if request_path.exists() { - break; - } - time::sleep(Duration::from_millis(50)).await; - } - assert!(request_path.exists(), "interview_request.json should exist"); + wait_for_exists(&request_path).await; // Verify the request contains valid Question JSON let request_data = fs::read_to_string(&request_path).await.unwrap(); @@ -205,10 +259,10 @@ mod tests { async fn timeout_returns_default() { let dir = tempfile::tempdir().unwrap(); let (request_path, response_path, claim_path) = interviewer_paths(dir.path()); - let interviewer = FileInterviewer::new(request_path, response_path, claim_path); + let interviewer = test_interviewer(request_path, response_path, claim_path); let mut question = Question::new("approve?", QuestionType::YesNo); - question.timeout_seconds = Some(0.1); + question.timeout_seconds = Some(0.02); question.default = Some(Answer::no()); let answer = interviewer.ask(question).await; @@ -220,41 +274,34 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let run_dir = dir.path().to_path_buf(); let (request_path, response_path, claim_path) = interviewer_paths(&run_dir); - let interviewer = - FileInterviewer::new(request_path.clone(), response_path, claim_path.clone()); + let interviewer = test_interviewer(request_path.clone(), response_path, claim_path.clone()); let question = Question::new("approve?", QuestionType::YesNo); let ask_handle = tokio::spawn(async move { interviewer.ask(question).await }); // Wait for request file to appear - for _ in 0..50 { - if request_path.exists() { - break; - } - time::sleep(Duration::from_millis(50)).await; - } - assert!(request_path.exists()); + wait_for_exists(&request_path).await; // Simulate attacher creating claim file std::fs::write(&claim_path, "12345\n").unwrap(); // Let the poll loop see the claim - time::sleep(Duration::from_millis(150)).await; + time::sleep(TEST_POLL_INTERVAL * 2).await; // Simulate attacher departing (deletes claim without writing response) std::fs::remove_file(&claim_path).unwrap(); // Should return timeout within REATTACH_WINDOW let started = time::Instant::now(); - let answer = time::timeout(Duration::from_secs(2), ask_handle) + let answer = time::timeout(Duration::from_millis(250), ask_handle) .await - .expect("should complete within 2s") + .expect("should complete quickly") .unwrap(); assert_eq!(answer.value, AnswerValue::Timeout); assert!( - started.elapsed() <= Duration::from_secs(1), + started.elapsed() <= Duration::from_millis(100), "should resolve well within the reattach window" ); } @@ -264,8 +311,7 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let run_dir = dir.path().to_path_buf(); let (request_path, response_path, claim_path) = interviewer_paths(&run_dir); - let interviewer = - FileInterviewer::new(request_path.clone(), response_path, claim_path.clone()); + let interviewer = test_interviewer(request_path.clone(), response_path, claim_path.clone()); let mut question = Question::new("approve?", QuestionType::YesNo); question.default = Some(Answer::no()); @@ -273,22 +319,16 @@ mod tests { let ask_handle = tokio::spawn(async move { interviewer.ask(question).await }); // Wait for request file - for _ in 0..50 { - if request_path.exists() { - break; - } - time::sleep(Duration::from_millis(50)).await; - } - assert!(request_path.exists()); + wait_for_exists(&request_path).await; // Simulate attacher creating then deleting claim std::fs::write(&claim_path, "12345\n").unwrap(); - time::sleep(Duration::from_millis(150)).await; + time::sleep(TEST_POLL_INTERVAL * 2).await; std::fs::remove_file(&claim_path).unwrap(); - let answer = time::timeout(Duration::from_secs(2), ask_handle) + let answer = time::timeout(Duration::from_millis(250), ask_handle) .await - .expect("should complete within 2s") + .expect("should complete quickly") .unwrap(); assert_eq!(answer.value, AnswerValue::No); @@ -299,7 +339,7 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let run_dir = dir.path().to_path_buf(); let (request_path, response_path, claim_path) = interviewer_paths(&run_dir); - let interviewer = FileInterviewer::new( + let interviewer = test_interviewer( request_path.clone(), response_path.clone(), claim_path.clone(), @@ -310,30 +350,24 @@ mod tests { let ask_handle = tokio::spawn(async move { interviewer.ask(question).await }); // Wait for request file - for _ in 0..50 { - if request_path.exists() { - break; - } - time::sleep(Duration::from_millis(50)).await; - } - assert!(request_path.exists()); + wait_for_exists(&request_path).await; // First attacher creates then releases claim std::fs::write(&claim_path, "12345\n").unwrap(); - time::sleep(Duration::from_millis(150)).await; + time::sleep(TEST_POLL_INTERVAL * 2).await; std::fs::remove_file(&claim_path).unwrap(); // Second attacher picks up and answers before reattach window expires - time::sleep(Duration::from_millis(50)).await; + time::sleep(TEST_POLL_INTERVAL * 2).await; std::fs::write(&claim_path, "12346\n").unwrap(); let answer = Answer::yes(); let response_json = serde_json::to_string_pretty(&answer).unwrap(); fs::write(response_path, response_json).await.unwrap(); - let result = time::timeout(Duration::from_secs(2), ask_handle) + let result = time::timeout(Duration::from_millis(250), ask_handle) .await - .expect("should complete within 2s") + .expect("should complete quickly") .unwrap(); assert_eq!(result.value, AnswerValue::Yes); @@ -343,10 +377,10 @@ mod tests { async fn timeout_without_default_returns_timeout() { let dir = tempfile::tempdir().unwrap(); let (request_path, response_path, claim_path) = interviewer_paths(dir.path()); - let interviewer = FileInterviewer::new(request_path, response_path, claim_path); + let interviewer = test_interviewer(request_path, response_path, claim_path); let mut question = Question::new("approve?", QuestionType::YesNo); - question.timeout_seconds = Some(0.1); + question.timeout_seconds = Some(0.02); let answer = interviewer.ask(question).await; assert_eq!(answer.value, AnswerValue::Timeout); diff --git a/lib/crates/fabro-interview/src/web.rs b/lib/crates/fabro-interview/src/web.rs index 474025ff5..4755e4a16 100644 --- a/lib/crates/fabro-interview/src/web.rs +++ b/lib/crates/fabro-interview/src/web.rs @@ -114,6 +114,16 @@ mod tests { use std::time::Duration; use tokio::time::sleep; + async fn wait_for_pending_count(interviewer: &WebInterviewer, expected: usize) { + for _ in 0..200 { + if interviewer.pending_questions().len() == expected { + return; + } + sleep(Duration::from_millis(1)).await; + } + panic!("pending question count did not reach {expected}"); + } + #[tokio::test] async fn ask_blocks_until_answer_submitted() { let interviewer = Arc::new(WebInterviewer::new()); @@ -124,10 +134,7 @@ mod tests { interviewer_clone.ask(q).await }); - // Give the ask task a moment to register the question - sleep(Duration::from_millis(50)).await; - - // Question should be pending + wait_for_pending_count(interviewer.as_ref(), 1).await; let pending = interviewer.pending_questions(); assert_eq!(pending.len(), 1); assert_eq!(pending[0].question.text, "approve?"); @@ -151,8 +158,7 @@ mod tests { interviewer_clone.ask(q).await }); - sleep(Duration::from_millis(50)).await; - + wait_for_pending_count(interviewer.as_ref(), 1).await; let pending = interviewer.pending_questions(); assert_eq!(pending.len(), 1); @@ -192,8 +198,7 @@ mod tests { i2.ask(q).await }); - sleep(Duration::from_millis(50)).await; - + wait_for_pending_count(interviewer.as_ref(), 2).await; let pending = interviewer.pending_questions(); assert_eq!(pending.len(), 2); @@ -245,8 +250,7 @@ mod tests { i_clone.ask(q).await }); - sleep(Duration::from_millis(50)).await; - + wait_for_pending_count(interviewer.as_ref(), 1).await; let pending = interviewer.pending_questions(); assert_eq!(pending.len(), 1); @@ -266,7 +270,7 @@ mod tests { interviewer_clone.ask(q).await }); - sleep(Duration::from_millis(50)).await; + wait_for_pending_count(interviewer.as_ref(), 1).await; { let mut inner = interviewer diff --git a/lib/crates/fabro-server/tests/it/api.rs b/lib/crates/fabro-server/tests/it/api.rs index dc1b9308c..c05cd5a6d 100644 --- a/lib/crates/fabro-server/tests/it/api.rs +++ b/lib/crates/fabro-server/tests/it/api.rs @@ -448,6 +448,70 @@ mod server_lifecycle { serde_json::from_slice(&bytes).unwrap() } + const POLL_INTERVAL: Duration = Duration::from_millis(10); + const POLL_ATTEMPTS: usize = 500; + + async fn run_json(app: &axum::Router, run_id: &str) -> serde_json::Value { + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::OK); + body_json(response.into_body()).await + } + + async fn wait_for_question_id(app: &axum::Router, run_id: &str) -> String { + for _ in 0..POLL_ATTEMPTS { + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/questions"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(req).await.unwrap(); + let body = body_json(response.into_body()).await; + let arr = body["data"].as_array().unwrap(); + if let Some(question_id) = arr + .first() + .and_then(|item| item["id"].as_str()) + .map(ToOwned::to_owned) + { + return question_id; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + panic!("question should have appeared"); + } + + async fn wait_for_run_status(app: &axum::Router, run_id: &str, expected: &[&str]) -> String { + for _ in 0..POLL_ATTEMPTS { + let body = run_json(app, run_id).await; + let status = body["status"].as_str().unwrap().to_string(); + if expected.iter().any(|candidate| *candidate == status) { + return status; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + panic!("run {run_id} did not reach any of {expected:?}"); + } + + async fn wait_for_run_status_not_in( + app: &axum::Router, + run_id: &str, + unexpected: &[&str], + ) -> String { + for _ in 0..POLL_ATTEMPTS { + let body = run_json(app, run_id).await; + let status = body["status"].as_str().unwrap().to_string(); + if unexpected.iter().all(|candidate| *candidate != status) { + return status; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + panic!("run {run_id} stayed in {unexpected:?}"); + } + const GATE_DOT: &str = r#"digraph GateTest { graph [goal="Test gate"] start [shape=Mdiamond] @@ -489,23 +553,7 @@ mod server_lifecycle { let run_id = body["id"].as_str().unwrap().to_string(); // 2. Poll for question to appear (run goes start -> work -> gate, then blocks) - let mut question_id = String::new(); - for _ in 0..500 { - tokio::time::sleep(Duration::from_millis(10)).await; - let req = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}/questions"))) - .body(Body::empty()) - .unwrap(); - let response = app.clone().oneshot(req).await.unwrap(); - let body = body_json(response.into_body()).await; - let arr = body["data"].as_array().unwrap(); - if !arr.is_empty() { - question_id = arr[0]["id"].as_str().unwrap().to_string(); - break; - } - } - assert!(!question_id.is_empty(), "question should have appeared"); + let question_id = wait_for_question_id(&app, &run_id).await; // 3. Submit answer selecting first option (Approve) let req = Request::builder() @@ -522,22 +570,7 @@ mod server_lifecycle { assert_eq!(response.status(), StatusCode::NO_CONTENT); // 4. Poll until completed - let mut final_status = String::new(); - for _ in 0..500 { - tokio::time::sleep(Duration::from_millis(10)).await; - let req = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}"))) - .body(Body::empty()) - .unwrap(); - let response = app.clone().oneshot(req).await.unwrap(); - let body = body_json(response.into_body()).await; - let status = body["status"].as_str().unwrap().to_string(); - if status == "completed" || status == "failed" { - final_status = status; - break; - } - } + let final_status = wait_for_run_status(&app, &run_id, &["completed", "failed"]).await; assert_eq!(final_status, "completed"); // 5. Verify context endpoint returns an object @@ -587,8 +620,7 @@ mod server_lifecycle { let body = body_json(response.into_body()).await; let run_id = body["id"].as_str().unwrap().to_string(); - // Wait briefly for scheduler to pick up the run - tokio::time::sleep(Duration::from_millis(100)).await; + wait_for_run_status_not_in(&app, &run_id, &["queued", "starting"]).await; // Cancel it let req = Request::builder() @@ -637,6 +669,57 @@ mod sse_events { start -> work -> exit }"#; + const POLL_INTERVAL: Duration = Duration::from_millis(10); + const POLL_ATTEMPTS: usize = 500; + + async fn body_json(body: Body) -> serde_json::Value { + let bytes = axum::body::to_bytes(body, usize::MAX).await.unwrap(); + serde_json::from_slice(&bytes).unwrap() + } + + async fn run_status(app: &axum::Router, run_id: &str) -> String { + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response.into_body()).await; + body["status"].as_str().unwrap().to_string() + } + + async fn wait_for_run_status_not_in( + app: &axum::Router, + run_id: &str, + unexpected: &[&str], + ) -> String { + for _ in 0..POLL_ATTEMPTS { + let status = run_status(app, run_id).await; + if unexpected.iter().all(|candidate| *candidate != status) { + return status; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + panic!("run {run_id} stayed in {unexpected:?}"); + } + + async fn wait_for_checkpoint(app: &axum::Router, run_id: &str) -> serde_json::Value { + for _ in 0..POLL_ATTEMPTS { + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/checkpoint"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(req).await.unwrap(); + if response.status() == StatusCode::OK { + return body_json(response.into_body()).await; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + panic!("checkpoint did not become available for {run_id}"); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn sse_stream_contains_expected_event_types() { let state = create_app_state(test_db().await); @@ -657,14 +740,10 @@ mod sse_events { .unwrap(); let response = app.clone().oneshot(req).await.unwrap(); assert_eq!(response.status(), StatusCode::CREATED); - let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) - .await - .unwrap(); - let body: serde_json::Value = serde_json::from_slice(&bytes).unwrap(); + let body = body_json(response.into_body()).await; let run_id = body["id"].as_str().unwrap().to_string(); - // Wait for scheduler to pick up the run before subscribing to SSE - tokio::time::sleep(Duration::from_millis(100)).await; + wait_for_run_status_not_in(&app, &run_id, &["queued", "starting"]).await; // Get SSE stream let req = Request::builder() @@ -708,13 +787,8 @@ mod sse_events { if let Some(json_str) = line.strip_prefix("data:") { let json_str = json_str.trim(); if let Ok(event) = serde_json::from_str::(json_str) { - // The event is serialized as a tagged enum, so the type is the first key - if let Some(obj) = event.as_object() { - for key in obj.keys() { - event_types.push(key.clone()); - } - } else if let Some(s) = event.as_str() { - event_types.push(s.to_string()); + if let Some(event_name) = event["event"].as_str() { + event_types.push(event_name.to_string()); } } } @@ -729,27 +803,13 @@ mod sse_events { assert!( event_types .iter() - .any(|t| t == "StageStarted" || t == "StageCompleted"), + .any(|t| t == "stage.started" || t == "stage.completed"), "should contain stage events, got: {event_types:?}" ); } // Pipeline is complete (SSE stream ended), verify checkpoint - // Small yield to let the spawned task update state - tokio::time::sleep(Duration::from_millis(10)).await; - - let req = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}/checkpoint"))) - .body(Body::empty()) - .unwrap(); - let response = app.clone().oneshot(req).await.unwrap(); - assert_eq!(response.status(), StatusCode::OK); - - let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) - .await - .unwrap(); - let cp_body: serde_json::Value = serde_json::from_slice(&bytes).unwrap(); + let cp_body = wait_for_checkpoint(&app, &run_id).await; // If run completed, checkpoint should have completed_nodes if !cp_body.is_null() { let completed = cp_body["completed_nodes"].as_array(); @@ -802,6 +862,29 @@ mod serve_dry_run { serde_json::from_slice(&bytes).unwrap() } + const POLL_INTERVAL: Duration = Duration::from_millis(10); + const POLL_ATTEMPTS: usize = 500; + + async fn wait_for_run_status(app: &axum::Router, run_id: &str, expected: &[&str]) -> String { + for _ in 0..POLL_ATTEMPTS { + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}"))) + .body(Body::empty()) + .unwrap(); + + let response = app.clone().oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response.into_body()).await; + let status = body["status"].as_str().unwrap().to_string(); + if expected.iter().any(|candidate| *candidate == status) { + return status; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + panic!("run {run_id} did not reach any of {expected:?}"); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn dry_run_serve_starts_and_runs_workflow() { let app = dry_run_app().await; @@ -823,21 +906,8 @@ mod serve_dry_run { let run_id = body["id"].as_str().unwrap().to_string(); assert!(!run_id.is_empty()); - // Wait for run to complete - tokio::time::sleep(Duration::from_millis(500)).await; - - // GET /runs/{id} to verify completion - let req = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}"))) - .body(Body::empty()) - .unwrap(); - - let response = app.oneshot(req).await.unwrap(); - assert_eq!(response.status(), StatusCode::OK); - - let body = body_json(response.into_body()).await; - assert_eq!(body["status"].as_str().unwrap(), "completed"); + let status = wait_for_run_status(&app, &run_id, &["completed", "failed"]).await; + assert_eq!(status, "completed"); } #[tokio::test] diff --git a/lib/crates/fabro-telemetry/src/buffer.rs b/lib/crates/fabro-telemetry/src/buffer.rs index 2c85b6e54..dea1fdb4e 100644 --- a/lib/crates/fabro-telemetry/src/buffer.rs +++ b/lib/crates/fabro-telemetry/src/buffer.rs @@ -150,6 +150,7 @@ mod tests { let (tx, rx) = mpsc::channel(); let mid_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); let final_flushes: Arc>>> = Arc::new(Mutex::new(Vec::new())); + let (notify_tx, notify_rx) = mpsc::channel(); let mid = mid_flushes.clone(); let fin = final_flushes.clone(); @@ -165,7 +166,8 @@ mod tests { }, move |tracks| { let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); - mid.lock().unwrap().push(events); + mid.lock().unwrap().push(events.clone()); + let _ = notify_tx.send(events); }, move |tracks| { let events: Vec = tracks.iter().map(|t| t.event.clone()).collect(); @@ -174,8 +176,10 @@ mod tests { ); }); - // Wait for time threshold to fire, then drop sender - std::thread::sleep(Duration::from_millis(150)); + let flushed = notify_rx + .recv_timeout(Duration::from_secs(1)) + .expect("time-based flush should have fired"); + assert_eq!(flushed, vec!["e1"]); drop(tx); handle.join().unwrap(); diff --git a/lib/crates/fabro-telemetry/src/spawn.rs b/lib/crates/fabro-telemetry/src/spawn.rs index 0b9bd8ef6..96f5bc2fb 100644 --- a/lib/crates/fabro-telemetry/src/spawn.rs +++ b/lib/crates/fabro-telemetry/src/spawn.rs @@ -150,6 +150,16 @@ mod tests { #[cfg(unix)] #[test] fn spawn_detached_unix_creates_marker_file() { + fn wait_for_file(path: &std::path::Path) { + for _ in 0..100 { + if path.exists() { + return; + } + std::thread::sleep(std::time::Duration::from_millis(10)); + } + panic!("detached process should have created {}", path.display()); + } + // Spawn a detached `touch ` and verify the file appears. let tmp = std::env::temp_dir().join("fabro-spawn-detached-test-marker"); let _ = std::fs::remove_file(&tmp); @@ -157,13 +167,7 @@ mod tests { let tmp_str = tmp.to_str().unwrap(); spawn_detached(&["touch", tmp_str], &[], &[]); - // Wait a bit for the detached process to complete. - std::thread::sleep(std::time::Duration::from_millis(500)); - - assert!( - tmp.exists(), - "detached process should have created the marker file" - ); + wait_for_file(&tmp); std::fs::remove_file(&tmp).ok(); } }