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 pub notice: bool,
62}
63
64pub async fn send(
68 homes: &HarnessHomes,
69 caller: &Caller,
70 to: &str,
71 body: &str,
72 options: SendOptions,
73) -> std::io::Result<Outcome> {
74 let SendOptions {
75 subject,
76 in_reply_to,
77 notify_when_idle,
78 queue,
79 idempotency_key,
80 notice,
81 } = options;
82 if let Some(machine) = remote_machine(to) {
83 let request = serde_json::json!({
84 "op": "send",
85 "from": caller.address.to_string(),
86 "from_name": caller.name,
87 "to": to,
88 "body": body,
89 "subject": subject,
90 "in_reply_to": in_reply_to,
91 "notify_when_idle": notify_when_idle,
92 "queue": queue,
93 "id": idempotency_key,
94 "notice": notice,
95 });
96 return Ok(remote_send(&machine, &request));
97 }
98 deliver(
99 homes,
100 caller,
101 to,
102 body,
103 subject,
104 in_reply_to,
105 notify_when_idle,
106 queue,
107 idempotency_key,
108 notice,
109 )
110 .await
111}
112
113fn remote_machine(to: &str) -> Option<String> {
115 let local = local_machine_name();
116 let machine = match MailAddress::parse(to) {
117 Ok(address) => address.machine,
118 Err(_) => to.rsplit_once('@')?.1.to_string(),
119 };
120 (machine != local).then_some(machine)
121}
122
123fn remote_send(machine: &str, request: &serde_json::Value) -> Outcome {
125 match crate::mailbox::teams_mail(machine, request) {
126 Ok(answer) => {
127 let mut outcome = Outcome::new(
128 answer["code"]
129 .as_i64()
130 .map_or(EXIT_FAILED, |code| code as i32),
131 answer["text"]
132 .as_str()
133 .or_else(|| answer["detail"].as_str())
134 .unwrap_or("the other machine answered nothing readable")
135 .to_string(),
136 );
137 outcome.receipt = answer.get("receipt").cloned();
138 outcome
139 }
140 Err(detail) if detail.contains("no mail grant") => Outcome::new(
141 EXIT_REFUSED,
142 format!(
143 "Not sent: machine {machine} does not accept mail from you (no mail grant). Only \
144 an operator can grant it; tell your user. Don't retry."
145 ),
146 ),
147 Err(detail) => Outcome::new(
148 EXIT_FAILED,
149 format!("Not sent to machine {machine}: {detail}"),
150 ),
151 }
152}
153
154#[allow(clippy::too_many_arguments)]
157async fn deliver(
158 homes: &HarnessHomes,
159 caller: &Caller,
160 to: &str,
161 body: &str,
162 subject: Option<String>,
163 in_reply_to: Option<String>,
164 notify_when_idle: bool,
165 queue: bool,
166 idempotency_key: Option<String>,
167 notice: bool,
168) -> std::io::Result<Outcome> {
169 use crate::mail_route::{deliver as route, door_for, Delivered, Refused};
170 let local_operator = to
175 .strip_prefix("operator:")
176 .and_then(|name| MailAddress::new(local_machine_name(), "operator", name).ok());
177 let (address, name) = match local_operator
178 .map(Ok)
179 .unwrap_or_else(|| MailAddress::parse(to))
180 {
181 Ok(address) if address.harness == "operator" || address.harness == "board" => {
183 let name = format!("{}@{}", address.session_id, address.machine);
184 (address, name)
185 }
186 Ok(address) if address.harness == crate::mail_agent::AGENT_HARNESS => {
188 if crate::mail_agent::load(&address.session_id).is_none() {
189 return Ok(Outcome::new(
190 EXIT_UNKNOWN,
191 format!(
192 "No agent named {} is declared on this machine. Nothing was sent. Run \
193 supercode agent show for the declared agents.",
194 address.session_id
195 ),
196 ));
197 }
198 let name = format!("agent {}@{}", address.session_id, address.machine);
199 (address, name)
200 }
201 _ => match LiveSessions::read(homes).resolve(to) {
202 Ok(session) => (session.address.clone(), session.name.clone()),
203 Err(Unresolved::Stale(message) | Unresolved::Unknown(message)) => {
204 return Ok(Outcome::new(EXIT_UNKNOWN, message));
205 }
206 },
207 };
208 if address == caller.address {
209 return Ok(Outcome::new(
210 EXIT_REFUSED,
211 "Not sent: that address is your own session.",
212 ));
213 }
214 let door = door_for(homes, &address).ok();
217 let message_id = match &idempotency_key {
218 Some(key)
221 if key.len() == 26
222 && key.starts_with("m-")
223 && key[2..]
224 .bytes()
225 .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()) =>
226 {
227 key.clone()
228 }
229 Some(key) => format!(
230 "m-{}",
231 &blake3::hash(format!("{}\0{key}", caller.address).as_bytes()).to_hex()[..24]
232 ),
233 None => crate::mailbox::new_message_id()?,
234 };
235 let sent_marker = sent_marker(&caller.address, &message_id);
236 let terminal = matches!(&door, Some(Door::Hook { pane: Some(_), .. }));
237 if let Ok(previous) = std::fs::read_to_string(&sent_marker) {
238 if terminal && !queue {
239 Mailbox::open(&mail_root(), &address)?.request_wake(&message_id)?;
240 }
241 let mut outcome = Outcome::new(
242 0,
243 format!(
244 "Already sent with --id {}: {previous} Nothing was sent again.",
245 idempotency_key.as_deref().unwrap_or_default()
246 ),
247 );
248 outcome.receipt = Some(delivery_receipt(&address, &message_id, terminal));
249 return Ok(outcome);
250 }
251 let in_reply_to = match in_reply_to {
256 Some(answered) if answered.len() < FULL_ID_LEN => {
257 let found = Mailbox::open(&mail_root(), &caller.address)
258 .and_then(|mailbox| mailbox.find_prefix(&answered))
259 .unwrap_or_default();
260 match found.as_slice() {
261 [only] => Some(only.envelope.id.clone()),
262 [] => {
263 return Ok(Outcome::new(
264 EXIT_REFUSED,
265 format!(
266 "Not sent: no message in your mailbox has an id starting {answered}. \
267 Name the message you answer as its envelope shows it."
268 ),
269 ))
270 }
271 many => {
272 let ids: Vec<&str> = many
273 .iter()
274 .map(|stored| stored.envelope.id.as_str())
275 .collect();
276 return Ok(Outcome::new(
277 EXIT_REFUSED,
278 format!(
279 "Not sent: {answered} names {} messages in your mailbox ({}). Give \
280 more of its id.",
281 ids.len(),
282 ids.join(", ")
283 ),
284 ));
285 }
286 }
287 }
288 other => other,
289 };
290 if let Some(answered) = in_reply_to.as_deref().filter(|id| id.starts_with("q-")) {
291 return crate::mail_question::reply(caller, &address, answered, body).await;
292 }
293 if let Some(answered) = in_reply_to.as_deref() {
294 if reply_chain_depth(&caller.address, &address, answered) >= REPLY_CHAIN_LIMIT {
295 return Ok(Outcome::new(
296 EXIT_REFUSED,
297 format!(
298 "Not sent: you and {name} have answered each other {REPLY_CHAIN_LIMIT} times \
299 in a row. If you are trading acknowledgements or status, stop; to go on, \
300 send a new message (without --re), or tell your user."
301 ),
302 ));
303 }
304 }
305 let (kind, reply_via) = if notice {
306 (MailKind::Notice, ReplyVia::None)
307 } else {
308 (MailKind::Peer, ReplyVia::Command)
309 };
310 let mut envelope = Envelope::new(
311 caller.address.clone(),
312 caller.name.clone(),
313 kind,
314 reply_via,
315 body,
316 )?;
317 envelope.subject = subject.as_deref().map(|value| {
318 value
319 .lines()
320 .next()
321 .unwrap_or("")
322 .chars()
323 .take(200)
324 .collect()
325 });
326 envelope.id = message_id.clone();
327 envelope.thread = crate::mailbox::thread_of_reply(in_reply_to.as_deref());
328 envelope.in_reply_to = in_reply_to;
329 if let Some(plan) = crate::mail_agent::plan(
330 &mut envelope,
331 &address,
332 &crate::mail_agent::Channel::default(),
333 )? {
334 return Ok(deliver_planned(homes, caller, &envelope, &plan, queue, notify_when_idle).await);
335 }
336 let Some(door) = door else {
337 return Ok(Outcome::new(
338 EXIT_UNKNOWN,
339 format!(
340 "{name} is no longer running. Nothing was sent. Run supercode message list for \
341 the live sessions."
342 ),
343 ));
344 };
345 let tier = door.name();
346 let idle_note = if notify_when_idle {
347 " Subscribed: one idle notice reaches you when its next turn ends."
348 } else {
349 ""
350 };
351 let (code, what) = match route(&envelope, &address, &door, !queue, notify_when_idle).await {
352 Err(detail) => {
353 return Ok(Outcome::new(
354 EXIT_FAILED,
355 format!("Not sent to {name}: {detail}"),
356 ))
357 }
358 Ok(Err(Refused::CannotQueueNative)) => {
359 return Ok(Outcome::new(
360 EXIT_REFUSED,
361 format!(
362 "Not sent: {name} is idle, and a Claude session always starts a turn when a \
363 message arrives, so --queue cannot hold it. Send without --queue to wake it."
364 ),
365 ))
366 }
367 Ok(Err(Refused::TooLong(bytes))) => {
368 return Ok(Outcome::new(
369 EXIT_REFUSED,
370 format!(
371 "Not sent: the message is {bytes} bytes; the limit is {}. Write it to a file \
372 and send its path instead.",
373 crate::mail_route::MAX_RELAYED_BYTES
374 ),
375 ))
376 }
377 Ok(Ok(Delivered::Steered)) => (0, "steered into its running turn"),
378 Ok(Ok(Delivered::Started)) => (0, "it was idle, so the message started a turn"),
379 Ok(Ok(Delivered::Native { busy: true })) => (0, "it will read it at its next tool call"),
380 Ok(Ok(Delivered::Native { busy: false })) => {
381 (0, "it was idle, so the message starts its next turn")
382 }
383 Ok(Ok(Delivered::Hooked)) => (
384 0,
385 "queued in its mailbox; terminal delivery waits for native transcript confirmation",
386 ),
387 Ok(Ok(Delivered::HookWoken)) => (
388 0,
389 "it was idle, so its pane was told to read its mailbox; its hook shows it the message in \
390 that turn",
391 ),
392 Ok(Ok(Delivered::Queued)) => (
393 0,
394 "it is idle and --queue leaves it so; the message waits in its mailbox",
395 ),
396 Ok(Ok(Delivered::Operator)) => (0, "filed in its mailbox"),
397 Ok(Ok(Delivered::Stored)) => (
398 EXIT_STORED,
399 "Stored (not failed): it has no delivery door, so it sees this only if it runs \
400 supercode message inbox (a Codex session gets one with: supercode message setup \
401 codex). Don't resend, and don't wait for a reply",
402 ),
403 };
404 record_send(&caller.address, &address, &message_id);
405 let shown_id = crate::mailbox::short_id(&message_id);
406 let text = if code == 0 {
407 format!(
408 "sent to {name} ({}, {tier}): {what}. Message id {shown_id}.{idle_note} Don't poll; \
409 carry on.",
410 address.harness
411 )
412 } else {
413 format!(
414 "{name} ({}): {what}. Message id {shown_id}.{idle_note}",
415 address.harness
416 )
417 };
418 if code == 0 && idempotency_key.is_some() {
419 if let Some(parent) = sent_marker.parent() {
420 std::fs::create_dir_all(parent).ok();
421 }
422 std::fs::write(&sent_marker, &text).ok();
423 }
424 let mut outcome = Outcome::new(code, text);
425 if code == 0 {
426 outcome.receipt = Some(delivery_receipt(&address, &message_id, terminal));
427 }
428 Ok(outcome)
429}
430
431pub async fn deliver_planned(
435 homes: &HarnessHomes,
436 caller: &Caller,
437 envelope: &Envelope,
438 plan: &crate::mail_agent::Plan,
439 queue: bool,
440 notify_when_idle: bool,
441) -> Outcome {
442 let mut reached: Vec<String> = Vec::new();
443 let mut copied: Vec<String> = Vec::new();
444 let mut failed: Vec<String> = Vec::new();
445 if let Some(agent) = &plan.agent {
446 if let Err(error) =
448 Mailbox::open(&mail_root(), agent).and_then(|mailbox| mailbox.deliver_read(envelope))
449 {
450 failed.push(format!("the agent's mailbox ({error})"));
451 }
452 }
453 for recipient in plan
454 .recipients
455 .iter()
456 .filter(|recipient| recipient.address != caller.address)
457 {
458 let label = recipient_label(homes, &recipient.address);
459 record_send(&caller.address, &recipient.address, &envelope.id);
460 if !recipient.wake {
461 match crate::mail_agent::file_unread(&recipient.address, envelope) {
462 Ok(()) => copied.push(label),
463 Err(error) => failed.push(format!("{label} ({error})")),
464 }
465 continue;
466 }
467 match deliver_woken(homes, envelope, &recipient.address, queue, notify_when_idle).await {
468 Ok(how) => reached.push(format!("{label} ({how})")),
469 Err(detail) => failed.push(format!("{label} ({detail})")),
470 }
471 }
472 let shown_id = crate::mailbox::short_id(&envelope.id);
473 let thread = plan
474 .thread
475 .as_ref()
476 .map(|thread| format!(" Thread {}.", crate::mailbox::short_id(&thread.id)))
477 .unwrap_or_default();
478 let mut text = if reached.is_empty() && copied.is_empty() {
479 format!("Not sent: nobody in its thread could be reached.{thread}")
480 } else {
481 let cc = if copied.is_empty() {
482 String::new()
483 } else {
484 format!(" CC (filed, not woken): {}.", copied.join(", "))
485 };
486 format!(
487 "sent to {}.{cc} Message id {shown_id}.{thread} Don't poll; carry on.",
488 if reached.is_empty() {
489 "no one woken".to_string()
490 } else {
491 reached.join(", ")
492 }
493 )
494 };
495 if !failed.is_empty() {
496 text.push_str(&format!(" Not delivered to: {}.", failed.join("; ")));
497 }
498 let code = if reached.is_empty() && copied.is_empty() {
499 EXIT_FAILED
500 } else {
501 0
502 };
503 let mut outcome = Outcome::new(code, text);
504 if code == 0 {
505 outcome.receipt = Some(
506 serde_json::json!({"delivered": true, "message_ids": [envelope.id], "receipt": "agent-mailbox"}),
507 );
508 }
509 outcome
510}
511
512fn recipient_label(homes: &HarnessHomes, address: &MailAddress) -> String {
514 if address.harness == "operator" || address.harness == "board" {
515 return format!("{}@{}", address.session_id, address.machine);
516 }
517 LiveSessions::read(homes)
518 .all()
519 .iter()
520 .find(|session| &session.address == address)
521 .map(|session| session.name.clone())
522 .unwrap_or_else(|| address.to_string())
523}
524
525async fn deliver_woken(
528 homes: &HarnessHomes,
529 envelope: &Envelope,
530 to: &MailAddress,
531 queue: bool,
532 notify_when_idle: bool,
533) -> Result<&'static str, String> {
534 use crate::mail_route::{deliver as route, door_for, Delivered, NoDoor};
535 let door = match door_for(homes, to) {
536 Ok(door) => door,
537 Err(NoDoor::OtherMachine(_)) => {
538 crate::mailbox::deliver_to(to, envelope).map_err(|error| error.to_string())?;
539 return Ok("filed on its machine");
540 }
541 Err(NoDoor::NotRunning) => {
542 crate::mail_agent::file_unread(to, envelope).map_err(|error| error.to_string())?;
545 if !crate::mail_agent::resumable(to) {
546 return Ok(
547 "stopped; not resumed, since the mailbox did not launch it: the \
548 message waits in its mailbox",
549 );
550 }
551 crate::mail_agent::resume(to).map_err(|error| format!("not resumed: {error}"))?;
552 let started = std::time::Instant::now();
553 loop {
554 if let Ok(door) = door_for(homes, to) {
555 break door;
556 }
557 if started.elapsed() > crate::mail_agent::RESUME_WAIT {
558 return Ok("resumed; the message waits in its mailbox");
559 }
560 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
561 }
562 }
563 };
564 match route(envelope, to, &door, !queue, notify_when_idle).await? {
565 Ok(Delivered::Steered) => Ok("steered into its running turn"),
566 Ok(Delivered::Started) => Ok("started a turn"),
567 Ok(Delivered::Native { busy: true }) => Ok("read at its next tool call"),
568 Ok(Delivered::Native { busy: false }) => Ok("starts its next turn"),
569 Ok(Delivered::Hooked | Delivered::HookWoken) => Ok("in its mailbox, its hook shows it"),
570 Ok(Delivered::Queued) => Ok("waits in its mailbox"),
571 Ok(Delivered::Operator) => Ok("filed"),
572 Ok(Delivered::Stored) => Ok("stored; it has no delivery door"),
573 Err(crate::mail_route::Refused::CannotQueueNative) => {
574 Err("idle, and --queue cannot hold a Claude session's mail".into())
575 }
576 Err(crate::mail_route::Refused::TooLong(bytes)) => Err(format!(
577 "{bytes} bytes; the limit is {}",
578 crate::mail_route::MAX_RELAYED_BYTES
579 )),
580 }
581}
582
583const REPLY_CHAIN_LIMIT: usize = 32;
588
589const FULL_ID_LEN: usize = 26;
591
592fn reply_chain_depth(a: &MailAddress, b: &MailAddress, answered: &str) -> usize {
596 let mut links = std::collections::HashMap::new();
597 for address in [a, b] {
598 let Ok(mailbox) = Mailbox::open(&mail_root(), address) else {
599 continue;
600 };
601 for stored in mailbox.list().unwrap_or_default() {
602 links.insert(
603 stored.envelope.id.clone(),
604 stored.envelope.in_reply_to.clone(),
605 );
606 }
607 }
608 let mut depth = 0;
609 let mut current = Some(answered.to_string());
610 while let Some(id) = current {
611 if depth > REPLY_CHAIN_LIMIT || !links.contains_key(&id) {
612 break;
613 }
614 depth += 1;
615 current = links.get(&id).cloned().flatten();
616 }
617 depth
618}
619
620fn send_log(sender: &MailAddress) -> PathBuf {
621 let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
622 mail_root().join("sent").join(&hash[..24]).join("log.jsonl")
623}
624
625fn record_send(sender: &MailAddress, to: &MailAddress, message_id: &str) {
626 let log = send_log(sender);
627 if let Some(parent) = log.parent() {
628 std::fs::create_dir_all(parent).ok();
629 }
630 if let Ok(mut file) = std::fs::OpenOptions::new()
631 .create(true)
632 .append(true)
633 .open(&log)
634 {
635 use std::io::Write as _;
636 let entry = serde_json::json!({"to": to.to_string(), "at_ms": now_ms(), "id": message_id});
637 writeln!(file, "{entry}").ok();
638 }
639}
640
641fn now_ms() -> u64 {
642 std::time::SystemTime::now()
643 .duration_since(std::time::UNIX_EPOCH)
644 .map(|elapsed| elapsed.as_millis() as u64)
645 .unwrap_or_default()
646}
647
648fn sent_marker(sender: &MailAddress, message_id: &str) -> PathBuf {
650 let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
651 mail_root().join("sent").join(&hash[..24]).join(message_id)
652}
653
654fn delivery_receipt(address: &MailAddress, id: &str, terminal: bool) -> serde_json::Value {
657 if !terminal {
658 return serde_json::json!({"delivered": true, "message_ids": [id], "receipt": "native-mail"});
659 }
660 let recipient = format!("{}:{}", address.harness, address.session_id);
661 if let Ok(entries) = std::fs::read_dir(mail_root().join("terminal-delivery")) {
662 for entry in entries.flatten() {
663 let Ok(bytes) = std::fs::read(entry.path().join("queue.json")) else {
664 continue;
665 };
666 let Ok(state) = serde_json::from_slice::<serde_json::Value>(&bytes) else {
667 continue;
668 };
669 if state["recipient"].as_str() != Some(recipient.as_str()) {
670 continue;
671 }
672 if let Some(receipt) = state["receipts"].as_array().and_then(|receipts| {
673 receipts.iter().find(|r| {
674 r["message_ids"]
675 .as_array()
676 .is_some_and(|ids| ids.iter().any(|value| value.as_str() == Some(id)))
677 })
678 }) {
679 return receipt.clone();
680 }
681 return serde_json::json!({"delivered": false, "queued": true, "message_ids": [id], "reason": state["active"]["error"]});
682 }
683 }
684 serde_json::json!({"delivered": false, "queued": true, "message_ids": [id], "reason": "awaiting terminal mailbox"})
685}