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