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