use std::io::{self, BufRead};
use std::path::{Path, PathBuf};
use crate::claude_attach::{perform_attach, AttachError, AttachRequest, ControlTransport};
pub const DETACH_SENTINELS: [&str; 2] = ["\x1b_cc-daemon-detach\x1b\\", "\x1b_cc-detach-msg;"];
pub const PROJECTS_DIR_ENV: &str = "FNO_CLAUDE_PROJECTS_DIR";
#[derive(Debug)]
pub enum DriveError {
UnsafeText,
Attach(AttachError),
Io(String),
NotDelivered,
}
impl std::fmt::Display for DriveError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
DriveError::UnsafeText => write!(f, "refused: text contains a detach sentinel"),
DriveError::Attach(e) => write!(f, "attach failed: {e}"),
DriveError::Io(s) => write!(f, "io error: {s}"),
DriveError::NotDelivered => {
write!(f, "injected but no confirming assistant turn appeared")
}
}
}
}
impl std::error::Error for DriveError {}
pub fn claude_projects_dir() -> PathBuf {
if let Some(v) = std::env::var_os(PROJECTS_DIR_ENV) {
return PathBuf::from(v);
}
let home = std::env::var_os("HOME")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("."));
home.join(".claude").join("projects")
}
fn is_session_uuid(s: &str) -> bool {
let parts: Vec<&str> = s.split('-').collect();
parts.len() == 5
&& [8, 4, 4, 4, 12]
.iter()
.zip(&parts)
.all(|(&n, p)| p.len() == n && p.bytes().all(|b| b.is_ascii_hexdigit()))
}
pub fn find_transcript(session_uuid: &str) -> Option<PathBuf> {
if !is_session_uuid(session_uuid) {
return None;
}
let base = claude_projects_dir();
let entries = std::fs::read_dir(&base).ok()?;
for entry in entries.flatten() {
let candidate = entry.path().join(format!("{session_uuid}.jsonl"));
if candidate.exists() {
return Some(candidate);
}
}
None
}
pub fn contains_detach_sentinel(text: &str) -> bool {
DETACH_SENTINELS.iter().any(|s| text.contains(s))
}
pub struct FnoMail<'a> {
pub from: &'a str,
pub harness: &'a str,
pub model: &'a str,
pub node: Option<&'a str>,
pub to: Option<&'a str>,
pub id: Option<&'a str>,
pub reply_to: Option<&'a str>,
}
pub fn fno_mail_open(m: &FnoMail) -> String {
let mut s = format!(
"<fno_mail from=\"{}\" harness=\"{}\" model=\"{}\"",
m.from, m.harness, m.model
);
if let Some(node) = m.node {
s.push_str(&format!(" node=\"{node}\""));
}
if let Some(to) = m.to {
s.push_str(&format!(" to=\"{to}\""));
}
if let Some(id) = m.id {
s.push_str(&format!(" id=\"{id}\""));
}
if let Some(reply_to) = m.reply_to {
s.push_str(&format!(" reply_to=\"{reply_to}\""));
}
s.push('>');
s
}
pub fn wrap_fno_mail(m: &FnoMail, text: &str) -> String {
format!("{}\n{}\n</fno_mail>", fno_mail_open(m), text)
}
pub fn build_reply_request(
short: &str,
text: &str,
auth: Option<&str>,
mail: Option<&FnoMail>,
) -> Result<String, DriveError> {
if contains_detach_sentinel(text) {
return Err(DriveError::UnsafeText);
}
let body = match mail {
Some(m) => wrap_fno_mail(m, text),
None => text.to_string(),
};
let mut obj = serde_json::Map::new();
obj.insert("op".into(), "reply".into());
obj.insert("short".into(), short.into());
obj.insert("text".into(), body.into());
if let Some(a) = auth {
obj.insert("auth".into(), a.into());
}
let mut line = serde_json::Value::Object(obj).to_string();
line.push('\n');
Ok(line)
}
pub fn inject_reply<T: ControlTransport>(
t: &mut T,
short: &str,
text: &str,
auth: Option<&str>,
mail: Option<&FnoMail>,
) -> Result<(), DriveError> {
let line = build_reply_request(short, text, auth, mail)?;
t.send_line(&line)
.map_err(|e| DriveError::Io(e.to_string()))
}
pub fn transcript_len(path: &Path) -> u64 {
std::fs::metadata(path).map(|m| m.len()).unwrap_or(0)
}
pub fn confirm_marker_after(path: &Path, marker: &str, since_byte: u64) -> io::Result<bool> {
let mut file = std::fs::File::open(path)?;
use std::io::Seek;
file.seek(io::SeekFrom::Start(since_byte))?;
let reader = io::BufReader::new(file);
for line in reader.lines() {
let line = line?;
if line.contains(marker) && line_is_assistant(&line) {
return Ok(true);
}
}
Ok(false)
}
fn line_is_assistant(line: &str) -> bool {
let Ok(v) = serde_json::from_str::<serde_json::Value>(line) else {
return false;
};
if v.get("type").and_then(serde_json::Value::as_str) == Some("assistant") {
return true;
}
v.get("message")
.and_then(|m| m.get("role"))
.and_then(serde_json::Value::as_str)
== Some("assistant")
}
pub struct DriveTurn<'a> {
pub text: &'a str,
pub marker: &'a str,
pub mail: Option<&'a FnoMail<'a>>,
}
pub fn drive_and_confirm<T: ControlTransport>(
transport: &mut T,
attach: &AttachRequest,
session_uuid: &str,
turn: &DriveTurn,
attempts: u32,
interval: std::time::Duration,
) -> Result<(), DriveError> {
perform_attach(transport, attach).map_err(DriveError::Attach)?;
let transcript = find_transcript(session_uuid)
.ok_or_else(|| DriveError::Io(format!("no transcript for session {session_uuid}")))?;
let baseline = transcript_len(&transcript);
inject_reply(
transport,
&attach.short,
turn.text,
attach.auth.as_deref(),
turn.mail,
)?;
for _ in 0..attempts.max(1) {
if confirm_marker_after(&transcript, turn.marker, baseline)
.map_err(|e| DriveError::Io(e.to_string()))?
{
return Ok(());
}
std::thread::sleep(interval);
}
Err(DriveError::NotDelivered)
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
struct Fake {
sent: Vec<String>,
}
impl ControlTransport for Fake {
fn send_line(&mut self, line: &str) -> io::Result<()> {
self.sent.push(line.to_string());
Ok(())
}
fn recv_line(&mut self) -> io::Result<Option<String>> {
Ok(None)
}
}
fn tmpdir(tag: &str) -> PathBuf {
let p = std::env::temp_dir().join(format!(
"fno-drive-{}-{}-{}",
tag,
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&p).unwrap();
p
}
const UUID: &str = "a1b2c3d4-1111-2222-3333-444455556666";
#[test]
fn is_session_uuid_validates_shape() {
assert!(is_session_uuid(UUID));
assert!(!is_session_uuid("a1b2c3d4")); assert!(!is_session_uuid("not-a-uuid"));
assert!(!is_session_uuid("g1b2c3d4-1111-2222-3333-444455556666")); }
fn from_orchestrator() -> FnoMail<'static> {
FnoMail {
from: "7d1f8bdc",
harness: "claude-code",
model: "opus-4.8",
node: Some("x-26df"),
to: None,
id: None,
reply_to: None,
}
}
#[test]
fn fno_mail_open_is_lowercase_quoted_attrs() {
assert_eq!(
fno_mail_open(&from_orchestrator()),
"<fno_mail from=\"7d1f8bdc\" harness=\"claude-code\" model=\"opus-4.8\" node=\"x-26df\">"
);
let directed = FnoMail {
from: "7d1f8bdc",
harness: "claude-code",
model: "opus-4.8",
node: None,
to: Some("claude-ee99ff00"),
id: None,
reply_to: None,
};
assert_eq!(
fno_mail_open(&directed),
"<fno_mail from=\"7d1f8bdc\" harness=\"claude-code\" model=\"opus-4.8\" to=\"claude-ee99ff00\">"
);
}
#[test]
fn fno_mail_open_renders_reply_to_last_when_present() {
let reply = FnoMail {
from: "7d1f8bdc",
harness: "claude-code",
model: "opus-4.8",
node: None,
to: Some("claude-e5f6a7b8"),
id: None,
reply_to: Some("msg-0091f3"),
};
assert_eq!(
fno_mail_open(&reply),
"<fno_mail from=\"7d1f8bdc\" harness=\"claude-code\" model=\"opus-4.8\" to=\"claude-e5f6a7b8\" reply_to=\"msg-0091f3\">"
);
}
#[test]
fn fno_mail_open_renders_id_last_but_one_before_reply_to() {
let with_id = FnoMail {
from: "7d1f8bdc",
harness: "claude-code",
model: "opus-4.8",
node: None,
to: Some("claude-e5f6a7b8"),
id: Some("msg-abc123"),
reply_to: Some("msg-0091f3"),
};
assert_eq!(
fno_mail_open(&with_id),
"<fno_mail from=\"7d1f8bdc\" harness=\"claude-code\" model=\"opus-4.8\" to=\"claude-e5f6a7b8\" id=\"msg-abc123\" reply_to=\"msg-0091f3\">"
);
let fresh = FnoMail {
from: "7d1f8bdc",
harness: "claude-code",
model: "opus-4.8",
node: None,
to: None,
id: Some("msg-abc123"),
reply_to: None,
};
assert_eq!(
fno_mail_open(&fresh),
"<fno_mail from=\"7d1f8bdc\" harness=\"claude-code\" model=\"opus-4.8\" id=\"msg-abc123\">"
);
}
#[test]
fn wrap_fno_mail_is_a_paired_envelope() {
let wrapped = wrap_fno_mail(&from_orchestrator(), "ship it");
assert_eq!(
wrapped,
"<fno_mail from=\"7d1f8bdc\" harness=\"claude-code\" model=\"opus-4.8\" node=\"x-26df\">\nship it\n</fno_mail>"
);
}
#[test]
fn build_reply_request_untagged_when_no_mail() {
let line = build_reply_request("a1b2c3d4", "hello world", Some("deadbeef"), None).unwrap();
assert!(line.ends_with('\n'));
let v: serde_json::Value = serde_json::from_str(line.trim()).unwrap();
assert_eq!(v["op"], "reply");
assert_eq!(v["short"], "a1b2c3d4");
assert_eq!(v["text"], "hello world");
assert_eq!(v["auth"], "deadbeef");
}
#[test]
fn build_reply_request_wraps_fno_mail() {
let from = from_orchestrator();
let line = build_reply_request("a1b2c3d4", "ship it MARKER42", None, Some(&from)).unwrap();
let v: serde_json::Value = serde_json::from_str(line.trim()).unwrap();
let text = v["text"].as_str().unwrap();
assert_eq!(
text,
"<fno_mail from=\"7d1f8bdc\" harness=\"claude-code\" model=\"opus-4.8\" node=\"x-26df\">\nship it MARKER42\n</fno_mail>"
);
}
#[test]
fn build_reply_request_omits_auth() {
let line = build_reply_request("a1b2c3d4", "hi", None, None).unwrap();
let v: serde_json::Value = serde_json::from_str(line.trim()).unwrap();
assert!(v.get("auth").is_none());
}
#[test]
fn build_reply_request_refuses_detach_sentinel() {
let evil = format!("hi {}", DETACH_SENTINELS[0]);
assert!(matches!(
build_reply_request("a1b2c3d4", &evil, None, None),
Err(DriveError::UnsafeText)
));
assert!(matches!(
build_reply_request("a1b2c3d4", DETACH_SENTINELS[1], None, None),
Err(DriveError::UnsafeText)
));
}
#[test]
fn inject_reply_writes_one_line_and_guards_sentinels() {
let mut t = Fake { sent: Vec::new() };
inject_reply(&mut t, "a1b2c3d4", "ping MARKER42", None, None).unwrap();
assert_eq!(t.sent.len(), 1);
assert!(t.sent[0].contains("\"op\":\"reply\""));
assert!(t.sent[0].contains("MARKER42"));
let mut t2 = Fake { sent: Vec::new() };
assert!(inject_reply(&mut t2, "a1b2c3d4", DETACH_SENTINELS[0], None, None).is_err());
assert!(t2.sent.is_empty(), "must not write unsafe text");
}
#[test]
fn find_transcript_by_uuid_across_project_dirs() {
let base = tmpdir("find");
std::env::set_var(PROJECTS_DIR_ENV, &base);
let proj = base.join("-Users-x-code-proj");
std::fs::create_dir_all(&proj).unwrap();
let t = proj.join(format!("{UUID}.jsonl"));
std::fs::write(&t, b"{}\n").unwrap();
assert_eq!(find_transcript(UUID), Some(t));
assert_eq!(
find_transcript("ffffffff-0000-0000-0000-000000000000"),
None
);
assert_eq!(find_transcript("bad-uuid"), None);
std::env::remove_var(PROJECTS_DIR_ENV);
std::fs::remove_dir_all(&base).ok();
}
#[test]
fn confirm_marker_after_counts_only_assistant_turns() {
let dir = tmpdir("confirm");
let path = dir.join("t.jsonl");
let mut f = std::fs::File::create(&path).unwrap();
writeln!(
f,
r#"{{"type":"user","message":{{"role":"user","content":"older"}}}}"#
)
.unwrap();
let baseline = transcript_len(&path);
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap();
writeln!(
f,
r#"{{"type":"user","message":{{"role":"user","content":"drive MARKER42"}}}}"#
)
.unwrap();
assert!(!confirm_marker_after(&path, "MARKER42", baseline).unwrap());
writeln!(
f,
r#"{{"type":"assistant","message":{{"role":"assistant","content":"done MARKER42"}}}}"#
)
.unwrap();
assert!(confirm_marker_after(&path, "MARKER42", baseline).unwrap());
let end = transcript_len(&path);
assert!(!confirm_marker_after(&path, "MARKER42", end).unwrap());
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn confirm_marker_absent_assistant_is_false() {
let dir = tmpdir("absent");
let path = dir.join("t.jsonl");
let mut f = std::fs::File::create(&path).unwrap();
writeln!(
f,
r#"{{"type":"assistant","message":{{"role":"assistant","content":"no token here"}}}}"#
)
.unwrap();
assert!(!confirm_marker_after(&path, "MARKER42", 0).unwrap());
std::fs::remove_dir_all(&dir).ok();
}
}