Skip to main content

supercode_harness/
mail_agent.rs

1//! An agent's one mailbox, its threads and their participants
2//! (docs/adr/0008-agent-mailbox.md).
3//!
4//! An agent is addressed `sc:<machine>:agent:<name>`. Its record names its main
5//! session; a root addressed to it goes there. A thread is a root and its
6//! replies, with a participant list: the session holding it, main (CC once the
7//! thread is delegated), and whoever else wrote in it. Every message in a
8//! thread is filed in each participant's mailbox but its sender's: a `to`
9//! participant through its door, woken; a `cc` participant filed unread, not
10//! woken. [`plan`] decides who receives one message; the senders
11//! (`mail_send`, `sessions.message`) deliver it.
12
13use std::collections::BTreeMap;
14use std::path::{Path, PathBuf};
15use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
16
17use serde::{Deserialize, Serialize};
18
19use crate::mailbox::{local_machine_name, mail_root, Envelope, MailAddress, MailKind, Mailbox};
20
21/// The pseudo-harness of an agent address.
22pub const AGENT_HARNESS: &str = "agent";
23
24/// Prefix of a thread that has no root message: a session someone started
25/// directly in an agent's folder.
26pub const SESSION_THREAD_PREFIX: &str = "s-";
27
28/// One agent, as the mailbox needs it.
29#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
30pub struct Agent {
31    /// Its name, unique on this machine; its address is `sc:<machine>:agent:<name>`.
32    pub name: String,
33    /// The session its work continues in: every root lands here.
34    pub main_session: MailAddress,
35    /// Where its sessions start (`delegate`'s working directory, and the folder
36    /// whose directly started sessions are threads).
37    #[serde(default, skip_serializing_if = "Option::is_none")]
38    pub folder: Option<PathBuf>,
39    /// The program `delegate` starts (`claude`).
40    #[serde(default, skip_serializing_if = "Option::is_none")]
41    pub harness: Option<String>,
42    /// Minutes a session the mailbox launched for it may stay idle before its
43    /// process ends; its transcript stays, and the next reply resumes it.
44    #[serde(default, skip_serializing_if = "Option::is_none")]
45    pub idle_minutes: Option<u64>,
46}
47
48impl Agent {
49    /// Its address.
50    pub fn address(&self) -> std::io::Result<MailAddress> {
51        agent_address(&self.name)
52    }
53}
54
55/// The address of the agent `name` on this machine.
56pub fn agent_address(name: &str) -> std::io::Result<MailAddress> {
57    MailAddress::new(local_machine_name(), AGENT_HARNESS, name)
58        .map_err(|error| std::io::Error::other(error.to_string()))
59}
60
61fn agents_dir() -> PathBuf {
62    mail_root().join("agents")
63}
64
65fn threads_dir() -> PathBuf {
66    mail_root().join("threads")
67}
68
69fn valid_name(name: &str) -> bool {
70    !name.is_empty()
71        && name
72            .chars()
73            .all(|character| character.is_ascii_alphanumeric() || "-_.".contains(character))
74}
75
76fn write_atomically(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
77    if let Some(parent) = path.parent() {
78        std::fs::create_dir_all(parent)?;
79    }
80    let temporary = path.with_extension(format!("tmp.{}", std::process::id()));
81    std::fs::write(&temporary, bytes)?;
82    std::fs::rename(temporary, path)
83}
84
85/// Record `agent`, replacing an earlier record of the same name.
86pub fn declare(agent: &Agent) -> std::io::Result<()> {
87    if !valid_name(&agent.name) {
88        return Err(std::io::Error::other(format!(
89            "`{}` is not an agent name: use letters, digits, `-`, `_` and `.`",
90            agent.name
91        )));
92    }
93    let bytes = serde_json::to_vec_pretty(agent).map_err(std::io::Error::other)?;
94    write_atomically(&agents_dir().join(format!("{}.json", agent.name)), &bytes)
95}
96
97/// The agent named `name`, when one is declared.
98pub fn load(name: &str) -> Option<Agent> {
99    if !valid_name(name) {
100        return None;
101    }
102    let bytes = std::fs::read(agents_dir().join(format!("{name}.json"))).ok()?;
103    serde_json::from_slice(&bytes).ok()
104}
105
106/// Every declared agent.
107pub fn all_agents() -> Vec<Agent> {
108    let Ok(entries) = std::fs::read_dir(agents_dir()) else {
109        return Vec::new();
110    };
111    let mut agents: Vec<Agent> = entries
112        .flatten()
113        .filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
114        .filter_map(|entry| std::fs::read(entry.path()).ok())
115        .filter_map(|bytes| serde_json::from_slice(&bytes).ok())
116        .collect();
117    agents.sort_by(|a, b| a.name.cmp(&b.name));
118    agents
119}
120
121fn owners_account_manager_file() -> PathBuf {
122    agents_dir().join("owners-account-manager")
123}
124
125/// Name `name` the owner's account manager: every line the owner writes to any
126/// agent is filed in its main session's mailbox as CC.
127pub fn set_owners_account_manager(name: &str) -> std::io::Result<()> {
128    write_atomically(
129        &owners_account_manager_file(),
130        format!("{name}\n").as_bytes(),
131    )
132}
133
134/// The owner's account manager, when one is named and declared.
135pub fn owners_account_manager() -> Option<Agent> {
136    let name = std::fs::read_to_string(owners_account_manager_file()).ok()?;
137    load(name.trim())
138}
139
140/// A participant's part in a thread.
141#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
142#[serde(rename_all = "snake_case")]
143pub enum Role {
144    /// Receives every line and is woken by it.
145    To,
146    /// Receives every line, filed unread, and is not woken.
147    Cc,
148}
149
150/// One participant of a thread.
151#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
152pub struct Participant {
153    /// Its address.
154    pub address: MailAddress,
155    /// Its part.
156    pub role: Role,
157}
158
159/// A root and its replies.
160#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
161pub struct Thread {
162    /// The root's message id (`s-<session id>` for a session started directly).
163    pub id: String,
164    /// The agent it belongs to.
165    pub agent: String,
166    /// The agent's session holding it: main until it is delegated.
167    pub holder: MailAddress,
168    /// Everyone who receives its lines.
169    pub participants: Vec<Participant>,
170    /// A channel's platform ids of its messages, each to its message id.
171    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
172    pub markers: BTreeMap<String, String>,
173    /// Where it is shown, when it has a surface (`rh2:<room>/<conversation>`).
174    #[serde(default, skip_serializing_if = "Option::is_none")]
175    pub surface: Option<String>,
176    /// When it was opened, epoch milliseconds.
177    pub created_at_ms: u64,
178    /// The mailbox launched its holder (`message delegate`), so it may resume
179    /// it the same way; a holder it did not launch is never resumed by mail.
180    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
181    pub holder_launched: bool,
182    /// The thread a side conversation was opened from: a thread session's
183    /// own mail to someone other than main. Its lines never reach the parent
184    /// thread's other participants.
185    #[serde(default, skip_serializing_if = "Option::is_none")]
186    pub parent: Option<String>,
187}
188
189impl Thread {
190    /// The part `address` has in it.
191    pub fn role_of(&self, address: &MailAddress) -> Option<Role> {
192        self.participants
193            .iter()
194            .find(|participant| &participant.address == address)
195            .map(|participant| participant.role)
196    }
197
198    /// Add `address` as `role`, or change its role.
199    pub fn join(&mut self, address: &MailAddress, role: Role) {
200        match self
201            .participants
202            .iter_mut()
203            .find(|participant| &participant.address == address)
204        {
205            Some(participant) => participant.role = role,
206            None => self.participants.push(Participant {
207                address: address.clone(),
208                role,
209            }),
210        }
211    }
212
213    /// Whether its holder is someone other than the agent's main session.
214    pub fn delegated(&self) -> bool {
215        load(&self.agent).is_some_and(|agent| agent.main_session != self.holder)
216    }
217}
218
219fn thread_path(id: &str) -> Option<PathBuf> {
220    valid_name(id).then(|| threads_dir().join(format!("{id}.json")))
221}
222
223/// The thread `id`, when the mailbox keeps one.
224pub fn thread(id: &str) -> Option<Thread> {
225    let bytes = std::fs::read(thread_path(id)?).ok()?;
226    serde_json::from_slice(&bytes).ok()
227}
228
229/// Every thread the mailbox keeps, oldest first.
230pub fn all_threads() -> Vec<Thread> {
231    let Ok(entries) = std::fs::read_dir(threads_dir()) else {
232        return Vec::new();
233    };
234    let mut threads: Vec<Thread> = entries
235        .flatten()
236        .filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
237        .filter_map(|entry| std::fs::read(entry.path()).ok())
238        .filter_map(|bytes| serde_json::from_slice(&bytes).ok())
239        .collect();
240    threads.sort_by(|a, b| (a.created_at_ms, &a.id).cmp(&(b.created_at_ms, &b.id)));
241    threads
242}
243
244/// The threads of the agent `name`, oldest first.
245pub fn threads_of(name: &str) -> Vec<Thread> {
246    all_threads()
247        .into_iter()
248        .filter(|thread| thread.agent == name)
249        .collect()
250}
251
252/// The thread `address` holds for an agent other than as its main session:
253/// a delegated session holds exactly one (its side conversations aside).
254pub fn thread_held_by(address: &MailAddress) -> Option<Thread> {
255    all_threads().into_iter().rev().find(|thread| {
256        &thread.holder == address
257            && thread.parent.is_none()
258            && load(&thread.agent).is_some_and(|agent| &agent.main_session != address)
259    })
260}
261
262/// The thread whose channel message has the platform id `marker`.
263pub fn thread_of_marker(marker: &str) -> Option<(Thread, String)> {
264    all_threads().into_iter().find_map(|thread| {
265        let id = thread.markers.get(marker)?.clone();
266        Some((thread, id))
267    })
268}
269
270/// The agent `address` is a session of (its main, or a thread's holder).
271pub fn agent_of_session(address: &MailAddress) -> Option<Agent> {
272    if let Some(agent) = all_agents()
273        .into_iter()
274        .find(|agent| &agent.main_session == address)
275    {
276        return Some(agent);
277    }
278    thread_held_by(address).and_then(|thread| load(&thread.agent))
279}
280
281/// Change the thread `id` under its lock, so two writers never lose each
282/// other's participants, and return it as saved. `create` makes it when it is
283/// not kept yet.
284pub fn update_thread(
285    id: &str,
286    create: impl FnOnce() -> Option<Thread>,
287    change: impl FnOnce(&mut Thread),
288) -> std::io::Result<Option<Thread>> {
289    let Some(path) = thread_path(id) else {
290        return Ok(None);
291    };
292    std::fs::create_dir_all(threads_dir())?;
293    let lock = path.with_extension("lock");
294    let started = Instant::now();
295    loop {
296        match std::fs::OpenOptions::new()
297            .write(true)
298            .create_new(true)
299            .open(&lock)
300        {
301            Ok(_) => break,
302            Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
303                // A lock older than the wait is left by a writer that died.
304                let stale = std::fs::metadata(&lock)
305                    .and_then(|meta| meta.modified())
306                    .ok()
307                    .and_then(|modified| modified.elapsed().ok())
308                    .is_some_and(|age| age > Duration::from_secs(10));
309                if stale {
310                    std::fs::remove_file(&lock).ok();
311                } else if started.elapsed() > Duration::from_secs(30) {
312                    // A live writer is never overridden; a send that cannot get the lock fails.
313                    return Err(std::io::Error::new(
314                        std::io::ErrorKind::TimedOut,
315                        format!("thread {id} is locked by another writer"),
316                    ));
317                }
318                std::thread::sleep(Duration::from_millis(20));
319            }
320            Err(error) => return Err(error),
321        }
322    }
323    let result = (|| {
324        let current = std::fs::read(&path)
325            .ok()
326            .and_then(|bytes| serde_json::from_slice::<Thread>(&bytes).ok())
327            .or_else(create);
328        let Some(mut thread) = current else {
329            return Ok(None);
330        };
331        change(&mut thread);
332        let bytes = serde_json::to_vec_pretty(&thread).map_err(std::io::Error::other)?;
333        write_atomically(&path, &bytes)?;
334        Ok(Some(thread))
335    })();
336    std::fs::remove_file(&lock).ok();
337    result
338}
339
340/// One receiver of a message.
341#[derive(Debug, Clone, PartialEq, Eq)]
342pub struct Recipient {
343    /// Where it is filed.
344    pub address: MailAddress,
345    /// Delivered through its door, waking it; else filed unread.
346    pub wake: bool,
347}
348
349/// Who receives one message, as [`plan`] decides it.
350#[derive(Debug, Clone, PartialEq, Eq)]
351pub struct Plan {
352    /// Every receiver, in the order they are delivered to.
353    pub recipients: Vec<Recipient>,
354    /// The agent whose own mailbox keeps a record of it (a root addressed to it).
355    pub agent: Option<MailAddress>,
356    /// The thread it is in, as saved.
357    pub thread: Option<Thread>,
358}
359
360fn now_ms() -> u64 {
361    SystemTime::now()
362        .duration_since(UNIX_EPOCH)
363        .map(|elapsed| elapsed.as_millis() as u64)
364        .unwrap_or_default()
365}
366
367/// What a channel says about one of its lines: the platform's ids, which the
368/// mailbox maps to its own.
369#[derive(Debug, Clone, Default, PartialEq, Eq)]
370pub struct Channel {
371    /// The line's own platform id (`rh2:<message id>`).
372    pub marker: Option<String>,
373    /// Its parent's platform id, when the platform marks it a reply.
374    pub reply_marker: Option<String>,
375    /// Where the thread is shown (`rh2:<room>/<conversation>`).
376    pub surface: Option<String>,
377}
378
379fn record_marker(thread: &mut Thread, channel: &Channel, id: &str) {
380    if let Some(marker) = &channel.marker {
381        thread.markers.insert(marker.clone(), id.to_string());
382    }
383}
384
385/// Who receives `envelope`, sent to `to`, when an agent's mailbox decides it;
386/// `None` when it is ordinary mail to one session. Sets the envelope's
387/// `thread` and `in_reply_to` where the mailbox knows them better than the
388/// sender: a channel names its parent by its platform id.
389pub fn plan(
390    envelope: &mut Envelope,
391    to: &MailAddress,
392    channel: &Channel,
393) -> std::io::Result<Option<Plan>> {
394    let sender = envelope.from.clone();
395    let id = envelope.id.clone();
396    // A channel's reply names its parent by the platform's id.
397    if envelope.in_reply_to.is_none() {
398        if let Some((thread, parent)) = channel.reply_marker.as_deref().and_then(thread_of_marker) {
399            envelope.in_reply_to = Some(parent);
400            envelope.thread = Some(thread.id);
401        }
402    }
403
404    // 1. A reply in a thread the mailbox keeps: to everyone in it but its sender.
405    let known = envelope
406        .in_reply_to
407        .as_ref()
408        .and(envelope.thread.as_deref())
409        .and_then(thread);
410    if let Some(known) = known {
411        let receiver = (to.harness != AGENT_HARNESS && to != &sender).then(|| to.clone());
412        let saved = update_thread(
413            &known.id,
414            || None,
415            |thread| {
416                if thread.role_of(&sender).is_none() {
417                    thread.join(&sender, Role::To);
418                }
419                if let Some(receiver) = &receiver {
420                    if thread.role_of(receiver).is_none() {
421                        thread.join(receiver, Role::To);
422                    }
423                }
424                record_marker(thread, channel, &id);
425            },
426        )?
427        .unwrap_or(known);
428        envelope.thread = Some(saved.id.clone());
429        let recipients = saved
430            .participants
431            .iter()
432            .filter(|participant| participant.address != sender)
433            .map(|participant| Recipient {
434                address: participant.address.clone(),
435                wake: participant.role == Role::To,
436            })
437            .collect();
438        return Ok(Some(Plan {
439            recipients,
440            agent: None,
441            thread: Some(saved),
442        }));
443    }
444
445    // 2. New mail from a session holding a delegated thread: a question to
446    //    main stays in that thread and wakes main alone; mail to anyone else
447    //    opens a side thread of the session, the receiver and main (CC), so
448    //    the thread's other participants (a principal's channel) never get a
449    //    third party's lines: what reaches them is the agent's call.
450    if envelope.in_reply_to.is_none() {
451        if let Some(held) = thread_held_by(&sender) {
452            let agent = load(&held.agent);
453            let own_agent = to.harness == AGENT_HARNESS && to.session_id == held.agent;
454            let to_main = own_agent || agent.as_ref().is_some_and(|a| &a.main_session == to);
455            if to_main {
456                let main = agent.map(|a| a.main_session).unwrap_or_else(|| to.clone());
457                let saved = update_thread(
458                    &held.id,
459                    || None,
460                    |thread| record_marker(thread, channel, &id),
461                )?
462                .unwrap_or(held);
463                envelope.thread = Some(saved.id.clone());
464                return Ok(Some(Plan {
465                    recipients: vec![Recipient {
466                        address: main,
467                        wake: true,
468                    }],
469                    agent: None,
470                    thread: Some(saved),
471                }));
472            }
473            let root = id.clone();
474            // Another agent is reached at its main session, which answers in this side thread.
475            let receiver = if to.harness == AGENT_HARNESS {
476                load(&to.session_id)
477                    .map(|other| other.main_session)
478                    .ok_or_else(|| {
479                        std::io::Error::other(format!(
480                            "no agent named {} is declared on this machine",
481                            to.session_id
482                        ))
483                    })?
484            } else {
485                to.clone()
486            };
487            let mut participants = vec![
488                Participant {
489                    address: sender.clone(),
490                    role: Role::To,
491                },
492                Participant {
493                    address: receiver,
494                    role: Role::To,
495                },
496            ];
497            if let Some(agent) = &agent {
498                participants.push(Participant {
499                    address: agent.main_session.clone(),
500                    role: Role::Cc,
501                });
502            }
503            let saved = update_thread(
504                &root,
505                || {
506                    Some(Thread {
507                        id: root.clone(),
508                        agent: held.agent.clone(),
509                        holder: sender.clone(),
510                        participants,
511                        markers: BTreeMap::new(),
512                        surface: None,
513                        created_at_ms: now_ms(),
514                        holder_launched: false,
515                        parent: Some(held.id.clone()),
516                    })
517                },
518                |thread| record_marker(thread, channel, &id),
519            )?;
520            envelope.thread = Some(root);
521            let recipients = saved
522                .as_ref()
523                .map(|thread| {
524                    thread
525                        .participants
526                        .iter()
527                        .filter(|p| p.address != sender)
528                        .map(|p| Recipient {
529                            address: p.address.clone(),
530                            wake: p.role == Role::To,
531                        })
532                        .collect()
533                })
534                .unwrap_or_default();
535            return Ok(Some(Plan {
536                recipients,
537                agent: None,
538                thread: saved,
539            }));
540        }
541    }
542
543    // 3. A voice front speaks for a thread's holder: what it passes on joins that thread,
544    //    so main has the call as CC.
545    if envelope.in_reply_to.is_none() && envelope.voice_for.as_ref() == Some(to) {
546        if let Some(held) = thread_held_by(to) {
547            envelope.thread = Some(held.id.clone());
548            let mut recipients = vec![Recipient {
549                address: to.clone(),
550                wake: true,
551            }];
552            recipients.extend(
553                held.participants
554                    .iter()
555                    .filter(|p| p.role == Role::Cc && p.address != sender && &p.address != to)
556                    .map(|p| Recipient {
557                        address: p.address.clone(),
558                        wake: false,
559                    }),
560            );
561            return Ok(Some(Plan {
562                recipients,
563                agent: None,
564                thread: Some(held),
565            }));
566        }
567    }
568
569    // 4. Mail to an agent: a root, to its main session.
570    if to.harness == AGENT_HARNESS {
571        let agent = load(&to.session_id).ok_or_else(|| {
572            std::io::Error::other(format!(
573                "no agent named {} is declared on this machine",
574                to.session_id
575            ))
576        })?;
577        let main = agent.main_session.clone();
578        if main == sender {
579            return Err(std::io::Error::other(format!(
580                "that is your own agent's address ({}); you are its main session",
581                agent.name
582            )));
583        }
584        // Its own sessions reach main, and open nothing.
585        if agent_of_session(&sender).is_some_and(|own| own.name == agent.name) {
586            return Ok(Some(Plan {
587                recipients: vec![Recipient {
588                    address: main,
589                    wake: true,
590                }],
591                agent: None,
592                thread: None,
593            }));
594        }
595        envelope.thread = None;
596        let root = id.clone();
597        let surface = channel.surface.clone();
598        let saved = update_thread(
599            &root,
600            || {
601                Some(Thread {
602                    id: root.clone(),
603                    agent: agent.name.clone(),
604                    holder: main.clone(),
605                    participants: vec![
606                        Participant {
607                            address: main.clone(),
608                            role: Role::To,
609                        },
610                        Participant {
611                            address: sender.clone(),
612                            role: Role::To,
613                        },
614                    ],
615                    markers: BTreeMap::new(),
616                    surface,
617                    created_at_ms: now_ms(),
618                    holder_launched: false,
619                    parent: None,
620                })
621            },
622            |thread| record_marker(thread, channel, &id),
623        )?;
624        return Ok(Some(Plan {
625            recipients: vec![Recipient {
626                address: main,
627                wake: true,
628            }],
629            agent: Some(to.clone()),
630            thread: saved,
631        }));
632    }
633
634    // 5. Main's new mail to anyone else opens a thread it holds, so the answer
635    //    comes back into it.
636    if envelope.in_reply_to.is_none() {
637        if let Some(agent) = all_agents()
638            .into_iter()
639            .find(|agent| agent.main_session == sender)
640        {
641            if to != &sender && to.harness != AGENT_HARNESS {
642                let root = id.clone();
643                let saved = update_thread(
644                    &root,
645                    || {
646                        Some(Thread {
647                            id: root.clone(),
648                            agent: agent.name.clone(),
649                            holder: sender.clone(),
650                            participants: vec![
651                                Participant {
652                                    address: sender.clone(),
653                                    role: Role::To,
654                                },
655                                Participant {
656                                    address: to.clone(),
657                                    role: Role::To,
658                                },
659                            ],
660                            markers: BTreeMap::new(),
661                            surface: None,
662                            created_at_ms: now_ms(),
663                            holder_launched: false,
664                            parent: None,
665                        })
666                    },
667                    |thread| record_marker(thread, channel, &id),
668                )?;
669                return Ok(Some(Plan {
670                    recipients: vec![Recipient {
671                        address: to.clone(),
672                        wake: true,
673                    }],
674                    agent: None,
675                    thread: saved,
676                }));
677            }
678        }
679    }
680    Ok(None)
681}
682
683/// Hand `thread` to `holder`: it becomes the thread's `to`, and main `cc`.
684/// `launched` says the mailbox is starting the holder itself.
685/// A line that reached an agent's main session through its own input (its terminal, or its DM in a Room: a channel's
686/// turn into that session) is a root of the agent (RFC 0020 decision 11, "every line is a root"): its thread, held by
687/// main, with its sender, opened when main first acts on it. `None` when `main` is no agent's main session or the
688/// line is main's own.
689pub fn root_of_main(main: &MailAddress, envelope: &Envelope) -> std::io::Result<Option<Thread>> {
690    let Some(agent) = agent_of_session(main).filter(|agent| &agent.main_session == main) else {
691        return Ok(None);
692    };
693    if &envelope.from == main || envelope.thread_id() != envelope.id {
694        return Ok(None);
695    }
696    let sender = envelope.from.clone();
697    update_thread(
698        &envelope.id,
699        || {
700            Some(Thread {
701                id: envelope.id.clone(),
702                agent: agent.name.clone(),
703                holder: main.clone(),
704                participants: vec![
705                    Participant {
706                        address: main.clone(),
707                        role: Role::To,
708                    },
709                    Participant {
710                        address: sender,
711                        role: Role::To,
712                    },
713                ],
714                markers: BTreeMap::new(),
715                surface: None,
716                created_at_ms: now_ms(),
717                holder_launched: false,
718                parent: None,
719            })
720        },
721        |_| {},
722    )
723}
724
725pub fn delegate_to(
726    thread_id: &str,
727    holder: &MailAddress,
728    launched: bool,
729) -> std::io::Result<Option<Thread>> {
730    let Some(current) = thread(thread_id) else {
731        return Ok(None);
732    };
733    let Some(agent) = load(&current.agent) else {
734        return Ok(None);
735    };
736    update_thread(
737        thread_id,
738        || None,
739        |thread| {
740            // Handed back to main: a holder the mailbox launched leaves the thread with it.
741            if holder == &agent.main_session && thread.holder_launched && &thread.holder != holder {
742                let previous = thread.holder.clone();
743                thread.participants.retain(|p| p.address != previous);
744            }
745            thread.holder = holder.clone();
746            thread.holder_launched = launched && holder != &agent.main_session;
747            thread.join(holder, Role::To);
748            if holder != &agent.main_session {
749                thread.join(&agent.main_session, Role::Cc);
750            }
751        },
752    )
753}
754
755/// How long a send waits for a resumed session to be reachable again.
756pub const RESUME_WAIT: Duration = Duration::from_secs(45);
757
758/// Whether mail may resume the stopped session at `address`: only a thread's
759/// holder the mailbox launched itself, which `open` resumes with the arguments
760/// it was launched with. Any other stopped session keeps its mail waiting.
761pub fn resumable(address: &MailAddress) -> bool {
762    all_threads()
763        .iter()
764        .any(|thread| &thread.holder == address && thread.holder_launched)
765        && recorded_resume_arguments(address).is_some()
766}
767
768/// The mode a stopped session resumes in: its own recorded permissions
769/// (`harness.v1.sessions.recorded_config`, ADR 0012), as `supercode open`
770/// arguments after `--`, which suppress open's unattended defaults. `None`
771/// when they are unrecorded or of a kind not carried, and then its mail waits.
772fn recorded_resume_arguments(address: &MailAddress) -> Option<Vec<String>> {
773    let query = crate::DiscoveryQuery {
774        harnesses: vec![crate::HarnessId::new(&address.harness)],
775        query: Some(address.session_id.clone()),
776        limit: Some(50),
777        ..Default::default()
778    };
779    let locator = crate::sdk::discover_sessions(&query)
780        .ok()?
781        .into_iter()
782        .find(|row| row.locator.session_id == address.session_id)?
783        .locator;
784    let config = crate::harness_service::recorded_config(&locator).ok()?;
785    if config.get("session_id").and_then(|v| v.as_str()) != Some(address.session_id.as_str()) {
786        return None;
787    }
788    resume_arguments(&config)
789}
790
791/// A recorded configuration as resume arguments: Codex's approval and sandbox
792/// policy (as the Room-notice resume carries them), Claude Code's permission
793/// mode, and the recorded model.
794fn resume_arguments(config: &serde_json::Value) -> Option<Vec<String>> {
795    let text = |key: &str| config.get(key).and_then(|v| v.as_str());
796    let mut args = Vec::new();
797    match text("harness")? {
798        "codex" => {
799            let approval = text("approval_policy")?;
800            if !["untrusted", "on-failure", "on-request", "never"].contains(&approval) {
801                return None;
802            }
803            let policy = config.get("sandbox_policy")?.as_object()?;
804            let kind = policy.get("type")?.as_str()?;
805            let allowed: &[&str] = match kind {
806                "read-only" | "danger-full-access" => &["type"],
807                "workspace-write" => &[
808                    "type",
809                    "writable_roots",
810                    "network_access",
811                    "exclude_tmpdir_env_var",
812                    "exclude_slash_tmp",
813                ],
814                _ => return None,
815            };
816            if policy.keys().any(|key| !allowed.contains(&key.as_str())) {
817                return None;
818            }
819            args.extend([
820                "--ask-for-approval".into(),
821                approval.into(),
822                "--sandbox".into(),
823                kind.into(),
824            ]);
825            if kind == "workspace-write" {
826                let roots = policy
827                    .get("writable_roots")
828                    .cloned()
829                    .unwrap_or(serde_json::json!([]));
830                if !roots
831                    .as_array()?
832                    .iter()
833                    .all(|root| root.as_str().is_some_and(|r| Path::new(r).is_absolute()))
834                {
835                    return None;
836                }
837                args.extend([
838                    "-c".into(),
839                    format!("sandbox_workspace_write.writable_roots={roots}"),
840                ]);
841                for key in [
842                    "network_access",
843                    "exclude_tmpdir_env_var",
844                    "exclude_slash_tmp",
845                ] {
846                    let value = match policy.get(key) {
847                        Some(value) => value.as_bool()?,
848                        None => key != "network_access",
849                    };
850                    args.extend([
851                        "-c".into(),
852                        format!("sandbox_workspace_write.{key}={value}"),
853                    ]);
854                }
855            }
856        }
857        "claude-code" => {
858            let mode = text("permission_mode")?;
859            if !["default", "acceptEdits", "plan", "bypassPermissions"].contains(&mode) {
860                return None;
861            }
862            args.extend(["--permission-mode".into(), mode.into()]);
863        }
864        _ => return None,
865    }
866    if let Some(model) = text("model").filter(|m| !m.is_empty()) {
867        args.extend(["--model".into(), model.into()]);
868    }
869    Some(args)
870}
871
872/// When the mailbox last resumed `address`: its idle time counts from then, not from its last message before it
873/// stopped, so a resumed session is not closed again at once.
874fn resumed_path(address: &MailAddress) -> PathBuf {
875    let hash = blake3::hash(address.to_string().as_bytes()).to_hex();
876    agents_dir().join("resumed").join(&hash[..24])
877}
878
879/// Resume the stopped session at `address` in a daemon pane, with its
880/// harness's own resume (`supercode open <session> --detach -- <its recorded
881/// mode>`): never open's unattended defaults.
882pub fn resume(address: &MailAddress) -> std::io::Result<()> {
883    let mut recorded = recorded_resume_arguments(address).ok_or_else(|| {
884        std::io::Error::other(format!(
885            "{address} has no recorded permissions to resume with; its mail waits"
886        ))
887    })?;
888    // A Claude session resumed in the agent's folder would derive the folder's name, which its main session (live in
889    // that folder) already holds; mail reaches a Claude session by its name, so the thread's next reply would wake main.
890    // The holder comes back under a name of its own: its agent's, and the start of its session id.
891    if address.harness == "claude-code" {
892        if let Some(thread) = all_threads()
893            .into_iter()
894            .find(|thread| &thread.holder == address)
895        {
896            let id: String = address.session_id.chars().take(6).collect();
897            recorded.extend(["--name".to_string(), format!("{}-{id}", thread.agent)]);
898        }
899    }
900    let program = crate::claude_relay::supercode_program().map_err(std::io::Error::other)?;
901    let marker = resumed_path(address);
902    if let Some(dir) = marker.parent() {
903        std::fs::create_dir_all(dir)?;
904    }
905    write_atomically(&marker, now_ms().to_string().as_bytes())?;
906    let output = std::process::Command::new(program)
907        // `open` names a saved session by its id; the mail address is not a name it resolves.
908        .args(["open", &address.session_id, "--detach", "--"])
909        .args(&recorded)
910        .stdin(std::process::Stdio::null())
911        .output()?;
912    if output.status.success() {
913        Ok(())
914    } else {
915        Err(std::io::Error::other(crate::mailbox::error_line(
916            &String::from_utf8_lossy(&output.stderr),
917        )))
918    }
919}
920
921/// File `envelope` unread in `address`'s mailbox without waking it.
922pub fn file_unread(address: &MailAddress, envelope: &Envelope) -> std::io::Result<()> {
923    Mailbox::open(&mail_root(), address)?
924        .deliver(envelope)
925        .map(|_| ())
926}
927
928/// Register as a thread every running session started directly in an agent's
929/// folder that is neither its main session nor a thread's holder: its terminal
930/// is a thread (`s-<session id>`), with main CC.
931pub fn register_folder_sessions(homes: &crate::HarnessHomes) {
932    let agents: Vec<Agent> = all_agents()
933        .into_iter()
934        .filter(|agent| agent.folder.is_some())
935        .collect();
936    if agents.is_empty() {
937        return;
938    }
939    let holders: Vec<MailAddress> = all_threads()
940        .into_iter()
941        .map(|thread| thread.holder)
942        .collect();
943    for session in crate::mail_route::LiveSessions::read(homes).all() {
944        let Some(cwd) = &session.cwd else { continue };
945        let Some(agent) = agents.iter().find(|agent| {
946            agent.folder.as_deref() == Some(cwd.as_path()) && agent.main_session != session.address
947        }) else {
948            continue;
949        };
950        if holders.contains(&session.address) {
951            continue;
952        }
953        let id = format!("{SESSION_THREAD_PREFIX}{}", session.address.session_id);
954        let holder = session.address.clone();
955        let main = agent.main_session.clone();
956        update_thread(
957            &id,
958            || {
959                Some(Thread {
960                    id: id.clone(),
961                    agent: agent.name.clone(),
962                    holder: holder.clone(),
963                    participants: vec![
964                        Participant {
965                            address: holder.clone(),
966                            role: Role::To,
967                        },
968                        Participant {
969                            address: main.clone(),
970                            role: Role::Cc,
971                        },
972                    ],
973                    markers: BTreeMap::new(),
974                    surface: None,
975                    created_at_ms: now_ms(),
976                    holder_launched: false,
977                    parent: None,
978                })
979            },
980            |_| {},
981        )
982        .ok();
983    }
984}
985
986fn typed_cursor_path(address: &MailAddress) -> PathBuf {
987    let hash = blake3::hash(address.to_string().as_bytes()).to_hex();
988    agents_dir().join("typed").join(&hash[..24])
989}
990
991/// Copy each line a person typed into an agent's session, since the last pass,
992/// to whoever is CC on it: a delegated thread's holder's lines to the thread's
993/// CC participants (decision 11), and every agent session's lines to the owner's
994/// account manager (decision 17), filed unread and not woken. A session seen for
995/// the first time starts from now, so its history is not copied.
996pub fn copy_typed_lines(homes: &crate::HarnessHomes) {
997    let agents = all_agents();
998    if agents.is_empty() {
999        return;
1000    }
1001    let threads = all_threads();
1002    let account_manager = owners_account_manager();
1003    let mut sessions: Vec<(MailAddress, String)> = agents
1004        .iter()
1005        .map(|agent| (agent.main_session.clone(), agent.name.clone()))
1006        .collect();
1007    for thread in &threads {
1008        if !sessions
1009            .iter()
1010            .any(|(address, _)| address == &thread.holder)
1011        {
1012            sessions.push((thread.holder.clone(), thread.agent.clone()));
1013        }
1014    }
1015    let now = now_ms();
1016    for (session, agent_name) in sessions {
1017        let cursor_path = typed_cursor_path(&session);
1018        let Some(cursor) = std::fs::read_to_string(&cursor_path)
1019            .ok()
1020            .and_then(|text| text.trim().parse::<u64>().ok())
1021        else {
1022            write_atomically(&cursor_path, now.to_string().as_bytes()).ok();
1023            continue;
1024        };
1025        // A transcript not written since the last pass has no new line.
1026        if let Some(path) = crate::mail_transcript::transcript(homes, &session) {
1027            let changed = std::fs::metadata(&path)
1028                .and_then(|meta| meta.modified())
1029                .ok()
1030                .and_then(|modified| modified.duration_since(UNIX_EPOCH).ok())
1031                .map(|since| since.as_millis() as u64);
1032            if changed.is_some_and(|changed| changed <= cursor) {
1033                continue;
1034            }
1035        }
1036        // A Claude session's lines are read from its transcript; a Codex session's were filed
1037        // in its mailbox by its UserPromptSubmit hook.
1038        let lines: Vec<(u64, String)> = if session.harness == "codex" {
1039            Mailbox::open(&mail_root(), &session)
1040                .and_then(|mailbox| mailbox.list())
1041                .unwrap_or_default()
1042                .into_iter()
1043                .filter(|stored| {
1044                    stored.envelope.kind == MailKind::User
1045                        && stored
1046                            .envelope
1047                            .id
1048                            .starts_with(crate::mail_transcript::TYPED_ID_PREFIX)
1049                        && stored.envelope.created_at_ms > cursor
1050                })
1051                .map(|stored| (stored.envelope.created_at_ms, stored.envelope.body))
1052                .collect()
1053        } else {
1054            let Some(mail) = crate::mail_transcript::read(homes, &session) else {
1055                continue;
1056            };
1057            // What supercode typed into the pane (a delivered user turn, filed under its own id) was copied to whoever
1058            // is CC when it was delivered: each such turn accounts for one typed line with its words.
1059            let mut typed_by_supercode: std::collections::HashMap<String, usize> =
1060                std::collections::HashMap::new();
1061            for stored in Mailbox::open(&mail_root(), &session)
1062                .and_then(|mailbox| mailbox.list())
1063                .unwrap_or_default()
1064            {
1065                if stored.envelope.kind == MailKind::User
1066                    && !stored
1067                        .envelope
1068                        .id
1069                        .starts_with(crate::mail_transcript::TYPED_ID_PREFIX)
1070                {
1071                    *typed_by_supercode
1072                        .entry(stored.envelope.body.trim().to_string())
1073                        .or_default() += 1;
1074                }
1075            }
1076            mail.typed
1077                .into_iter()
1078                .filter(|line| !line.withdrawn && line.sent_at_ms > cursor)
1079                .filter(|line| match typed_by_supercode.get_mut(line.text.trim()) {
1080                    Some(left) if *left > 0 => {
1081                        *left -= 1;
1082                        false
1083                    }
1084                    _ => true,
1085                })
1086                .map(|line| (line.sent_at_ms, line.text))
1087                .collect()
1088        };
1089        let Some(newest) = lines.iter().map(|(at, _)| *at).max() else {
1090            continue;
1091        };
1092        let held = threads
1093            .iter()
1094            .rev()
1095            .find(|thread| thread.holder == session && thread.delegated());
1096        let mut receivers: Vec<MailAddress> = held
1097            .map(|thread| {
1098                thread
1099                    .participants
1100                    .iter()
1101                    .filter(|p| p.role == Role::Cc && p.address != session)
1102                    .map(|p| p.address.clone())
1103                    .collect()
1104            })
1105            .unwrap_or_default();
1106        if let Some(account_manager) = &account_manager {
1107            if agent_name != account_manager.name
1108                && !receivers.contains(&account_manager.main_session)
1109            {
1110                receivers.push(account_manager.main_session.clone());
1111            }
1112        }
1113        let Ok(user) = crate::mail_transcript::user_address(&session.machine) else {
1114            continue;
1115        };
1116        for (sent_at_ms, text) in lines {
1117            let copy = Envelope {
1118                id: crate::mail_transcript::typed_line_id(&session, sent_at_ms, &text),
1119                created_at_ms: sent_at_ms,
1120                from: user.clone(),
1121                from_name: format!("user@{} (typed into {session})", session.machine),
1122                sender_identity: None,
1123                kind: MailKind::Typed,
1124                reply_via: crate::mailbox::ReplyVia::None,
1125                in_reply_to: None,
1126                in_reply_to_inferred: false,
1127                thread: held.map(|thread| thread.id.clone()),
1128                native_from: None,
1129                voice_for: None,
1130                subject: None,
1131                body: text,
1132            };
1133            for receiver in &receivers {
1134                // A copy that cannot be filed is said, never dropped in silence.
1135                if let Err(error) = file_unread(receiver, &copy) {
1136                    eprintln!("supercode mail: the line typed into {session} was not copied to {receiver}: {error}");
1137                }
1138            }
1139        }
1140        write_atomically(&cursor_path, newest.to_string().as_bytes()).ok();
1141    }
1142}
1143
1144/// End the process of every thread session the mailbox launched that has been
1145/// idle longer than its agent's `idle_minutes` (decision 9), through `supercode
1146/// close`. Its transcript stays; the next reply in its thread resumes it.
1147pub fn close_idle_threads(homes: &crate::HarnessHomes) {
1148    let launched: Vec<(Thread, u64)> = all_threads()
1149        .into_iter()
1150        .filter(|thread| thread.holder_launched)
1151        .filter_map(|thread| {
1152            let minutes = load(&thread.agent)?.idle_minutes?;
1153            Some((thread, minutes))
1154        })
1155        .collect();
1156    if launched.is_empty() {
1157        return;
1158    }
1159    let live = crate::mail_route::LiveSessions::read(homes);
1160    let now = now_ms();
1161    for (thread, minutes) in launched {
1162        let Some(session) = live
1163            .all()
1164            .iter()
1165            .find(|session| session.address == thread.holder)
1166        else {
1167            continue;
1168        };
1169        if session.status != "idle" {
1170            continue;
1171        }
1172        let Some(last) = session.last_message_at_ms(homes) else {
1173            continue;
1174        };
1175        let resumed = std::fs::read_to_string(resumed_path(&thread.holder))
1176            .ok()
1177            .and_then(|text| text.trim().parse::<u64>().ok())
1178            .unwrap_or(0);
1179        let last = last.max(resumed);
1180        if now.saturating_sub(last) < minutes.saturating_mul(60_000) {
1181            continue;
1182        }
1183        let Ok(program) = crate::claude_relay::supercode_program() else {
1184            return;
1185        };
1186        std::process::Command::new(program)
1187            .args(["close", &thread.holder.to_string()])
1188            .stdin(std::process::Stdio::null())
1189            .stdout(std::process::Stdio::null())
1190            .stderr(std::process::Stdio::null())
1191            .status()
1192            .ok();
1193    }
1194}