use super::inbox_dir;
use crate::prompt::Clock;
use std::io;
use std::path::{Path, PathBuf};
const MESSAGE_EXT: &str = "md";
const SEQ_WIDTH: usize = 3;
const FIRST_SEQ: u32 = 1;
#[derive(Debug, thiserror::Error)]
pub enum DepositError {
#[error("inbox i/o at {path}: {source}")]
Io {
path: PathBuf,
#[source]
source: io::Error,
},
}
fn io_err(path: &Path, source: io::Error) -> DepositError {
DepositError::Io {
path: path.to_path_buf(),
source,
}
}
pub fn deposit(
workspace: &Path,
agent_id: &str,
sender: &str,
content: &str,
clock: &dyn Clock,
) -> Result<PathBuf, DepositError> {
let dir = inbox_dir(workspace, agent_id);
std::fs::create_dir_all(&dir).map_err(|e| io_err(&dir, e))?;
let seq = next_sequence(&dir, sender).map_err(|e| io_err(&dir, e))?;
let filename = message_filename(sender, seq);
let body = render(sender, &clock.now_iso8601(), content);
atomic_create(&dir, &filename, body.as_bytes())?;
Ok(dir.join(filename))
}
fn message_filename(sender: &str, seq: u32) -> String {
format!("{sender}-{seq:0width$}.{MESSAGE_EXT}", width = SEQ_WIDTH)
}
pub(super) fn next_sequence(dir: &Path, sender: &str) -> io::Result<u32> {
let prefix = format!("{sender}-");
let suffix = format!(".{MESSAGE_EXT}");
let mut max: Option<u32> = None;
for entry in std::fs::read_dir(dir)?.flatten() {
let name = entry.file_name();
let Some(name) = name.to_str() else { continue };
if let Some(seq) = parse_seq(name, &prefix, &suffix) {
max = Some(max.map_or(seq, |m| m.max(seq)));
}
}
Ok(max.map_or(FIRST_SEQ, |m| m + 1))
}
fn parse_seq(name: &str, prefix: &str, suffix: &str) -> Option<u32> {
let mid = name.strip_prefix(prefix)?.strip_suffix(suffix)?;
mid.parse::<u32>().ok()
}
fn render(sender: &str, deposited_at: &str, content: &str) -> String {
format!("---\nfrom: {sender}\ndeposited_at: {deposited_at}\n---\n{content}")
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Epitaph {
FinalResponse,
Stopped,
BudgetExhausted,
Died,
}
impl Epitaph {
pub fn as_str(self) -> &'static str {
match self {
Epitaph::FinalResponse => "final-response",
Epitaph::Stopped => "stopped",
Epitaph::BudgetExhausted => "budget-exhausted",
Epitaph::Died => "died",
}
}
}
pub fn deposit_result(
workspace: &Path,
parent_id: &str,
child_id: &str,
epitaph: Epitaph,
terminal_ref: &str,
terminal_response: Option<&str>,
clock: &dyn Clock,
) -> Result<PathBuf, DepositError> {
let dir = inbox_dir(workspace, parent_id);
std::fs::create_dir_all(&dir).map_err(|e| io_err(&dir, e))?;
let seq = next_sequence(&dir, child_id).map_err(|e| io_err(&dir, e))?;
let filename = message_filename(child_id, seq);
let body = render_result(
child_id,
&clock.now_iso8601(),
epitaph,
terminal_ref,
terminal_response,
);
atomic_create(&dir, &filename, body.as_bytes())?;
Ok(dir.join(filename))
}
fn render_result(
child_id: &str,
deposited_at: &str,
epitaph: Epitaph,
terminal_ref: &str,
terminal_response: Option<&str>,
) -> String {
let head = format!(
"---\nfrom: {child_id}\ndeposited_at: {deposited_at}\n\
epitaph: {ep}\nterminal_ref: {terminal_ref}\n---\n",
ep = epitaph.as_str(),
);
match terminal_response {
Some(body) => format!("{head}{body}"),
None => head,
}
}
pub(super) fn atomic_create(dir: &Path, filename: &str, bytes: &[u8]) -> Result<(), DepositError> {
let final_path = dir.join(filename);
let tmp_path = dir.join(format!(".{filename}.tmp"));
std::fs::write(&tmp_path, bytes).map_err(|e| io_err(&tmp_path, e))?;
std::fs::rename(&tmp_path, &final_path).map_err(|e| io_err(&final_path, e))?;
Ok(())
}