use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
type Result<T> = std::result::Result<T, Box<dyn std::error::Error>>;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Task {
pub id: String,
pub prompt: String,
pub added_at: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub repository: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub adapter: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub outcome: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum State {
Pending,
Running,
Done,
}
impl State {
pub fn dir(self) -> &'static str {
match self {
State::Pending => "pending",
State::Running => "running",
State::Done => "done",
}
}
pub fn name(self) -> &'static str {
self.dir()
}
}
fn root(workspace_ostraka: &Path) -> PathBuf {
workspace_ostraka.join("tasks")
}
fn place(workspace_ostraka: &Path, state: State, id: &str) -> PathBuf {
root(workspace_ostraka).join(state.dir()).join(id)
}
pub fn add(
workspace_ostraka: &Path,
prompt: &str,
repository: Option<&str>,
adapter: Option<&str>,
) -> Result<Task> {
let dir = place(workspace_ostraka, State::Pending, "");
std::fs::create_dir_all(&dir)?;
let now = ostraka_core::clock::now_rfc3339();
let stamp = now.replace([':', '-'], "").replace('.', "");
for n in 0..1000 {
let id = format!("k{stamp}-{n}");
let path = dir.join(format!("{id}.json"));
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&path)
{
Ok(mut file) => {
let task = Task {
id,
prompt: prompt.to_string(),
added_at: now,
repository: repository.map(str::to_string),
adapter: adapter.map(str::to_string),
run_id: None,
outcome: None,
};
use std::io::Write;
file.write_all(serde_json::to_string_pretty(&task)?.as_bytes())?;
return Ok(task);
}
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => continue,
Err(e) => return Err(e.into()),
}
}
Err("a thousand tasks in one second is not a task list".into())
}
pub fn list(workspace_ostraka: &Path, state: State) -> Vec<Task> {
let dir = root(workspace_ostraka).join(state.dir());
let Ok(entries) = std::fs::read_dir(&dir) else {
return Vec::new();
};
let mut paths: Vec<PathBuf> = entries
.filter_map(|e| e.ok().map(|e| e.path()))
.filter(|p| p.extension().is_some_and(|e| e == "json"))
.collect();
paths.sort();
paths
.iter()
.filter_map(|p| std::fs::read_to_string(p).ok())
.filter_map(|text| serde_json::from_str(&text).ok())
.collect()
}
pub fn claim(workspace_ostraka: &Path) -> Result<Option<Task>> {
let running = root(workspace_ostraka).join(State::Running.dir());
std::fs::create_dir_all(&running)?;
for task in list(workspace_ostraka, State::Pending) {
let from = place(
workspace_ostraka,
State::Pending,
&format!("{}.json", task.id),
);
let to = running.join(format!("{}.json", task.id));
match std::fs::rename(&from, &to) {
Ok(()) => return Ok(Some(task)),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
Err(e) => return Err(e.into()),
}
}
Ok(None)
}
pub fn finish(workspace_ostraka: &Path, mut task: Task, run_id: &str, outcome: &str) -> Result<()> {
task.run_id = Some(run_id.to_string());
task.outcome = Some(outcome.to_string());
let done = root(workspace_ostraka).join(State::Done.dir());
std::fs::create_dir_all(&done)?;
std::fs::write(
done.join(format!("{}.json", task.id)),
serde_json::to_string_pretty(&task)?,
)?;
let running = place(
workspace_ostraka,
State::Running,
&format!("{}.json", task.id),
);
match std::fs::remove_file(&running) {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(format!(
"task {} finished but could not be taken off the running list ({}): \
it will be listed twice until that file is removed — {e}",
task.id,
running.display()
)
.into()),
}
}
pub fn release(workspace_ostraka: &Path, id: &str) -> Result<()> {
let from = place(workspace_ostraka, State::Running, &format!("{id}.json"));
if !from.is_file() {
return Err(format!("no task {id:?} is running").into());
}
let pending = root(workspace_ostraka).join(State::Pending.dir());
std::fs::create_dir_all(&pending)?;
std::fs::rename(&from, pending.join(format!("{id}.json")))?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn finishing_twice_is_not_an_error_but_a_failed_removal_is() {
let dir = std::env::temp_dir().join(format!("ostraka-finish-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("scratch");
let task = add(&dir, "do a thing", None, None).expect("added");
let claimed = claim(&dir).expect("claims").expect("a task was waiting");
assert_eq!(claimed.id, task.id);
assert!(place(&dir, State::Running, &format!("{}.json", task.id)).is_file());
finish(&dir, claimed.clone(), "r-1", "approved").expect("finishes");
assert!(!place(&dir, State::Running, &format!("{}.json", task.id)).is_file());
assert!(place(&dir, State::Done, &format!("{}.json", task.id)).is_file());
finish(&dir, claimed, "r-1", "approved").expect("finishing twice is not an error");
let _ = std::fs::remove_dir_all(&dir);
}
fn scratch(name: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!("ostraka-tasks-{}-{name}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("scratch");
dir
}
#[test]
fn a_task_written_down_is_there_to_be_taken() {
let dir = scratch("add");
let task = add(&dir, "write a file", None, None).expect("adds");
assert_eq!(list(&dir, State::Pending).len(), 1);
assert_eq!(list(&dir, State::Pending)[0].prompt, "write a file");
let taken = claim(&dir).expect("claims").expect("one was there");
assert_eq!(taken.id, task.id);
assert!(list(&dir, State::Pending).is_empty(), "it is still pending");
assert_eq!(list(&dir, State::Running).len(), 1);
finish(&dir, taken, "t1-2026", "approved").expect("finishes");
assert!(list(&dir, State::Running).is_empty());
let done = list(&dir, State::Done);
assert_eq!(done.len(), 1);
assert_eq!(done[0].run_id.as_deref(), Some("t1-2026"));
assert_eq!(done[0].outcome.as_deref(), Some("approved"));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn two_claimants_never_take_the_same_task() {
let dir = scratch("race");
for i in 0..8 {
add(&dir, &format!("task {i}"), None, None).expect("adds");
}
let taken: Vec<String> = std::thread::scope(|s| {
let handles: Vec<_> = (0..4)
.map(|_| {
let dir = dir.clone();
s.spawn(move || {
let mut mine = Vec::new();
while let Ok(Some(task)) = claim(&dir) {
mine.push(task.id);
}
mine
})
})
.collect();
handles
.into_iter()
.flat_map(|h| h.join().expect("thread"))
.collect()
});
assert_eq!(taken.len(), 8, "a task was lost or taken twice: {taken:?}");
let mut unique = taken.clone();
unique.sort();
unique.dedup();
assert_eq!(
unique.len(),
8,
"the same task was claimed twice: {taken:?}"
);
assert!(list(&dir, State::Pending).is_empty());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_task_nobody_reported_on_can_be_put_back() {
let dir = scratch("release");
add(&dir, "a task", None, None).expect("adds");
let taken = claim(&dir).expect("claims").expect("one");
assert!(release(&dir, &taken.id).is_ok());
assert_eq!(list(&dir, State::Pending).len(), 1);
assert!(list(&dir, State::Running).is_empty());
assert!(release(&dir, &taken.id).is_err());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn an_empty_list_is_an_answer_rather_than_an_error() {
let dir = scratch("empty");
assert!(list(&dir, State::Pending).is_empty());
assert!(claim(&dir).expect("claims").is_none());
let _ = std::fs::remove_dir_all(&dir);
}
}