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