use std::{
io::Write,
path::{Path, PathBuf},
};
use anyhow::Context;
use serde::{Deserialize, Serialize};
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct Entry {
pub delivery_id: String,
pub message_id: String,
pub conversation_id: String,
#[serde(default)]
pub seq: i64,
#[serde(default)]
pub from_address: String,
#[serde(default)]
pub created_at: String,
#[serde(default)]
pub confirmed: bool,
pub spooled_at: String,
}
pub const KEEP_HOURS: i64 = 48;
pub fn spool_path(dir: &Path, session: &str) -> PathBuf {
dir.join("inbox")
.join(format!("{}.jsonl", crate::context::binding_key(session)))
}
fn ensure_dir(path: &Path) -> anyhow::Result<()> {
if let Some(dir) = path.parent() {
std::fs::create_dir_all(dir).with_context(|| format!("creating {}", dir.display()))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::set_permissions(dir, std::fs::Permissions::from_mode(0o700));
}
}
Ok(())
}
pub fn append(path: &Path, entries: &[Entry]) -> anyhow::Result<Vec<Entry>> {
ensure_dir(path)?;
let existing = read(path);
let mut fresh = Vec::new();
for entry in entries {
if existing.iter().any(|e| e.delivery_id == entry.delivery_id)
|| fresh
.iter()
.any(|e: &Entry| e.delivery_id == entry.delivery_id)
{
continue;
}
fresh.push(entry.clone());
}
if fresh.is_empty() {
return Ok(fresh);
}
repair_tail(path)?;
let mut file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(path)
.with_context(|| format!("opening {}", path.display()))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = file.set_permissions(std::fs::Permissions::from_mode(0o600));
}
for entry in &fresh {
let line = serde_json::to_string(entry)?;
writeln!(file, "{line}")?;
}
file.sync_all().context("fsync of the inbox spool")?;
Ok(fresh)
}
fn repair_tail(path: &Path) -> anyhow::Result<()> {
let Ok(text) = std::fs::read(path) else {
return Ok(());
};
if text.is_empty() || text.ends_with(b"\n") {
return Ok(());
}
let mut file = std::fs::OpenOptions::new()
.append(true)
.open(path)
.with_context(|| format!("opening {}", path.display()))?;
file.write_all(b"\n")?;
file.sync_all()?;
Ok(())
}
pub fn read(path: &Path) -> Vec<Entry> {
let Ok(text) = std::fs::read_to_string(path) else {
return Vec::new();
};
text.lines()
.filter(|l| !l.trim().is_empty())
.filter_map(|l| serde_json::from_str::<Entry>(l).ok())
.collect()
}
pub fn rewrite(path: &Path, entries: &[Entry]) -> anyhow::Result<()> {
ensure_dir(path)?;
let cutoff = chrono::Utc::now() - chrono::Duration::hours(KEEP_HOURS);
let kept: Vec<String> = entries
.iter()
.filter(|e| {
if !e.confirmed {
return true;
}
match chrono::DateTime::parse_from_rfc3339(&e.spooled_at) {
Ok(at) => at.with_timezone(&chrono::Utc) > cutoff,
Err(_) => true,
}
})
.filter_map(|e| serde_json::to_string(e).ok())
.collect();
crate::context::write_private(path, &format!("{}\n", kept.join("\n")))
}
pub fn unconfirmed(entries: &[Entry]) -> Vec<String> {
entries
.iter()
.filter(|e| !e.confirmed)
.map(|e| e.delivery_id.clone())
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
fn entry(id: &str) -> Entry {
Entry {
delivery_id: id.to_owned(),
message_id: "m".into(),
conversation_id: "c".into(),
seq: 1,
from_address: "dani/review".into(),
created_at: "2026-09-20T00:00:00Z".into(),
confirmed: false,
spooled_at: chrono::Utc::now().to_rfc3339(),
}
}
#[test]
fn a_redelivered_reference_is_not_spooled_twice() {
let dir = std::env::temp_dir().join(format!("acs-spool-{}", uuid::Uuid::new_v4()));
let path = dir.join("s.jsonl");
assert_eq!(append(&path, &[entry("a"), entry("b")]).unwrap().len(), 2);
assert_eq!(
append(&path, &[entry("b"), entry("c")]).unwrap().len(),
1,
"'b' was already held; only 'c' is new"
);
assert_eq!(read(&path).len(), 3);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_truncated_line_does_not_lose_the_rest() {
let dir = std::env::temp_dir().join(format!("acs-spool-{}", uuid::Uuid::new_v4()));
let path = dir.join("s.jsonl");
append(&path, &[entry("a")]).unwrap();
{
use std::io::Write as _;
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap();
write!(f, "{{\"delivery_id\": \"hal").unwrap();
}
append(&path, &[entry("b")]).unwrap();
let held = read(&path);
assert_eq!(held.len(), 2);
assert!(held.iter().any(|e| e.delivery_id == "b"));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn compaction_keeps_what_the_bus_has_not_been_told() {
let dir = std::env::temp_dir().join(format!("acs-spool-{}", uuid::Uuid::new_v4()));
let path = dir.join("s.jsonl");
let mut old = entry("old");
old.confirmed = true;
old.spooled_at =
(chrono::Utc::now() - chrono::Duration::hours(KEEP_HOURS + 1)).to_rfc3339();
let mut pending = entry("pending");
pending.spooled_at = old.spooled_at.clone();
rewrite(&path, &[old, pending]).unwrap();
let held = read(&path);
assert_eq!(held.len(), 1);
assert_eq!(
held[0].delivery_id, "pending",
"an unconfirmed entry is never dropped, however old"
);
let _ = std::fs::remove_dir_all(&dir);
}
}