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 #[cfg(test)]
321 pub(crate) fn with(entries: &[(&str, &str, &'static str)]) -> Self {
322 Self {
323 sessions: entries
324 .iter()
325 .filter_map(|(harness, id, door)| {
326 let door = match *door {
327 "native" => Door::Native(Box::new(ClaudePeerSession {
328 pid: 0,
329 session_id: id.to_string(),
330 cwd: None,
331 name: id.to_string(),
332 socket_path: Default::default(),
333 status: None,
334 updated_at_ms: None,
335 version: None,
336 tmux: None,
337 })),
338 "hook" => Door::Hook,
339 "stored" => Door::Stored,
340 other => panic!("no test door {other}"),
341 };
342 Some(LiveSession {
343 address: MailAddress::new("test", *harness, *id).ok()?,
344 name: format!("{id}@test"),
345 status: "unknown".into(),
346 pid: None,
347 cwd: None,
348 tmux: None,
349 door,
350 })
351 })
352 .collect(),
353 }
354 }
355}
356
357pub fn has_message_tools(pid: u32) -> bool {
361 let Ok(output) = std::process::Command::new("ps")
362 .args(["-A", "-o", "ppid=,command="])
363 .output()
364 else {
365 return false;
366 };
367 String::from_utf8_lossy(&output.stdout).lines().any(|line| {
368 let line = line.trim_start();
369 let Some((ppid, command)) = line.split_once(' ') else {
370 return false;
371 };
372 ppid.parse::<u32>() == Ok(pid) && command.trim_end().ends_with(" message mcp")
373 })
374}
375
376pub fn codex_name(session_id: &str) -> String {
378 format!("codex-{}", session_id.chars().take(8).collect::<String>())
379}
380
381#[derive(Debug, Clone, PartialEq, Eq)]
383pub struct Caller {
384 pub address: MailAddress,
386 pub name: String,
388}
389
390pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
392 replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
393 shell tool.";
394
395pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
400 let machine = local_machine_name();
401 let registry = read_registry(®istry_dir(homes));
402 let hosted = crate::runtime_mail::controlled_runtimes();
403 for &pid in pids {
404 if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
405 let short: String = record.source.session_id.chars().take(8).collect();
406 return Ok(Caller {
407 address: MailAddress::new(
408 &machine,
409 &record.source.harness,
410 &record.source.session_id,
411 )
412 .map_err(|error| error.to_string())?,
413 name: format!("{}-{short}@{machine}", record.source.harness),
414 });
415 }
416 if let Some(session) = registry.iter().find(|session| session.pid == pid) {
417 if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
418 if !claimed.is_empty() && claimed != session.session_id {
419 return Err(CALLER_UNRESOLVED.into());
420 }
421 }
422 return Ok(Caller {
423 address: MailAddress::new(&machine, "claude-code", &session.session_id)
424 .map_err(|error| error.to_string())?,
425 name: format!("{}@{machine}", session.name),
426 });
427 }
428 if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
429 return Ok(Caller {
430 address: MailAddress::new(&machine, "codex", &session_id)
431 .map_err(|error| error.to_string())?,
432 name: format!("{}@{machine}", codex_name(&session_id)),
433 });
434 }
435 }
436 Err(CALLER_UNRESOLVED.into())
437}
438
439pub fn process_ancestry() -> Vec<u32> {
441 let parents = parent_pids();
442 let mut chain = vec![std::process::id()];
443 let mut current = std::process::id();
444 while let Some(&parent) = parents.get(¤t) {
445 if parent <= 1 || chain.contains(&parent) {
446 break;
447 }
448 chain.push(parent);
449 current = parent;
450 }
451 chain
452}
453
454#[cfg(not(windows))]
456fn parent_pids() -> std::collections::HashMap<u32, u32> {
457 let Ok(output) = std::process::Command::new("ps")
458 .args(["-axo", "pid=,ppid="])
459 .output()
460 else {
461 return Default::default();
462 };
463 String::from_utf8_lossy(&output.stdout)
464 .lines()
465 .filter_map(|line| {
466 let mut fields = line.split_whitespace();
467 Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
468 })
469 .collect()
470}
471
472#[cfg(windows)]
474fn parent_pids() -> std::collections::HashMap<u32, u32> {
475 use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
476 use windows_sys::Win32::System::Diagnostics::ToolHelp::{
477 CreateToolhelp32Snapshot, Process32FirstW, Process32NextW, PROCESSENTRY32W,
478 TH32CS_SNAPPROCESS,
479 };
480
481 let mut parents = std::collections::HashMap::new();
482 let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPPROCESS, 0) };
483 if snapshot == INVALID_HANDLE_VALUE {
484 return parents;
485 }
486 let mut entry: PROCESSENTRY32W = unsafe { std::mem::zeroed() };
487 entry.dwSize = std::mem::size_of::<PROCESSENTRY32W>() as u32;
488 let mut has_entry = unsafe { Process32FirstW(snapshot, &mut entry) } != 0;
489 while has_entry {
490 parents.insert(entry.th32ProcessID, entry.th32ParentProcessID);
491 has_entry = unsafe { Process32NextW(snapshot, &mut entry) } != 0;
492 }
493 unsafe {
494 CloseHandle(snapshot);
495 }
496 parents
497}
498
499pub fn codex_hooks_path() -> std::path::PathBuf {
501 std::env::var_os("CODEX_HOME")
502 .map(std::path::PathBuf::from)
503 .or_else(|| {
504 supercode_interchange::user_home()
505 .map(std::path::PathBuf::into_os_string)
506 .map(|home| std::path::PathBuf::from(home).join(".codex"))
507 })
508 .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
509 .join("hooks.json")
510}
511
512pub fn codex_user_hook_installed() -> bool {
515 std::fs::read_to_string(codex_hooks_path())
516 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
517}
518
519pub fn codex_project_hook_installed(cwd: &Path) -> bool {
522 cwd.ancestors().any(|directory| {
523 std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
524 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
525 })
526}
527
528#[derive(Debug, Clone, PartialEq, Eq)]
530pub enum Delivered {
531 Steered,
533 Started,
535 Native {
538 busy: bool,
540 },
541 Hooked,
543 Queued,
545 Stored,
547 Operator,
549}
550
551#[derive(Debug, Clone, PartialEq, Eq)]
553pub enum Refused {
554 CannotQueueNative,
557 TooLong(usize),
559}
560
561pub const MAX_RELAYED_BYTES: usize = 100_000;
564
565pub async fn deliver(
569 envelope: &Envelope,
570 to: &MailAddress,
571 door: &Door,
572 wake: bool,
573 notify_when_idle: bool,
574) -> Result<Result<Delivered, Refused>, String> {
575 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
576 let mut final_reply = false;
580 let delivered = match door {
581 Door::Runtime(record) => {
582 let mut sent = envelope.clone();
583 let answers = matches!(envelope.reply_via, ReplyVia::Command);
584 if answers {
585 sent.reply_via = ReplyVia::FinalMessage {
586 destination: envelope.from_name.clone(),
587 };
588 }
589 match deliver_to_runtime(record, sent.render(), wake).await? {
590 RuntimeDelivery::Steered => {
591 mailbox
592 .deliver_read(&sent)
593 .map_err(|error| error.to_string())?;
594 final_reply = answers;
595 Delivered::Steered
596 }
597 RuntimeDelivery::Started => {
598 mailbox
599 .deliver_read(&sent)
600 .map_err(|error| error.to_string())?;
601 final_reply = answers;
602 Delivered::Started
603 }
604 RuntimeDelivery::NotWoken => {
605 mailbox
606 .deliver(envelope)
607 .map_err(|error| error.to_string())?;
608 Delivered::Queued
609 }
610 }
611 }
612 Door::Native(session) => {
613 let busy = session.status != Some(ClaudePeerStatus::Idle);
614 if !wake && !busy {
615 return Ok(Err(Refused::CannotQueueNative));
616 }
617 let mut envelope = envelope.clone();
621 if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
622 envelope.reply_via = ReplyVia::Tool;
623 }
624 let envelope = &envelope;
625 let text = envelope.render();
626 if text.len() > MAX_RELAYED_BYTES {
627 return Ok(Err(Refused::TooLong(text.len())));
628 }
629 match send_through_relay(
630 &envelope.from,
631 &envelope.from_name,
632 &session.name,
633 text,
634 &envelope.id,
635 )
636 .await
637 {
638 RelayReceipt::Delivered { .. } => {
639 mailbox
640 .deliver_read(envelope)
641 .map_err(|error| error.to_string())?;
642 Delivered::Native { busy }
643 }
644 RelayReceipt::Failed { detail } => return Err(detail),
645 }
646 }
647 Door::Hook => {
648 mailbox
649 .deliver(envelope)
650 .map_err(|error| error.to_string())?;
651 Delivered::Hooked
652 }
653 Door::Stored => {
654 mailbox
655 .deliver(envelope)
656 .map_err(|error| error.to_string())?;
657 Delivered::Stored
658 }
659 Door::Operator => {
660 mailbox
661 .deliver(envelope)
662 .map_err(|error| error.to_string())?;
663 Delivered::Operator
664 }
665 };
666 let notice = notify_when_idle && !matches!(door, Door::Operator);
667 if notice || final_reply {
668 let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
669 subscription.notice = notice;
670 subscription.final_reply = final_reply;
671 mailbox
672 .subscribe_idle(&subscription)
673 .map_err(|error| error.to_string())?;
674 if let Ok(program) = supercode_program() {
676 crate::claude_relay::ensure_machine_daemon(&program)
677 .await
678 .ok();
679 }
680 }
681 Ok(Ok(delivered))
682}
683
684#[derive(Debug, Clone, Copy, PartialEq, Eq)]
686pub enum UserTurn {
687 Steered,
689 Started,
691 Typed,
693 Waiting,
697}
698
699impl UserTurn {
700 pub const fn as_str(self) -> &'static str {
702 match self {
703 Self::Steered => "steered",
704 Self::Started => "started",
705 Self::Typed => "typed",
706 Self::Waiting => "waiting",
707 }
708 }
709}
710
711pub fn daemon_pane(session: &LiveSession) -> Option<String> {
713 let name = session.tmux.as_deref()?.split(':').next()?;
714 let rest = name.strip_prefix("p_")?;
715 (!rest.is_empty()
716 && rest
717 .chars()
718 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
719 .then(|| name.to_string())
720}
721
722pub async fn deliver_user_turn(
728 homes: &HarnessHomes,
729 envelope: &Envelope,
730 to: &MailAddress,
731) -> Result<UserTurn, String> {
732 let session = LiveSessions::read(homes)
733 .sessions
734 .into_iter()
735 .find(|session| &session.address == to)
736 .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
737 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
738 if let Door::Runtime(record) = &session.door {
739 let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
740 mailbox
741 .deliver_read(envelope)
742 .map_err(|error| error.to_string())?;
743 return Ok(match delivered {
744 RuntimeDelivery::Steered => UserTurn::Steered,
745 _ => UserTurn::Started,
746 });
747 }
748 let Some(pane) = daemon_pane(&session) else {
749 let name = session.name.split('@').next().unwrap_or(&session.name);
750 return Err(format!(
751 "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
752 pane with `supercode open {name}`; nothing was sent."
753 ));
754 };
755 mailbox
756 .deliver(envelope)
757 .map_err(|error| error.to_string())?;
758 let typed = type_user_turns(&mailbox, &pane).await;
759 Ok(if typed.contains(&envelope.id) {
760 UserTurn::Typed
761 } else {
762 UserTurn::Waiting
763 })
764}
765
766pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
770 let mut typed = Vec::new();
771 for stored in mailbox.user_turns().unwrap_or_default() {
772 match submit_when_composer_empty(pane, &stored.envelope.body).await {
773 Ok(true) => {
774 mailbox.mark_read(&stored).ok();
775 typed.push(stored.envelope.id.clone());
776 }
777 Ok(false) => break,
778 Err(error) => {
779 eprintln!(
780 "supercode: the user's turn {} for {} waits: {error}",
781 stored.envelope.id,
782 mailbox.address()
783 );
784 break;
785 }
786 }
787 }
788 typed
789}
790
791async fn submit_when_composer_empty(pane: &str, text: &str) -> Result<bool, String> {
794 let entry = crate::teams_entry().map_err(|error| error.to_string())?;
795 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
796 .ok()
797 .filter(|value| !value.trim().is_empty())
798 .unwrap_or_else(|| "node".into());
799 let output = tokio::process::Command::new(node)
800 .arg(entry)
801 .args(["input", pane, text, "--when-composer-empty"])
802 .stdin(std::process::Stdio::null())
803 .output()
804 .await
805 .map_err(|error| error.to_string())?;
806 if !output.status.success() {
807 return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
808 &output.stderr,
809 )));
810 }
811 let answer: serde_json::Value =
812 serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
813 Ok(answer["delivered"] == true)
814}