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    let (address, name) = match MailAddress::parse(to) {
153        Ok(address) if address.harness == "operator" => {
154            let name = format!("{}@{}", address.session_id, address.machine);
155            (address, name)
156        }
157        _ => match LiveSessions::read(homes).resolve(to) {
158            Ok(session) => (session.address.clone(), session.name.clone()),
159            Err(Unresolved::Stale(message) | Unresolved::Unknown(message)) => {
160                return Ok(Outcome::new(EXIT_UNKNOWN, message));
161            }
162        },
163    };
164    if address == caller.address {
165        return Ok(Outcome::new(
166            EXIT_REFUSED,
167            "Not sent: that address is your own session.",
168        ));
169    }
170    let door = match door_for(homes, &address) {
171        Ok(door) => door,
172        Err(_) => {
173            return Ok(Outcome::new(
174                EXIT_UNKNOWN,
175                format!(
176                    "{name} is no longer running. Nothing was sent. Run supercode message list for \
177                     the live sessions."
178                ),
179            ))
180        }
181    };
182    let message_id = match &idempotency_key {
183        Some(key) => format!(
184            "m-{}",
185            &blake3::hash(format!("{}\0{key}", caller.address).as_bytes()).to_hex()[..24]
186        ),
187        None => crate::mailbox::new_message_id()?,
188    };
189    let sent_marker = sent_marker(&caller.address, &message_id);
190    if let Ok(previous) = std::fs::read_to_string(&sent_marker) {
191        return Ok(Outcome::new(
192            0,
193            format!(
194                "Already sent with --id {}: {previous} Nothing was sent again.",
195                idempotency_key.as_deref().unwrap_or_default()
196            ),
197        ));
198    }
199    if let Some(answered) = in_reply_to.as_deref() {
200        if reply_chain_depth(&caller.address, &address, answered) >= REPLY_CHAIN_LIMIT {
201            return Ok(Outcome::new(
202                EXIT_REFUSED,
203                format!(
204                    "Not sent: you and {name} have answered each other {REPLY_CHAIN_LIMIT} times \
205                     in a row. If you are trading acknowledgements or status, stop; to go on, \
206                     send a new message (without --re), or tell your user."
207                ),
208            ));
209        }
210    }
211    let mut envelope = Envelope::new(
212        caller.address.clone(),
213        caller.name.clone(),
214        MailKind::Peer,
215        ReplyVia::Command,
216        body,
217    )?;
218    envelope.id = message_id.clone();
219    envelope.in_reply_to = in_reply_to;
220    let tier = door.name();
221    let idle_note = if notify_when_idle {
222        " Subscribed: one idle notice reaches you when its next turn ends."
223    } else {
224        ""
225    };
226    let (code, what) = match route(&envelope, &address, &door, !queue, notify_when_idle).await {
227        Err(detail) => {
228            return Ok(Outcome::new(
229                EXIT_FAILED,
230                format!("Not sent to {name}: {detail}"),
231            ))
232        }
233        Ok(Err(Refused::CannotQueueNative)) => {
234            return Ok(Outcome::new(
235                EXIT_REFUSED,
236                format!(
237                    "Not sent: {name} is idle, and a Claude session always starts a turn when a \
238                     message arrives, so --queue cannot hold it. Send without --queue to wake it."
239                ),
240            ))
241        }
242        Ok(Err(Refused::TooLong(bytes))) => {
243            return Ok(Outcome::new(
244                EXIT_REFUSED,
245                format!(
246                    "Not sent: the message is {bytes} bytes; the limit is {}. Write it to a file \
247                     and send its path instead.",
248                    crate::mail_route::MAX_RELAYED_BYTES
249                ),
250            ))
251        }
252        Ok(Ok(Delivered::Steered)) => (0, "steered into its running turn"),
253        Ok(Ok(Delivered::Started)) => (0, "it was idle, so the message started a turn"),
254        Ok(Ok(Delivered::Native { busy: true })) => (0, "it will read it at its next tool call"),
255        Ok(Ok(Delivered::Native { busy: false })) => {
256            (0, "it was idle, so the message starts its next turn")
257        }
258        Ok(Ok(Delivered::Hooked)) => (
259            0,
260            "its hook points it at the message at its next tool call, or when its turn ends; an \
261             idle session sees it on its next turn",
262        ),
263        Ok(Ok(Delivered::Queued)) => (
264            0,
265            "it is idle and --queue leaves it so; the message waits in its mailbox",
266        ),
267        Ok(Ok(Delivered::Operator)) => (0, "filed in its mailbox"),
268        Ok(Ok(Delivered::Stored)) => (
269            EXIT_STORED,
270            "Stored (not failed): it has no delivery door, so it sees this only if it runs \
271             supercode message inbox (a Codex session gets one with: supercode message setup \
272             codex). Don't resend, and don't wait for a reply",
273        ),
274    };
275    record_send(&caller.address, &address, &message_id);
276    let text = if code == 0 {
277        format!(
278            "sent to {name} ({}, {tier}): {what}. Message id {message_id}.{idle_note} Don't poll; \
279             carry on.",
280            address.harness
281        )
282    } else {
283        format!(
284            "{name} ({}): {what}. Message id {message_id}.{idle_note}",
285            address.harness
286        )
287    };
288    if code == 0 && idempotency_key.is_some() {
289        if let Some(parent) = sent_marker.parent() {
290            std::fs::create_dir_all(parent).ok();
291        }
292        std::fs::write(&sent_marker, &text).ok();
293    }
294    Ok(Outcome::new(code, text))
295}
296
297/// Replies in one unbroken chain between two sessions beyond which another
298/// reply is refused: two agents trading acknowledgements never stop on their
299/// own, and every message of such a loop answers the one before it. A new
300/// message (no `--re`) starts a chain, so volume alone is never refused.
301const REPLY_CHAIN_LIMIT: usize = 32;
302
303/// How many messages deep the reply chain ending at `answered` runs, between
304/// `a` and `b`: each message's `in_reply_to`, followed back through both
305/// sessions' mailboxes.
306fn reply_chain_depth(a: &MailAddress, b: &MailAddress, answered: &str) -> usize {
307    let mut links = std::collections::HashMap::new();
308    for address in [a, b] {
309        let Ok(mailbox) = Mailbox::open(&mail_root(), address) else {
310            continue;
311        };
312        for stored in mailbox.list().unwrap_or_default() {
313            links.insert(
314                stored.envelope.id.clone(),
315                stored.envelope.in_reply_to.clone(),
316            );
317        }
318    }
319    let mut depth = 0;
320    let mut current = Some(answered.to_string());
321    while let Some(id) = current {
322        if depth > REPLY_CHAIN_LIMIT || !links.contains_key(&id) {
323            break;
324        }
325        depth += 1;
326        current = links.get(&id).cloned().flatten();
327    }
328    depth
329}
330
331fn send_log(sender: &MailAddress) -> PathBuf {
332    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
333    mail_root().join("sent").join(&hash[..24]).join("log.jsonl")
334}
335
336fn record_send(sender: &MailAddress, to: &MailAddress, message_id: &str) {
337    let log = send_log(sender);
338    if let Some(parent) = log.parent() {
339        std::fs::create_dir_all(parent).ok();
340    }
341    if let Ok(mut file) = std::fs::OpenOptions::new()
342        .create(true)
343        .append(true)
344        .open(&log)
345    {
346        use std::io::Write as _;
347        let entry = serde_json::json!({"to": to.to_string(), "at_ms": now_ms(), "id": message_id});
348        writeln!(file, "{entry}").ok();
349    }
350}
351
352fn now_ms() -> u64 {
353    std::time::SystemTime::now()
354        .duration_since(std::time::UNIX_EPOCH)
355        .map(|elapsed| elapsed.as_millis() as u64)
356        .unwrap_or_default()
357}
358
359/// Marker recording that `message_id` was sent by `sender`, for `--id`.
360fn sent_marker(sender: &MailAddress, message_id: &str) -> PathBuf {
361    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
362    mail_root().join("sent").join(&hash[..24]).join(message_id)
363}