use super::{transcript, transfer};
use crate::prompt::Error;
use crate::template::GitRunner;
use std::path::{Path, PathBuf};
use std::time::SystemTime;
const MESSAGE_EXT: &str = "md";
pub(super) fn drain(
worktree: &Path,
inbox: &Path,
conv_id: &str,
git: &dyn GitRunner,
) -> Result<Delivery, Error> {
recover_strays(worktree, conv_id, git)?;
let mut delivery = Delivery {
delivered: 0,
left: Vec::new(),
};
for msg in pending(inbox)? {
let body = std::fs::read_to_string(&msg.path).map_err(Error::Io)?;
if transfer::terminal_ref_of(&body).is_some() {
delivery.left.push(SeenDeposit {
name: msg.name,
mtime: msg.mtime,
});
continue;
}
transcript::deliver_message(worktree, conv_id, &msg.sender, &msg.path, git)?;
delivery.delivered += 1;
}
Ok(delivery)
}
#[derive(Debug)]
pub(super) struct Delivery {
pub(super) delivered: usize,
pub(super) left: Vec<SeenDeposit>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(super) struct SeenDeposit {
name: String,
mtime: SystemTime,
}
impl SeenDeposit {
#[cfg(test)]
pub(super) fn new(name: String, mtime: SystemTime) -> Self {
SeenDeposit { name, mtime }
}
pub(super) fn matches(&self, pending: &Pending) -> bool {
self.name == pending.name && self.mtime == pending.mtime
}
}
#[derive(Debug)]
pub(super) struct Pending {
mtime: SystemTime,
pub(super) name: String,
pub(super) path: PathBuf,
pub(super) sender: String,
}
pub(super) fn pending(inbox: &Path) -> Result<Vec<Pending>, Error> {
let rd = match std::fs::read_dir(inbox) {
Ok(rd) => rd,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => return Err(Error::Io(e)),
};
let mut out = Vec::new();
for entry in rd {
let entry = entry.map_err(Error::Io)?;
let name = entry.file_name().to_string_lossy().into_owned();
let Some(sender) = sender_of(&name) else {
continue;
};
let mtime = entry
.metadata()
.and_then(|m| m.modified())
.unwrap_or(SystemTime::UNIX_EPOCH);
out.push(Pending {
mtime,
name,
path: entry.path(),
sender,
});
}
out.sort_by(|a, b| (a.mtime, &a.name).cmp(&(b.mtime, &b.name)));
Ok(out)
}
fn sender_of(name: &str) -> Option<String> {
let stem = name.strip_suffix(&format!(".{MESSAGE_EXT}"))?;
let (sender, seq) = stem.rsplit_once('-')?;
seq.parse::<u32>().ok()?;
(!sender.is_empty() && !sender.starts_with('.')).then(|| sender.to_string())
}
fn recover_strays(worktree: &Path, conv_id: &str, git: &dyn GitRunner) -> Result<(), Error> {
let status = git
.run_capture(
worktree,
&["status", "--porcelain", "--", transcript::MESSAGES_DIR],
)
.map_err(|source| Error::Git {
op: "drain status",
source,
})?;
if status.trim().is_empty() {
return Ok(());
}
git.run(worktree, &["add", transcript::MESSAGES_DIR])
.map_err(|source| Error::Git {
op: "drain recover add",
source,
})?;
let msg = format!("transcript: recover delivered stray [{conv_id}]");
git.run(worktree, &["commit", "-m", msg.as_str()])
.map_err(|source| Error::Git {
op: "drain recover commit",
source,
})
}
#[cfg(test)]
mod tests;