1use std::path::PathBuf;
8
9use crate::mail_route::{Caller, LiveSessions, Unresolved};
10use crate::mailbox::{
11 local_machine_name, mail_root, Envelope, MailAddress, MailKind, Mailbox, ReplyVia,
12};
13use crate::HarnessHomes;
14
15pub const EXIT_REFUSED: i32 = 2;
17pub const EXIT_UNKNOWN: i32 = 3;
19pub const EXIT_STORED: i32 = 4;
21pub const EXIT_FAILED: i32 = 5;
23
24pub struct Outcome {
26 pub code: i32,
29 pub text: String,
31}
32
33impl Outcome {
34 pub fn new(code: i32, text: impl Into<String>) -> Self {
36 Self {
37 code,
38 text: text.into(),
39 }
40 }
41}
42
43#[derive(Debug, Clone, Default)]
45pub struct SendOptions {
46 pub subject: Option<String>,
48 pub in_reply_to: Option<String>,
50 pub notify_when_idle: bool,
52 pub queue: bool,
54 pub idempotency_key: Option<String>,
56}
57
58pub async fn send(
62 homes: &HarnessHomes,
63 caller: &Caller,
64 to: &str,
65 body: &str,
66 options: SendOptions,
67) -> std::io::Result<Outcome> {
68 let SendOptions {
69 subject,
70 in_reply_to,
71 notify_when_idle,
72 queue,
73 idempotency_key,
74 } = options;
75 if let Some(machine) = remote_machine(to) {
76 let request = serde_json::json!({
77 "op": "send",
78 "from": caller.address.to_string(),
79 "from_name": caller.name,
80 "to": to,
81 "body": body,
82 "subject": subject,
83 "in_reply_to": in_reply_to,
84 "notify_when_idle": notify_when_idle,
85 "queue": queue,
86 "id": idempotency_key,
87 });
88 return Ok(remote_send(&machine, &request));
89 }
90 deliver(
91 homes,
92 caller,
93 to,
94 body,
95 subject,
96 in_reply_to,
97 notify_when_idle,
98 queue,
99 idempotency_key,
100 )
101 .await
102}
103
104fn remote_machine(to: &str) -> Option<String> {
106 let local = local_machine_name();
107 let machine = match MailAddress::parse(to) {
108 Ok(address) => address.machine,
109 Err(_) => to.rsplit_once('@')?.1.to_string(),
110 };
111 (machine != local).then_some(machine)
112}
113
114fn remote_send(machine: &str, request: &serde_json::Value) -> Outcome {
116 match crate::mailbox::teams_mail(machine, request) {
117 Ok(answer) => Outcome::new(
118 answer["code"]
119 .as_i64()
120 .map_or(EXIT_FAILED, |code| code as i32),
121 answer["text"]
122 .as_str()
123 .or_else(|| answer["detail"].as_str())
124 .unwrap_or("the other machine answered nothing readable")
125 .to_string(),
126 ),
127 Err(detail) if detail.contains("no mail grant") => Outcome::new(
128 EXIT_REFUSED,
129 format!(
130 "Not sent: machine {machine} does not accept mail from you (no mail grant). Only \
131 an operator can grant it; tell your user. Don't retry."
132 ),
133 ),
134 Err(detail) => Outcome::new(
135 EXIT_FAILED,
136 format!("Not sent to machine {machine}: {detail}"),
137 ),
138 }
139}
140
141#[allow(clippy::too_many_arguments)]
144async fn deliver(
145 homes: &HarnessHomes,
146 caller: &Caller,
147 to: &str,
148 body: &str,
149 subject: Option<String>,
150 in_reply_to: Option<String>,
151 notify_when_idle: bool,
152 queue: bool,
153 idempotency_key: Option<String>,
154) -> std::io::Result<Outcome> {
155 use crate::mail_route::{deliver as route, door_for, Delivered, Refused};
156 let local_operator = to
161 .strip_prefix("operator:")
162 .and_then(|name| MailAddress::new(local_machine_name(), "operator", name).ok());
163 let (address, name) = match local_operator
164 .map(Ok)
165 .unwrap_or_else(|| MailAddress::parse(to))
166 {
167 Ok(address) if address.harness == "operator" || address.harness == "board" => {
169 let name = format!("{}@{}", address.session_id, address.machine);
170 (address, name)
171 }
172 _ => match LiveSessions::read(homes).resolve(to) {
173 Ok(session) => (session.address.clone(), session.name.clone()),
174 Err(Unresolved::Stale(message) | Unresolved::Unknown(message)) => {
175 return Ok(Outcome::new(EXIT_UNKNOWN, message));
176 }
177 },
178 };
179 if address == caller.address {
180 return Ok(Outcome::new(
181 EXIT_REFUSED,
182 "Not sent: that address is your own session.",
183 ));
184 }
185 let door = match door_for(homes, &address) {
186 Ok(door) => door,
187 Err(_) => {
188 return Ok(Outcome::new(
189 EXIT_UNKNOWN,
190 format!(
191 "{name} is no longer running. Nothing was sent. Run supercode message list for \
192 the live sessions."
193 ),
194 ))
195 }
196 };
197 let message_id = match &idempotency_key {
198 Some(key)
201 if key.len() == 26
202 && key.starts_with("m-")
203 && key[2..]
204 .bytes()
205 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()) =>
206 {
207 key.clone()
208 }
209 Some(key) => format!(
210 "m-{}",
211 &blake3::hash(format!("{}\0{key}", caller.address).as_bytes()).to_hex()[..24]
212 ),
213 None => crate::mailbox::new_message_id()?,
214 };
215 let sent_marker = sent_marker(&caller.address, &message_id);
216 if let Ok(previous) = std::fs::read_to_string(&sent_marker) {
217 return Ok(Outcome::new(
218 0,
219 format!(
220 "Already sent with --id {}: {previous} Nothing was sent again.",
221 idempotency_key.as_deref().unwrap_or_default()
222 ),
223 ));
224 }
225 let in_reply_to = match in_reply_to {
230 Some(answered) if answered.len() < FULL_ID_LEN => {
231 let found = Mailbox::open(&mail_root(), &caller.address)
232 .and_then(|mailbox| mailbox.find_prefix(&answered))
233 .unwrap_or_default();
234 match found.as_slice() {
235 [only] => Some(only.envelope.id.clone()),
236 [] => {
237 return Ok(Outcome::new(
238 EXIT_REFUSED,
239 format!(
240 "Not sent: no message in your mailbox has an id starting {answered}. \
241 Name the message you answer as its envelope shows it."
242 ),
243 ))
244 }
245 many => {
246 let ids: Vec<&str> = many
247 .iter()
248 .map(|stored| stored.envelope.id.as_str())
249 .collect();
250 return Ok(Outcome::new(
251 EXIT_REFUSED,
252 format!(
253 "Not sent: {answered} names {} messages in your mailbox ({}). Give \
254 more of its id.",
255 ids.len(),
256 ids.join(", ")
257 ),
258 ));
259 }
260 }
261 }
262 other => other,
263 };
264 if let Some(answered) = in_reply_to.as_deref().filter(|id| id.starts_with("q-")) {
265 return crate::mail_question::reply(caller, &address, answered, body).await;
266 }
267 if let Some(answered) = in_reply_to.as_deref() {
268 if reply_chain_depth(&caller.address, &address, answered) >= REPLY_CHAIN_LIMIT {
269 return Ok(Outcome::new(
270 EXIT_REFUSED,
271 format!(
272 "Not sent: you and {name} have answered each other {REPLY_CHAIN_LIMIT} times \
273 in a row. If you are trading acknowledgements or status, stop; to go on, \
274 send a new message (without --re), or tell your user."
275 ),
276 ));
277 }
278 }
279 let mut envelope = Envelope::new(
280 caller.address.clone(),
281 caller.name.clone(),
282 MailKind::Peer,
283 ReplyVia::Command,
284 body,
285 )?;
286 envelope.subject = subject.as_deref().map(|value| {
287 value
288 .lines()
289 .next()
290 .unwrap_or("")
291 .chars()
292 .take(200)
293 .collect()
294 });
295 envelope.id = message_id.clone();
296 envelope.thread = crate::mailbox::thread_of_reply(in_reply_to.as_deref());
297 envelope.in_reply_to = in_reply_to;
298 let tier = door.name();
299 let idle_note = if notify_when_idle {
300 " Subscribed: one idle notice reaches you when its next turn ends."
301 } else {
302 ""
303 };
304 let (code, what) = match route(&envelope, &address, &door, !queue, notify_when_idle).await {
305 Err(detail) => {
306 return Ok(Outcome::new(
307 EXIT_FAILED,
308 format!("Not sent to {name}: {detail}"),
309 ))
310 }
311 Ok(Err(Refused::CannotQueueNative)) => {
312 return Ok(Outcome::new(
313 EXIT_REFUSED,
314 format!(
315 "Not sent: {name} is idle, and a Claude session always starts a turn when a \
316 message arrives, so --queue cannot hold it. Send without --queue to wake it."
317 ),
318 ))
319 }
320 Ok(Err(Refused::TooLong(bytes))) => {
321 return Ok(Outcome::new(
322 EXIT_REFUSED,
323 format!(
324 "Not sent: the message is {bytes} bytes; the limit is {}. Write it to a file \
325 and send its path instead.",
326 crate::mail_route::MAX_RELAYED_BYTES
327 ),
328 ))
329 }
330 Ok(Ok(Delivered::Steered)) => (0, "steered into its running turn"),
331 Ok(Ok(Delivered::Started)) => (0, "it was idle, so the message started a turn"),
332 Ok(Ok(Delivered::Native { busy: true })) => (0, "it will read it at its next tool call"),
333 Ok(Ok(Delivered::Native { busy: false })) => {
334 (0, "it was idle, so the message starts its next turn")
335 }
336 Ok(Ok(Delivered::Hooked)) => (
337 0,
338 "its hook points it at the message at its next tool call, or when its turn ends; an \
339 idle session sees it on its next turn",
340 ),
341 Ok(Ok(Delivered::HookWoken)) => (
342 0,
343 "it was idle, so its pane was told to read its mailbox; its hook shows it the message in \
344 that turn",
345 ),
346 Ok(Ok(Delivered::Queued)) => (
347 0,
348 "it is idle and --queue leaves it so; the message waits in its mailbox",
349 ),
350 Ok(Ok(Delivered::Operator)) => (0, "filed in its mailbox"),
351 Ok(Ok(Delivered::Stored)) => (
352 EXIT_STORED,
353 "Stored (not failed): it has no delivery door, so it sees this only if it runs \
354 supercode message inbox (a Codex session gets one with: supercode message setup \
355 codex). Don't resend, and don't wait for a reply",
356 ),
357 };
358 record_send(&caller.address, &address, &message_id);
359 let shown_id = crate::mailbox::short_id(&message_id);
360 let text = if code == 0 {
361 format!(
362 "sent to {name} ({}, {tier}): {what}. Message id {shown_id}.{idle_note} Don't poll; \
363 carry on.",
364 address.harness
365 )
366 } else {
367 format!(
368 "{name} ({}): {what}. Message id {shown_id}.{idle_note}",
369 address.harness
370 )
371 };
372 if code == 0 && idempotency_key.is_some() {
373 if let Some(parent) = sent_marker.parent() {
374 std::fs::create_dir_all(parent).ok();
375 }
376 std::fs::write(&sent_marker, &text).ok();
377 }
378 Ok(Outcome::new(code, text))
379}
380
381const REPLY_CHAIN_LIMIT: usize = 32;
386
387const FULL_ID_LEN: usize = 26;
389
390fn reply_chain_depth(a: &MailAddress, b: &MailAddress, answered: &str) -> usize {
394 let mut links = std::collections::HashMap::new();
395 for address in [a, b] {
396 let Ok(mailbox) = Mailbox::open(&mail_root(), address) else {
397 continue;
398 };
399 for stored in mailbox.list().unwrap_or_default() {
400 links.insert(
401 stored.envelope.id.clone(),
402 stored.envelope.in_reply_to.clone(),
403 );
404 }
405 }
406 let mut depth = 0;
407 let mut current = Some(answered.to_string());
408 while let Some(id) = current {
409 if depth > REPLY_CHAIN_LIMIT || !links.contains_key(&id) {
410 break;
411 }
412 depth += 1;
413 current = links.get(&id).cloned().flatten();
414 }
415 depth
416}
417
418fn send_log(sender: &MailAddress) -> PathBuf {
419 let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
420 mail_root().join("sent").join(&hash[..24]).join("log.jsonl")
421}
422
423fn record_send(sender: &MailAddress, to: &MailAddress, message_id: &str) {
424 let log = send_log(sender);
425 if let Some(parent) = log.parent() {
426 std::fs::create_dir_all(parent).ok();
427 }
428 if let Ok(mut file) = std::fs::OpenOptions::new()
429 .create(true)
430 .append(true)
431 .open(&log)
432 {
433 use std::io::Write as _;
434 let entry = serde_json::json!({"to": to.to_string(), "at_ms": now_ms(), "id": message_id});
435 writeln!(file, "{entry}").ok();
436 }
437}
438
439fn now_ms() -> u64 {
440 std::time::SystemTime::now()
441 .duration_since(std::time::UNIX_EPOCH)
442 .map(|elapsed| elapsed.as_millis() as u64)
443 .unwrap_or_default()
444}
445
446fn sent_marker(sender: &MailAddress, message_id: &str) -> PathBuf {
448 let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
449 mail_root().join("sent").join(&hash[..24]).join(message_id)
450}