use std::io::Write;
use std::path::{Path, PathBuf};
use anyhow::{bail, Context, Result};
use serde_json::Value;
pub(crate) const CHANNEL_DIR_NAME: &str = "csift-channel";
pub(crate) fn channel_dir(sidecar_dir: &Path) -> PathBuf {
sidecar_dir.join(CHANNEL_DIR_NAME)
}
pub(crate) fn messages_dir(root: &Path) -> PathBuf {
root.join("messages")
}
pub(crate) fn message_path(root: &Path, id: &str) -> Result<PathBuf> {
validate_message_id(id)?;
Ok(messages_dir(root).join(format!("{id}.json")))
}
pub(crate) fn outbox_path(root: &Path) -> PathBuf {
root.join("outbox.jsonl")
}
pub(crate) fn inbox_path(root: &Path, lane: &str) -> Result<PathBuf> {
validate_lane_id(lane)?;
Ok(root.join("inbox").join(format!("{lane}.jsonl")))
}
pub(crate) fn ledger_path(root: &Path, lane: &str) -> Result<PathBuf> {
validate_lane_id(lane)?;
Ok(root.join("ledger").join(format!("{lane}.jsonl")))
}
pub(crate) fn armed_path(root: &Path, lane: &str) -> Result<PathBuf> {
validate_lane_id(lane)?;
Ok(root.join("armed").join(format!("{lane}.json")))
}
pub(crate) fn is_lane_id(s: &str) -> bool {
crate::path::is_uuid(s) || crate::path::is_subagent_id(s)
}
pub(crate) fn validate_lane_id(s: &str) -> Result<()> {
if !is_lane_id(s) {
bail!(
"not a lane id: `{s}` - a lane is a top-level session uuid, a bare `a<16 hex>` \
agent id, or a teammate id `a<Name>-<16 hex>`"
);
}
Ok(())
}
pub(crate) fn is_message_id(s: &str) -> bool {
s.len() == 16
&& s.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
}
pub(crate) fn validate_message_id(s: &str) -> Result<()> {
if !is_message_id(s) {
bail!("not a message id: `{s}` - a message id is 16 lowercase hex characters");
}
Ok(())
}
pub(crate) fn append_line(path: &Path, line: &str) -> Result<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)
.with_context(|| format!("creating channel directory {}", parent.display()))?;
}
let mut buf = String::with_capacity(line.len() + 1);
buf.push_str(line);
buf.push('\n');
let mut f = std::fs::OpenOptions::new()
.append(true)
.create(true)
.open(path)
.with_context(|| format!("opening {} for append", path.display()))?;
f.write_all(buf.as_bytes())
.with_context(|| format!("appending to {}", path.display()))
}
pub(crate) fn write_atomic(path: &Path, contents: &str) -> Result<()> {
let parent = path
.parent()
.ok_or_else(|| anyhow::anyhow!("no parent directory for {}", path.display()))?;
std::fs::create_dir_all(parent)
.with_context(|| format!("creating channel directory {}", parent.display()))?;
let stem = path
.file_name()
.and_then(|s| s.to_str())
.unwrap_or("channel");
static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let tmp = parent.join(format!(".{stem}.tmp-{}-{seq}", std::process::id()));
{
let mut f =
std::fs::File::create(&tmp).with_context(|| format!("creating {}", tmp.display()))?;
f.write_all(contents.as_bytes())
.with_context(|| format!("writing {}", tmp.display()))?;
}
std::fs::rename(&tmp, path).with_context(|| {
format!(
"renaming {} over {} (atomic marker rewrite)",
tmp.display(),
path.display()
)
})
}
pub(crate) fn read_jsonl(path: &Path) -> Result<(Vec<Value>, usize)> {
let Ok(raw) = std::fs::read_to_string(path) else {
return Ok((Vec::new(), 0));
};
let mut out = Vec::new();
let mut skipped = 0usize;
for line in raw.lines() {
if line.trim().is_empty() {
continue;
}
match serde_json::from_str::<Value>(line) {
Ok(v) => out.push(v),
Err(_) => skipped += 1,
}
}
Ok((out, skipped))
}
pub(crate) fn read_json_object(path: &Path) -> Option<Value> {
let raw = std::fs::read_to_string(path).ok()?;
let v: Value = serde_json::from_str(&raw).ok()?;
v.is_object().then_some(v)
}