mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Speed up sleep-heavy test suites
This commit is contained in:
parent
e7e6ae00b6
commit
17eb572f87
13 changed files with 393 additions and 220 deletions
|
|
@ -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()]);
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -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]
|
||||
");
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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([
|
||||
|
|
|
|||
|
|
@ -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"]);
|
||||
|
||||
|
|
|
|||
|
|
@ -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"]);
|
||||
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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::<serde_json::Value>(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]
|
||||
|
|
|
|||
|
|
@ -150,6 +150,7 @@ mod tests {
|
|||
let (tx, rx) = mpsc::channel();
|
||||
let mid_flushes: Arc<Mutex<Vec<Vec<String>>>> = Arc::new(Mutex::new(Vec::new()));
|
||||
let final_flushes: Arc<Mutex<Vec<Vec<String>>>> = 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<String> = 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<String> = 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();
|
||||
|
||||
|
|
|
|||
|
|
@ -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 <marker>` 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();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue