1use std::path::PathBuf;
8use std::time::Duration;
9
10use crate::mail_route::{Caller, LiveSessions, Unresolved};
11use crate::mailbox::{local_machine_name, mail_root, Envelope, MailAddress, MailKind, ReplyVia};
12use crate::HarnessHomes;
13
14pub const EXIT_REFUSED: i32 = 2;
16pub const EXIT_UNKNOWN: i32 = 3;
18pub const EXIT_STORED: i32 = 4;
20pub const EXIT_FAILED: i32 = 5;
22
23pub struct Outcome {
25 pub code: i32,
28 pub text: String,
30}
31
32impl Outcome {
33 pub fn new(code: i32, text: impl Into<String>) -> Self {
35 Self {
36 code,
37 text: text.into(),
38 }
39 }
40}
41
42#[derive(Debug, Clone, Default)]
44pub struct SendOptions {
45 pub in_reply_to: Option<String>,
47 pub notify_when_idle: bool,
49 pub queue: bool,
51 pub idempotency_key: Option<String>,
53}
54
55pub async fn send(
59 homes: &HarnessHomes,
60 caller: &Caller,
61 to: &str,
62 body: &str,
63 options: SendOptions,
64) -> std::io::Result<Outcome> {
65 let SendOptions {
66 in_reply_to,
67 notify_when_idle,
68 queue,
69 idempotency_key,
70 } = options;
71 if let Some(machine) = remote_machine(to) {
72 let request = serde_json::json!({
73 "op": "send",
74 "from": caller.address.to_string(),
75 "from_name": caller.name,
76 "to": to,
77 "body": body,
78 "in_reply_to": in_reply_to,
79 "notify_when_idle": notify_when_idle,
80 "queue": queue,
81 "id": idempotency_key,
82 });
83 return Ok(remote_send(&machine, &request));
84 }
85 deliver(
86 homes,
87 caller,
88 to,
89 body,
90 in_reply_to,
91 notify_when_idle,
92 queue,
93 idempotency_key,
94 )
95 .await
96}
97
98fn remote_machine(to: &str) -> Option<String> {
100 let local = local_machine_name();
101 let machine = match MailAddress::parse(to) {
102 Ok(address) => address.machine,
103 Err(_) => to.rsplit_once('@')?.1.to_string(),
104 };
105 (machine != local).then_some(machine)
106}
107
108fn remote_send(machine: &str, request: &serde_json::Value) -> Outcome {
110 match crate::mailbox::teams_mail(machine, request) {
111 Ok(answer) => Outcome::new(
112 answer["code"]
113 .as_i64()
114 .map_or(EXIT_FAILED, |code| code as i32),
115 answer["text"]
116 .as_str()
117 .or_else(|| answer["detail"].as_str())
118 .unwrap_or("the other machine answered nothing readable")
119 .to_string(),
120 ),
121 Err(detail) if detail.contains("no mail grant") => Outcome::new(
122 EXIT_REFUSED,
123 format!(
124 "Not sent: machine {machine} does not accept mail from you (no mail grant). Only \
125 an operator can grant it; tell your user. Don't retry."
126 ),
127 ),
128 Err(detail) => Outcome::new(
129 EXIT_FAILED,
130 format!("Not sent to machine {machine}: {detail}"),
131 ),
132 }
133}
134
135#[allow(clippy::too_many_arguments)]
138async fn deliver(
139 homes: &HarnessHomes,
140 caller: &Caller,
141 to: &str,
142 body: &str,
143 in_reply_to: Option<String>,
144 notify_when_idle: bool,
145 queue: bool,
146 idempotency_key: Option<String>,
147) -> std::io::Result<Outcome> {
148 use crate::mail_route::{deliver as route, door_for, Delivered, Refused};
149 let (address, name) = match MailAddress::parse(to) {
152 Ok(address) if address.harness == "operator" => {
153 let name = format!("{}@{}", address.session_id, address.machine);
154 (address, name)
155 }
156 _ => match LiveSessions::read(homes).resolve(to) {
157 Ok(session) => (session.address.clone(), session.name.clone()),
158 Err(Unresolved::Stale(message) | Unresolved::Unknown(message)) => {
159 return Ok(Outcome::new(EXIT_UNKNOWN, message));
160 }
161 },
162 };
163 if address == caller.address {
164 return Ok(Outcome::new(
165 EXIT_REFUSED,
166 "Not sent: that address is your own session.",
167 ));
168 }
169 let door = match door_for(homes, &address) {
170 Ok(door) => door,
171 Err(_) => {
172 return Ok(Outcome::new(
173 EXIT_UNKNOWN,
174 format!(
175 "{name} is no longer running. Nothing was sent. Run supercode message list for \
176 the live sessions."
177 ),
178 ))
179 }
180 };
181 let message_id = match &idempotency_key {
182 Some(key) => format!(
183 "m-{}",
184 &blake3::hash(format!("{}\0{key}", caller.address).as_bytes()).to_hex()[..24]
185 ),
186 None => crate::mailbox::new_message_id()?,
187 };
188 let sent_marker = sent_marker(&caller.address, &message_id);
189 if let Ok(previous) = std::fs::read_to_string(&sent_marker) {
190 return Ok(Outcome::new(
191 0,
192 format!(
193 "Already sent with --id {}: {previous} Nothing was sent again.",
194 idempotency_key.as_deref().unwrap_or_default()
195 ),
196 ));
197 }
198 if recent_sends(&caller.address, &address) >= LOOP_LIMIT {
199 return Ok(Outcome::new(
200 EXIT_REFUSED,
201 format!(
202 "Not sent: this would be message {} from you to {name} in {} minutes. If you two \
203 are trading acknowledgements or status, stop; tell your user if the exchange must \
204 go on.",
205 LOOP_LIMIT + 1,
206 LOOP_WINDOW.as_secs() / 60
207 ),
208 ));
209 }
210 let mut envelope = Envelope::new(
211 caller.address.clone(),
212 caller.name.clone(),
213 MailKind::Peer,
214 ReplyVia::Command,
215 body,
216 )?;
217 envelope.id = message_id.clone();
218 envelope.in_reply_to = in_reply_to;
219 let tier = door.name();
220 let idle_note = if notify_when_idle {
221 " Subscribed: one idle notice reaches you when its next turn ends."
222 } else {
223 ""
224 };
225 let (code, what) = match route(&envelope, &address, &door, !queue, notify_when_idle).await {
226 Err(detail) => {
227 return Ok(Outcome::new(
228 EXIT_FAILED,
229 format!("Not sent to {name}: {detail}"),
230 ))
231 }
232 Ok(Err(Refused::CannotQueueNative)) => {
233 return Ok(Outcome::new(
234 EXIT_REFUSED,
235 format!(
236 "Not sent: {name} is idle, and a Claude session always starts a turn when a \
237 message arrives, so --queue cannot hold it. Send without --queue to wake it."
238 ),
239 ))
240 }
241 Ok(Err(Refused::TooLong(bytes))) => {
242 return Ok(Outcome::new(
243 EXIT_REFUSED,
244 format!(
245 "Not sent: the message is {bytes} bytes; the limit is {}. Write it to a file \
246 and send its path instead.",
247 crate::mail_route::MAX_RELAYED_BYTES
248 ),
249 ))
250 }
251 Ok(Ok(Delivered::Steered)) => (0, "steered into its running turn"),
252 Ok(Ok(Delivered::Started)) => (0, "it was idle, so the message started a turn"),
253 Ok(Ok(Delivered::Native { busy: true })) => (0, "it will read it at its next tool call"),
254 Ok(Ok(Delivered::Native { busy: false })) => {
255 (0, "it was idle, so the message starts its next turn")
256 }
257 Ok(Ok(Delivered::Hooked)) => (
258 0,
259 "its hook points it at the message at its next tool call, or when its turn ends; an \
260 idle session sees it on its next turn",
261 ),
262 Ok(Ok(Delivered::Queued)) => (
263 0,
264 "it is idle and --queue leaves it so; the message waits in its mailbox",
265 ),
266 Ok(Ok(Delivered::Operator)) => (0, "filed in its mailbox"),
267 Ok(Ok(Delivered::Stored)) => (
268 EXIT_STORED,
269 "Stored (not failed): it has no delivery door, so it sees this only if it runs \
270 supercode message inbox (a Codex session gets one with: supercode message setup \
271 codex). Don't resend, and don't wait for a reply",
272 ),
273 };
274 record_send(&caller.address, &address, &message_id);
275 let text = if code == 0 {
276 format!(
277 "sent to {name} ({}, {tier}): {what}. Message id {message_id}.{idle_note} Don't poll; \
278 carry on.",
279 address.harness
280 )
281 } else {
282 format!(
283 "{name} ({}): {what}. Message id {message_id}.{idle_note}",
284 address.harness
285 )
286 };
287 if code == 0 && idempotency_key.is_some() {
288 if let Some(parent) = sent_marker.parent() {
289 std::fs::create_dir_all(parent).ok();
290 }
291 std::fs::write(&sent_marker, &text).ok();
292 }
293 Ok(Outcome::new(code, text))
294}
295
296const LOOP_LIMIT: usize = 8;
299const LOOP_WINDOW: Duration = Duration::from_secs(10 * 60);
300
301fn send_log(sender: &MailAddress) -> PathBuf {
302 let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
303 mail_root().join("sent").join(&hash[..24]).join("log.jsonl")
304}
305
306fn recent_sends(sender: &MailAddress, to: &MailAddress) -> usize {
308 let now = now_ms();
309 let window = LOOP_WINDOW.as_millis() as u64;
310 std::fs::read_to_string(send_log(sender))
311 .unwrap_or_default()
312 .lines()
313 .filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
314 .filter(|entry| entry["to"].as_str() == Some(to.to_string().as_str()))
315 .filter(|entry| {
316 entry["at_ms"]
317 .as_u64()
318 .is_some_and(|at| now.saturating_sub(at) < window)
319 })
320 .count()
321}
322
323fn record_send(sender: &MailAddress, to: &MailAddress, message_id: &str) {
324 let log = send_log(sender);
325 if let Some(parent) = log.parent() {
326 std::fs::create_dir_all(parent).ok();
327 }
328 if let Ok(mut file) = std::fs::OpenOptions::new()
329 .create(true)
330 .append(true)
331 .open(&log)
332 {
333 use std::io::Write as _;
334 let entry = serde_json::json!({"to": to.to_string(), "at_ms": now_ms(), "id": message_id});
335 writeln!(file, "{entry}").ok();
336 }
337}
338
339fn now_ms() -> u64 {
340 std::time::SystemTime::now()
341 .duration_since(std::time::UNIX_EPOCH)
342 .map(|elapsed| elapsed.as_millis() as u64)
343 .unwrap_or_default()
344}
345
346fn sent_marker(sender: &MailAddress, message_id: &str) -> PathBuf {
348 let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
349 mail_root().join("sent").join(&hash[..24]).join(message_id)
350}