1use 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
35pub const CODEX_HOOK_ARGUMENTS: &str = "message hook codex";
37
38#[derive(Debug, Clone, PartialEq, Eq)]
40pub enum Door {
41 Runtime(Box<LiveRuntimeRecord>),
43 Native(Box<ClaudePeerSession>),
45 Hook,
47 Stored,
49 Operator,
51}
52
53impl Door {
54 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#[derive(Debug, Clone, PartialEq, Eq)]
68pub enum NoDoor {
69 OtherMachine(String),
71 NotRunning,
73}
74
75pub 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#[derive(Debug, Clone, PartialEq, Eq)]
93pub struct LiveSession {
94 pub address: MailAddress,
96 pub name: String,
98 pub status: String,
100 pub door: Door,
102 pub pid: Option<u32>,
105 pub cwd: Option<std::path::PathBuf>,
107 pub tmux: Option<String>,
109}
110
111#[derive(Debug, Default)]
115pub struct LiveSessions {
116 sessions: Vec<LiveSession>,
117}
118
119#[derive(Debug, Clone, PartialEq, Eq)]
121pub enum Unresolved {
122 Stale(String),
124 Unknown(String),
126}
127
128impl LiveSessions {
129 pub fn read(homes: &HarnessHomes) -> Self {
133 let machine = local_machine_name();
134 let registry = read_registry(®istry_dir(homes));
135 let registered = |record: &crate::live_runtime::LiveRuntimeRecord| {
138 registry.iter().find(|session| {
139 session.session_id == record.source.session_id
140 || session.session_id == record.runtime_session_id
141 })
142 };
143 let mut sessions: Vec<LiveSession> = crate::runtime_mail::controlled_runtimes()
144 .into_iter()
145 .filter_map(|record| {
146 let registered = registered(&record);
147 if registered.is_some_and(|session| session.name.starts_with(RELAY_NAME_PREFIX)) {
148 return None;
149 }
150 let address =
151 MailAddress::new(&machine, &record.source.harness, &record.source.session_id)
152 .ok()?;
153 let short: String = record.source.session_id.chars().take(8).collect();
154 let name = match registered {
155 Some(session) if !session.name.is_empty() => session.name.clone(),
156 _ => format!("{}-{short}", record.source.harness),
157 };
158 Some(LiveSession {
159 name: format!("{name}@{machine}"),
160 address,
161 status: "hosted".into(),
162 pid: Some(record.pid),
163 cwd: Some(record.source.workspace.clone()),
164 tmux: None,
165 door: Door::Runtime(Box::new(record)),
166 })
167 })
168 .collect();
169 let controlled = |address: &MailAddress, sessions: &[LiveSession]| {
170 sessions.iter().any(|session| &session.address == address)
171 };
172 for session in registry {
173 if session.name.starts_with(RELAY_NAME_PREFIX) {
174 continue;
175 }
176 let Ok(address) = MailAddress::new(&machine, "claude-code", &session.session_id) else {
177 continue;
178 };
179 if controlled(&address, &sessions) {
180 continue;
181 }
182 sessions.push(LiveSession {
183 address,
184 name: format!("{}@{machine}", session.name),
185 status: session
186 .status
187 .as_ref()
188 .map(|status| status.as_str().to_string())
189 .unwrap_or_else(|| "unknown".into()),
190 pid: Some(session.pid),
191 cwd: session.cwd.clone(),
192 tmux: session.tmux.clone(),
193 door: Door::Native(Box::new(session)),
194 });
195 }
196 let user_hook = codex_user_hook_installed();
197 for (path, status) in crate::codex_peer::live_rollouts(&homes.codex) {
198 let Some((session_id, None)) = crate::codex_peer::rollout_session(&path) else {
199 continue;
200 };
201 let Ok(address) = MailAddress::new(&machine, "codex", &session_id) else {
202 continue;
203 };
204 if controlled(&address, &sessions) {
205 continue;
206 }
207 let cwd = crate::codex_peer::rollout_cwd(&path);
208 let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
209 sessions.push(LiveSession {
210 name: format!("{}@{machine}", codex_name(&session_id)),
211 address,
212 status: status.as_str().to_string(),
213 pid: None,
214 cwd,
215 tmux: None,
216 door: if hooked { Door::Hook } else { Door::Stored },
217 });
218 }
219 Self { sessions }
220 }
221
222 pub fn all(&self) -> &[LiveSession] {
224 &self.sessions
225 }
226
227 pub fn door(&self, harness: &str, session_id: &str) -> Option<&'static str> {
230 self.sessions
231 .iter()
232 .find(|session| {
233 session.address.harness == harness && session.address.session_id == session_id
234 })
235 .map(|session| session.door.name())
236 }
237
238 pub fn resolve(&self, to: &str) -> Result<&LiveSession, Unresolved> {
241 let machine = local_machine_name();
242 if let Ok(address) = MailAddress::parse(to) {
243 return self
244 .sessions
245 .iter()
246 .find(|session| session.address == address)
247 .ok_or_else(|| {
248 Unresolved::Stale(format!(
249 "{to} is no longer running. Nothing was sent. Run supercode message list \
250 for the live sessions."
251 ))
252 });
253 }
254 let wanted = match to.split_once('@') {
255 Some((name, at)) if at == machine => name.to_string(),
256 Some((_, at)) => {
257 return Err(Unresolved::Unknown(format!(
258 "Not sent: {to} is on machine {at}, not this one. Nothing was sent."
259 )))
260 }
261 None => to.to_string(),
262 };
263 let short = |session: &LiveSession| {
264 session
265 .name
266 .split('@')
267 .next()
268 .unwrap_or_default()
269 .to_string()
270 };
271 let matching: Vec<&LiveSession> = self
272 .sessions
273 .iter()
274 .filter(|session| short(session) == wanted)
275 .collect();
276 if let [only] = matching.as_slice() {
277 return Ok(only);
278 }
279 let hint = if matching.len() > 1 {
280 format!(
281 " {} sessions are named {wanted}; use its address.",
282 matching.len()
283 )
284 } else {
285 let near: Vec<String> = self
286 .sessions
287 .iter()
288 .filter(|session| {
289 let name = short(session);
290 name.contains(&wanted)
291 || wanted.contains(&name)
292 || name
293 .chars()
294 .zip(wanted.chars())
295 .take_while(|(a, b)| a == b)
296 .count()
297 >= 4
298 })
299 .take(3)
300 .map(|session| {
301 format!(
302 "{} ({}, {})",
303 session.name, session.address.harness, session.status
304 )
305 })
306 .collect();
307 if near.is_empty() {
308 String::new()
309 } else {
310 format!(" Did you mean: {}?", near.join(", "))
311 }
312 };
313 Err(Unresolved::Unknown(format!(
314 "No session named \"{wanted}\" is reachable.{hint} Run supercode message list. Nothing \
315 was sent."
316 )))
317 }
318}
319
320pub fn has_message_tools(pid: u32) -> bool {
324 let Ok(output) = std::process::Command::new("ps")
325 .args(["-A", "-o", "ppid=,command="])
326 .output()
327 else {
328 return false;
329 };
330 String::from_utf8_lossy(&output.stdout).lines().any(|line| {
331 let line = line.trim_start();
332 let Some((ppid, command)) = line.split_once(' ') else {
333 return false;
334 };
335 ppid.parse::<u32>() == Ok(pid) && command.trim_end().ends_with(" message mcp")
336 })
337}
338
339pub fn codex_name(session_id: &str) -> String {
341 format!("codex-{}", session_id.chars().take(8).collect::<String>())
342}
343
344#[derive(Debug, Clone, PartialEq, Eq)]
346pub struct Caller {
347 pub address: MailAddress,
349 pub name: String,
351}
352
353pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
355 replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
356 shell tool.";
357
358pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
363 let machine = local_machine_name();
364 let registry = read_registry(®istry_dir(homes));
365 let hosted = crate::runtime_mail::controlled_runtimes();
366 for &pid in pids {
367 if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
368 let short: String = record.source.session_id.chars().take(8).collect();
369 return Ok(Caller {
370 address: MailAddress::new(
371 &machine,
372 &record.source.harness,
373 &record.source.session_id,
374 )
375 .map_err(|error| error.to_string())?,
376 name: format!("{}-{short}@{machine}", record.source.harness),
377 });
378 }
379 if let Some(session) = registry.iter().find(|session| session.pid == pid) {
380 if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
381 if !claimed.is_empty() && claimed != session.session_id {
382 return Err(format!(
383 "{CALLER_UNRESOLVED} (CLAUDE_CODE_SESSION_ID names {claimed}, but the Claude \
384 process {pid} above this command is session {})",
385 session.session_id
386 ));
387 }
388 }
389 return Ok(Caller {
390 address: MailAddress::new(&machine, "claude-code", &session.session_id)
391 .map_err(|error| error.to_string())?,
392 name: format!("{}@{machine}", session.name),
393 });
394 }
395 if let Some(thread) = std::env::var("CODEX_THREAD_ID")
399 .ok()
400 .filter(|thread| !thread.is_empty())
401 .filter(|thread| crate::codex_peer::holds_session(pid, thread))
402 {
403 return Ok(Caller {
404 address: MailAddress::new(&machine, "codex", &thread)
405 .map_err(|error| error.to_string())?,
406 name: format!("{}@{machine}", codex_name(&thread)),
407 });
408 }
409 if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
410 return Ok(Caller {
411 address: MailAddress::new(&machine, "codex", &session_id)
412 .map_err(|error| error.to_string())?,
413 name: format!("{}@{machine}", codex_name(&session_id)),
414 });
415 }
416 }
417 #[cfg(windows)]
418 if let Some(session) = msys_cut_claim(®istry, pids) {
419 return Ok(Caller {
420 address: MailAddress::new(&machine, "claude-code", &session.session_id)
421 .map_err(|error| error.to_string())?,
422 name: format!("{}@{machine}", session.name),
423 });
424 }
425 Err(format!(
426 "{CALLER_UNRESOLVED} (looked for this command's processes {pids:?} among {} Claude \
427 sessions in {} and {} hosted runtimes)",
428 registry.len(),
429 registry_dir(homes).display(),
430 hosted.len()
431 ))
432}
433
434#[cfg(windows)]
441fn msys_cut_claim<'a>(
442 registry: &'a [ClaudePeerSession],
443 pids: &[u32],
444) -> Option<&'a ClaudePeerSession> {
445 let table = process_table();
446 let [.., shell, cut] = pids else {
447 return None;
448 };
449 if table.contains_key(cut) {
450 return None;
451 }
452 let name = table.get(shell)?.1.to_ascii_lowercase();
453 if !matches!(name.as_str(), "sh.exe" | "bash.exe" | "dash.exe") {
454 return None;
455 }
456 let claimed = std::env::var("CLAUDE_CODE_SESSION_ID").ok()?;
457 registry
458 .iter()
459 .find(|session| !claimed.is_empty() && session.session_id == claimed)
460}
461
462pub fn process_ancestry() -> Vec<u32> {
464 ancestry_of(std::process::id())
465}
466
467pub fn ancestry_of(pid: u32) -> Vec<u32> {
469 let parents = parent_pids();
470 let mut chain = vec![pid];
471 let mut current = pid;
472 while let Some(&parent) = parents.get(¤t) {
473 if parent <= 1 || chain.contains(&parent) {
474 break;
475 }
476 chain.push(parent);
477 current = parent;
478 }
479 chain
480}
481
482#[cfg(not(windows))]
484fn parent_pids() -> std::collections::HashMap<u32, u32> {
485 let Ok(output) = std::process::Command::new("ps")
486 .args(["-axo", "pid=,ppid="])
487 .output()
488 else {
489 return Default::default();
490 };
491 String::from_utf8_lossy(&output.stdout)
492 .lines()
493 .filter_map(|line| {
494 let mut fields = line.split_whitespace();
495 Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
496 })
497 .collect()
498}
499
500#[cfg(windows)]
502fn parent_pids() -> std::collections::HashMap<u32, u32> {
503 process_table()
504 .into_iter()
505 .map(|(pid, (parent, _))| (pid, parent))
506 .collect()
507}
508
509#[cfg(windows)]
511fn process_table() -> std::collections::HashMap<u32, (u32, String)> {
512 use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
513 use windows_sys::Win32::System::Diagnostics::ToolHelp::{
514 CreateToolhelp32Snapshot, Process32FirstW, Process32NextW, PROCESSENTRY32W,
515 TH32CS_SNAPPROCESS,
516 };
517
518 let mut parents = std::collections::HashMap::new();
519 let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPPROCESS, 0) };
520 if snapshot == INVALID_HANDLE_VALUE {
521 return parents;
522 }
523 let mut entry: PROCESSENTRY32W = unsafe { std::mem::zeroed() };
524 entry.dwSize = std::mem::size_of::<PROCESSENTRY32W>() as u32;
525 let mut has_entry = unsafe { Process32FirstW(snapshot, &mut entry) } != 0;
526 while has_entry {
527 let length = entry
528 .szExeFile
529 .iter()
530 .position(|&unit| unit == 0)
531 .unwrap_or(entry.szExeFile.len());
532 let name = String::from_utf16_lossy(&entry.szExeFile[..length]);
533 parents.insert(entry.th32ProcessID, (entry.th32ParentProcessID, name));
534 has_entry = unsafe { Process32NextW(snapshot, &mut entry) } != 0;
535 }
536 unsafe {
537 CloseHandle(snapshot);
538 }
539 parents
540}
541
542pub fn codex_hooks_path() -> std::path::PathBuf {
544 std::env::var_os("CODEX_HOME")
545 .map(std::path::PathBuf::from)
546 .or_else(|| {
547 supercode_interchange::user_home()
548 .map(std::path::PathBuf::into_os_string)
549 .map(|home| std::path::PathBuf::from(home).join(".codex"))
550 })
551 .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
552 .join("hooks.json")
553}
554
555pub fn codex_user_hook_installed() -> bool {
558 std::fs::read_to_string(codex_hooks_path())
559 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
560}
561
562pub fn codex_project_hook_installed(cwd: &Path) -> bool {
565 cwd.ancestors().any(|directory| {
566 std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
567 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
568 })
569}
570
571#[derive(Debug, Clone, PartialEq, Eq)]
573pub enum Delivered {
574 Steered,
576 Started,
578 Native {
581 busy: bool,
583 },
584 Hooked,
586 Queued,
588 Stored,
590 Operator,
592}
593
594#[derive(Debug, Clone, PartialEq, Eq)]
596pub enum Refused {
597 CannotQueueNative,
600 TooLong(usize),
602}
603
604pub const MAX_RELAYED_BYTES: usize = 100_000;
607
608pub async fn deliver(
612 envelope: &Envelope,
613 to: &MailAddress,
614 door: &Door,
615 wake: bool,
616 notify_when_idle: bool,
617) -> Result<Result<Delivered, Refused>, String> {
618 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
619 let mut final_reply = false;
623 let delivered = match door {
624 Door::Runtime(record) => {
625 let mut sent = envelope.clone();
626 let answers = matches!(envelope.reply_via, ReplyVia::Command);
627 if answers {
628 sent.reply_via = ReplyVia::FinalMessage {
629 destination: envelope.from_name.clone(),
630 };
631 }
632 match deliver_to_runtime(record, sent.render(), wake).await? {
633 RuntimeDelivery::Steered => {
634 mailbox
635 .deliver_read(&sent)
636 .map_err(|error| error.to_string())?;
637 final_reply = answers;
638 Delivered::Steered
639 }
640 RuntimeDelivery::Started => {
641 mailbox
642 .deliver_read(&sent)
643 .map_err(|error| error.to_string())?;
644 final_reply = answers;
645 Delivered::Started
646 }
647 RuntimeDelivery::NotWoken => {
648 mailbox
649 .deliver(envelope)
650 .map_err(|error| error.to_string())?;
651 Delivered::Queued
652 }
653 }
654 }
655 Door::Native(session) => {
656 let busy = session.status != Some(ClaudePeerStatus::Idle);
657 if !wake && !busy {
658 return Ok(Err(Refused::CannotQueueNative));
659 }
660 let mut envelope = envelope.clone();
664 if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
665 envelope.reply_via = ReplyVia::Tool;
666 }
667 let envelope = &envelope;
668 let text = envelope.render();
669 if text.len() > MAX_RELAYED_BYTES {
670 return Ok(Err(Refused::TooLong(text.len())));
671 }
672 match send_through_relay(
673 &envelope.from,
674 &envelope.from_name,
675 &session.name,
676 text,
677 &envelope.id,
678 )
679 .await
680 {
681 RelayReceipt::Delivered { .. } => {
682 mailbox
683 .deliver_read(envelope)
684 .map_err(|error| error.to_string())?;
685 Delivered::Native { busy }
686 }
687 RelayReceipt::Failed { detail } => return Err(detail),
688 }
689 }
690 Door::Hook => {
691 mailbox
692 .deliver(envelope)
693 .map_err(|error| error.to_string())?;
694 Delivered::Hooked
695 }
696 Door::Stored => {
697 mailbox
698 .deliver(envelope)
699 .map_err(|error| error.to_string())?;
700 Delivered::Stored
701 }
702 Door::Operator => {
703 mailbox
704 .deliver(envelope)
705 .map_err(|error| error.to_string())?;
706 Delivered::Operator
707 }
708 };
709 let notice = notify_when_idle && !matches!(door, Door::Operator);
710 if notice || final_reply {
711 let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
712 subscription.notice = notice;
713 subscription.final_reply = final_reply;
714 mailbox
715 .subscribe_idle(&subscription)
716 .map_err(|error| error.to_string())?;
717 if let Ok(program) = supercode_program() {
719 crate::claude_relay::ensure_machine_daemon(&program)
720 .await
721 .ok();
722 }
723 }
724 Ok(Ok(delivered))
725}
726
727#[derive(Debug, Clone, Copy, PartialEq, Eq)]
729pub enum UserTurn {
730 Steered,
732 Started,
734 Typed,
736 Waiting,
740}
741
742impl UserTurn {
743 pub const fn as_str(self) -> &'static str {
745 match self {
746 Self::Steered => "steered",
747 Self::Started => "started",
748 Self::Typed => "typed",
749 Self::Waiting => "waiting",
750 }
751 }
752}
753
754pub fn daemon_pane(session: &LiveSession) -> Option<String> {
756 let name = session.tmux.as_deref()?.split(':').next()?;
757 let rest = name.strip_prefix("p_")?;
758 (!rest.is_empty()
759 && rest
760 .chars()
761 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
762 .then(|| name.to_string())
763}
764
765pub async fn deliver_user_turn(
771 homes: &HarnessHomes,
772 envelope: &Envelope,
773 to: &MailAddress,
774) -> Result<UserTurn, String> {
775 let session = LiveSessions::read(homes)
776 .sessions
777 .into_iter()
778 .find(|session| &session.address == to)
779 .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
780 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
781 if let Door::Runtime(record) = &session.door {
782 let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
783 mailbox
784 .deliver_read(envelope)
785 .map_err(|error| error.to_string())?;
786 return Ok(match delivered {
787 RuntimeDelivery::Steered => UserTurn::Steered,
788 _ => UserTurn::Started,
789 });
790 }
791 let Some(pane) = daemon_pane(&session) else {
792 let name = session.name.split('@').next().unwrap_or(&session.name);
793 return Err(format!(
794 "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
795 pane with `supercode open {name}`; nothing was sent."
796 ));
797 };
798 mailbox
799 .deliver(envelope)
800 .map_err(|error| error.to_string())?;
801 let typed = type_user_turns(&mailbox, &pane).await;
802 Ok(if typed.contains(&envelope.id) {
803 UserTurn::Typed
804 } else {
805 UserTurn::Waiting
806 })
807}
808
809pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
813 let mut typed = Vec::new();
814 for waiting in mailbox.user_turns().unwrap_or_default() {
815 let Ok(Some(stored)) = mailbox.claim_user_turn(&waiting) else {
819 break;
820 };
821 match submit_when_composer_empty(pane, &stored.envelope.body).await {
822 Ok(true) => {
823 mailbox.acknowledge(&stored).ok();
824 typed.push(stored.envelope.id.clone());
825 }
826 Ok(false) => {
827 mailbox.release(&stored).ok();
828 break;
829 }
830 Err(error) => {
831 mailbox.release(&stored).ok();
832 eprintln!(
833 "supercode: the user's turn {} for {} waits: {error}",
834 stored.envelope.id,
835 mailbox.address()
836 );
837 break;
838 }
839 }
840 }
841 typed
842}
843
844async fn submit_when_composer_empty(pane: &str, text: &str) -> Result<bool, String> {
847 let entry = crate::teams_entry().map_err(|error| error.to_string())?;
848 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
849 .ok()
850 .filter(|value| !value.trim().is_empty())
851 .unwrap_or_else(|| "node".into());
852 let output = tokio::process::Command::new(node)
853 .arg(entry)
854 .args(["input", pane, text, "--when-composer-empty"])
855 .stdin(std::process::Stdio::null())
856 .output()
857 .await
858 .map_err(|error| error.to_string())?;
859 if !output.status.success() {
860 return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
861 &output.stderr,
862 )));
863 }
864 let answer: serde_json::Value =
865 serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
866 Ok(answer["delivered"] == true)
867}