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        // a board is an operator too: its dispatcher reads its mailbox (answers to the questions it asks)
162        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        // A key that is itself a whole message id is the message's id: a sender that chose it beforehand (a board
193        // blocking a card on the answer to the question it is about to send) knows what the answer will name.
194        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    // An answered message is named as the agent was shown it, by a prefix of its id: the one in
220    // the caller's mailbox it names.
221    // A prefix names the one message it matches there, or nothing is sent; a whole id is kept as
222    // given (it may answer a message filed elsewhere).
223    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() {
259        if reply_chain_depth(&caller.address, &address, answered) >= REPLY_CHAIN_LIMIT {
260            return Ok(Outcome::new(
261                EXIT_REFUSED,
262                format!(
263                    "Not sent: you and {name} have answered each other {REPLY_CHAIN_LIMIT} times \
264                     in a row. If you are trading acknowledgements or status, stop; to go on, \
265                     send a new message (without --re), or tell your user."
266                ),
267            ));
268        }
269    }
270    let mut envelope = Envelope::new(
271        caller.address.clone(),
272        caller.name.clone(),
273        MailKind::Peer,
274        ReplyVia::Command,
275        body,
276    )?;
277    envelope.id = message_id.clone();
278    envelope.thread = crate::mailbox::thread_of_reply(in_reply_to.as_deref());
279    envelope.in_reply_to = in_reply_to;
280    let tier = door.name();
281    let idle_note = if notify_when_idle {
282        " Subscribed: one idle notice reaches you when its next turn ends."
283    } else {
284        ""
285    };
286    let (code, what) = match route(&envelope, &address, &door, !queue, notify_when_idle).await {
287        Err(detail) => {
288            return Ok(Outcome::new(
289                EXIT_FAILED,
290                format!("Not sent to {name}: {detail}"),
291            ))
292        }
293        Ok(Err(Refused::CannotQueueNative)) => {
294            return Ok(Outcome::new(
295                EXIT_REFUSED,
296                format!(
297                    "Not sent: {name} is idle, and a Claude session always starts a turn when a \
298                     message arrives, so --queue cannot hold it. Send without --queue to wake it."
299                ),
300            ))
301        }
302        Ok(Err(Refused::TooLong(bytes))) => {
303            return Ok(Outcome::new(
304                EXIT_REFUSED,
305                format!(
306                    "Not sent: the message is {bytes} bytes; the limit is {}. Write it to a file \
307                     and send its path instead.",
308                    crate::mail_route::MAX_RELAYED_BYTES
309                ),
310            ))
311        }
312        Ok(Ok(Delivered::Steered)) => (0, "steered into its running turn"),
313        Ok(Ok(Delivered::Started)) => (0, "it was idle, so the message started a turn"),
314        Ok(Ok(Delivered::Native { busy: true })) => (0, "it will read it at its next tool call"),
315        Ok(Ok(Delivered::Native { busy: false })) => {
316            (0, "it was idle, so the message starts its next turn")
317        }
318        Ok(Ok(Delivered::Hooked)) => (
319            0,
320            "its hook points it at the message at its next tool call, or when its turn ends; an \
321             idle session sees it on its next turn",
322        ),
323        Ok(Ok(Delivered::HookWoken)) => (
324            0,
325            "it was idle, so its pane was told to read its mailbox; its hook shows it the message in \
326             that turn",
327        ),
328        Ok(Ok(Delivered::Queued)) => (
329            0,
330            "it is idle and --queue leaves it so; the message waits in its mailbox",
331        ),
332        Ok(Ok(Delivered::Operator)) => (0, "filed in its mailbox"),
333        Ok(Ok(Delivered::Stored)) => (
334            EXIT_STORED,
335            "Stored (not failed): it has no delivery door, so it sees this only if it runs \
336             supercode message inbox (a Codex session gets one with: supercode message setup \
337             codex). Don't resend, and don't wait for a reply",
338        ),
339    };
340    record_send(&caller.address, &address, &message_id);
341    let shown_id = crate::mailbox::short_id(&message_id);
342    let text = if code == 0 {
343        format!(
344            "sent to {name} ({}, {tier}): {what}. Message id {shown_id}.{idle_note} Don't poll; \
345             carry on.",
346            address.harness
347        )
348    } else {
349        format!(
350            "{name} ({}): {what}. Message id {shown_id}.{idle_note}",
351            address.harness
352        )
353    };
354    if code == 0 && idempotency_key.is_some() {
355        if let Some(parent) = sent_marker.parent() {
356            std::fs::create_dir_all(parent).ok();
357        }
358        std::fs::write(&sent_marker, &text).ok();
359    }
360    Ok(Outcome::new(code, text))
361}
362
363/// Replies in one unbroken chain between two sessions beyond which another
364/// reply is refused: two agents trading acknowledgements never stop on their
365/// own, and every message of such a loop answers the one before it. A new
366/// message (no `--re`) starts a chain, so volume alone is never refused.
367const REPLY_CHAIN_LIMIT: usize = 32;
368
369/// Length of a whole message id: its kind, `-`, and 24 hex digits.
370const FULL_ID_LEN: usize = 26;
371
372/// How many messages deep the reply chain ending at `answered` runs, between
373/// `a` and `b`: each message's `in_reply_to`, followed back through both
374/// sessions' mailboxes.
375fn reply_chain_depth(a: &MailAddress, b: &MailAddress, answered: &str) -> usize {
376    let mut links = std::collections::HashMap::new();
377    for address in [a, b] {
378        let Ok(mailbox) = Mailbox::open(&mail_root(), address) else {
379            continue;
380        };
381        for stored in mailbox.list().unwrap_or_default() {
382            links.insert(
383                stored.envelope.id.clone(),
384                stored.envelope.in_reply_to.clone(),
385            );
386        }
387    }
388    let mut depth = 0;
389    let mut current = Some(answered.to_string());
390    while let Some(id) = current {
391        if depth > REPLY_CHAIN_LIMIT || !links.contains_key(&id) {
392            break;
393        }
394        depth += 1;
395        current = links.get(&id).cloned().flatten();
396    }
397    depth
398}
399
400fn send_log(sender: &MailAddress) -> PathBuf {
401    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
402    mail_root().join("sent").join(&hash[..24]).join("log.jsonl")
403}
404
405fn record_send(sender: &MailAddress, to: &MailAddress, message_id: &str) {
406    let log = send_log(sender);
407    if let Some(parent) = log.parent() {
408        std::fs::create_dir_all(parent).ok();
409    }
410    if let Ok(mut file) = std::fs::OpenOptions::new()
411        .create(true)
412        .append(true)
413        .open(&log)
414    {
415        use std::io::Write as _;
416        let entry = serde_json::json!({"to": to.to_string(), "at_ms": now_ms(), "id": message_id});
417        writeln!(file, "{entry}").ok();
418    }
419}
420
421fn now_ms() -> u64 {
422    std::time::SystemTime::now()
423        .duration_since(std::time::UNIX_EPOCH)
424        .map(|elapsed| elapsed.as_millis() as u64)
425        .unwrap_or_default()
426}
427
428/// Marker recording that `message_id` was sent by `sender`, for `--id`.
429fn sent_marker(sender: &MailAddress, message_id: &str) -> PathBuf {
430    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
431    mail_root().join("sent").join(&hash[..24]).join(message_id)
432}