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    Hook,
47    /// A running session with no door.
48    Stored,
49    /// An operator's mailbox.
50    Operator,
51}
52
53impl Door {
54    /// Stable name of the door, as `message list` and receipts show it.
55    pub fn name(&self) -> &'static str {
56        match self {
57            Self::Runtime(_) => "runtime",
58            Self::Native(_) => "native",
59            Self::Hook => "hook",
60            Self::Stored => "stored",
61            Self::Operator => "operator",
62        }
63    }
64}
65
66/// Why no door reaches an address on this machine.
67#[derive(Debug, Clone, PartialEq, Eq)]
68pub enum NoDoor {
69    /// The address names another machine.
70    OtherMachine(String),
71    /// No running session has that address.
72    NotRunning,
73}
74
75/// The door that reaches `to` on this machine.
76pub fn door_for(homes: &HarnessHomes, to: &MailAddress) -> Result<Door, NoDoor> {
77    if to.machine != local_machine_name() {
78        return Err(NoDoor::OtherMachine(to.machine.clone()));
79    }
80    if to.harness == "operator" {
81        return Ok(Door::Operator);
82    }
83    LiveSessions::read(homes)
84        .sessions
85        .into_iter()
86        .find(|session| &session.address == to)
87        .map(|session| session.door)
88        .ok_or(NoDoor::NotRunning)
89}
90
91/// One running session on this machine, as the router reaches it.
92#[derive(Debug, Clone, PartialEq, Eq)]
93pub struct LiveSession {
94    /// Its address.
95    pub address: MailAddress,
96    /// The name an agent knows it by (`volter-2c@machine`).
97    pub name: String,
98    /// Its status as its harness reports it (`busy`, `idle`, `hosted`, …).
99    pub status: String,
100    /// The door that reaches it.
101    pub door: Door,
102    /// The process running it, when the harness names one (a hosted
103    /// runtime's host, a Claude session's own process).
104    pub pid: Option<u32>,
105    /// Its working directory.
106    pub cwd: Option<std::path::PathBuf>,
107    /// The tmux session and pane it runs in, when Claude records one.
108    pub tmux: Option<String>,
109    /// Its transcript, when the harness names the file (a Codex rollout).
110    pub transcript: Option<std::path::PathBuf>,
111}
112
113impl LiveSession {
114    /// When its transcript last recorded a message (a prompt, a reply, a tool
115    /// call or result), in epoch milliseconds: what the session last did,
116    /// read from the harness's own records. A queued or injected notice is
117    /// not a message, so this is not the transcript's mtime.
118    pub fn last_message_at_ms(&self, homes: &HarnessHomes) -> Option<u64> {
119        let harness = self.address.harness.as_str();
120        let path = match &self.transcript {
121            Some(path) => path.clone(),
122            None if harness == "claude-code" => {
123                let file = format!("{}.jsonl", self.address.session_id);
124                std::fs::read_dir(&homes.claude_code)
125                    .ok()?
126                    .flatten()
127                    .map(|project| project.path().join(&file))
128                    .find(|path| path.is_file())?
129            }
130            None => return None,
131        };
132        last_message_at_ms(&path, harness)
133    }
134}
135
136impl LiveSession {
137    /// What a session waiting on its user is asking, read from its transcript: the newest tool call with no result
138    /// yet. An `AskUserQuestion` gives its questions, headers and option labels; any other tool its name and a
139    /// shortened input. `None` when the session is not waiting or its transcript shows nothing pending.
140    pub fn pending_request(&self, homes: &HarnessHomes) -> Option<serde_json::Value> {
141        if self.status != "waiting" || self.address.harness != "claude-code" {
142            return None;
143        }
144        let file = format!("{}.jsonl", self.address.session_id);
145        let path = std::fs::read_dir(&homes.claude_code)
146            .ok()?
147            .flatten()
148            .map(|project| project.path().join(&file))
149            .find(|path| path.is_file())?;
150        pending_request(&path)
151    }
152}
153
154/// The newest tool call in a Claude transcript that has no result yet, as [`LiveSession::pending_request`] shows it.
155pub fn pending_request(path: &Path) -> Option<serde_json::Value> {
156    use std::io::{Read, Seek, SeekFrom};
157    let mut file = std::fs::File::open(path).ok()?;
158    let length = file.metadata().ok()?.len();
159    let start = length.saturating_sub(512 * 1024);
160    file.seek(SeekFrom::Start(start)).ok()?;
161    let mut bytes = Vec::new();
162    file.read_to_end(&mut bytes).ok()?;
163    let text = String::from_utf8_lossy(&bytes);
164    let mut answered = std::collections::HashSet::new();
165    for line in text.lines().rev() {
166        let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
167            continue;
168        };
169        if record["isSidechain"] == true {
170            continue;
171        }
172        let Some(content) = record
173            .pointer("/message/content")
174            .and_then(serde_json::Value::as_array)
175        else {
176            continue;
177        };
178        match record["type"].as_str() {
179            Some("user") => {
180                for item in content.iter().filter(|item| item["type"] == "tool_result") {
181                    if let Some(id) = item["tool_use_id"].as_str() {
182                        answered.insert(id.to_string());
183                    }
184                }
185            }
186            Some("assistant") => {
187                let Some(call) = content.iter().rev().find(|item| item["type"] == "tool_use")
188                else {
189                    continue;
190                };
191                if call["id"].as_str().is_some_and(|id| answered.contains(id)) {
192                    return None;
193                }
194                let tool = call["name"].as_str().unwrap_or_default();
195                if tool == "AskUserQuestion" {
196                    let questions = call["input"]["questions"]
197                        .as_array()
198                        .map(|questions| {
199                            questions
200                                .iter()
201                                .map(|question| {
202                                    serde_json::json!({
203                                        "question": question["question"],
204                                        "header": question["header"],
205                                        "options": question["options"].as_array().map(|options| options.iter().map(|option| option["label"].clone()).collect::<Vec<_>>()).unwrap_or_default(),
206                                    })
207                                })
208                                .collect::<Vec<_>>()
209                        })
210                        .unwrap_or_default();
211                    return Some(serde_json::json!({"tool": tool, "questions": questions}));
212                }
213                let input = call["input"].to_string();
214                let input: String = input.chars().take(500).collect();
215                return Some(serde_json::json!({"tool": tool, "input": input}));
216            }
217            _ => {}
218        }
219    }
220    None
221}
222
223/// The newest message record's timestamp in a Claude or Codex transcript.
224fn last_message_at_ms(path: &Path, harness: &str) -> Option<u64> {
225    use std::io::{Read, Seek, SeekFrom};
226    // A tool result can be large; widen the window once before giving up.
227    for window in [256 * 1024_u64, 8 * 1024 * 1024] {
228        let mut file = std::fs::File::open(path).ok()?;
229        let length = file.metadata().ok()?.len();
230        let start = length.saturating_sub(window);
231        file.seek(SeekFrom::Start(start)).ok()?;
232        let mut bytes = Vec::new();
233        file.read_to_end(&mut bytes).ok()?;
234        let text = String::from_utf8_lossy(&bytes);
235        let mut lines = text.lines().rev().collect::<Vec<_>>();
236        if start > 0 {
237            lines.pop(); // the first line may begin mid-record
238        }
239        for line in lines {
240            let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
241                continue;
242            };
243            let message = match harness {
244                "codex" => record["type"] == "response_item",
245                _ => {
246                    matches!(record["type"].as_str(), Some("user" | "assistant"))
247                        && record["isSidechain"] != true
248                }
249            };
250            if !message {
251                continue;
252            }
253            if let Some(at) = record["timestamp"]
254                .as_str()
255                .and_then(supercode_interchange::sidecar::rfc3339_to_ms)
256                .and_then(|at| u64::try_from(at).ok())
257            {
258                return Some(at);
259            }
260        }
261        if start == 0 {
262            return None;
263        }
264    }
265    None
266}
267
268/// Every running session on this machine with its door, read once: what
269/// `message list` shows, what names resolve against, and what discovery
270/// projects onto each discovered session as `delivery`.
271#[derive(Debug, Default)]
272pub struct LiveSessions {
273    sessions: Vec<LiveSession>,
274}
275
276/// Why a receiver named by an agent was not found.
277#[derive(Debug, Clone, PartialEq, Eq)]
278pub enum Unresolved {
279    /// The address names a session that is not running.
280    Stale(String),
281    /// No running session has that name (the text suggests near names).
282    Unknown(String),
283}
284
285impl LiveSessions {
286    /// Read every running session now. Sessions supercode controls come
287    /// first, and a degraded-tier row for the same session is dropped: the
288    /// runtime is the default tier.
289    pub fn read(homes: &HarnessHomes) -> Self {
290        let machine = local_machine_name();
291        let registry = read_registry(&registry_dir(homes));
292        // A hosted Claude session is also in Claude's registry, which knows
293        // its name; a relay is not a session and is not listed.
294        let registered = |record: &crate::live_runtime::LiveRuntimeRecord| {
295            registry.iter().find(|session| {
296                session.session_id == record.source.session_id
297                    || session.session_id == record.runtime_session_id
298            })
299        };
300        let mut sessions: Vec<LiveSession> = crate::runtime_mail::controlled_runtimes()
301            .into_iter()
302            .filter_map(|record| {
303                let registered = registered(&record);
304                if registered.is_some_and(|session| session.name.starts_with(RELAY_NAME_PREFIX)) {
305                    return None;
306                }
307                let address =
308                    MailAddress::new(&machine, &record.source.harness, &record.source.session_id)
309                        .ok()?;
310                let short: String = record.source.session_id.chars().take(8).collect();
311                let name = match registered {
312                    Some(session) if !session.name.is_empty() => session.name.clone(),
313                    _ => format!("{}-{short}", record.source.harness),
314                };
315                Some(LiveSession {
316                    name: format!("{name}@{machine}"),
317                    address,
318                    status: "hosted".into(),
319                    pid: Some(record.pid),
320                    cwd: Some(record.source.workspace.clone()),
321                    tmux: None,
322                    transcript: None,
323                    door: Door::Runtime(Box::new(record)),
324                })
325            })
326            .collect();
327        let controlled = |address: &MailAddress, sessions: &[LiveSession]| {
328            sessions.iter().any(|session| &session.address == address)
329        };
330        for session in registry {
331            if session.name.starts_with(RELAY_NAME_PREFIX) {
332                continue;
333            }
334            let Ok(address) = MailAddress::new(&machine, "claude-code", &session.session_id) else {
335                continue;
336            };
337            if controlled(&address, &sessions) {
338                continue;
339            }
340            sessions.push(LiveSession {
341                address,
342                name: format!("{}@{machine}", session.name),
343                status: session
344                    .status
345                    .as_ref()
346                    .map(|status| status.as_str().to_string())
347                    .unwrap_or_else(|| "unknown".into()),
348                pid: Some(session.pid),
349                cwd: session.cwd.clone(),
350                tmux: session.tmux.clone(),
351                transcript: None,
352                door: Door::Native(Box::new(session)),
353            });
354        }
355        let user_hook = codex_user_hook_installed();
356        for (path, status) in crate::codex_peer::live_rollouts(&homes.codex) {
357            let Some((session_id, None)) = crate::codex_peer::rollout_session(&path) else {
358                continue;
359            };
360            let Ok(address) = MailAddress::new(&machine, "codex", &session_id) else {
361                continue;
362            };
363            if controlled(&address, &sessions) {
364                continue;
365            }
366            let cwd = crate::codex_peer::rollout_cwd(&path);
367            let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
368            sessions.push(LiveSession {
369                name: format!("{}@{machine}", codex_name(&session_id)),
370                address,
371                status: status.as_str().to_string(),
372                pid: None,
373                cwd,
374                tmux: None,
375                transcript: Some(path),
376                door: if hooked { Door::Hook } else { Door::Stored },
377            });
378        }
379        Self { sessions }
380    }
381
382    /// Every running session.
383    pub fn all(&self) -> &[LiveSession] {
384        &self.sessions
385    }
386
387    /// The door that reaches `harness`'s session `session_id`, when it is
388    /// running.
389    pub fn door(&self, harness: &str, session_id: &str) -> Option<&'static str> {
390        self.sessions
391            .iter()
392            .find(|session| {
393                session.address.harness == harness && session.address.session_id == session_id
394            })
395            .map(|session| session.door.name())
396    }
397
398    /// The running session an agent means by `to`: an address, `name@machine`
399    /// on this machine, or a name unique here.
400    pub fn resolve(&self, to: &str) -> Result<&LiveSession, Unresolved> {
401        let machine = local_machine_name();
402        if let Ok(address) = MailAddress::parse(to) {
403            return self
404                .sessions
405                .iter()
406                .find(|session| session.address == address)
407                .ok_or_else(|| {
408                    Unresolved::Stale(format!(
409                        "{to} is no longer running. Nothing was sent. Run supercode message list \
410                         for the live sessions."
411                    ))
412                });
413        }
414        let wanted = match to.split_once('@') {
415            Some((name, at)) if at == machine => name.to_string(),
416            Some((_, at)) => {
417                return Err(Unresolved::Unknown(format!(
418                    "Not sent: {to} is on machine {at}, not this one. Nothing was sent."
419                )))
420            }
421            None => to.to_string(),
422        };
423        let short = |session: &LiveSession| {
424            session
425                .name
426                .split('@')
427                .next()
428                .unwrap_or_default()
429                .to_string()
430        };
431        let matching: Vec<&LiveSession> = self
432            .sessions
433            .iter()
434            .filter(|session| short(session) == wanted)
435            .collect();
436        if let [only] = matching.as_slice() {
437            return Ok(only);
438        }
439        let hint = if matching.len() > 1 {
440            format!(
441                " {} sessions are named {wanted}; use its address.",
442                matching.len()
443            )
444        } else {
445            let near: Vec<String> = self
446                .sessions
447                .iter()
448                .filter(|session| {
449                    let name = short(session);
450                    name.contains(&wanted)
451                        || wanted.contains(&name)
452                        || name
453                            .chars()
454                            .zip(wanted.chars())
455                            .take_while(|(a, b)| a == b)
456                            .count()
457                            >= 4
458                })
459                .take(3)
460                .map(|session| {
461                    format!(
462                        "{} ({}, {})",
463                        session.name, session.address.harness, session.status
464                    )
465                })
466                .collect();
467            if near.is_empty() {
468                String::new()
469            } else {
470                format!(" Did you mean: {}?", near.join(", "))
471            }
472        };
473        Err(Unresolved::Unknown(format!(
474            "No session named \"{wanted}\" is reachable.{hint} Run supercode message list. Nothing \
475             was sent."
476        )))
477    }
478}
479
480/// Whether the session process `pid` has supercode's messaging tools: a
481/// harness starts each MCP server as a child of the session, so the tools are
482/// loaded exactly when a live `supercode message mcp` is one of its children.
483pub fn has_message_tools(pid: u32) -> bool {
484    let Ok(output) = std::process::Command::new("ps")
485        .args(["-A", "-o", "ppid=,command="])
486        .output()
487    else {
488        return false;
489    };
490    String::from_utf8_lossy(&output.stdout).lines().any(|line| {
491        let line = line.trim_start();
492        let Some((ppid, command)) = line.split_once(' ') else {
493            return false;
494        };
495        ppid.parse::<u32>() == Ok(pid) && command.trim_end().ends_with(" message mcp")
496    })
497}
498
499/// Display name of a Codex conversation: `codex-` and the start of its id.
500pub fn codex_name(session_id: &str) -> String {
501    format!("codex-{}", session_id.chars().take(8).collect::<String>())
502}
503
504/// The session a process belongs to: the sender a message is from.
505#[derive(Debug, Clone, PartialEq, Eq)]
506pub struct Caller {
507    /// Its address, where replies go.
508    pub address: MailAddress,
509    /// The name it is known by.
510    pub name: String,
511}
512
513/// Why no session could be found behind a process.
514pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
515    replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
516    shell tool.";
517
518/// The session behind a process, from its ancestry `pids` (nearest first):
519/// the first that owns a hosted runtime, a Claude session or a Codex
520/// conversation. Never declared by the caller; an environment id only
521/// corroborates, and a mismatch refuses.
522pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
523    let machine = local_machine_name();
524    let registry = read_registry(&registry_dir(homes));
525    let hosted = crate::runtime_mail::controlled_runtimes();
526    for &pid in pids {
527        if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
528            let short: String = record.source.session_id.chars().take(8).collect();
529            return Ok(Caller {
530                address: MailAddress::new(
531                    &machine,
532                    &record.source.harness,
533                    &record.source.session_id,
534                )
535                .map_err(|error| error.to_string())?,
536                name: format!("{}-{short}@{machine}", record.source.harness),
537            });
538        }
539        if let Some(session) = registry.iter().find(|session| session.pid == pid) {
540            if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
541                if !claimed.is_empty() && claimed != session.session_id {
542                    return Err(format!(
543                        "{CALLER_UNRESOLVED} (CLAUDE_CODE_SESSION_ID names {claimed}, but the Claude \
544                         process {pid} above this command is session {})",
545                        session.session_id
546                    ));
547                }
548            }
549            return Ok(Caller {
550                address: MailAddress::new(&machine, "claude-code", &session.session_id)
551                    .map_err(|error| error.to_string())?,
552                name: format!("{}@{machine}", session.name),
553            });
554        }
555        // Codex names the thread a command runs for (`CODEX_THREAD_ID`); the
556        // claim stands when this process holds that thread's rollout, as
557        // Codex's shared app-server daemon does for every thread it runs.
558        if let Some(thread) = std::env::var("CODEX_THREAD_ID")
559            .ok()
560            .filter(|thread| !thread.is_empty())
561            .filter(|thread| crate::codex_peer::holds_session(pid, thread))
562        {
563            return Ok(Caller {
564                address: MailAddress::new(&machine, "codex", &thread)
565                    .map_err(|error| error.to_string())?,
566                name: format!("{}@{machine}", codex_name(&thread)),
567            });
568        }
569        if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
570            return Ok(Caller {
571                address: MailAddress::new(&machine, "codex", &session_id)
572                    .map_err(|error| error.to_string())?,
573                name: format!("{}@{machine}", codex_name(&session_id)),
574            });
575        }
576    }
577    #[cfg(windows)]
578    if let Some(session) = msys_cut_claim(&registry, pids) {
579        return Ok(Caller {
580            address: MailAddress::new(&machine, "claude-code", &session.session_id)
581                .map_err(|error| error.to_string())?,
582            name: format!("{}@{machine}", session.name),
583        });
584    }
585    Err(format!(
586        "{CALLER_UNRESOLVED} (looked for this command's processes {pids:?} among {} Claude \
587         sessions in {} and {} hosted runtimes)",
588        registry.len(),
589        registry_dir(homes).display(),
590        hosted.len()
591    ))
592}
593
594/// Git Bash (MSYS) runs a command in a forked process that hands over to it. When the command is
595/// itself an MSYS program (a `sh` script, as npm's command shims are), the forked process exits once
596/// it has handed over, so the Windows parent chain is cut just above that shell and never reaches
597/// the Claude session running it. When the ancestry ends at exactly such a cut (its last pid is gone,
598/// the one before it is an MSYS shell), the session Claude names in `CLAUDE_CODE_SESSION_ID` stands
599/// if it is a live Claude session on this machine.
600#[cfg(windows)]
601fn msys_cut_claim<'a>(
602    registry: &'a [ClaudePeerSession],
603    pids: &[u32],
604) -> Option<&'a ClaudePeerSession> {
605    let table = process_table();
606    let [.., shell, cut] = pids else {
607        return None;
608    };
609    if table.contains_key(cut) {
610        return None;
611    }
612    let name = table.get(shell)?.1.to_ascii_lowercase();
613    if !matches!(name.as_str(), "sh.exe" | "bash.exe" | "dash.exe") {
614        return None;
615    }
616    let claimed = std::env::var("CLAUDE_CODE_SESSION_ID").ok()?;
617    registry
618        .iter()
619        .find(|session| !claimed.is_empty() && session.session_id == claimed)
620}
621
622/// This process's ancestry, nearest first (itself included).
623pub fn process_ancestry() -> Vec<u32> {
624    ancestry_of(std::process::id())
625}
626
627/// A process's ancestry, nearest first (itself included).
628pub fn ancestry_of(pid: u32) -> Vec<u32> {
629    let parents = parent_pids();
630    let mut chain = vec![pid];
631    let mut current = pid;
632    while let Some(&parent) = parents.get(&current) {
633        if parent <= 1 || chain.contains(&parent) {
634            break;
635        }
636        chain.push(parent);
637        current = parent;
638    }
639    chain
640}
641
642/// Every process's parent: pid → parent pid, from `ps`.
643#[cfg(not(windows))]
644fn parent_pids() -> std::collections::HashMap<u32, u32> {
645    let Ok(output) = std::process::Command::new("ps")
646        .args(["-axo", "pid=,ppid="])
647        .output()
648    else {
649        return Default::default();
650    };
651    String::from_utf8_lossy(&output.stdout)
652        .lines()
653        .filter_map(|line| {
654            let mut fields = line.split_whitespace();
655            Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
656        })
657        .collect()
658}
659
660/// Every process's parent: pid → parent pid, from a process snapshot.
661#[cfg(windows)]
662fn parent_pids() -> std::collections::HashMap<u32, u32> {
663    process_table()
664        .into_iter()
665        .map(|(pid, (parent, _))| (pid, parent))
666        .collect()
667}
668
669/// Every process: pid → (parent pid, executable file name), from a process snapshot.
670#[cfg(windows)]
671fn process_table() -> std::collections::HashMap<u32, (u32, String)> {
672    use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
673    use windows_sys::Win32::System::Diagnostics::ToolHelp::{
674        CreateToolhelp32Snapshot, Process32FirstW, Process32NextW, PROCESSENTRY32W,
675        TH32CS_SNAPPROCESS,
676    };
677
678    let mut parents = std::collections::HashMap::new();
679    let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPPROCESS, 0) };
680    if snapshot == INVALID_HANDLE_VALUE {
681        return parents;
682    }
683    let mut entry: PROCESSENTRY32W = unsafe { std::mem::zeroed() };
684    entry.dwSize = std::mem::size_of::<PROCESSENTRY32W>() as u32;
685    let mut has_entry = unsafe { Process32FirstW(snapshot, &mut entry) } != 0;
686    while has_entry {
687        let length = entry
688            .szExeFile
689            .iter()
690            .position(|&unit| unit == 0)
691            .unwrap_or(entry.szExeFile.len());
692        let name = String::from_utf16_lossy(&entry.szExeFile[..length]);
693        parents.insert(entry.th32ProcessID, (entry.th32ParentProcessID, name));
694        has_entry = unsafe { Process32NextW(snapshot, &mut entry) } != 0;
695    }
696    unsafe {
697        CloseHandle(snapshot);
698    }
699    parents
700}
701
702/// Codex's user-level hooks file.
703pub fn codex_hooks_path() -> std::path::PathBuf {
704    std::env::var_os("CODEX_HOME")
705        .map(std::path::PathBuf::from)
706        .or_else(|| {
707            supercode_interchange::user_home()
708                .map(std::path::PathBuf::into_os_string)
709                .map(|home| std::path::PathBuf::from(home).join(".codex"))
710        })
711        .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
712        .join("hooks.json")
713}
714
715/// Whether supercode's mail hook is in Codex's user hooks file. (Codex runs
716/// it only once its user has trusted it.)
717pub fn codex_user_hook_installed() -> bool {
718    std::fs::read_to_string(codex_hooks_path())
719        .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
720}
721
722/// Whether supercode's mail hook is in a project's Codex hooks file, in the
723/// session's directory or one above it.
724pub fn codex_project_hook_installed(cwd: &Path) -> bool {
725    cwd.ancestors().any(|directory| {
726        std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
727            .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
728    })
729}
730
731/// How a delivery ended.
732#[derive(Debug, Clone, PartialEq, Eq)]
733pub enum Delivered {
734    /// The receiver has it now: steered into a running turn.
735    Steered,
736    /// The receiver has it now: it started a turn in an idle session.
737    Started,
738    /// Claude reported it in the receiver's inbox; `busy` says whether the
739    /// receiver reads it at its next tool call (true) or it starts a turn.
740    Native {
741        /// True when the receiver was in a turn.
742        busy: bool,
743    },
744    /// Filed in the receiver's mailbox, which a hook points it at.
745    Hooked,
746    /// Filed in the receiver's mailbox, waiting to be read.
747    Queued,
748    /// Filed with no door to show it.
749    Stored,
750    /// Filed in an operator's mailbox.
751    Operator,
752}
753
754/// A delivery refused before anything was sent.
755#[derive(Debug, Clone, PartialEq, Eq)]
756pub enum Refused {
757    /// `--queue` cannot hold a message for an idle Claude session: Claude
758    /// starts a turn for every message it delivers.
759    CannotQueueNative,
760    /// The relayed text would not fit one relay turn.
761    TooLong(usize),
762}
763
764/// Largest message relayed into a Claude session. The relay copies it into a
765/// model turn byte for byte, so it must fit comfortably in one.
766pub const MAX_RELAYED_BYTES: usize = 100_000;
767
768/// Deliver `envelope` to `to` through `door`. With `wake` false an idle
769/// receiver is not started. With `notify_when_idle` the sender gets one idle
770/// notice after the receiver's next turn ends.
771pub async fn deliver(
772    envelope: &Envelope,
773    to: &MailAddress,
774    door: &Door,
775    wake: bool,
776    notify_when_idle: bool,
777) -> Result<Result<Delivered, Refused>, String> {
778    let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
779    // A runtime answers with its turn's final message, which the watcher
780    // sends back: not every runtime can run a command (a hosted agent may
781    // have no shell), and each can end a turn.
782    let mut final_reply = false;
783    let delivered = match door {
784        Door::Runtime(record) => {
785            let mut sent = envelope.clone();
786            let answers = matches!(envelope.reply_via, ReplyVia::Command);
787            if answers {
788                sent.reply_via = ReplyVia::FinalMessage {
789                    destination: envelope.from_name.clone(),
790                };
791            }
792            match deliver_to_runtime(record, sent.render(), wake).await? {
793                RuntimeDelivery::Steered => {
794                    mailbox
795                        .deliver_read(&sent)
796                        .map_err(|error| error.to_string())?;
797                    final_reply = answers;
798                    Delivered::Steered
799                }
800                RuntimeDelivery::Started => {
801                    mailbox
802                        .deliver_read(&sent)
803                        .map_err(|error| error.to_string())?;
804                    final_reply = answers;
805                    Delivered::Started
806                }
807                RuntimeDelivery::NotWoken => {
808                    mailbox
809                        .deliver(envelope)
810                        .map_err(|error| error.to_string())?;
811                    Delivered::Queued
812                }
813            }
814        }
815        Door::Native(session) => {
816            let busy = session.status != Some(ClaudePeerStatus::Idle);
817            if !wake && !busy {
818                return Ok(Err(Refused::CannotQueueNative));
819            }
820            // The same envelope every door delivers; Claude wraps it in its
821            // own, which names the relay. Its reply line names the door this
822            // session really has: its tools when it runs supercode's MCP server.
823            let mut envelope = envelope.clone();
824            if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
825                envelope.reply_via = ReplyVia::Tool;
826            }
827            let envelope = &envelope;
828            let text = envelope.render();
829            if text.len() > MAX_RELAYED_BYTES {
830                return Ok(Err(Refused::TooLong(text.len())));
831            }
832            match send_through_relay(
833                &envelope.from,
834                &envelope.from_name,
835                &session.name,
836                text,
837                &envelope.id,
838            )
839            .await
840            {
841                RelayReceipt::Delivered { .. } => {
842                    mailbox
843                        .deliver_read(envelope)
844                        .map_err(|error| error.to_string())?;
845                    Delivered::Native { busy }
846                }
847                RelayReceipt::Failed { detail } => return Err(detail),
848            }
849        }
850        Door::Hook => {
851            mailbox
852                .deliver(envelope)
853                .map_err(|error| error.to_string())?;
854            Delivered::Hooked
855        }
856        Door::Stored => {
857            mailbox
858                .deliver(envelope)
859                .map_err(|error| error.to_string())?;
860            Delivered::Stored
861        }
862        Door::Operator => {
863            mailbox
864                .deliver(envelope)
865                .map_err(|error| error.to_string())?;
866            Delivered::Operator
867        }
868    };
869    let notice = notify_when_idle && !matches!(door, Door::Operator);
870    if notice || final_reply {
871        let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
872        subscription.notice = notice;
873        subscription.final_reply = final_reply;
874        mailbox
875            .subscribe_idle(&subscription)
876            .map_err(|error| error.to_string())?;
877        // The watcher that settles it runs beside the machine daemon.
878        if let Ok(program) = supercode_program() {
879            crate::claude_relay::ensure_machine_daemon(&program)
880                .await
881                .ok();
882        }
883    }
884    Ok(Ok(delivered))
885}
886
887/// How the user's own turn reached its session.
888#[derive(Debug, Clone, Copy, PartialEq, Eq)]
889pub enum UserTurn {
890    /// A hosted runtime was in a turn; the words were steered into it.
891    Steered,
892    /// A hosted runtime was idle; the words started a turn.
893    Started,
894    /// Typed into the session's pane as its own submitted turn.
895    Typed,
896    /// The pane's composer holds a draft, or its program is not at its
897    /// composer: the turn waits in the session's mailbox and is typed once
898    /// the composer is empty.
899    Waiting,
900}
901
902impl UserTurn {
903    /// Stable wire spelling.
904    pub const fn as_str(self) -> &'static str {
905        match self {
906            Self::Steered => "steered",
907            Self::Started => "started",
908            Self::Typed => "typed",
909            Self::Waiting => "waiting",
910        }
911    }
912}
913
914/// The daemon pane a session runs in, when it runs in one.
915pub fn daemon_pane(session: &LiveSession) -> Option<String> {
916    let name = session.tmux.as_deref()?.split(':').next()?;
917    let rest = name.strip_prefix("p_")?;
918    (!rest.is_empty()
919        && rest
920            .chars()
921            .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
922    .then(|| name.to_string())
923}
924
925/// Deliver the user's own turn to `to` through the one door that carries the
926/// user's authority: a hosted runtime's own input, or the session's pane.
927/// A session with neither (running outside supercode) is refused, never
928/// reached as a peer instead. Callers hold the owner's authority already:
929/// this is reached only through owner doors (the machine's own harness.v1).
930pub async fn deliver_user_turn(
931    homes: &HarnessHomes,
932    envelope: &Envelope,
933    to: &MailAddress,
934) -> Result<UserTurn, String> {
935    let session = LiveSessions::read(homes)
936        .sessions
937        .into_iter()
938        .find(|session| &session.address == to)
939        .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
940    let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
941    if let Door::Runtime(record) = &session.door {
942        let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
943        mailbox
944            .deliver_read(envelope)
945            .map_err(|error| error.to_string())?;
946        return Ok(match delivered {
947            RuntimeDelivery::Steered => UserTurn::Steered,
948            _ => UserTurn::Started,
949        });
950    }
951    let Some(pane) = daemon_pane(&session) else {
952        let name = session.name.split('@').next().unwrap_or(&session.name);
953        return Err(format!(
954            "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
955             pane with `supercode open {name}`; nothing was sent."
956        ));
957    };
958    mailbox
959        .deliver(envelope)
960        .map_err(|error| error.to_string())?;
961    let typed = type_user_turns(&mailbox, &pane).await;
962    Ok(if typed.contains(&envelope.id) {
963        UserTurn::Typed
964    } else {
965        UserTurn::Waiting
966    })
967}
968
969/// Type the user's waiting turns into `pane`, oldest first, each only when
970/// the pane's composer is empty. Stops at the first that has to wait.
971/// Returns the ids typed.
972pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
973    let mut typed = Vec::new();
974    for waiting in mailbox.user_turns().unwrap_or_default() {
975        // Claimed before it is typed: the machine's watcher and a direct delivery both type waiting turns, and
976        // two typers filling one composer sent a turn's text twice in one prompt. A turn another typer holds is
977        // that typer's, and so is every turn after it.
978        let Ok(Some(stored)) = mailbox.claim_user_turn(&waiting) else {
979            break;
980        };
981        match submit_when_composer_empty(pane, &stored.envelope.body).await {
982            Ok(true) => {
983                mailbox.acknowledge(&stored).ok();
984                typed.push(stored.envelope.id.clone());
985            }
986            Ok(false) => {
987                mailbox.release(&stored).ok();
988                break;
989            }
990            Err(error) => {
991                mailbox.release(&stored).ok();
992                eprintln!(
993                    "supercode: the user's turn {} for {} waits: {error}",
994                    stored.envelope.id,
995                    mailbox.address()
996                );
997                break;
998            }
999        }
1000    }
1001    typed
1002}
1003
1004/// The daemon's pane door: type `text` as a submitted turn if the composer
1005/// is empty now. `Ok(false)` when it holds a draft or is not on screen.
1006async fn submit_when_composer_empty(pane: &str, text: &str) -> Result<bool, String> {
1007    let entry = crate::teams_entry().map_err(|error| error.to_string())?;
1008    let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1009        .ok()
1010        .filter(|value| !value.trim().is_empty())
1011        .unwrap_or_else(|| "node".into());
1012    let output = tokio::process::Command::new(node)
1013        .arg(entry)
1014        .args(["input", pane, text, "--when-composer-empty"])
1015        // A message is the session's own door, not a session acting on a pane: no caller rides it.
1016        .env("SUPERCODE_CALLER", "")
1017        .stdin(std::process::Stdio::null())
1018        .output()
1019        .await
1020        .map_err(|error| error.to_string())?;
1021    if !output.status.success() {
1022        return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
1023            &output.stderr,
1024        )));
1025    }
1026    let answer: serde_json::Value =
1027        serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
1028    Ok(answer["delivered"] == true)
1029}