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 Ok(output) = std::process::Command::new("ps")
442 .args(["-axo", "pid=,ppid="])
443 .output()
444 else {
445 return vec![std::process::id()];
446 };
447 let parents: std::collections::HashMap<u32, u32> = String::from_utf8_lossy(&output.stdout)
448 .lines()
449 .filter_map(|line| {
450 let mut fields = line.split_whitespace();
451 Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
452 })
453 .collect();
454 let mut chain = vec![std::process::id()];
455 let mut current = std::process::id();
456 while let Some(&parent) = parents.get(¤t) {
457 if parent <= 1 || chain.contains(&parent) {
458 break;
459 }
460 chain.push(parent);
461 current = parent;
462 }
463 chain
464}
465
466pub fn codex_hooks_path() -> std::path::PathBuf {
468 std::env::var_os("CODEX_HOME")
469 .map(std::path::PathBuf::from)
470 .or_else(|| {
471 supercode_interchange::user_home()
472 .map(std::path::PathBuf::into_os_string)
473 .map(|home| std::path::PathBuf::from(home).join(".codex"))
474 })
475 .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
476 .join("hooks.json")
477}
478
479pub fn codex_user_hook_installed() -> bool {
482 std::fs::read_to_string(codex_hooks_path())
483 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
484}
485
486pub fn codex_project_hook_installed(cwd: &Path) -> bool {
489 cwd.ancestors().any(|directory| {
490 std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
491 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
492 })
493}
494
495#[derive(Debug, Clone, PartialEq, Eq)]
497pub enum Delivered {
498 Steered,
500 Started,
502 Native {
505 busy: bool,
507 },
508 Hooked,
510 Queued,
512 Stored,
514 Operator,
516}
517
518#[derive(Debug, Clone, PartialEq, Eq)]
520pub enum Refused {
521 CannotQueueNative,
524 TooLong(usize),
526}
527
528pub const MAX_RELAYED_BYTES: usize = 100_000;
531
532pub async fn deliver(
536 envelope: &Envelope,
537 to: &MailAddress,
538 door: &Door,
539 wake: bool,
540 notify_when_idle: bool,
541) -> Result<Result<Delivered, Refused>, String> {
542 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
543 let mut final_reply = false;
547 let delivered = match door {
548 Door::Runtime(record) => {
549 let mut sent = envelope.clone();
550 let answers = matches!(envelope.reply_via, ReplyVia::Command);
551 if answers {
552 sent.reply_via = ReplyVia::FinalMessage {
553 destination: envelope.from_name.clone(),
554 };
555 }
556 match deliver_to_runtime(record, sent.render(), wake).await? {
557 RuntimeDelivery::Steered => {
558 mailbox
559 .deliver_read(&sent)
560 .map_err(|error| error.to_string())?;
561 final_reply = answers;
562 Delivered::Steered
563 }
564 RuntimeDelivery::Started => {
565 mailbox
566 .deliver_read(&sent)
567 .map_err(|error| error.to_string())?;
568 final_reply = answers;
569 Delivered::Started
570 }
571 RuntimeDelivery::NotWoken => {
572 mailbox
573 .deliver(envelope)
574 .map_err(|error| error.to_string())?;
575 Delivered::Queued
576 }
577 }
578 }
579 Door::Native(session) => {
580 let busy = session.status != Some(ClaudePeerStatus::Idle);
581 if !wake && !busy {
582 return Ok(Err(Refused::CannotQueueNative));
583 }
584 let mut envelope = envelope.clone();
588 if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
589 envelope.reply_via = ReplyVia::Tool;
590 }
591 let envelope = &envelope;
592 let text = envelope.render();
593 if text.len() > MAX_RELAYED_BYTES {
594 return Ok(Err(Refused::TooLong(text.len())));
595 }
596 match send_through_relay(
597 &envelope.from,
598 &envelope.from_name,
599 &session.name,
600 text,
601 &envelope.id,
602 )
603 .await
604 {
605 RelayReceipt::Delivered { .. } => {
606 mailbox
607 .deliver_read(envelope)
608 .map_err(|error| error.to_string())?;
609 Delivered::Native { busy }
610 }
611 RelayReceipt::Failed { detail } => return Err(detail),
612 }
613 }
614 Door::Hook => {
615 mailbox
616 .deliver(envelope)
617 .map_err(|error| error.to_string())?;
618 Delivered::Hooked
619 }
620 Door::Stored => {
621 mailbox
622 .deliver(envelope)
623 .map_err(|error| error.to_string())?;
624 Delivered::Stored
625 }
626 Door::Operator => {
627 mailbox
628 .deliver(envelope)
629 .map_err(|error| error.to_string())?;
630 Delivered::Operator
631 }
632 };
633 let notice = notify_when_idle && !matches!(door, Door::Operator);
634 if notice || final_reply {
635 let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
636 subscription.notice = notice;
637 subscription.final_reply = final_reply;
638 mailbox
639 .subscribe_idle(&subscription)
640 .map_err(|error| error.to_string())?;
641 if let Ok(program) = supercode_program() {
643 crate::claude_relay::ensure_machine_daemon(&program)
644 .await
645 .ok();
646 }
647 }
648 Ok(Ok(delivered))
649}
650
651#[derive(Debug, Clone, Copy, PartialEq, Eq)]
653pub enum UserTurn {
654 Steered,
656 Started,
658 Typed,
660 Waiting,
664}
665
666impl UserTurn {
667 pub const fn as_str(self) -> &'static str {
669 match self {
670 Self::Steered => "steered",
671 Self::Started => "started",
672 Self::Typed => "typed",
673 Self::Waiting => "waiting",
674 }
675 }
676}
677
678pub fn daemon_pane(session: &LiveSession) -> Option<String> {
680 let name = session.tmux.as_deref()?.split(':').next()?;
681 let rest = name.strip_prefix("p_")?;
682 (!rest.is_empty()
683 && rest
684 .chars()
685 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
686 .then(|| name.to_string())
687}
688
689pub async fn deliver_user_turn(
695 homes: &HarnessHomes,
696 envelope: &Envelope,
697 to: &MailAddress,
698) -> Result<UserTurn, String> {
699 let session = LiveSessions::read(homes)
700 .sessions
701 .into_iter()
702 .find(|session| &session.address == to)
703 .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
704 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
705 if let Door::Runtime(record) = &session.door {
706 let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
707 mailbox
708 .deliver_read(envelope)
709 .map_err(|error| error.to_string())?;
710 return Ok(match delivered {
711 RuntimeDelivery::Steered => UserTurn::Steered,
712 _ => UserTurn::Started,
713 });
714 }
715 let Some(pane) = daemon_pane(&session) else {
716 let name = session.name.split('@').next().unwrap_or(&session.name);
717 return Err(format!(
718 "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
719 pane with `supercode open {name}`; nothing was sent."
720 ));
721 };
722 mailbox
723 .deliver(envelope)
724 .map_err(|error| error.to_string())?;
725 let typed = type_user_turns(&mailbox, &pane).await;
726 Ok(if typed.contains(&envelope.id) {
727 UserTurn::Typed
728 } else {
729 UserTurn::Waiting
730 })
731}
732
733pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
737 let mut typed = Vec::new();
738 for stored in mailbox.user_turns().unwrap_or_default() {
739 match submit_when_composer_empty(pane, &stored.envelope.body).await {
740 Ok(true) => {
741 mailbox.mark_read(&stored).ok();
742 typed.push(stored.envelope.id.clone());
743 }
744 Ok(false) => break,
745 Err(error) => {
746 eprintln!(
747 "supercode: the user's turn {} for {} waits: {error}",
748 stored.envelope.id,
749 mailbox.address()
750 );
751 break;
752 }
753 }
754 }
755 typed
756}
757
758async fn submit_when_composer_empty(pane: &str, text: &str) -> Result<bool, String> {
761 let entry = crate::teams_entry().map_err(|error| error.to_string())?;
762 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
763 .ok()
764 .filter(|value| !value.trim().is_empty())
765 .unwrap_or_else(|| "node".into());
766 let output = tokio::process::Command::new(node)
767 .arg(entry)
768 .args(["input", pane, text, "--when-composer-empty"])
769 .stdin(std::process::Stdio::null())
770 .output()
771 .await
772 .map_err(|error| error.to_string())?;
773 if !output.status.success() {
774 return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
775 &output.stderr,
776 )));
777 }
778 let answer: serde_json::Value =
779 serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
780 Ok(answer["delivered"] == true)
781}