use std::path::Path;
use basis::RunUsage;
use serde_json::{Value, json};
use crate::{
data_dir::{AgentPaths, DataDir, canonical_workspace, valid_task_handle, workspace_key},
inbox, lock,
state::{MessageRecord, TaskMeta, load_meta, read_terminal},
};
const PROMPT_BUDGET: usize = 64;
#[derive(Debug, Clone)]
pub struct TaskSummary {
pub task: String,
pub state: String,
pub started_ms: u64,
pub last_activity_ms: u64,
pub prompt: String,
pub agent_id: String,
pub usage: RunUsage,
}
impl TaskSummary {
pub fn payload(&self) -> Value {
let mut payload = json!({
"task": self.task,
"state": self.state,
"started_ms": self.started_ms,
"last_activity_ms": self.last_activity_ms,
"prompt": self.prompt,
"continuable": !self.agent_id.is_empty(),
});
if self.usage != RunUsage::default() {
payload["usage"] = json!(self.usage);
}
payload
}
}
pub(crate) fn workspace_tasks(
data: &DataDir,
workspace: &Path,
) -> Result<Option<Vec<TaskSummary>>, String> {
let canonical = canonical_workspace(workspace)
.map_err(|error| format!("resolve workspace {}: {error}", workspace.display()))?;
let key = workspace_key(&canonical);
match data.described_workspace(&key) {
None => return Ok(None),
Some(described) if described != canonical => {
return Err(format!(
"workspace key collision: {key} describes {}, not {}",
described.display(),
canonical.display()
));
}
Some(_) => {}
}
Ok(Some(scan(data, &key)?))
}
fn scan(data: &DataDir, key: &str) -> Result<Vec<TaskSummary>, String> {
let agents = data.agents_dir(key);
let entries = match std::fs::read_dir(&agents) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(error) => return Err(format!("scan workspace agents: {error}")),
};
let mut summaries = Vec::new();
for entry in entries {
let entry = entry.map_err(|error| format!("scan workspace agents: {error}"))?;
let task = format!("{key}/{}", entry.file_name().to_string_lossy());
let Some(paths) = data.agent_dir(&task).filter(AgentPaths::exists) else {
continue;
};
let Ok(meta) = load_meta(&paths) else {
continue;
};
let messages = inbox::load(&paths).unwrap_or_default();
let state = task_state(&paths).unwrap_or_else(|_| "unknown".to_string());
summaries.push(TaskSummary {
state,
started_ms: meta.created_ms,
last_activity_ms: last_activity_ms(&meta, &messages),
prompt: first_line(&meta.prompt),
agent_id: meta.agent_id,
usage: meta.usage,
task,
});
}
summaries.sort_by(|left, right| {
right
.last_activity_ms
.cmp(&left.last_activity_ms)
.then_with(|| left.task.cmp(&right.task))
});
Ok(summaries)
}
fn last_activity_ms(meta: &TaskMeta, messages: &[MessageRecord]) -> u64 {
let sent = messages
.iter()
.map(|message| message.created_ms)
.max()
.unwrap_or_default();
meta.created_ms.max(meta.updated_ms).max(sent)
}
fn task_state(paths: &AgentPaths) -> Result<String, String> {
match read_terminal(paths)? {
Some(terminal) => Ok(terminal["state"].as_str().unwrap_or("unknown").to_string()),
None => Ok(probe_state(lock::is_held(&paths.attach_lock())).to_string()),
}
}
pub fn probe_state(attached: bool) -> &'static str {
if attached { "running" } else { "resumable" }
}
pub(crate) fn latest_conversation(summaries: &[TaskSummary]) -> Option<&TaskSummary> {
summaries
.iter()
.find(|summary| !summary.agent_id.is_empty())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum NamedError {
InvalidReference(String),
NotFound(String),
}
pub(crate) fn named<'a>(
summaries: &'a [TaskSummary],
workspace_key: &str,
handle: &str,
) -> Result<&'a TaskSummary, NamedError> {
let Some((key, _)) = valid_task_handle(handle) else {
return Err(NamedError::InvalidReference(format!(
"`{handle}` is not a task handle"
)));
};
if key != workspace_key {
return Err(NamedError::InvalidReference(format!(
"task {handle} belongs to another workspace; list it where it was started"
)));
}
summaries
.iter()
.find(|summary| summary.task == handle)
.ok_or_else(|| NamedError::NotFound(format!("no task directory for {handle}")))
}
pub(crate) fn claimed_continuation(
data: &DataDir,
key: &str,
agent_id: &str,
) -> Result<Option<String>, String> {
let agents = data.agents_dir(key);
let entries = match std::fs::read_dir(&agents) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(error) => return Err(format!("scan workspace agents: {error}")),
};
for entry in entries {
let entry = entry.map_err(|error| format!("scan workspace agents: {error}"))?;
let task = format!("{key}/{}", entry.file_name().to_string_lossy());
let Some(paths) = data.agent_dir(&task).filter(AgentPaths::exists) else {
continue;
};
let Ok(meta) = load_meta(&paths) else {
continue;
};
if meta.continues.as_deref() != Some(agent_id) {
continue;
}
if read_terminal(&paths)?.is_some() {
continue;
}
return Ok(Some(task));
}
Ok(None)
}
fn first_line(prompt: &str) -> String {
let line = prompt.lines().next().unwrap_or_default().trim();
match line.char_indices().nth(PROMPT_BUDGET) {
Some((end, _)) => format!("{}…", &line[..end]),
None => line.to_string(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::state::{MessageState, RunOptions, save_meta, write_terminal};
fn summary(task: &str, state: &str, started_ms: u64, agent_id: &str) -> TaskSummary {
TaskSummary {
task: task.to_string(),
state: state.to_string(),
started_ms,
last_activity_ms: started_ms,
prompt: "fix the failing test".to_string(),
agent_id: agent_id.to_string(),
usage: RunUsage::default(),
}
}
fn meta(created_ms: u64, updated_ms: u64) -> TaskMeta {
let mut meta = TaskMeta::new(
"w/t".to_string(),
None,
true,
"/repo".to_string(),
"fix the failing test".to_string(),
RunOptions::default(),
None,
);
meta.created_ms = created_ms;
meta.updated_ms = updated_ms;
meta
}
fn message(created_ms: u64) -> MessageRecord {
MessageRecord {
id: format!("m{created_ms}"),
body: "more".to_string(),
state: MessageState::Pending,
created_ms,
reply: None,
}
}
#[test]
fn a_json_row_says_whether_it_can_be_continued() {
let started = summary("w/t", "succeeded", 5, "agent-1").payload();
assert_eq!(started["task"], "w/t");
assert_eq!(started["state"], "succeeded");
assert_eq!(started["started_ms"], 5);
assert_eq!(started["continuable"], true);
assert!(
started.get("usage").is_none(),
"a task that reported nothing claims no measurement: {started}"
);
let never_attached = summary("w/t2", "resumable", 5, "").payload();
assert_eq!(
never_attached["continuable"], false,
"an agent nobody attached to has no conversation yet"
);
}
#[test]
fn a_json_row_carries_what_the_task_spent() {
let mut spent = summary("w/t", "succeeded", 5, "agent-1");
spent.usage = RunUsage {
input_tokens: 900,
output_tokens: 100,
..RunUsage::default()
};
assert_eq!(spent.payload()["usage"]["input_tokens"], 900);
}
#[test]
fn continue_takes_the_first_task_in_the_list_that_has_a_conversation() {
let summaries = vec![
summary("w/untouched", "resumable", 300, ""),
summary("w/middle", "succeeded", 200, "agent-middle"),
summary("w/oldest", "succeeded", 100, "agent-oldest"),
];
assert_eq!(
latest_conversation(&summaries).map(|summary| summary.task.as_str()),
Some("w/middle")
);
assert!(
latest_conversation(&summaries[..1]).is_none(),
"a workspace whose only task never ran has nothing to continue"
);
}
#[test]
fn activity_is_the_latest_thing_either_writer_recorded() {
assert_eq!(
last_activity_ms(&meta(100, 300), &[]),
300,
"a task nobody has written to since its last turn"
);
assert_eq!(
last_activity_ms(&meta(100, 300), &[message(200), message(700)]),
700,
"a message sent after the last turn is the newer fact"
);
assert_eq!(
last_activity_ms(&meta(100, 900), &[message(700)]),
900,
"and so is a turn run after the last message"
);
assert_eq!(
last_activity_ms(&meta(100, 0), &[]),
100,
"a record written before basis kept this clock falls back to its start"
);
}
#[test]
fn a_json_row_carries_both_clocks() {
let mut worked = summary("w/t", "succeeded", 0, "agent-1");
worked.last_activity_ms = 7_140_000;
let payload = worked.payload();
assert_eq!(payload["started_ms"], 0, "the birthday survives in --json");
assert_eq!(payload["last_activity_ms"], 7_140_000);
}
#[test]
fn a_handle_from_another_workspace_is_refused_rather_than_searched_for() {
let here = "0123456789abcdef";
let elsewhere = format!("fedcba9876543210/{:032x}", 1);
let summaries = vec![summary(&format!("{here}/{:032x}", 1), "succeeded", 1, "a")];
let error = named(&summaries, here, &elsewhere).expect_err("refused");
assert!(
matches!(error, NamedError::InvalidReference(_)),
"{error:?}"
);
assert!(
format!("{error:?}").contains("another workspace"),
"{error:?}"
);
let malformed = named(&summaries, here, "not-a-handle").expect_err("refused");
assert!(
matches!(malformed, NamedError::InvalidReference(_)),
"{malformed:?}"
);
assert!(format!("{malformed:?}").contains("not a task handle"));
let found = named(&summaries, here, &summaries[0].task).expect("in this workspace");
assert_eq!(found.task, summaries[0].task);
}
#[test]
fn a_handle_that_fits_the_grammar_but_names_nothing_is_not_found_not_invalid() {
let here = "0123456789abcdef";
let missing = format!("{here}/{:032x}", 99);
let summaries = vec![summary(&format!("{here}/{:032x}", 1), "succeeded", 1, "a")];
let error = named(&summaries, here, &missing).expect_err("refused");
assert!(matches!(error, NamedError::NotFound(_)), "{error:?}");
assert!(
format!("{error:?}").contains("no task directory"),
"{error:?}"
);
}
#[test]
fn a_still_open_continuation_is_a_claim_a_settled_one_releases() {
let dir = tempfile::tempdir().unwrap();
let data = DataDir::from_path(dir.path()).unwrap();
let key = "0123456789abcdef";
let claimant = format!("{key}/{:032x}", 1);
let paths = data.agent_dir(&claimant).unwrap();
std::fs::create_dir_all(paths.dir()).unwrap();
let meta = TaskMeta::new(
claimant.clone(),
None,
true,
"/repo".to_string(),
"continue it".to_string(),
RunOptions::default(),
None,
)
.continuing(Some("conversation-1".to_string()));
save_meta(&paths, &meta).unwrap();
assert_eq!(
claimed_continuation(&data, key, "conversation-1").unwrap(),
Some(claimant.clone()),
"an unsettled claimant is still holding its claim"
);
assert_eq!(
claimed_continuation(&data, key, "some-other-conversation").unwrap(),
None,
"a claim on one conversation says nothing about another"
);
write_terminal(&paths, &json!({"state": "succeeded", "result": "done"})).unwrap();
assert_eq!(
claimed_continuation(&data, key, "conversation-1").unwrap(),
None,
"a settled claimant's claim released the moment it settled"
);
}
#[test]
fn a_row_carries_one_bounded_line_of_the_prompt() {
assert_eq!(first_line("fix the test\nthen push"), "fix the test");
assert!(first_line(&"x".repeat(200)).ends_with('…'));
assert_eq!(
first_line(&"界".repeat(100)).chars().count(),
PROMPT_BUDGET + 1
);
}
}