mod common;
use std::time::{Duration, Instant};
use common::{Judges, fixture, home_lock};
use magi::daemon::{self, Opts, Stop};
use magi::queue::{Queue, Source, Task, TaskStatus};
fn write_config(dir: &std::path::Path, config: &magi::config::Config) -> std::path::PathBuf {
let path = dir.join("magi.toml");
std::fs::write(&path, toml::to_string(config).expect("serialize config")).unwrap();
path
}
fn daemon_config(fx: &common::Fixture) -> magi::config::Config {
let mut config = fx.config.clone();
config.graph.candidates = 1;
config.graph.reviewers = 1;
config.disk.min_free_bytes = 0;
config.disk.auto_fold = false;
config.disk.cache_limit_bytes = 0;
config
}
fn block_impl_a(fx: &mut common::Fixture, block_dir: &std::path::Path) {
std::fs::create_dir_all(block_dir).expect("block dir");
for a in &mut fx.config.agents {
a.env
.insert("MOCK_BLOCK_SEAT".to_owned(), "impl-A".to_owned());
a.env.insert(
"MOCK_BLOCK_DIR".to_owned(),
block_dir.to_string_lossy().into_owned(),
);
}
}
async fn wait_until(mut cond: impl FnMut() -> bool, timeout: Duration, what: &str) {
let start = Instant::now();
while !cond() {
if start.elapsed() > timeout {
panic!("timed out after {:?} waiting for: {what}", start.elapsed());
}
tokio::time::sleep(Duration::from_millis(15)).await;
}
}
const MARKER_WAIT: Duration = Duration::from_secs(60);
common::e2e! {
async fn an_urgent_task_starts_alongside_an_already_running_task() {
let home = home_lock().await;
let mut fx = fixture(home, Judges::Unanimous, false);
let block_dir = fx.tmp.path().join("block");
block_impl_a(&mut fx, &block_dir);
let started_marker = block_dir.join("started-impl-A");
let release_marker = block_dir.join("release-impl-A");
let config = daemon_config(&fx);
let config_path = write_config(fx.tmp.path(), &config);
let queue = Queue::open();
let mut ordinary = Task::new(
"ordinary task".to_owned(),
"create note.txt".to_owned(),
fx.repo.clone(),
Source::Human,
);
queue.put(&mut ordinary).expect("file the ordinary task");
let stop = Stop::new();
let opts = Opts {
repo: fx.repo.clone(),
config: Some(config_path),
once: false,
poll: Duration::from_millis(20),
worktrees_root: Some(fx.tmp.path().join("wt")),
..Opts::default()
};
let serve = tokio::spawn(daemon::serve_until(opts, stop.clone()));
wait_until(
|| started_marker.exists(),
MARKER_WAIT,
"the ordinary task's implement seat never signalled it started",
)
.await;
assert_eq!(
queue.get(&ordinary.id).unwrap().status,
TaskStatus::Running,
"the ordinary task is claimed and running before the urgent task exists"
);
let mut urgent = Task::new(
"urgent task".to_owned(),
"create note.txt".to_owned(),
fx.repo.clone(),
Source::Human,
);
urgent.urgent = true;
queue.put(&mut urgent).expect("file the urgent task");
wait_until(
|| {
queue
.get(&urgent.id)
.is_ok_and(|t| t.status == TaskStatus::Running)
},
MARKER_WAIT,
"the urgent task never joined the blocked ordinary run in flight",
)
.await;
assert_eq!(
queue.get(&ordinary.id).unwrap().status,
TaskStatus::Running,
"the existing run was not touched, let alone restarted, by the urgent dispatch"
);
assert_eq!(queue.get(&ordinary.id).unwrap().runs.len(), 1);
std::fs::write(&release_marker, b"go").expect("release impl-A");
wait_until(
|| {
queue
.get(&ordinary.id)
.is_ok_and(|t| t.status == TaskStatus::Done)
&& queue
.get(&urgent.id)
.is_ok_and(|t| t.status == TaskStatus::Done)
},
MARKER_WAIT,
"both tasks never reached a terminal status",
)
.await;
stop.stop();
tokio::time::timeout(MARKER_WAIT, serve)
.await
.expect("the daemon loop did not stop in time")
.expect("the daemon task panicked")
.expect("the daemon loop returned an error");
let after_ordinary = queue.get(&ordinary.id).unwrap();
assert_eq!(
after_ordinary.status,
TaskStatus::Done,
"{after_ordinary:?}"
);
assert_eq!(
after_ordinary.runs.len(),
1,
"the ordinary run was never recompeted, only ever the one run: {after_ordinary:?}"
);
let after_urgent = queue.get(&urgent.id).unwrap();
assert_eq!(after_urgent.status, TaskStatus::Done, "{after_urgent:?}");
assert_eq!(after_urgent.runs.len(), 1);
}
}
common::e2e! {
async fn an_ordinary_task_still_waits_for_the_ordinary_slot_to_free() {
let home = home_lock().await;
let mut fx = fixture(home, Judges::Unanimous, false);
let block_dir = fx.tmp.path().join("block");
block_impl_a(&mut fx, &block_dir);
let started_marker = block_dir.join("started-impl-A");
let release_marker = block_dir.join("release-impl-A");
let config = daemon_config(&fx);
let config_path = write_config(fx.tmp.path(), &config);
let queue = Queue::open();
let mut ordinary = Task::new(
"ordinary task".to_owned(),
"create note.txt".to_owned(),
fx.repo.clone(),
Source::Human,
);
queue.put(&mut ordinary).expect("file the ordinary task");
let stop = Stop::new();
let opts = Opts {
repo: fx.repo.clone(),
config: Some(config_path),
once: false,
poll: Duration::from_millis(20),
worktrees_root: Some(fx.tmp.path().join("wt")),
..Opts::default()
};
let serve = tokio::spawn(daemon::serve_until(opts, stop.clone()));
wait_until(
|| started_marker.exists(),
MARKER_WAIT,
"the ordinary task's implement seat never signalled it started",
)
.await;
let mut plain = Task::new(
"second ordinary task".to_owned(),
"create note.txt".to_owned(),
fx.repo.clone(),
Source::Human,
);
queue.put(&mut plain).expect("file the second task");
for _ in 0..10 {
tokio::time::sleep(Duration::from_millis(20)).await;
assert_eq!(
queue.get(&plain.id).unwrap().status,
TaskStatus::Queued,
"an ordinary task must not dispatch while the one ordinary slot \
is already spent on the blocked task"
);
}
std::fs::write(&release_marker, b"go").expect("release impl-A");
wait_until(
|| {
queue
.get(&ordinary.id)
.is_ok_and(|t| t.status == TaskStatus::Done)
},
MARKER_WAIT,
"the ordinary task never finished",
)
.await;
wait_until(
|| {
queue
.get(&plain.id)
.is_ok_and(|t| t.status == TaskStatus::Done)
},
MARKER_WAIT,
"the second ordinary task never started once the slot freed",
)
.await;
stop.stop();
tokio::time::timeout(MARKER_WAIT, serve)
.await
.expect("the daemon loop did not stop in time")
.expect("the daemon task panicked")
.expect("the daemon loop returned an error");
let after_ordinary = queue.get(&ordinary.id).unwrap();
assert_eq!(
after_ordinary.status,
TaskStatus::Done,
"{after_ordinary:?}"
);
let after_plain = queue.get(&plain.id).unwrap();
assert_eq!(after_plain.status, TaskStatus::Done, "{after_plain:?}");
}
}