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