use crate::tasks;
use crate::workspace::Workspace;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
type Outcome = Result<bool, Box<dyn std::error::Error>>;
pub const WORKERS: usize = 1;
pub fn run(workspace: &Workspace, args: crate::run::Args, workers: usize, json: bool) -> Outcome {
if workers == 0 {
return Err("--workers 0 would take nothing from the list".into());
}
ostraka_adapter::interrupt::clear();
let interrupts = AtomicUsize::new(0);
let _ = ctrlc::set_handler(move || {
if interrupts.fetch_add(1, Ordering::SeqCst) == 0 {
eprintln!("\nstopping — press Ctrl-C again to give up on what is going");
ostraka_adapter::interrupt::request();
} else {
std::process::exit(130);
}
});
let say = Mutex::new(());
let taken = AtomicUsize::new(0);
let approved = AtomicUsize::new(0);
let refused = AtomicUsize::new(0);
let failed = AtomicUsize::new(0);
std::thread::scope(|scope| {
for _ in 0..workers {
scope.spawn(|| {
loop {
if ostraka_adapter::interrupt::requested() {
return;
}
let claimed = match tasks::claim(&workspace.ostraka()) {
Ok(Some(task)) => task,
Ok(None) => return,
Err(e) => {
let _guard = say.lock();
eprintln!("could not take from the list: {e}");
return;
}
};
taken.fetch_add(1, Ordering::SeqCst);
let mut mine = args.clone();
mine.prompt = claimed.prompt.clone();
if claimed.repository.is_some() {
mine.repository = claimed.repository.clone();
}
if claimed.adapter.is_some() {
mine.adapter = claimed.adapter.clone();
}
match crate::run::execute(workspace, &mine, None) {
Ok(report) => {
let outcome = if report.approved() {
approved.fetch_add(1, Ordering::SeqCst);
"approved"
} else {
refused.fetch_add(1, Ordering::SeqCst);
"refused"
};
let run_id = report.record.run_id.clone();
if let Err(e) =
tasks::finish(&workspace.ostraka(), claimed, &run_id, outcome)
{
let _guard = say.lock();
eprintln!("{run_id} finished but the list did not record it: {e}");
}
if !json {
let _guard = say.lock();
println!("{outcome:<9} {run_id}");
}
}
Err(e) => {
failed.fetch_add(1, Ordering::SeqCst);
let _guard = say.lock();
eprintln!("could not run {}: {e}", claimed.id);
}
}
}
});
}
});
let taken = taken.load(Ordering::SeqCst);
let approved = approved.load(Ordering::SeqCst);
let refused = refused.load(Ordering::SeqCst);
let failed = failed.load(Ordering::SeqCst);
let stopped = ostraka_adapter::interrupt::requested();
if json {
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"taken": taken,
"approved": approved,
"refused": refused,
"failed": failed,
"workers": workers,
"stopped": stopped,
}))?
);
} else if taken == 0 {
println!("nothing on the list");
} else {
let tail = if stopped { ", stopped" } else { "" };
println!(
"{taken} taken — {approved} approved, {refused} refused, {failed} could not run{tail}"
);
}
Ok(failed == 0 && !stopped)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn no_workers_is_refused_rather_than_quietly_doing_nothing() {
let dir = std::env::temp_dir().join(format!("ostraka-drain-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("scratch");
let workspace = Workspace::at(&dir);
let args = crate::run::Args::for_task(String::new());
assert!(run(&workspace, args, 0, false).is_err());
let _ = std::fs::remove_dir_all(&dir);
}
}