Skip to main content

supercode_harness/
mail_send.rs

1//! One send for every sender: `supercode message send`, the native agent's
2//! `send_message`, and a request arriving through another machine's mail
3//! door. It resolves the receiver, routes to another machine through Teams or
4//! to the router's door here, guards against loops and repeats, and words
5//! the outcome for the sending agent.
6
7use 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
15/// Exit code of a refused send.
16pub const EXIT_REFUSED: i32 = 2;
17/// Exit code of an unknown or stale receiver.
18pub const EXIT_UNKNOWN: i32 = 3;
19/// Exit code of a message stored with no door to show it.
20pub const EXIT_STORED: i32 = 4;
21/// Exit code of a failed send.
22pub const EXIT_FAILED: i32 = 5;
23
24/// What a delivery came to: the exit code, and the text the sending agent reads.
25pub struct Outcome {
26    /// Exit code for scripts: 0 sent, [`EXIT_REFUSED`], [`EXIT_UNKNOWN`],
27    /// [`EXIT_STORED`] or [`EXIT_FAILED`].
28    pub code: i32,
29    /// What the sending agent reads.
30    pub text: String,
31}
32
33impl Outcome {
34    /// An outcome with this code and text.
35    pub fn new(code: i32, text: impl Into<String>) -> Self {
36        Self {
37            code,
38            text: text.into(),
39        }
40    }
41}
42
43/// Options of one send.
44#[derive(Debug, Clone, Default)]
45pub struct SendOptions {
46    /// Id of the message this one answers.
47    pub in_reply_to: Option<String>,
48    /// Also send one notice when the receiver's next turn ends.
49    pub notify_when_idle: bool,
50    /// Only enqueue; never start a turn in an idle receiver.
51    pub queue: bool,
52    /// Idempotency key: repeating a send with it does not send twice.
53    pub idempotency_key: Option<String>,
54}
55
56/// Send `body` from `caller` to `to` (an address, `name@machine`, or a name
57/// unique on this machine): through Teams when `to` is on another machine,
58/// else through the router's door here. The one send every sender uses.
59pub 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
99/// The machine `to` names when it is not this one.
100fn 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
109/// Hand a request to another machine's mail door through Teams.
110fn 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/// Deliver from `caller` to `to` on this machine, through the one router
137/// every sender uses (`crate::mail_route`).
138#[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    // An operator (a board, a script) is not a listed session; its address is
151    // taken as given.
152    // `operator:<name>` is that operator on this machine: a tool files what it posts on the
153    // session's behalf without knowing the machine's name.
154    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" => {
162            let name = format!("{}@{}", address.session_id, address.machine);
163            (address, name)
164        }
165        _ => match LiveSessions::read(homes).resolve(to) {
166            Ok(session) => (session.address.clone(), session.name.clone()),
167            Err(Unresolved::Stale(message) | Unresolved::Unknown(message)) => {
168                return Ok(Outcome::new(EXIT_UNKNOWN, message));
169            }
170        },
171    };
172    if address == caller.address {
173        return Ok(Outcome::new(
174            EXIT_REFUSED,
175            "Not sent: that address is your own session.",
176        ));
177    }
178    let door = match door_for(homes, &address) {
179        Ok(door) => door,
180        Err(_) => {
181            return Ok(Outcome::new(
182                EXIT_UNKNOWN,
183                format!(
184                    "{name} is no longer running. Nothing was sent. Run supercode message list for \
185                     the live sessions."
186                ),
187            ))
188        }
189    };
190    let message_id = match &idempotency_key {
191        // A key that is itself a whole message id is the message's id: a sender that chose it beforehand (a board
192        // blocking a card on the answer to the question it is about to send) knows what the answer will name.
193        Some(key)
194            if key.len() == 26
195                && key.starts_with("m-")
196                && key[2..]
197                    .bytes()
198                    .all(|byte| byte.is_ascii_hexdigit() && !byte.is_ascii_uppercase()) =>
199        {
200            key.clone()
201        }
202        Some(key) => format!(
203            "m-{}",
204            &blake3::hash(format!("{}\0{key}", caller.address).as_bytes()).to_hex()[..24]
205        ),
206        None => crate::mailbox::new_message_id()?,
207    };
208    let sent_marker = sent_marker(&caller.address, &message_id);
209    if let Ok(previous) = std::fs::read_to_string(&sent_marker) {
210        return Ok(Outcome::new(
211            0,
212            format!(
213                "Already sent with --id {}: {previous} Nothing was sent again.",
214                idempotency_key.as_deref().unwrap_or_default()
215            ),
216        ));
217    }
218    // An answered message is named as the agent was shown it, by a prefix of its id: the one in
219    // the caller's mailbox it names.
220    // A prefix names the one message it matches there, or nothing is sent; a whole id is kept as
221    // given (it may answer a message filed elsewhere).
222    let in_reply_to = match in_reply_to {
223        Some(answered) if answered.len() < FULL_ID_LEN => {
224            let found = Mailbox::open(&mail_root(), &caller.address)
225                .and_then(|mailbox| mailbox.find_prefix(&answered))
226                .unwrap_or_default();
227            match found.as_slice() {
228                [only] => Some(only.envelope.id.clone()),
229                [] => {
230                    return Ok(Outcome::new(
231                        EXIT_REFUSED,
232                        format!(
233                            "Not sent: no message in your mailbox has an id starting {answered}. \
234                             Name the message you answer as its envelope shows it."
235                        ),
236                    ))
237                }
238                many => {
239                    let ids: Vec<&str> = many
240                        .iter()
241                        .map(|stored| stored.envelope.id.as_str())
242                        .collect();
243                    return Ok(Outcome::new(
244                        EXIT_REFUSED,
245                        format!(
246                            "Not sent: {answered} names {} messages in your mailbox ({}). Give \
247                             more of its id.",
248                            ids.len(),
249                            ids.join(", ")
250                        ),
251                    ));
252                }
253            }
254        }
255        other => other,
256    };
257    if let Some(answered) = in_reply_to.as_deref() {
258        if reply_chain_depth(&caller.address, &address, answered) >= REPLY_CHAIN_LIMIT {
259            return Ok(Outcome::new(
260                EXIT_REFUSED,
261                format!(
262                    "Not sent: you and {name} have answered each other {REPLY_CHAIN_LIMIT} times \
263                     in a row. If you are trading acknowledgements or status, stop; to go on, \
264                     send a new message (without --re), or tell your user."
265                ),
266            ));
267        }
268    }
269    let mut envelope = Envelope::new(
270        caller.address.clone(),
271        caller.name.clone(),
272        MailKind::Peer,
273        ReplyVia::Command,
274        body,
275    )?;
276    envelope.id = message_id.clone();
277    envelope.thread = crate::mailbox::thread_of_reply(in_reply_to.as_deref());
278    envelope.in_reply_to = in_reply_to;
279    let tier = door.name();
280    let idle_note = if notify_when_idle {
281        " Subscribed: one idle notice reaches you when its next turn ends."
282    } else {
283        ""
284    };
285    let (code, what) = match route(&envelope, &address, &door, !queue, notify_when_idle).await {
286        Err(detail) => {
287            return Ok(Outcome::new(
288                EXIT_FAILED,
289                format!("Not sent to {name}: {detail}"),
290            ))
291        }
292        Ok(Err(Refused::CannotQueueNative)) => {
293            return Ok(Outcome::new(
294                EXIT_REFUSED,
295                format!(
296                    "Not sent: {name} is idle, and a Claude session always starts a turn when a \
297                     message arrives, so --queue cannot hold it. Send without --queue to wake it."
298                ),
299            ))
300        }
301        Ok(Err(Refused::TooLong(bytes))) => {
302            return Ok(Outcome::new(
303                EXIT_REFUSED,
304                format!(
305                    "Not sent: the message is {bytes} bytes; the limit is {}. Write it to a file \
306                     and send its path instead.",
307                    crate::mail_route::MAX_RELAYED_BYTES
308                ),
309            ))
310        }
311        Ok(Ok(Delivered::Steered)) => (0, "steered into its running turn"),
312        Ok(Ok(Delivered::Started)) => (0, "it was idle, so the message started a turn"),
313        Ok(Ok(Delivered::Native { busy: true })) => (0, "it will read it at its next tool call"),
314        Ok(Ok(Delivered::Native { busy: false })) => {
315            (0, "it was idle, so the message starts its next turn")
316        }
317        Ok(Ok(Delivered::Hooked)) => (
318            0,
319            "its hook points it at the message at its next tool call, or when its turn ends; an \
320             idle session sees it on its next turn",
321        ),
322        Ok(Ok(Delivered::HookWoken)) => (
323            0,
324            "it was idle, so its pane was told to read its mailbox; its hook shows it the message in \
325             that turn",
326        ),
327        Ok(Ok(Delivered::Queued)) => (
328            0,
329            "it is idle and --queue leaves it so; the message waits in its mailbox",
330        ),
331        Ok(Ok(Delivered::Operator)) => (0, "filed in its mailbox"),
332        Ok(Ok(Delivered::Stored)) => (
333            EXIT_STORED,
334            "Stored (not failed): it has no delivery door, so it sees this only if it runs \
335             supercode message inbox (a Codex session gets one with: supercode message setup \
336             codex). Don't resend, and don't wait for a reply",
337        ),
338    };
339    record_send(&caller.address, &address, &message_id);
340    let shown_id = crate::mailbox::short_id(&message_id);
341    let text = if code == 0 {
342        format!(
343            "sent to {name} ({}, {tier}): {what}. Message id {shown_id}.{idle_note} Don't poll; \
344             carry on.",
345            address.harness
346        )
347    } else {
348        format!(
349            "{name} ({}): {what}. Message id {shown_id}.{idle_note}",
350            address.harness
351        )
352    };
353    if code == 0 && idempotency_key.is_some() {
354        if let Some(parent) = sent_marker.parent() {
355            std::fs::create_dir_all(parent).ok();
356        }
357        std::fs::write(&sent_marker, &text).ok();
358    }
359    Ok(Outcome::new(code, text))
360}
361
362/// Replies in one unbroken chain between two sessions beyond which another
363/// reply is refused: two agents trading acknowledgements never stop on their
364/// own, and every message of such a loop answers the one before it. A new
365/// message (no `--re`) starts a chain, so volume alone is never refused.
366const REPLY_CHAIN_LIMIT: usize = 32;
367
368/// Length of a whole message id: its kind, `-`, and 24 hex digits.
369const FULL_ID_LEN: usize = 26;
370
371/// How many messages deep the reply chain ending at `answered` runs, between
372/// `a` and `b`: each message's `in_reply_to`, followed back through both
373/// sessions' mailboxes.
374fn reply_chain_depth(a: &MailAddress, b: &MailAddress, answered: &str) -> usize {
375    let mut links = std::collections::HashMap::new();
376    for address in [a, b] {
377        let Ok(mailbox) = Mailbox::open(&mail_root(), address) else {
378            continue;
379        };
380        for stored in mailbox.list().unwrap_or_default() {
381            links.insert(
382                stored.envelope.id.clone(),
383                stored.envelope.in_reply_to.clone(),
384            );
385        }
386    }
387    let mut depth = 0;
388    let mut current = Some(answered.to_string());
389    while let Some(id) = current {
390        if depth > REPLY_CHAIN_LIMIT || !links.contains_key(&id) {
391            break;
392        }
393        depth += 1;
394        current = links.get(&id).cloned().flatten();
395    }
396    depth
397}
398
399fn send_log(sender: &MailAddress) -> PathBuf {
400    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
401    mail_root().join("sent").join(&hash[..24]).join("log.jsonl")
402}
403
404fn record_send(sender: &MailAddress, to: &MailAddress, message_id: &str) {
405    let log = send_log(sender);
406    if let Some(parent) = log.parent() {
407        std::fs::create_dir_all(parent).ok();
408    }
409    if let Ok(mut file) = std::fs::OpenOptions::new()
410        .create(true)
411        .append(true)
412        .open(&log)
413    {
414        use std::io::Write as _;
415        let entry = serde_json::json!({"to": to.to_string(), "at_ms": now_ms(), "id": message_id});
416        writeln!(file, "{entry}").ok();
417    }
418}
419
420fn now_ms() -> u64 {
421    std::time::SystemTime::now()
422        .duration_since(std::time::UNIX_EPOCH)
423        .map(|elapsed| elapsed.as_millis() as u64)
424        .unwrap_or_default()
425}
426
427/// Marker recording that `message_id` was sent by `sender`, for `--id`.
428fn sent_marker(sender: &MailAddress, message_id: &str) -> PathBuf {
429    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
430    mail_root().join("sent").join(&hash[..24]).join(message_id)
431}