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