use std::io::{self, BufRead, Read, Seek};
use std::path::Path;
use std::time::Duration;
use crate::claude_attach::{perform_attach, AttachRequest, UnixControlTransport};
use crate::claude_drive::{contains_detach_sentinel, find_transcript, transcript_len, DriveError};
use crate::claude_roster::{read_control_key, ClaudeRoster};
pub const DEFAULT_ATTEMPTS: u32 = 40;
pub const DEFAULT_INTERVAL_MS: u64 = 250;
const CR_SETTLE_MS: u64 = 800;
const CR_RESUBMIT_EVERY: u32 = 8;
#[derive(Debug, PartialEq, Clone, Copy)]
pub enum MailInjectProvider {
Claude,
Codex,
}
const PROVIDER_AXIS_TOMBSTONE: &str = concat!(
"--provider was split at the axis rename: the CLI binary is --harness/-H; ",
"a model vendor is only routable at spawn ",
"(`fno agents spawn --provider <vendor> --model <m>`).",
);
#[derive(Debug, PartialEq)]
pub struct MailInjectArgs {
pub session: String,
pub provider: MailInjectProvider,
pub attempts: u32,
pub interval_ms: u64,
}
pub fn parse_args(rest: &[String]) -> Result<MailInjectArgs, (i32, String)> {
let mut session: Option<String> = None;
let mut provider = MailInjectProvider::Claude;
let mut attempts = DEFAULT_ATTEMPTS;
let mut interval_ms = DEFAULT_INTERVAL_MS;
let mut it = rest.iter();
while let Some(a) = it.next() {
match a.as_str() {
"--session" => {
session = Some(
it.next()
.ok_or((2, "mail-inject: --session needs a value".to_string()))?
.to_string(),
);
}
"--harness" | "-H" => {
provider = match it.next().map(String::as_str) {
Some("claude") => MailInjectProvider::Claude,
Some("codex") => MailInjectProvider::Codex,
_ => {
return Err((
2,
"mail-inject: --harness must be claude or codex".to_string(),
))
}
};
}
"--provider" => return Err((2, PROVIDER_AXIS_TOMBSTONE.to_string())),
"--attempts" => {
attempts = it.next().and_then(|v| v.parse().ok()).ok_or((
2,
"mail-inject: --attempts needs a positive integer".to_string(),
))?;
}
"--interval-ms" => {
interval_ms = it.next().and_then(|v| v.parse().ok()).ok_or((
2,
"mail-inject: --interval-ms needs a positive integer".to_string(),
))?;
}
other => {
return Err((2, format!("mail-inject: unknown flag: {other}")));
}
}
}
let session = session.ok_or((2, "mail-inject: --session is required".to_string()))?;
Ok(MailInjectArgs {
session,
provider,
attempts,
interval_ms,
})
}
pub fn outcome_json(delivered: bool, reason: &str) -> String {
serde_json::json!({ "delivered": delivered, "reason": reason }).to_string()
}
pub fn outcome_exit(delivered: bool) -> i32 {
i32::from(!delivered)
}
fn emit(delivered: bool, reason: &str) -> i32 {
println!("{}", outcome_json(delivered, reason));
outcome_exit(delivered)
}
const PASTE_BEGIN: &str = "\x1b[200~";
const PASTE_END: &str = "\x1b[201~";
fn inject_with_submit<T: crate::claude_attach::ControlTransport>(
transport: &mut T,
text: &str,
settle: Duration,
) -> Result<(), DriveError> {
if contains_detach_sentinel(text) {
return Err(DriveError::UnsafeText);
}
transport
.send_line(&format!("{PASTE_BEGIN}{text}{PASTE_END}"))
.map_err(|e| DriveError::Io(e.to_string()))?;
std::thread::sleep(settle);
transport
.send_line("\r")
.map_err(|e| DriveError::Io(e.to_string()))
}
fn confirm_with_cr_retry<T: crate::claude_attach::ControlTransport>(
transport: &mut T,
attempts: u32,
interval: Duration,
mut confirmed: impl FnMut() -> bool,
) -> Result<(), &'static str> {
for i in 0..attempts.max(1) {
if confirmed() {
return Ok(());
}
std::thread::sleep(interval);
if (i + 1) % CR_RESUBMIT_EVERY == 0 {
let _ = transport.send_line("\r");
}
}
Err("not-confirmed")
}
fn escaped_marker(marker: &str) -> String {
let s = serde_json::to_string(marker).unwrap_or_default();
s.strip_prefix('"')
.and_then(|s| s.strip_suffix('"'))
.unwrap_or("")
.to_string()
}
fn confirm_content_after(path: &Path, marker: &str, since_byte: u64) -> io::Result<bool> {
let escaped = escaped_marker(marker);
if escaped.is_empty() {
return Ok(false);
}
let mut file = std::fs::File::open(path)?;
file.seek(io::SeekFrom::Start(since_byte))?;
for line in io::BufReader::new(file).lines() {
if line?.contains(&escaped) {
return Ok(true);
}
}
Ok(false)
}
pub fn deliver_via_control_sock(
session: &str,
text: &str,
attempts: u32,
interval_ms: u64,
) -> Result<(), &'static str> {
let roster = ClaudeRoster::load_default().map_err(|_| "not-live")?;
let worker = roster.find(session).ok_or("not-live")?;
let sock = worker.resolve_control_sock().ok_or("not-live")?;
let short = worker.short_id().to_string();
let auth = read_control_key();
let transcript = find_transcript(&worker.session_id).ok_or("no-transcript")?;
let mut transport = UnixControlTransport::connect(&sock).map_err(|_| "io-error")?;
if perform_attach(
&mut transport,
&AttachRequest::for_frame_stream(short.clone(), auth.clone()),
)
.is_err()
{
return Err("attach-failed");
}
let baseline = transcript_len(&transcript);
let marker = text.lines().next().unwrap_or(text);
inject_with_submit(&mut transport, text, Duration::from_millis(CR_SETTLE_MS)).map_err(|e| {
match e {
DriveError::UnsafeText => "unsafe-text",
_ => "io-error",
}
})?;
confirm_with_cr_retry(
&mut transport,
attempts,
Duration::from_millis(interval_ms),
|| confirm_content_after(&transcript, marker, baseline).unwrap_or(false),
)
}
pub async fn run_mail_inject(rest: &[String]) -> i32 {
let args = match parse_args(rest) {
Ok(a) => a,
Err((code, msg)) => {
eprintln!("{msg}");
return code;
}
};
let mut text = String::new();
if let Err(e) = std::io::stdin().read_to_string(&mut text) {
eprintln!("mail-inject: reading stdin: {e}");
return emit(false, "io-error");
}
let result = match args.provider {
MailInjectProvider::Claude => {
deliver_via_control_sock(&args.session, &text, args.attempts, args.interval_ms)
}
MailInjectProvider::Codex => {
crate::codex_inject::deliver_via_codex_daemon(&args.session, &text).await
}
};
match result {
Ok(()) => emit(true, "delivered"),
Err(reason) => emit(false, reason),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::claude_attach::ControlTransport;
use crate::claude_drive::DETACH_SENTINELS;
use std::fs::{File, OpenOptions};
use std::io::{self, Write};
use std::path::PathBuf;
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 argv(parts: &[&str]) -> Vec<String> {
parts.iter().map(|s| s.to_string()).collect()
}
fn tmp_transcript(tag: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!("mailinj-{}-{}", tag, std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
dir.join("t.jsonl")
}
#[test]
fn inject_with_submit_bracketed_pastes_then_separate_cr() {
let mut t = Fake { sent: Vec::new() };
let envelope = "<fno_mail from=\"a1b2c3d4\" node=\"x-178e\">\nhi MARKER\n</fno_mail>";
inject_with_submit(&mut t, envelope, Duration::ZERO).unwrap();
assert_eq!(
t.sent,
vec![
format!("{PASTE_BEGIN}{envelope}{PASTE_END}"),
"\r".to_string()
]
);
assert!(t.sent[0].contains(envelope), "envelope pasted verbatim");
assert!(
!t.sent[0].contains("\"op\""),
"envelope must be raw bytes, not a JSON op"
);
assert!(
!t.sent[0].contains("auth"),
"raw paste must never carry the control auth key"
);
}
#[test]
fn inject_with_submit_refuses_unsafe_envelope_and_writes_nothing() {
let mut t = Fake { sent: Vec::new() };
let err = inject_with_submit(&mut t, DETACH_SENTINELS[0], Duration::ZERO);
assert!(matches!(err, Err(DriveError::UnsafeText)));
assert!(t.sent.is_empty(), "unsafe envelope must not paste or CR");
}
#[test]
fn busy_recipient_gets_raw_paste_then_retried_crs() {
let mut t = Fake { sent: Vec::new() };
inject_with_submit(&mut t, "hi MARKER", Duration::ZERO).unwrap();
let attempts = 2 * CR_RESUBMIT_EVERY; let r = confirm_with_cr_retry(&mut t, attempts, Duration::ZERO, || false);
assert_eq!(r, Err("not-confirmed"));
assert_eq!(t.sent.len() as u32, 2 + attempts / CR_RESUBMIT_EVERY);
for line in &t.sent[1..] {
assert_eq!(line, "\r");
}
}
#[test]
fn confirm_stops_on_landing_without_extra_cr() {
let mut t = Fake { sent: Vec::new() };
let mut calls = 0;
let r = confirm_with_cr_retry(&mut t, 40, Duration::ZERO, || {
calls += 1;
calls >= 2
});
assert_eq!(r, Ok(()));
assert!(
t.sent.is_empty(),
"landing before a resubmit window sends no CR"
);
}
#[test]
fn content_confirm_rejects_growth_and_accepts_the_landed_envelope() {
let path = tmp_transcript("content");
let mut f = File::create(&path).unwrap();
writeln!(
f,
r#"{{"type":"user","message":{{"role":"user","content":"older"}}}}"#
)
.unwrap();
let baseline = transcript_len(&path);
let marker = "<fno_mail from=\"a1b2c3d4\" node=\"x-178e\">";
let mut f = OpenOptions::new().append(true).open(&path).unwrap();
writeln!(
f,
r#"{{"type":"assistant","message":{{"role":"assistant","content":"streaming something else"}}}}"#
)
.unwrap();
assert!(
!confirm_content_after(&path, marker, baseline).unwrap(),
"growth without the marker must not confirm"
);
writeln!(
f,
r#"{{"type":"user","message":{{"role":"user","content":"{}\nhi\n</fno_mail>"}}}}"#,
escaped_marker(marker)
)
.unwrap();
assert!(
confirm_content_after(&path, marker, baseline).unwrap(),
"the landed envelope confirms delivery"
);
std::fs::remove_dir_all(path.parent().unwrap()).ok();
}
#[test]
fn parse_args_requires_session() {
assert_eq!(parse_args(&[]).unwrap_err().0, 2);
assert_eq!(
parse_args(&argv(&["--attempts", "5"])).unwrap_err().0,
2,
"no --session is an error even with other flags"
);
}
#[test]
fn parse_args_defaults_and_overrides() {
let a = parse_args(&argv(&["--session", "a1b2c3d4"])).unwrap();
assert_eq!(a.session, "a1b2c3d4");
assert_eq!(a.provider, MailInjectProvider::Claude);
assert_eq!(a.attempts, DEFAULT_ATTEMPTS);
assert_eq!(a.interval_ms, DEFAULT_INTERVAL_MS);
let b = parse_args(&argv(&[
"--session",
"a1b2c3d4-1111-2222-3333-444455556666",
"--attempts",
"3",
"--interval-ms",
"10",
]))
.unwrap();
assert_eq!(b.session, "a1b2c3d4-1111-2222-3333-444455556666");
assert_eq!(b.attempts, 3);
assert_eq!(b.interval_ms, 10);
}
#[test]
fn parse_args_harness_defaults_claude_and_accepts_codex() {
let d = parse_args(&argv(&["--session", "x"])).unwrap();
assert_eq!(d.provider, MailInjectProvider::Claude);
let c = parse_args(&argv(&["--session", "x", "--harness", "codex"])).unwrap();
assert_eq!(c.provider, MailInjectProvider::Codex);
let h = parse_args(&argv(&["--session", "x", "-H", "codex"])).unwrap();
assert_eq!(h.provider, MailInjectProvider::Codex);
assert_eq!(
parse_args(&argv(&["--session", "x", "--harness", "gemini"]))
.unwrap_err()
.0,
2
);
}
#[test]
fn parse_args_provider_is_the_axis_rename_tombstone() {
let err = parse_args(&argv(&["--session", "x", "--provider", "codex"])).unwrap_err();
assert_eq!(err.0, 2);
assert!(
err.1.contains("--harness/-H"),
"tombstone points at --harness: {err:?}"
);
}
#[test]
fn parse_args_rejects_unknown_flag_and_missing_value() {
assert_eq!(parse_args(&argv(&["--nope"])).unwrap_err().0, 2);
assert_eq!(parse_args(&argv(&["--session"])).unwrap_err().0, 2);
assert_eq!(
parse_args(&argv(&["--session", "x", "--attempts", "notnum"]))
.unwrap_err()
.0,
2
);
}
#[test]
fn outcome_json_is_the_python_contract() {
let v: serde_json::Value = serde_json::from_str(&outcome_json(true, "delivered")).unwrap();
assert_eq!(v["delivered"], true);
assert_eq!(v["reason"], "delivered");
let w: serde_json::Value = serde_json::from_str(&outcome_json(false, "not-live")).unwrap();
assert_eq!(w["delivered"], false);
assert_eq!(w["reason"], "not-live");
}
#[test]
fn outcome_exit_maps_delivered_to_zero() {
assert_eq!(outcome_exit(true), 0);
assert_eq!(outcome_exit(false), 1);
}
}