Skip to main content

supercode_harness/
mail_route.rs

1//! One route for every message: which door reaches a receiver, and delivery
2//! through it.
3//!
4//! Every caller that delivers mail — `supercode message send`,
5//! `harness.v1.sessions.message`, and idle notices — goes through
6//! [`door_for`] and [`deliver`], so the choice of door is made in one place.
7//! The doors, best first:
8//!
9//! 1. **runtime** — a session supercode controls (a hosted runtime of any
10//!    harness): its own `steer`/`send_input` ([`crate::runtime_mail`]). This
11//!    is the default tier.
12//! 2. **native** — a Claude Code session supercode does not control: a Claude
13//!    relay sends with Claude's own `SendMessage` ([`crate::claude_relay`]).
14//! 3. **hook** — a Codex session supercode does not control, with
15//!    supercode's hooks in its hooks file: filed in its mailbox, and the hook
16//!    points the session at it.
17//! 4. **stored** — no door: filed in its mailbox, seen only when the session
18//!    reads it.
19//!
20//! Tiers 2–4 are the degradation tier, for sessions supercode does not
21//! control. An operator (a board, a script) is not a session: its mail is
22//! filed in its mailbox and read with `sessions.inbox`.
23
24use std::path::Path;
25
26use crate::claude_peer::{read_registry, registry_dir, ClaudePeerSession, ClaudePeerStatus};
27use crate::claude_relay::{send_through_relay, supercode_program, RelayReceipt, RELAY_NAME_PREFIX};
28use crate::live_runtime::LiveRuntimeRecord;
29use crate::mailbox::{
30    local_machine_name, mail_root, Envelope, IdleSubscription, MailAddress, Mailbox, ReplyVia,
31};
32use crate::runtime_mail::{deliver_to_runtime, RuntimeDelivery};
33use crate::HarnessHomes;
34
35/// Arguments a Codex hook entry carries, which is how its presence is found.
36pub const CODEX_HOOK_ARGUMENTS: &str = "message hook codex";
37
38/// The door that reaches one receiver.
39#[derive(Debug, Clone, PartialEq, Eq)]
40pub enum Door {
41    /// A session supercode controls.
42    Runtime(Box<LiveRuntimeRecord>),
43    /// A Claude Code session supercode does not control.
44    Native(Box<ClaudePeerSession>),
45    /// A Codex session supercode does not control, with supercode's hooks:
46    /// they show it its mailbox at its next tool call or when its turn ends.
47    /// An idle one in a daemon pane is woken through its pane.
48    Hook {
49        /// The daemon pane it runs in.
50        pane: Option<String>,
51        /// Whether its last turn has ended.
52        idle: bool,
53    },
54    /// A running session with no door.
55    Stored,
56    /// An operator's mailbox.
57    Operator,
58}
59
60impl Door {
61    /// Stable name of the door, as `message list` and receipts show it.
62    pub fn name(&self) -> &'static str {
63        match self {
64            Self::Runtime(_) => "runtime",
65            Self::Native(_) => "native",
66            Self::Hook { .. } => "hook",
67            Self::Stored => "stored",
68            Self::Operator => "operator",
69        }
70    }
71}
72
73/// Why no door reaches an address on this machine.
74#[derive(Debug, Clone, PartialEq, Eq)]
75pub enum NoDoor {
76    /// The address names another machine.
77    OtherMachine(String),
78    /// No running session has that address.
79    NotRunning,
80}
81
82/// The door that reaches `to` on this machine.
83pub fn door_for(homes: &HarnessHomes, to: &MailAddress) -> Result<Door, NoDoor> {
84    if to.machine != local_machine_name() {
85        return Err(NoDoor::OtherMachine(to.machine.clone()));
86    }
87    if to.harness == "operator" || to.harness == "board" {
88        return Ok(Door::Operator);
89    }
90    LiveSessions::read(homes)
91        .sessions
92        .into_iter()
93        .find(|session| &session.address == to)
94        .map(|session| session.door)
95        .ok_or(NoDoor::NotRunning)
96}
97
98/// One running session on this machine, as the router reaches it.
99#[derive(Debug, Clone, PartialEq, Eq)]
100pub struct LiveSession {
101    /// Its address.
102    pub address: MailAddress,
103    /// The name an agent knows it by (`volter-2c@machine`).
104    pub name: String,
105    /// Its status as its harness reports it (`busy`, `idle`, `hosted`, …).
106    pub status: String,
107    /// The door that reaches it.
108    pub door: Door,
109    /// The process running it, when the harness names one (a hosted
110    /// runtime's host, a Claude session's own process).
111    pub pid: Option<u32>,
112    /// Its working directory.
113    pub cwd: Option<std::path::PathBuf>,
114    /// The tmux session and pane it runs in, when Claude records one.
115    pub tmux: Option<String>,
116    /// Its transcript, when the harness names the file (a Codex rollout).
117    pub transcript: Option<std::path::PathBuf>,
118}
119
120impl LiveSession {
121    /// When its transcript last recorded a message (a prompt, a reply, a tool
122    /// call or result), in epoch milliseconds: what the session last did,
123    /// read from the harness's own records. A queued or injected notice is
124    /// not a message, so this is not the transcript's mtime.
125    pub fn last_message_at_ms(&self, homes: &HarnessHomes) -> Option<u64> {
126        let harness = self.address.harness.as_str();
127        let path = match &self.transcript {
128            Some(path) => path.clone(),
129            None if harness == "claude-code" => {
130                let file = format!("{}.jsonl", self.address.session_id);
131                std::fs::read_dir(&homes.claude_code)
132                    .ok()?
133                    .flatten()
134                    .map(|project| project.path().join(&file))
135                    .find(|path| path.is_file())?
136            }
137            None => return None,
138        };
139        last_message_at_ms(&path, harness)
140    }
141}
142
143impl LiveSession {
144    /// What a session waiting on its user is asking, read from its transcript: the newest tool call with no result
145    /// yet. An `AskUserQuestion` gives its questions, headers and option labels; any other tool its name and a
146    /// shortened input. `None` when the session is not waiting or its transcript shows nothing pending.
147    pub fn pending_request(&self, homes: &HarnessHomes) -> Option<serde_json::Value> {
148        if self.status != "waiting" {
149            return None;
150        }
151        // A Codex prompt writes nothing to its rollout: what it asks is read off its pane.
152        if self.address.harness == "codex" {
153            let prompt = daemon_pane(self).and_then(|pane| pane_prompt(&pane))?;
154            return Some(serde_json::json!({"prompt": prompt, "source": "screen"}));
155        }
156        if self.address.harness != "claude-code" {
157            return None;
158        }
159        let file = format!("{}.jsonl", self.address.session_id);
160        let from_transcript = std::fs::read_dir(&homes.claude_code)
161            .ok()
162            .and_then(|projects| {
163                projects
164                    .flatten()
165                    .map(|project| project.path().join(&file))
166                    .find(|path| path.is_file())
167            })
168            .and_then(|path| pending_request(&path));
169        // Claude Code writes a question's tool call to its transcript only once it is answered
170        // (2.1.280): until then, what it asks is on its pane.
171        from_transcript.or_else(|| {
172            let prompt = daemon_pane(self).and_then(|pane| pane_prompt(&pane))?;
173            Some(serde_json::json!({"prompt": prompt, "source": "screen"}))
174        })
175    }
176}
177
178/// The newest tool call in a Claude transcript that has no result yet, as [`LiveSession::pending_request`] shows it.
179pub fn pending_request(path: &Path) -> Option<serde_json::Value> {
180    use std::io::{Read, Seek, SeekFrom};
181    let mut file = std::fs::File::open(path).ok()?;
182    let length = file.metadata().ok()?.len();
183    let start = length.saturating_sub(512 * 1024);
184    file.seek(SeekFrom::Start(start)).ok()?;
185    let mut bytes = Vec::new();
186    file.read_to_end(&mut bytes).ok()?;
187    let text = String::from_utf8_lossy(&bytes);
188    let mut answered = std::collections::HashSet::new();
189    for line in text.lines().rev() {
190        let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
191            continue;
192        };
193        if record["isSidechain"] == true {
194            continue;
195        }
196        let Some(content) = record
197            .pointer("/message/content")
198            .and_then(serde_json::Value::as_array)
199        else {
200            continue;
201        };
202        match record["type"].as_str() {
203            Some("user") => {
204                for item in content.iter().filter(|item| item["type"] == "tool_result") {
205                    if let Some(id) = item["tool_use_id"].as_str() {
206                        answered.insert(id.to_string());
207                    }
208                }
209            }
210            Some("assistant") => {
211                let Some(call) = content.iter().rev().find(|item| item["type"] == "tool_use")
212                else {
213                    continue;
214                };
215                if call["id"].as_str().is_some_and(|id| answered.contains(id)) {
216                    return None;
217                }
218                let tool = call["name"].as_str().unwrap_or_default();
219                if tool == "AskUserQuestion" {
220                    let questions = call["input"]["questions"]
221                        .as_array()
222                        .map(|questions| {
223                            questions
224                                .iter()
225                                .map(|question| {
226                                    serde_json::json!({
227                                        "question": question["question"],
228                                        "header": question["header"],
229                                        "options": question["options"].as_array().map(|options| options.iter().map(|option| option["label"].clone()).collect::<Vec<_>>()).unwrap_or_default(),
230                                    })
231                                })
232                                .collect::<Vec<_>>()
233                        })
234                        .unwrap_or_default();
235                    return Some(serde_json::json!({"tool": tool, "questions": questions}));
236                }
237                let input = call["input"].to_string();
238                let input: String = input.chars().take(500).collect();
239                return Some(serde_json::json!({"tool": tool, "input": input}));
240            }
241            _ => {}
242        }
243    }
244    None
245}
246
247/// The newest message record's timestamp in a Claude or Codex transcript.
248fn last_message_at_ms(path: &Path, harness: &str) -> Option<u64> {
249    use std::io::{Read, Seek, SeekFrom};
250    // A tool result can be large; widen the window once before giving up.
251    for window in [256 * 1024_u64, 8 * 1024 * 1024] {
252        let mut file = std::fs::File::open(path).ok()?;
253        let length = file.metadata().ok()?.len();
254        let start = length.saturating_sub(window);
255        file.seek(SeekFrom::Start(start)).ok()?;
256        let mut bytes = Vec::new();
257        file.read_to_end(&mut bytes).ok()?;
258        let text = String::from_utf8_lossy(&bytes);
259        let mut lines = text.lines().rev().collect::<Vec<_>>();
260        if start > 0 {
261            lines.pop(); // the first line may begin mid-record
262        }
263        for line in lines {
264            let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
265                continue;
266            };
267            let message = match harness {
268                "codex" => record["type"] == "response_item",
269                _ => {
270                    matches!(record["type"].as_str(), Some("user" | "assistant"))
271                        && record["isSidechain"] != true
272                }
273            };
274            if !message {
275                continue;
276            }
277            if let Some(at) = record["timestamp"]
278                .as_str()
279                .and_then(supercode_interchange::sidecar::rfc3339_to_ms)
280                .and_then(|at| u64::try_from(at).ok())
281            {
282                return Some(at);
283            }
284        }
285        if start == 0 {
286            return None;
287        }
288    }
289    None
290}
291
292/// Every running session on this machine with its door, read once: what
293/// `message list` shows, what names resolve against, and what discovery
294/// projects onto each discovered session as `delivery`.
295#[derive(Debug, Default)]
296pub struct LiveSessions {
297    sessions: Vec<LiveSession>,
298    /// The Codex subagent threads running inside them, each with the conversation that spawned it (Codex's own
299    /// `parent_thread_id`). Not sessions to message: work inside their parent, read only by the session tree.
300    threads: Vec<CodexThread>,
301}
302
303/// A running Codex subagent thread and the conversation Codex recorded as its parent.
304#[derive(Debug, Clone)]
305pub struct CodexThread {
306    /// Its address.
307    pub address: MailAddress,
308    /// The conversation that spawned it, as Codex recorded it.
309    pub parent: MailAddress,
310    /// Its rollout's state (`busy`, `idle`, …).
311    pub status: String,
312    /// Its working directory.
313    pub cwd: Option<std::path::PathBuf>,
314}
315
316/// Why a receiver named by an agent was not found.
317#[derive(Debug, Clone, PartialEq, Eq)]
318pub enum Unresolved {
319    /// The address names a session that is not running.
320    Stale(String),
321    /// No running session has that name (the text suggests near names).
322    Unknown(String),
323}
324
325impl LiveSessions {
326    /// Read every running session now. Sessions supercode controls come
327    /// first, and a degraded-tier row for the same session is dropped: the
328    /// runtime is the default tier.
329    pub fn read(homes: &HarnessHomes) -> Self {
330        crate::slow_log::timed("read live sessions", || Self::read_now(homes))
331    }
332
333    fn read_now(homes: &HarnessHomes) -> Self {
334        let machine = local_machine_name();
335        let registry = read_registry(&registry_dir(homes));
336        // A hosted Claude session is also in Claude's registry, which knows
337        // its name; a relay is not a session and is not listed.
338        let registered = |record: &crate::live_runtime::LiveRuntimeRecord| {
339            registry.iter().find(|session| {
340                session.session_id == record.source.session_id
341                    || session.session_id == record.runtime_session_id
342            })
343        };
344        let records = crate::runtime_mail::controlled_runtimes();
345        // A runtime supercode hosts says whether a turn runs: its label is that, not only its door.
346        let turns = crate::runtime_mail::runtime_turn_states(&records);
347        let mut sessions: Vec<LiveSession> = records
348            .into_iter()
349            .zip(turns)
350            .filter_map(|(record, turn)| {
351                let registered = registered(&record);
352                if registered.is_some_and(|session| session.name.starts_with(RELAY_NAME_PREFIX)) {
353                    return None;
354                }
355                let address =
356                    MailAddress::new(&machine, &record.source.harness, &record.source.session_id)
357                        .ok()?;
358                let short: String = record.source.session_id.chars().take(8).collect();
359                let name = match registered {
360                    Some(session) if !session.name.is_empty() => session.name.clone(),
361                    _ => format!("{}-{short}", record.source.harness),
362                };
363                Some(LiveSession {
364                    name: format!("{name}@{machine}"),
365                    address,
366                    // `hosted` only when its runtime did not answer: no label rather than a false one.
367                    status: match turn {
368                        Some(crate::frontend::FrontendTurnState::Busy) => "busy".into(),
369                        Some(crate::frontend::FrontendTurnState::Idle) => "idle".into(),
370                        None => "hosted".into(),
371                    },
372                    pid: Some(record.pid),
373                    cwd: Some(record.source.workspace.clone()),
374                    tmux: None,
375                    transcript: None,
376                    door: Door::Runtime(Box::new(record)),
377                })
378            })
379            .collect();
380        let controlled = |address: &MailAddress, sessions: &[LiveSession]| {
381            sessions.iter().any(|session| &session.address == address)
382        };
383        for session in registry {
384            if session.name.starts_with(RELAY_NAME_PREFIX) {
385                continue;
386            }
387            let Ok(address) = MailAddress::new(&machine, "claude-code", &session.session_id) else {
388                continue;
389            };
390            if controlled(&address, &sessions) {
391                continue;
392            }
393            sessions.push(LiveSession {
394                address,
395                name: format!("{}@{machine}", session.name),
396                status: session
397                    .status
398                    .as_ref()
399                    .map(|status| status.as_str().to_string())
400                    .unwrap_or_else(|| "unknown".into()),
401                pid: Some(session.pid),
402                cwd: session.cwd.clone(),
403                tmux: session.tmux.clone(),
404                transcript: None,
405                door: Door::Native(Box::new(session)),
406            });
407        }
408        let user_hook = codex_user_hook_installed();
409        let rollouts = crate::slow_log::timed("codex live rollouts", || {
410            crate::codex_peer::live_rollouts(&homes.codex)
411        });
412        let panes = crate::slow_log::timed("codex session panes", crate::codex_peer::session_panes);
413        // A shared native app-server can retain a rollout after its terminal has gone.
414        // Pane absence rules out the pane, not the native writer holding the thread.
415        let live_panes = (!panes.is_empty())
416            .then(|| live_daemon_panes(panes.values()))
417            .flatten();
418        let mut threads = Vec::new();
419        for (path, status) in rollouts {
420            let Some((session_id, parent)) = crate::codex_peer::rollout_session(&path) else {
421                continue;
422            };
423            if let Some(parent) = parent {
424                if let (Ok(address), Ok(parent)) = (
425                    MailAddress::new(&machine, "codex", &session_id),
426                    MailAddress::new(&machine, "codex", &parent),
427                ) {
428                    threads.push(CodexThread {
429                        address,
430                        parent,
431                        status: status.as_str().to_string(),
432                        cwd: crate::codex_peer::rollout_cwd(&path),
433                    });
434                }
435                continue;
436            }
437            let Ok(address) = MailAddress::new(&machine, "codex", &session_id) else {
438                continue;
439            };
440            if controlled(&address, &sessions) {
441                continue;
442            }
443            let cwd = crate::codex_peer::rollout_cwd(&path);
444            let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
445            let pane = panes
446                .get(&session_id)
447                .filter(|pane| {
448                    live_panes
449                        .as_ref()
450                        .is_none_or(|live| live.contains_key(*pane))
451                })
452                .cloned();
453            // A Codex turn waiting on an approval or a choice writes nothing to its rollout, which
454            // reads that turn as working; in a daemon pane its screen shows the prompt.
455            let status = match status {
456                crate::codex_peer::CodexPeerStatus::Busy
457                | crate::codex_peer::CodexPeerStatus::Running
458                    if pane.as_ref().is_some_and(|p| {
459                        live_panes.as_ref().and_then(|live| live.get(p)) == Some(&true)
460                    }) =>
461                {
462                    "idle".to_string()
463                }
464
465                crate::codex_peer::CodexPeerStatus::Busy
466                | crate::codex_peer::CodexPeerStatus::Running
467                    if pane.as_deref().and_then(pane_prompt).is_some() =>
468                {
469                    "waiting".to_string()
470                }
471                status => status.as_str().to_string(),
472            };
473            let door = if hooked {
474                Door::Hook {
475                    pane: pane.clone(),
476                    idle: status == "idle",
477                }
478            } else {
479                Door::Stored
480            };
481            sessions.push(LiveSession {
482                name: format!("{}@{machine}", codex_name(&session_id)),
483                address,
484                status,
485                pid: None,
486                cwd,
487                // The daemon pane it runs in, as a Claude session's tmux names its own.
488                tmux: pane,
489                transcript: Some(path),
490                door,
491            });
492        }
493        // A conversation resumed in a daemon pane and idle since holds no rollout open; its pane and the
494        // process's own `resume <id>` still name it, so it is reachable as an idle session there.
495        for (session_id, pane) in &panes {
496            let Ok(address) = MailAddress::new(&machine, "codex", session_id) else {
497                continue;
498            };
499            if controlled(&address, &sessions)
500                || !live_panes
501                    .as_ref()
502                    .is_some_and(|live| live.contains_key(pane))
503            {
504                continue;
505            }
506            let Some(path) = crate::codex_peer::rollout_of_session(&homes.codex, session_id) else {
507                continue;
508            };
509            let cwd = crate::codex_peer::rollout_cwd(&path);
510            let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
511            sessions.push(LiveSession {
512                name: format!("{}@{machine}", codex_name(session_id)),
513                address,
514                status: "idle".to_string(),
515                pid: None,
516                cwd,
517                tmux: Some(pane.clone()),
518                transcript: Some(path),
519                door: if hooked {
520                    Door::Hook {
521                        pane: Some(pane.clone()),
522                        idle: true,
523                    }
524                } else {
525                    Door::Stored
526                },
527            });
528        }
529        remember_names(&sessions);
530        Self { sessions, threads }
531    }
532
533    /// Every running session.
534    pub fn all(&self) -> &[LiveSession] {
535        &self.sessions
536    }
537
538    /// The running Codex subagent threads, each with its parent conversation.
539    pub fn codex_threads(&self) -> &[CodexThread] {
540        &self.threads
541    }
542
543    /// The door that reaches `harness`'s session `session_id`, when it is
544    /// running.
545    pub fn door(&self, harness: &str, session_id: &str) -> Option<&'static str> {
546        self.sessions
547            .iter()
548            .find(|session| {
549                session.address.harness == harness && session.address.session_id == session_id
550            })
551            .map(|session| session.door.name())
552    }
553
554    /// The running session an agent means by `to`: an address, `name@machine`
555    /// on this machine, or a name unique here.
556    pub fn resolve(&self, to: &str) -> Result<&LiveSession, Unresolved> {
557        let machine = local_machine_name();
558        if let Ok(address) = MailAddress::parse(to) {
559            return self
560                .sessions
561                .iter()
562                .find(|session| session.address == address)
563                .ok_or_else(|| {
564                    Unresolved::Stale(format!(
565                        "{to} is no longer running. Nothing was sent. Run supercode message list \
566                         for the live sessions."
567                    ))
568                });
569        }
570        let wanted = match to.split_once('@') {
571            Some((name, at)) if at == machine => name.to_string(),
572            Some((_, at)) => {
573                return Err(Unresolved::Unknown(format!(
574                    "Not sent: {to} is on machine {at}, not this one. Nothing was sent."
575                )))
576            }
577            None => to.to_string(),
578        };
579        let short = |session: &LiveSession| {
580            session
581                .name
582                .split('@')
583                .next()
584                .unwrap_or_default()
585                .to_string()
586        };
587        let matching: Vec<&LiveSession> = self
588            .sessions
589            .iter()
590            .filter(|session| short(session) == wanted)
591            .collect();
592        if let [only] = matching.as_slice() {
593            return Ok(only);
594        }
595        // A name no running session carries now still means the session that last carried it: a session resumed
596        // under another name (a restart, a rename) keeps every name it was reached by.
597        if matching.is_empty() {
598            if let Some(address) = remembered_address(&wanted) {
599                return self
600                    .sessions
601                    .iter()
602                    .find(|session| session.address == address)
603                    .ok_or_else(|| {
604                        Unresolved::Stale(format!(
605                            "{wanted} was {address}, which is not running. Nothing was sent. Run \
606                             supercode message list for the live sessions."
607                        ))
608                    });
609            }
610        }
611        let hint = if matching.len() > 1 {
612            format!(
613                " {} sessions are named {wanted}; use its address.",
614                matching.len()
615            )
616        } else {
617            let near: Vec<String> = self
618                .sessions
619                .iter()
620                .filter(|session| {
621                    let name = short(session);
622                    name.contains(&wanted)
623                        || wanted.contains(&name)
624                        || name
625                            .chars()
626                            .zip(wanted.chars())
627                            .take_while(|(a, b)| a == b)
628                            .count()
629                            >= 4
630                })
631                .take(3)
632                .map(|session| {
633                    format!(
634                        "{} ({}, {})",
635                        session.name, session.address.harness, session.status
636                    )
637                })
638                .collect();
639            if near.is_empty() {
640                String::new()
641            } else {
642                format!(" Did you mean: {}?", near.join(", "))
643            }
644        };
645        Err(Unresolved::Unknown(format!(
646            "No session named \"{wanted}\" is reachable.{hint} Run supercode message list. Nothing \
647             was sent."
648        )))
649    }
650}
651
652/// Every name a session was listed by on this machine, each with the session it named, when it began to and when it
653/// was last listed by it: `<mail root>/names.json`. Names are not addresses: Claude derives a new one for a session
654/// each time it starts unless one is given, and a title is not a name. The ledger lets a name an agent was given
655/// reach the same session after it runs under another, within [`NAME_KEPT_MS`] of its last listing, and lets a
656/// resume give a session back the name it had.
657#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
658struct NameRecord {
659    address: String,
660    since_ms: u64,
661    /// Last listed by it (renewed at most every [`NAME_SEEN_EVERY_MS`]); 0 in a record from before this was kept.
662    #[serde(default)]
663    seen_ms: u64,
664}
665
666impl NameRecord {
667    fn last_seen_ms(&self) -> u64 {
668        self.seen_ms.max(self.since_ms)
669    }
670}
671
672/// How long a name no running session carries still reaches the session that last carried it.
673pub const NAME_KEPT_MS: u64 = 7 * 24 * 60 * 60 * 1000;
674/// How often a name's last listing is written again while its session runs.
675const NAME_SEEN_EVERY_MS: u64 = 60 * 60 * 1000;
676
677fn names_path() -> std::path::PathBuf {
678    crate::mailbox::mail_root().join("names.json")
679}
680
681/// The ledger, empty when there is none yet; `None` when it exists and cannot be read, so it is never rewritten
682/// from nothing. With `set_aside` (a writer, under the ledger's lock), a ledger that reads but does not parse is
683/// moved aside to `names.json.unreadable-<ms>`, kept for whoever wants its history, and logged; the ledger starts
684/// again empty rather than going dark for good.
685fn read_names(set_aside: bool) -> Option<std::collections::BTreeMap<String, NameRecord>> {
686    let path = names_path();
687    match std::fs::read(&path) {
688        Ok(bytes) => match serde_json::from_slice(&bytes) {
689            Ok(names) => Some(names),
690            Err(error) if set_aside => {
691                let aside = path.with_extension(format!("json.unreadable-{}", now_ms()));
692                let moved = std::fs::rename(&path, &aside).is_ok();
693                tracing::warn!(
694                    "the names ledger {} does not parse ({error}); {}",
695                    path.display(),
696                    if moved {
697                        format!("moved aside to {} and started again", aside.display())
698                    } else {
699                        "it could not be moved aside and is left as it is".to_string()
700                    }
701                );
702                moved.then(Default::default)
703            }
704            Err(_) => None,
705        },
706        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Some(Default::default()),
707        Err(_) => None,
708    }
709}
710
711fn now_ms() -> u64 {
712    std::time::SystemTime::now()
713        .duration_since(std::time::UNIX_EPOCH)
714        .map(|elapsed| elapsed.as_millis() as u64)
715        .unwrap_or_default()
716}
717
718/// Record each running session's name. A name two running sessions carry is left as it was (it names neither).
719/// Writes take the ledger's lock and re-read it under the lock, so concurrent listings merge rather than overwrite,
720/// and only when a name is new, names another session, or its last listing is an hour old.
721fn remember_names(sessions: &[LiveSession]) {
722    let mut carried: std::collections::BTreeMap<&str, Vec<String>> = Default::default();
723    for session in sessions {
724        let name = session.name.split('@').next().unwrap_or_default();
725        if !name.is_empty() && !name.starts_with(RELAY_NAME_PREFIX) {
726            carried
727                .entry(name)
728                .or_default()
729                .push(session.address.to_string());
730        }
731    }
732    let now = now_ms();
733    let wanted = |names: &std::collections::BTreeMap<String, NameRecord>, name: &str, address: &str| {
734        let current = names.get(name).is_some_and(|record| {
735            record.address == address && now.saturating_sub(record.last_seen_ms()) < NAME_SEEN_EVERY_MS
736        });
737        !current || latest_name_in(names, address).as_deref() != Some(name)
738    };
739    let unique = || {
740        carried
741            .iter()
742            .filter_map(|(name, addresses)| match addresses.as_slice() {
743                [address] => Some((*name, address.as_str())),
744                _ => None,
745            })
746    };
747    // Most listings change nothing: decided without the lock (a ledger that does not parse is set aside under it).
748    let names = read_names(false).unwrap_or_default();
749    if !unique().any(|(name, address)| wanted(&names, name, address)) {
750        return;
751    }
752    let path = names_path();
753    if let Some(parent) = path.parent() {
754        std::fs::create_dir_all(parent).ok();
755    }
756    let Ok(_lock) = NamesLock::acquire(&path.with_extension("json.lock")) else {
757        return;
758    };
759    let Some(mut names) = read_names(true) else { return };
760    let mut changed = false;
761    for (name, address) in unique() {
762        if !wanted(&names, name, address) {
763            continue;
764        }
765        let renewed = names
766            .get(name)
767            .is_some_and(|record| record.address == address)
768            && latest_name_in(&names, address).as_deref() == Some(name);
769        let since_ms = if renewed { names[name].since_ms } else { now };
770        names.insert(
771            name.to_string(),
772            NameRecord {
773                address: address.to_string(),
774                since_ms,
775                seen_ms: now,
776            },
777        );
778        changed = true;
779    }
780    if changed {
781        let staged = path.with_extension(format!("json.{}.tmp", std::process::id()));
782        if std::fs::write(&staged, serde_json::to_vec(&names).unwrap_or_default()).is_ok() {
783            if std::fs::rename(&staged, &path).is_err() {
784                std::fs::remove_file(&staged).ok();
785            }
786        }
787    }
788}
789
790/// The ledger's write lock: held while it is re-read, changed and replaced; released when dropped.
791struct NamesLock {
792    _file: std::fs::File,
793}
794
795impl NamesLock {
796    fn acquire(path: &Path) -> std::io::Result<Self> {
797        let file = std::fs::OpenOptions::new()
798            .create(true)
799            .truncate(false)
800            .write(true)
801            .open(path)?;
802        #[cfg(unix)]
803        {
804            use std::os::unix::io::AsRawFd;
805            // SAFETY: flock on a descriptor this function owns; the lock is released when the file is dropped.
806            if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) } != 0 {
807                return Err(std::io::Error::last_os_error());
808            }
809        }
810        Ok(Self { _file: file })
811    }
812}
813
814fn latest_name_in(
815    names: &std::collections::BTreeMap<String, NameRecord>,
816    address: &str,
817) -> Option<String> {
818    names
819        .iter()
820        .filter(|(_, record)| record.address == address)
821        .max_by_key(|(_, record)| record.since_ms)
822        .map(|(name, _)| name.clone())
823}
824
825/// The session a name last meant on this machine, from the names ledger, when it was listed by that name within
826/// [`NAME_KEPT_MS`].
827pub fn remembered_address(name: &str) -> Option<MailAddress> {
828    let names = read_names(false)?;
829    let record = names.get(name)?;
830    if now_ms().saturating_sub(record.last_seen_ms()) > NAME_KEPT_MS {
831        return None;
832    }
833    MailAddress::parse(&record.address).ok()
834}
835
836/// The name `address` was last listed by on this machine: the name a resume gives it back.
837pub fn remembered_name(address: &MailAddress) -> Option<String> {
838    latest_name_in(&read_names(false)?, &address.to_string())
839}
840
841/// Whether the session process `pid` has supercode's messaging tools: a
842/// harness starts each MCP server as a child of the session, so the tools are
843/// loaded exactly when a live `supercode message mcp` is one of its children.
844pub fn has_message_tools(pid: u32) -> bool {
845    let Ok(output) = std::process::Command::new("ps")
846        .args(["-A", "-o", "ppid=,command="])
847        .output()
848    else {
849        return false;
850    };
851    String::from_utf8_lossy(&output.stdout).lines().any(|line| {
852        let line = line.trim_start();
853        let Some((ppid, command)) = line.split_once(' ') else {
854            return false;
855        };
856        ppid.parse::<u32>() == Ok(pid) && command.trim_end().ends_with(" message mcp")
857    })
858}
859
860/// Display name of a Codex conversation: `codex-` and the start of its id.
861pub fn codex_name(session_id: &str) -> String {
862    format!("codex-{}", session_id.chars().take(8).collect::<String>())
863}
864
865/// The session a process belongs to: the sender a message is from.
866#[derive(Debug, Clone, PartialEq, Eq)]
867pub struct Caller {
868    /// Its address, where replies go.
869    pub address: MailAddress,
870    /// The name it is known by.
871    pub name: String,
872}
873
874/// Why no session could be found behind a process.
875pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
876    replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
877    shell tool.";
878
879/// The session behind a process, from its ancestry `pids` (nearest first):
880/// the first that owns a hosted runtime, a Claude session or a Codex
881/// conversation. Never declared by the caller; an environment id only
882/// corroborates, and a mismatch refuses.
883pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
884    crate::slow_log::timed("resolve the caller", || resolve_caller_now(homes, pids))
885}
886
887fn resolve_caller_now(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
888    let machine = local_machine_name();
889    let registry = read_registry(&registry_dir(homes));
890    let hosted = crate::runtime_mail::controlled_runtimes();
891    for &pid in pids {
892        if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
893            let short: String = record.source.session_id.chars().take(8).collect();
894            return Ok(Caller {
895                address: MailAddress::new(
896                    &machine,
897                    &record.source.harness,
898                    &record.source.session_id,
899                )
900                .map_err(|error| error.to_string())?,
901                name: format!("{}-{short}@{machine}", record.source.harness),
902            });
903        }
904        if let Some(session) = registry.iter().find(|session| session.pid == pid) {
905            if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
906                if !claimed.is_empty() && claimed != session.session_id {
907                    return Err(format!(
908                        "{CALLER_UNRESOLVED} (CLAUDE_CODE_SESSION_ID names {claimed}, but the Claude \
909                         process {pid} above this command is session {})",
910                        session.session_id
911                    ));
912                }
913            }
914            return Ok(Caller {
915                address: MailAddress::new(&machine, "claude-code", &session.session_id)
916                    .map_err(|error| error.to_string())?,
917                name: format!("{}@{machine}", session.name),
918            });
919        }
920        // Codex names the thread a command runs for (`CODEX_THREAD_ID`); the
921        // claim stands when this process holds that thread's rollout, as
922        // Codex's shared app-server daemon does for every thread it runs.
923        if let Some(thread) = std::env::var("CODEX_THREAD_ID")
924            .ok()
925            .filter(|thread| !thread.is_empty())
926            .filter(|thread| crate::codex_peer::holds_session(pid, thread))
927        {
928            return Ok(Caller {
929                address: MailAddress::new(&machine, "codex", &thread)
930                    .map_err(|error| error.to_string())?,
931                name: format!("{}@{machine}", codex_name(&thread)),
932            });
933        }
934        if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
935            return Ok(Caller {
936                address: MailAddress::new(&machine, "codex", &session_id)
937                    .map_err(|error| error.to_string())?,
938                name: format!("{}@{machine}", codex_name(&session_id)),
939            });
940        }
941    }
942    #[cfg(windows)]
943    if let Some(session) = msys_cut_claim(&registry, pids) {
944        return Ok(Caller {
945            address: MailAddress::new(&machine, "claude-code", &session.session_id)
946                .map_err(|error| error.to_string())?,
947            name: format!("{}@{machine}", session.name),
948        });
949    }
950    Err(format!(
951        "{CALLER_UNRESOLVED} (looked for this command's processes {pids:?} among {} Claude \
952         sessions in {} and {} hosted runtimes)",
953        registry.len(),
954        registry_dir(homes).display(),
955        hosted.len()
956    ))
957}
958
959/// Git Bash (MSYS) runs a command in a forked process that hands over to it. When the command is
960/// itself an MSYS program (a `sh` script, as npm's command shims are), the forked process exits once
961/// it has handed over, so the Windows parent chain is cut just above that shell and never reaches
962/// the Claude session running it. When the ancestry ends at exactly such a cut (its last pid is gone,
963/// the one before it is an MSYS shell), the session Claude names in `CLAUDE_CODE_SESSION_ID` stands
964/// if it is a live Claude session on this machine.
965#[cfg(windows)]
966fn msys_cut_claim<'a>(
967    registry: &'a [ClaudePeerSession],
968    pids: &[u32],
969) -> Option<&'a ClaudePeerSession> {
970    let table = process_table();
971    let [.., shell, cut] = pids else {
972        return None;
973    };
974    if table.contains_key(cut) {
975        return None;
976    }
977    let name = table.get(shell)?.1.to_ascii_lowercase();
978    if !matches!(name.as_str(), "sh.exe" | "bash.exe" | "dash.exe") {
979        return None;
980    }
981    let claimed = std::env::var("CLAUDE_CODE_SESSION_ID").ok()?;
982    registry
983        .iter()
984        .find(|session| !claimed.is_empty() && session.session_id == claimed)
985}
986
987/// This process's ancestry, nearest first (itself included).
988pub fn process_ancestry() -> Vec<u32> {
989    crate::slow_log::timed("read the process ancestry", || {
990        ancestry_of(std::process::id())
991    })
992}
993
994/// A process's ancestry, nearest first (itself included).
995pub fn ancestry_of(pid: u32) -> Vec<u32> {
996    ancestry_in(pid, &parent_pids())
997}
998
999/// Every process's parent now (pid → parent pid): one read for a caller that walks several
1000/// processes' ancestries ([`ancestry_in`]) instead of reading the table once per process.
1001pub fn parent_table() -> std::collections::HashMap<u32, u32> {
1002    parent_pids()
1003}
1004
1005/// `pid`'s ancestry, nearest first (itself included), in a table [`parent_table`] read.
1006pub fn ancestry_in(pid: u32, parents: &std::collections::HashMap<u32, u32>) -> Vec<u32> {
1007    let mut chain = vec![pid];
1008    let mut current = pid;
1009    while let Some(&parent) = parents.get(&current) {
1010        if parent <= 1 || chain.contains(&parent) {
1011            break;
1012        }
1013        chain.push(parent);
1014        current = parent;
1015    }
1016    chain
1017}
1018
1019/// Every process's parent: pid → parent pid, from `ps`.
1020#[cfg(not(windows))]
1021fn parent_pids() -> std::collections::HashMap<u32, u32> {
1022    crate::slow_log::timed("ps parent table", parent_pids_now)
1023}
1024
1025#[cfg(not(windows))]
1026fn parent_pids_now() -> std::collections::HashMap<u32, u32> {
1027    let Ok(output) = std::process::Command::new("ps")
1028        .args(["-axo", "pid=,ppid="])
1029        .output()
1030    else {
1031        return Default::default();
1032    };
1033    String::from_utf8_lossy(&output.stdout)
1034        .lines()
1035        .filter_map(|line| {
1036            let mut fields = line.split_whitespace();
1037            Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
1038        })
1039        .collect()
1040}
1041
1042/// Every process's parent: pid → parent pid, from a process snapshot.
1043#[cfg(windows)]
1044fn parent_pids() -> std::collections::HashMap<u32, u32> {
1045    process_table()
1046        .into_iter()
1047        .map(|(pid, (parent, _))| (pid, parent))
1048        .collect()
1049}
1050
1051/// Every process: pid → (parent pid, executable file name), from a process snapshot.
1052#[cfg(windows)]
1053fn process_table() -> std::collections::HashMap<u32, (u32, String)> {
1054    use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
1055    use windows_sys::Win32::System::Diagnostics::ToolHelp::{
1056        CreateToolhelp32Snapshot, Process32FirstW, Process32NextW, PROCESSENTRY32W,
1057        TH32CS_SNAPPROCESS,
1058    };
1059
1060    let mut parents = std::collections::HashMap::new();
1061    let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPPROCESS, 0) };
1062    if snapshot == INVALID_HANDLE_VALUE {
1063        return parents;
1064    }
1065    let mut entry: PROCESSENTRY32W = unsafe { std::mem::zeroed() };
1066    entry.dwSize = std::mem::size_of::<PROCESSENTRY32W>() as u32;
1067    let mut has_entry = unsafe { Process32FirstW(snapshot, &mut entry) } != 0;
1068    while has_entry {
1069        let length = entry
1070            .szExeFile
1071            .iter()
1072            .position(|&unit| unit == 0)
1073            .unwrap_or(entry.szExeFile.len());
1074        let name = String::from_utf16_lossy(&entry.szExeFile[..length]);
1075        parents.insert(entry.th32ProcessID, (entry.th32ParentProcessID, name));
1076        has_entry = unsafe { Process32NextW(snapshot, &mut entry) } != 0;
1077    }
1078    unsafe {
1079        CloseHandle(snapshot);
1080    }
1081    parents
1082}
1083
1084/// Codex's user-level hooks file.
1085pub fn codex_hooks_path() -> std::path::PathBuf {
1086    std::env::var_os("CODEX_HOME")
1087        .map(std::path::PathBuf::from)
1088        .or_else(|| {
1089            supercode_interchange::user_home()
1090                .map(std::path::PathBuf::into_os_string)
1091                .map(|home| std::path::PathBuf::from(home).join(".codex"))
1092        })
1093        .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
1094        .join("hooks.json")
1095}
1096
1097/// Whether supercode's mail hook is in Codex's user hooks file. (Codex runs
1098/// it only once its user has trusted it.)
1099pub fn codex_user_hook_installed() -> bool {
1100    std::fs::read_to_string(codex_hooks_path())
1101        .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
1102}
1103
1104/// Whether supercode's mail hook is in a project's Codex hooks file, in the
1105/// session's directory or one above it.
1106pub fn codex_project_hook_installed(cwd: &Path) -> bool {
1107    cwd.ancestors().any(|directory| {
1108        std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
1109            .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
1110    })
1111}
1112
1113/// How a delivery ended.
1114#[derive(Debug, Clone, PartialEq, Eq)]
1115pub enum Delivered {
1116    /// The receiver has it now: steered into a running turn.
1117    Steered,
1118    /// The receiver has it now: it started a turn in an idle session.
1119    Started,
1120    /// Claude reported it in the receiver's inbox; `busy` says whether the
1121    /// receiver reads it at its next tool call (true) or it starts a turn.
1122    Native {
1123        /// True when the receiver was in a turn.
1124        busy: bool,
1125    },
1126    /// Filed in the receiver's mailbox, which a hook points it at.
1127    Hooked,
1128    /// Filed in an idle receiver's mailbox, and its pane told to read it: a
1129    /// turn has started, in which its hook shows it the message.
1130    HookWoken,
1131    /// Filed in the receiver's mailbox, waiting to be read.
1132    Queued,
1133    /// Filed with no door to show it.
1134    Stored,
1135    /// Filed in an operator's mailbox.
1136    Operator,
1137    /// Its receiver already had it (the message is read in its mailbox): nothing was handed over again.
1138    Already,
1139}
1140
1141/// A delivery refused before anything was sent.
1142#[derive(Debug, Clone, PartialEq, Eq)]
1143pub enum Refused {
1144    /// `--queue` cannot hold a message for an idle Claude session: Claude
1145    /// starts a turn for every message it delivers.
1146    CannotQueueNative,
1147    /// The relayed text would not fit one relay turn.
1148    TooLong(usize),
1149}
1150
1151/// Largest message relayed into a Claude session. The relay copies it into a
1152/// model turn byte for byte, so it must fit comfortably in one.
1153pub const MAX_RELAYED_BYTES: usize = 100_000;
1154
1155/// Deliver `envelope` to `to` through `door`. With `wake` false an idle
1156/// receiver is not started. With `notify_when_idle` the sender gets one idle
1157/// notice after the receiver's next turn ends.
1158pub async fn deliver(
1159    envelope: &Envelope,
1160    to: &MailAddress,
1161    door: &Door,
1162    wake: bool,
1163    notify_when_idle: bool,
1164) -> Result<Result<Delivered, Refused>, String> {
1165    let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
1166    // A runtime answers with its turn's final message, which the watcher
1167    // sends back: not every runtime can run a command (a hosted agent may
1168    // have no shell), and each can end a turn.
1169    let mut final_reply = false;
1170    // A message its receiver already has (its own claim, or an earlier hand-over to its door) is never handed over
1171    // again, whichever attempt asks (idempotent by its id). A message read any other way is still handed over.
1172    if matches!(door, Door::Native(_) | Door::Runtime(_))
1173        && mailbox
1174            .find(&envelope.id)
1175            .map_err(|error| error.to_string())?
1176            .is_some_and(|stored| stored.state == crate::mailbox::MailState::Read)
1177        && mailbox.delivered_to_recipient(&envelope.id)
1178    {
1179        return Ok(Ok(Delivered::Already));
1180    }
1181    let delivered = match door {
1182        Door::Runtime(record) => {
1183            let mut sent = envelope.clone();
1184            let answers = matches!(envelope.reply_via, ReplyVia::Command);
1185            if answers {
1186                sent.reply_via = ReplyVia::FinalMessage {
1187                    destination: envelope.from_name.clone(),
1188                };
1189            }
1190            match deliver_to_runtime(record, sent.render(), wake).await? {
1191                RuntimeDelivery::Steered => {
1192                    mailbox
1193                        .deliver_read(&sent)
1194                        .map_err(|error| error.to_string())?;
1195                    final_reply = answers;
1196                    Delivered::Steered
1197                }
1198                RuntimeDelivery::Started => {
1199                    mailbox
1200                        .deliver_read(&sent)
1201                        .map_err(|error| error.to_string())?;
1202                    final_reply = answers;
1203                    Delivered::Started
1204                }
1205                RuntimeDelivery::NotWoken => {
1206                    mailbox
1207                        .deliver(envelope)
1208                        .map_err(|error| error.to_string())?;
1209                    Delivered::Queued
1210                }
1211            }
1212        }
1213        Door::Native(session) => {
1214            let busy = session.status != Some(ClaudePeerStatus::Idle);
1215            if !wake && !busy {
1216                return Ok(Err(Refused::CannotQueueNative));
1217            }
1218            // The same envelope every door delivers; Claude wraps it in its
1219            // own, which names the relay. Its reply line names the door this
1220            // session really has: its tools when it runs supercode's MCP server.
1221            let mut envelope = envelope.clone();
1222            if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
1223                envelope.reply_via = ReplyVia::Tool;
1224            }
1225            let envelope = &envelope;
1226            let text = envelope.render();
1227            if text.len() > MAX_RELAYED_BYTES {
1228                return Ok(Err(Refused::TooLong(text.len())));
1229            }
1230            // Filed before it is handed over: a session that reads its mailbox on its first tool call (to reply by
1231            // the message's id) finds it there. A relay that fails takes back the copy this delivery filed.
1232            let filed_here = mailbox
1233                .find(&envelope.id)
1234                .map_err(|error| error.to_string())?
1235                .is_none();
1236            let filed = mailbox
1237                .deliver_read(envelope)
1238                .map_err(|error| error.to_string())?;
1239            match send_through_relay(
1240                &envelope.from,
1241                &envelope.from_name,
1242                session,
1243                text,
1244                &envelope.id,
1245            )
1246            .await
1247            {
1248                RelayReceipt::Delivered { .. } => Delivered::Native { busy },
1249                RelayReceipt::Failed { detail } => {
1250                    if filed_here {
1251                        std::fs::remove_file(&filed).ok();
1252                    }
1253                    return Err(detail);
1254                }
1255            }
1256        }
1257        Door::Hook { .. } => {
1258            mailbox
1259                .deliver(envelope)
1260                .map_err(|error| error.to_string())?;
1261            if wake {
1262                mailbox
1263                    .request_wake(&envelope.id)
1264                    .map_err(|error| error.to_string())?;
1265            }
1266            Delivered::Hooked
1267        }
1268        Door::Stored => {
1269            mailbox
1270                .deliver(envelope)
1271                .map_err(|error| error.to_string())?;
1272            Delivered::Stored
1273        }
1274        Door::Operator => {
1275            mailbox
1276                .deliver(envelope)
1277                .map_err(|error| error.to_string())?;
1278            Delivered::Operator
1279        }
1280    };
1281    // a hand-over to the recipient's own door (its relay, its runtime) is its delivery, recorded as such
1282    if matches!(
1283        delivered,
1284        Delivered::Steered | Delivered::Started | Delivered::Native { .. }
1285    ) {
1286        mailbox
1287            .record_claim(
1288                &envelope.id,
1289                &crate::mailbox::Claim::new("door", None, true),
1290            )
1291            .ok();
1292    }
1293    let notice = notify_when_idle && !matches!(door, Door::Operator);
1294    if notice || final_reply {
1295        let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
1296        subscription.notice = notice;
1297        subscription.final_reply = final_reply;
1298        mailbox
1299            .subscribe_idle(&subscription)
1300            .map_err(|error| error.to_string())?;
1301        // The watcher that settles it runs beside the machine daemon.
1302        if let Ok(program) = supercode_program() {
1303            crate::claude_relay::ensure_machine_daemon(&program)
1304                .await
1305                .ok();
1306        }
1307    }
1308    Ok(Ok(delivered))
1309}
1310
1311/// How the user's own turn reached its session.
1312#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1313pub enum UserTurn {
1314    /// A hosted runtime was in a turn; the words were steered into it.
1315    Steered,
1316    /// A hosted runtime was idle; the words started a turn.
1317    Started,
1318    /// Typed into the session's pane as its own submitted turn.
1319    Typed,
1320    /// The pane's composer holds a draft, or its program is not at its
1321    /// composer: the turn waits in the session's mailbox and is typed once
1322    /// the composer is empty.
1323    Waiting,
1324}
1325
1326impl UserTurn {
1327    /// Stable wire spelling.
1328    pub const fn as_str(self) -> &'static str {
1329        match self {
1330            Self::Steered => "steered",
1331            Self::Started => "started",
1332            Self::Typed => "typed",
1333            Self::Waiting => "waiting",
1334        }
1335    }
1336}
1337
1338/// The choice or permission prompt a daemon pane's screen shows, as its lines (the question above
1339/// it and its options), or `None`. The terminal substrate's own check (`classifyAttentionSignal`):
1340/// a selector (❯ or ›) on a numbered option is a prompt blocked on an answer; the idle input line
1341/// has no number after its selector.
1342pub fn pane_prompt(pane: &str) -> Option<Vec<String>> {
1343    let screen = pane_screen(pane)?;
1344    let lines: Vec<&str> = screen.lines().map(str::trim_end).collect();
1345    let is_option = |line: &str| {
1346        let line = line.trim_start();
1347        // Windows consoles draw Claude's selector as `>`.
1348        let rest = line
1349            .strip_prefix('❯')
1350            .or_else(|| line.strip_prefix('›'))
1351            .or_else(|| line.strip_prefix('>'))
1352            .unwrap_or(line)
1353            .trim_start();
1354        let digits = rest.chars().take_while(char::is_ascii_digit).count();
1355        digits > 0 && matches!(rest[digits..].chars().next(), Some('.' | ')'))
1356    };
1357    let selected = lines.iter().position(|line| {
1358        let line = line.trim_start();
1359        (line.starts_with('❯') || line.starts_with('›') || line.starts_with('>')) && is_option(line)
1360    })?;
1361    // The prompt: the options' block of non-empty lines and the block above it (the question).
1362    let block_start = |end: usize| {
1363        lines[..end]
1364            .iter()
1365            .rposition(|line| line.trim().is_empty())
1366            .map_or(0, |blank| blank + 1)
1367    };
1368    let options = block_start(selected);
1369    let above = lines[..options]
1370        .iter()
1371        .rposition(|line| !line.trim().is_empty())
1372        .map(|last| block_start(last));
1373    let start = above.unwrap_or(options);
1374    let end = lines[selected..]
1375        .iter()
1376        .position(|line| line.trim().is_empty())
1377        .map_or(lines.len(), |blank| selected + blank);
1378    Some(
1379        lines[start..end]
1380            .iter()
1381            .map(|line| line.trim().to_string())
1382            .filter(|line| !line.is_empty())
1383            .take(20)
1384            .collect(),
1385    )
1386}
1387
1388/// Read this machine's actual panes once for the session list. Unavailability is unknown.
1389/// Only `wanted` are read: the daemon captures each pane it is asked about, so asking about every
1390/// pane cost a screen capture of all of them on every live-session read (t_23e47b73).
1391fn live_daemon_panes<'a>(
1392    wanted: impl Iterator<Item = &'a String>,
1393) -> Option<std::collections::HashMap<String, bool>> {
1394    let wanted: Vec<&str> = wanted.map(String::as_str).collect();
1395    crate::slow_log::timed("teams panes ls --presence-only", || {
1396        live_daemon_panes_now(&wanted.join(","))
1397    })
1398}
1399
1400fn live_daemon_panes_now(wanted: &str) -> Option<std::collections::HashMap<String, bool>> {
1401    let entry = crate::teams_entry().ok()?;
1402    let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1403        .ok()
1404        .filter(|value| !value.trim().is_empty())
1405        .unwrap_or_else(|| "node".into());
1406    let output = std::process::Command::new(node)
1407        .arg(entry)
1408        .args(["panes", "ls", "--presence-only", "--panes", wanted])
1409        .stdin(std::process::Stdio::null())
1410        .output()
1411        .ok()?;
1412    if !output.status.success() {
1413        return None;
1414    }
1415    let rows: Vec<serde_json::Value> = serde_json::from_slice(&output.stdout).ok()?;
1416    rows.iter()
1417        .map(|row| {
1418            Some((
1419                row["pane"].as_str()?.to_owned(),
1420                row["idleComposer"].as_bool()?,
1421            ))
1422        })
1423        .collect()
1424}
1425
1426/// Each live daemon pane's root process (pane → pid), the daemon's own record (`teams panes ls
1427/// --processes`). A daemon older than `panes.processes` refuses it; its full listing carries the
1428/// same pane and pid. No daemon, no panes.
1429pub(crate) fn daemon_pane_roots() -> std::collections::HashMap<String, u32> {
1430    daemon_panes()
1431        .into_iter()
1432        .map(|(pane, (pid, _))| (pane, pid))
1433        .collect()
1434}
1435
1436/// Each live daemon pane's root process and the harness the daemon launched in it (pane → (pid,
1437/// kind)), from the same listing as [`daemon_pane_roots`]. The kind is the daemon's launch record
1438/// (`codex`, `claude-code`, …); a daemon older than that field, or a pane opened without a kind,
1439/// names none.
1440pub(crate) fn daemon_panes() -> std::collections::HashMap<String, (u32, Option<String>)> {
1441    crate::slow_log::timed("teams panes ls --processes", daemon_panes_now)
1442}
1443
1444fn daemon_panes_now() -> std::collections::HashMap<String, (u32, Option<String>)> {
1445    let Ok(entry) = crate::teams_entry() else {
1446        return Default::default();
1447    };
1448    let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1449        .ok()
1450        .filter(|value| !value.trim().is_empty())
1451        .unwrap_or_else(|| "node".into());
1452    let listed = |args: &[&str]| {
1453        std::process::Command::new(&node)
1454            .arg(&entry)
1455            .args(args)
1456            .stdin(std::process::Stdio::null())
1457            .stderr(std::process::Stdio::null())
1458            .output()
1459            .ok()
1460            .filter(|output| output.status.success())
1461    };
1462    let Some(output) = listed(&["panes", "ls", "--processes"]).or_else(|| listed(&["panes", "ls"]))
1463    else {
1464        return Default::default();
1465    };
1466    serde_json::from_slice::<Vec<serde_json::Value>>(&output.stdout)
1467        .unwrap_or_default()
1468        .into_iter()
1469        .filter_map(|row| {
1470            let pid = u32::try_from(row["pid"].as_u64()?).ok()?;
1471            let kind = row["kind"]
1472                .as_str()
1473                .map(str::trim)
1474                .filter(|kind| !kind.is_empty())
1475                .map(str::to_ascii_lowercase);
1476            Some((row["pane"].as_str()?.to_string(), (pid, kind)))
1477        })
1478        .collect()
1479}
1480
1481/// A daemon pane's screen: tmux's capture where the daemon's host is tmux, else the machine daemon's
1482/// own capture (`teams panes capture`), as on Windows, where panes are ConPTY.
1483fn pane_screen(pane: &str) -> Option<String> {
1484    #[cfg(unix)]
1485    {
1486        let output = std::process::Command::new("tmux")
1487            .args(["capture-pane", "-p", "-t", &format!("{pane}:")])
1488            .output()
1489            .ok()?;
1490        output
1491            .status
1492            .success()
1493            .then(|| String::from_utf8_lossy(&output.stdout).into_owned())
1494    }
1495    #[cfg(not(unix))]
1496    {
1497        let entry = crate::teams_entry().ok()?;
1498        let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1499            .ok()
1500            .filter(|value| !value.trim().is_empty())
1501            .unwrap_or_else(|| "node".into());
1502        let output = std::process::Command::new(node)
1503            .arg(entry)
1504            .args(["panes", "capture", pane, "--lines", "60"])
1505            .stdin(std::process::Stdio::null())
1506            .output()
1507            .ok()?;
1508        output
1509            .status
1510            .success()
1511            .then(|| String::from_utf8_lossy(&output.stdout).into_owned())
1512    }
1513}
1514
1515/// The daemon pane a session runs in, when it runs in one.
1516pub fn daemon_pane(session: &LiveSession) -> Option<String> {
1517    let name = session.tmux.as_deref()?.split(':').next()?;
1518    let rest = name.strip_prefix("p_")?;
1519    (!rest.is_empty()
1520        && rest
1521            .chars()
1522            .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
1523    .then(|| name.to_string())
1524}
1525
1526/// Whether `to` is running with a door that carries its user's own turn: a hosted runtime's input or a pane
1527/// this daemon holds. A session without one can be reached only as a peer.
1528pub fn has_user_door(homes: &HarnessHomes, to: &MailAddress) -> bool {
1529    LiveSessions::read(homes)
1530        .sessions
1531        .into_iter()
1532        .find(|session| &session.address == to)
1533        .is_some_and(|session| {
1534            matches!(session.door, Door::Runtime(_)) || daemon_pane(&session).is_some()
1535        })
1536}
1537
1538/// Deliver the user's own turn to `to` through the one door that carries the
1539/// user's authority: a hosted runtime's own input, or the session's pane.
1540/// A session with neither (running outside supercode) is refused, never
1541/// reached as a peer instead. Callers hold the owner's authority already:
1542/// this is reached only through owner doors (the machine's own harness.v1).
1543pub async fn deliver_user_turn(
1544    homes: &HarnessHomes,
1545    envelope: &Envelope,
1546    to: &MailAddress,
1547) -> Result<UserTurn, String> {
1548    let session = LiveSessions::read(homes)
1549        .sessions
1550        .into_iter()
1551        .find(|session| &session.address == to)
1552        .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
1553    let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
1554    if let Door::Runtime(record) = &session.door {
1555        let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
1556        mailbox
1557            .deliver_read(envelope)
1558            .map_err(|error| error.to_string())?;
1559        return Ok(match delivered {
1560            RuntimeDelivery::Steered => UserTurn::Steered,
1561            _ => UserTurn::Started,
1562        });
1563    }
1564    let Some(pane) = daemon_pane(&session) else {
1565        let name = session.name.split('@').next().unwrap_or(&session.name);
1566        return Err(format!(
1567            "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
1568             pane with `supercode open {name}`; nothing was sent."
1569        ));
1570    };
1571    mailbox
1572        .deliver(envelope)
1573        .map_err(|error| error.to_string())?;
1574    let typed = type_user_turns(&mailbox, &pane).await;
1575    Ok(if typed.contains(&envelope.id) {
1576        UserTurn::Typed
1577    } else {
1578        UserTurn::Waiting
1579    })
1580}
1581
1582/// Type the user's waiting turns into `pane`, oldest first, each only when
1583/// the pane's composer is empty. Stops at the first that has to wait.
1584/// Returns the ids typed.
1585pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
1586    let mut typed = Vec::new();
1587    for waiting in mailbox.user_turns().unwrap_or_default() {
1588        // Claimed before it is typed: the machine's watcher and a direct delivery both type waiting turns, and
1589        // two typers filling one composer sent a turn's text twice in one prompt. A turn another typer holds is
1590        // that typer's, and so is every turn after it.
1591        let Ok(Some(stored)) = mailbox.claim_user_turn(&waiting) else {
1592            break;
1593        };
1594        match submit_mail_batch(
1595            pane,
1596            &mailbox.address().harness,
1597            &stored.envelope.body,
1598            &[stored.envelope.id.clone()],
1599        )
1600        .await
1601        {
1602            Ok(true) => {
1603                mailbox.acknowledge(&stored).ok();
1604                typed.push(stored.envelope.id.clone());
1605            }
1606            Ok(false) => {
1607                mailbox.release(&stored).ok();
1608                break;
1609            }
1610            Err(error) => {
1611                mailbox.release(&stored).ok();
1612                eprintln!(
1613                    "supercode: the user's turn {} for {} waits: {error}",
1614                    stored.envelope.id,
1615                    mailbox.address()
1616                );
1617                break;
1618            }
1619        }
1620    }
1621    typed
1622}
1623
1624/// The daemon drains a recipient's waiting envelopes as one paste. The machine
1625/// retains the batch and its ids until the native transcript confirms receipt.
1626pub(crate) async fn wake_hook_mailbox(
1627    mailbox: &Mailbox,
1628    pane: &str,
1629    ids: &[String],
1630) -> Result<bool, String> {
1631    let messages = ids
1632        .iter()
1633        .map(|id| mailbox.find(id))
1634        .collect::<Result<Vec<_>, _>>()
1635        .map_err(|error| error.to_string())?;
1636    let messages: Vec<_> = messages.into_iter().flatten().collect();
1637    if messages.is_empty() {
1638        return Ok(false);
1639    }
1640    let mut delivered = true;
1641    for message in messages {
1642        delivered &= submit_mail_batch(
1643            pane,
1644            &mailbox.address().harness,
1645            &message.envelope.render(),
1646            &[message.envelope.id.clone()],
1647        )
1648        .await?;
1649    }
1650    Ok(delivered)
1651}
1652
1653/// `harness` is the addressed session's: the daemon refuses a pane running another program
1654/// (`harness-mismatch`). A Teams package older than that check has no `--harness` and dies on it;
1655/// that one is asked again without it, as before the check.
1656async fn submit_mail_batch(
1657    pane: &str,
1658    harness: &str,
1659    text: &str,
1660    ids: &[String],
1661) -> Result<bool, String> {
1662    let entry = crate::teams_entry().map_err(|error| error.to_string())?;
1663    let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1664        .ok()
1665        .filter(|value| !value.trim().is_empty())
1666        .unwrap_or_else(|| "node".into());
1667    let ids = ids.join(",");
1668    let harness = format!("--harness={harness}");
1669    let mut output = None;
1670    for checked in [true, false] {
1671        let mut args: Vec<&str> = vec![
1672            "input",
1673            pane,
1674            text,
1675            "--when-composer-empty",
1676            "--message-ids",
1677            ids.as_str(),
1678        ];
1679        if checked {
1680            args.push(harness.as_str());
1681        }
1682        let answer = tokio::process::Command::new(&node)
1683            .arg(&entry)
1684            .args(&args)
1685            // A message is the session's own door, not a session acting on a pane: the pane doors read that from this
1686            // process's place under the daemon's `message watch`, not from anything named here.
1687            .env_remove("SUPERCODE_CALLER")
1688            .stdin(std::process::Stdio::null())
1689            .output()
1690            .await
1691            .map_err(|error| error.to_string())?;
1692        let older = checked
1693            && !answer.status.success()
1694            && String::from_utf8_lossy(&answer.stderr).contains("unknown flag: --harness");
1695        output = Some(answer);
1696        if !older {
1697            break;
1698        }
1699    }
1700    let output = output.expect("one attempt always runs");
1701    if !output.status.success() {
1702        return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
1703            &output.stderr,
1704        )));
1705    }
1706    let answer: serde_json::Value =
1707        serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
1708    Ok(answer["delivered"] == true)
1709}