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;
8use std::time::Duration;
9
10use crate::mail_route::{Caller, LiveSessions, Unresolved};
11use crate::mailbox::{local_machine_name, mail_root, Envelope, MailAddress, MailKind, ReplyVia};
12use crate::HarnessHomes;
13
14/// Exit code of a refused send.
15pub const EXIT_REFUSED: i32 = 2;
16/// Exit code of an unknown or stale receiver.
17pub const EXIT_UNKNOWN: i32 = 3;
18/// Exit code of a message stored with no door to show it.
19pub const EXIT_STORED: i32 = 4;
20/// Exit code of a failed send.
21pub const EXIT_FAILED: i32 = 5;
22
23/// What a delivery came to: the exit code, and the text the sending agent reads.
24pub struct Outcome {
25    /// Exit code for scripts: 0 sent, [`EXIT_REFUSED`], [`EXIT_UNKNOWN`],
26    /// [`EXIT_STORED`] or [`EXIT_FAILED`].
27    pub code: i32,
28    /// What the sending agent reads.
29    pub text: String,
30}
31
32impl Outcome {
33    /// An outcome with this code and text.
34    pub fn new(code: i32, text: impl Into<String>) -> Self {
35        Self {
36            code,
37            text: text.into(),
38        }
39    }
40}
41
42/// Options of one send.
43#[derive(Debug, Clone, Default)]
44pub struct SendOptions {
45    /// Id of the message this one answers.
46    pub in_reply_to: Option<String>,
47    /// Also send one notice when the receiver's next turn ends.
48    pub notify_when_idle: bool,
49    /// Only enqueue; never start a turn in an idle receiver.
50    pub queue: bool,
51    /// Idempotency key: repeating a send with it does not send twice.
52    pub idempotency_key: Option<String>,
53}
54
55/// Send `body` from `caller` to `to` (an address, `name@machine`, or a name
56/// unique on this machine): through Teams when `to` is on another machine,
57/// else through the router's door here. The one send every sender uses.
58pub async fn send(
59    homes: &HarnessHomes,
60    caller: &Caller,
61    to: &str,
62    body: &str,
63    options: SendOptions,
64) -> std::io::Result<Outcome> {
65    let SendOptions {
66        in_reply_to,
67        notify_when_idle,
68        queue,
69        idempotency_key,
70    } = options;
71    if let Some(machine) = remote_machine(to) {
72        let request = serde_json::json!({
73            "op": "send",
74            "from": caller.address.to_string(),
75            "from_name": caller.name,
76            "to": to,
77            "body": body,
78            "in_reply_to": in_reply_to,
79            "notify_when_idle": notify_when_idle,
80            "queue": queue,
81            "id": idempotency_key,
82        });
83        return Ok(remote_send(&machine, &request));
84    }
85    deliver(
86        homes,
87        caller,
88        to,
89        body,
90        in_reply_to,
91        notify_when_idle,
92        queue,
93        idempotency_key,
94    )
95    .await
96}
97
98/// The machine `to` names when it is not this one.
99fn remote_machine(to: &str) -> Option<String> {
100    let local = local_machine_name();
101    let machine = match MailAddress::parse(to) {
102        Ok(address) => address.machine,
103        Err(_) => to.rsplit_once('@')?.1.to_string(),
104    };
105    (machine != local).then_some(machine)
106}
107
108/// Hand a request to another machine's mail door through Teams.
109fn remote_send(machine: &str, request: &serde_json::Value) -> Outcome {
110    match crate::mailbox::teams_mail(machine, request) {
111        Ok(answer) => Outcome::new(
112            answer["code"]
113                .as_i64()
114                .map_or(EXIT_FAILED, |code| code as i32),
115            answer["text"]
116                .as_str()
117                .or_else(|| answer["detail"].as_str())
118                .unwrap_or("the other machine answered nothing readable")
119                .to_string(),
120        ),
121        Err(detail) if detail.contains("no mail grant") => Outcome::new(
122            EXIT_REFUSED,
123            format!(
124                "Not sent: machine {machine} does not accept mail from you (no mail grant). Only \
125                 an operator can grant it; tell your user. Don't retry."
126            ),
127        ),
128        Err(detail) => Outcome::new(
129            EXIT_FAILED,
130            format!("Not sent to machine {machine}: {detail}"),
131        ),
132    }
133}
134
135/// Deliver from `caller` to `to` on this machine, through the one router
136/// every sender uses (`crate::mail_route`).
137#[allow(clippy::too_many_arguments)]
138async fn deliver(
139    homes: &HarnessHomes,
140    caller: &Caller,
141    to: &str,
142    body: &str,
143    in_reply_to: Option<String>,
144    notify_when_idle: bool,
145    queue: bool,
146    idempotency_key: Option<String>,
147) -> std::io::Result<Outcome> {
148    use crate::mail_route::{deliver as route, door_for, Delivered, Refused};
149    // An operator (a board, a script) is not a listed session; its address is
150    // taken as given.
151    let (address, name) = match MailAddress::parse(to) {
152        Ok(address) if address.harness == "operator" => {
153            let name = format!("{}@{}", address.session_id, address.machine);
154            (address, name)
155        }
156        _ => match LiveSessions::read(homes).resolve(to) {
157            Ok(session) => (session.address.clone(), session.name.clone()),
158            Err(Unresolved::Stale(message) | Unresolved::Unknown(message)) => {
159                return Ok(Outcome::new(EXIT_UNKNOWN, message));
160            }
161        },
162    };
163    if address == caller.address {
164        return Ok(Outcome::new(
165            EXIT_REFUSED,
166            "Not sent: that address is your own session.",
167        ));
168    }
169    let door = match door_for(homes, &address) {
170        Ok(door) => door,
171        Err(_) => {
172            return Ok(Outcome::new(
173                EXIT_UNKNOWN,
174                format!(
175                    "{name} is no longer running. Nothing was sent. Run supercode message list for \
176                     the live sessions."
177                ),
178            ))
179        }
180    };
181    let message_id = match &idempotency_key {
182        Some(key) => format!(
183            "m-{}",
184            &blake3::hash(format!("{}\0{key}", caller.address).as_bytes()).to_hex()[..24]
185        ),
186        None => crate::mailbox::new_message_id()?,
187    };
188    let sent_marker = sent_marker(&caller.address, &message_id);
189    if let Ok(previous) = std::fs::read_to_string(&sent_marker) {
190        return Ok(Outcome::new(
191            0,
192            format!(
193                "Already sent with --id {}: {previous} Nothing was sent again.",
194                idempotency_key.as_deref().unwrap_or_default()
195            ),
196        ));
197    }
198    if recent_sends(&caller.address, &address) >= LOOP_LIMIT {
199        return Ok(Outcome::new(
200            EXIT_REFUSED,
201            format!(
202                "Not sent: this would be message {} from you to {name} in {} minutes. If you two \
203                 are trading acknowledgements or status, stop; tell your user if the exchange must \
204                 go on.",
205                LOOP_LIMIT + 1,
206                LOOP_WINDOW.as_secs() / 60
207            ),
208        ));
209    }
210    let mut envelope = Envelope::new(
211        caller.address.clone(),
212        caller.name.clone(),
213        MailKind::Peer,
214        ReplyVia::Command,
215        body,
216    )?;
217    envelope.id = message_id.clone();
218    envelope.in_reply_to = in_reply_to;
219    let tier = door.name();
220    let idle_note = if notify_when_idle {
221        " Subscribed: one idle notice reaches you when its next turn ends."
222    } else {
223        ""
224    };
225    let (code, what) = match route(&envelope, &address, &door, !queue, notify_when_idle).await {
226        Err(detail) => {
227            return Ok(Outcome::new(
228                EXIT_FAILED,
229                format!("Not sent to {name}: {detail}"),
230            ))
231        }
232        Ok(Err(Refused::CannotQueueNative)) => {
233            return Ok(Outcome::new(
234                EXIT_REFUSED,
235                format!(
236                    "Not sent: {name} is idle, and a Claude session always starts a turn when a \
237                     message arrives, so --queue cannot hold it. Send without --queue to wake it."
238                ),
239            ))
240        }
241        Ok(Err(Refused::TooLong(bytes))) => {
242            return Ok(Outcome::new(
243                EXIT_REFUSED,
244                format!(
245                    "Not sent: the message is {bytes} bytes; the limit is {}. Write it to a file \
246                     and send its path instead.",
247                    crate::mail_route::MAX_RELAYED_BYTES
248                ),
249            ))
250        }
251        Ok(Ok(Delivered::Steered)) => (0, "steered into its running turn"),
252        Ok(Ok(Delivered::Started)) => (0, "it was idle, so the message started a turn"),
253        Ok(Ok(Delivered::Native { busy: true })) => (0, "it will read it at its next tool call"),
254        Ok(Ok(Delivered::Native { busy: false })) => {
255            (0, "it was idle, so the message starts its next turn")
256        }
257        Ok(Ok(Delivered::Hooked)) => (
258            0,
259            "its hook points it at the message at its next tool call, or when its turn ends; an \
260             idle session sees it on its next turn",
261        ),
262        Ok(Ok(Delivered::Queued)) => (
263            0,
264            "it is idle and --queue leaves it so; the message waits in its mailbox",
265        ),
266        Ok(Ok(Delivered::Operator)) => (0, "filed in its mailbox"),
267        Ok(Ok(Delivered::Stored)) => (
268            EXIT_STORED,
269            "Stored (not failed): it has no delivery door, so it sees this only if it runs \
270             supercode message inbox (a Codex session gets one with: supercode message setup \
271             codex). Don't resend, and don't wait for a reply",
272        ),
273    };
274    record_send(&caller.address, &address, &message_id);
275    let text = if code == 0 {
276        format!(
277            "sent to {name} ({}, {tier}): {what}. Message id {message_id}.{idle_note} Don't poll; \
278             carry on.",
279            address.harness
280        )
281    } else {
282        format!(
283            "{name} ({}): {what}. Message id {message_id}.{idle_note}",
284            address.harness
285        )
286    };
287    if code == 0 && idempotency_key.is_some() {
288        if let Some(parent) = sent_marker.parent() {
289            std::fs::create_dir_all(parent).ok();
290        }
291        std::fs::write(&sent_marker, &text).ok();
292    }
293    Ok(Outcome::new(code, text))
294}
295
296/// Messages to one peer within [`LOOP_WINDOW`] beyond which a send is
297/// refused: two agents trading acknowledgements never stop on their own.
298const LOOP_LIMIT: usize = 8;
299const LOOP_WINDOW: Duration = Duration::from_secs(10 * 60);
300
301fn send_log(sender: &MailAddress) -> PathBuf {
302    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
303    mail_root().join("sent").join(&hash[..24]).join("log.jsonl")
304}
305
306/// Sends from `sender` to `to` within the loop window.
307fn recent_sends(sender: &MailAddress, to: &MailAddress) -> usize {
308    let now = now_ms();
309    let window = LOOP_WINDOW.as_millis() as u64;
310    std::fs::read_to_string(send_log(sender))
311        .unwrap_or_default()
312        .lines()
313        .filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
314        .filter(|entry| entry["to"].as_str() == Some(to.to_string().as_str()))
315        .filter(|entry| {
316            entry["at_ms"]
317                .as_u64()
318                .is_some_and(|at| now.saturating_sub(at) < window)
319        })
320        .count()
321}
322
323fn record_send(sender: &MailAddress, to: &MailAddress, message_id: &str) {
324    let log = send_log(sender);
325    if let Some(parent) = log.parent() {
326        std::fs::create_dir_all(parent).ok();
327    }
328    if let Ok(mut file) = std::fs::OpenOptions::new()
329        .create(true)
330        .append(true)
331        .open(&log)
332    {
333        use std::io::Write as _;
334        let entry = serde_json::json!({"to": to.to_string(), "at_ms": now_ms(), "id": message_id});
335        writeln!(file, "{entry}").ok();
336    }
337}
338
339fn now_ms() -> u64 {
340    std::time::SystemTime::now()
341        .duration_since(std::time::UNIX_EPOCH)
342        .map(|elapsed| elapsed.as_millis() as u64)
343        .unwrap_or_default()
344}
345
346/// Marker recording that `message_id` was sent by `sender`, for `--id`.
347fn sent_marker(sender: &MailAddress, message_id: &str) -> PathBuf {
348    let hash = blake3::hash(sender.to_string().as_bytes()).to_hex();
349    mail_root().join("sent").join(&hash[..24]).join(message_id)
350}