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 if matching.is_empty() {
280 if let Some(worker) = self
283 .sessions
284 .iter()
285 .filter(|session| orchestrator_of(session).as_deref() == Some(wanted.as_str()))
286 .max_by_key(|session| last_active(session))
287 {
288 return Ok(worker);
289 }
290 }
291 let hint = if matching.len() > 1 {
292 format!(
293 " {} sessions are named {wanted}; use its address.",
294 matching.len()
295 )
296 } else {
297 let near: Vec<String> = self
298 .sessions
299 .iter()
300 .filter(|session| {
301 let name = short(session);
302 name.contains(&wanted)
303 || wanted.contains(&name)
304 || name
305 .chars()
306 .zip(wanted.chars())
307 .take_while(|(a, b)| a == b)
308 .count()
309 >= 4
310 })
311 .take(3)
312 .map(|session| {
313 format!(
314 "{} ({}, {})",
315 session.name, session.address.harness, session.status
316 )
317 })
318 .collect();
319 if near.is_empty() {
320 String::new()
321 } else {
322 format!(" Did you mean: {}?", near.join(", "))
323 }
324 };
325 Err(Unresolved::Unknown(format!(
326 "No session named \"{wanted}\" is reachable.{hint} Run supercode message list. Nothing \
327 was sent."
328 )))
329 }
330
331 #[cfg(test)]
333 pub(crate) fn with(entries: &[(&str, &str, &'static str)]) -> Self {
334 Self {
335 sessions: entries
336 .iter()
337 .filter_map(|(harness, id, door)| {
338 let door = match *door {
339 "native" => Door::Native(Box::new(ClaudePeerSession {
340 pid: 0,
341 session_id: id.to_string(),
342 cwd: None,
343 name: id.to_string(),
344 socket_path: Default::default(),
345 status: None,
346 updated_at_ms: None,
347 version: None,
348 tmux: None,
349 })),
350 "hook" => Door::Hook,
351 "stored" => Door::Stored,
352 other => panic!("no test door {other}"),
353 };
354 Some(LiveSession {
355 address: MailAddress::new("test", *harness, *id).ok()?,
356 name: format!("{id}@test"),
357 status: "unknown".into(),
358 pid: None,
359 cwd: None,
360 tmux: None,
361 door,
362 })
363 })
364 .collect(),
365 }
366 }
367}
368
369pub fn orchestrator_of(session: &LiveSession) -> Option<String> {
374 let cwd = session.cwd.as_ref()?;
375 if !cwd.join("orchestrator.lock").is_file() {
376 return None;
377 }
378 cwd.file_name()?.to_str().map(str::to_string)
379}
380
381fn last_active(session: &LiveSession) -> u64 {
382 match &session.door {
383 Door::Native(peer) => peer.updated_at_ms.unwrap_or(0),
384 _ => 0,
385 }
386}
387
388pub fn codex_name(session_id: &str) -> String {
390 format!("codex-{}", session_id.chars().take(8).collect::<String>())
391}
392
393#[derive(Debug, Clone, PartialEq, Eq)]
395pub struct Caller {
396 pub address: MailAddress,
398 pub name: String,
400}
401
402pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
404 replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
405 shell tool.";
406
407pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
412 let machine = local_machine_name();
413 let registry = read_registry(®istry_dir(homes));
414 let hosted = crate::runtime_mail::controlled_runtimes();
415 for &pid in pids {
416 if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
417 let short: String = record.source.session_id.chars().take(8).collect();
418 return Ok(Caller {
419 address: MailAddress::new(
420 &machine,
421 &record.source.harness,
422 &record.source.session_id,
423 )
424 .map_err(|error| error.to_string())?,
425 name: format!("{}-{short}@{machine}", record.source.harness),
426 });
427 }
428 if let Some(session) = registry.iter().find(|session| session.pid == pid) {
429 if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
430 if !claimed.is_empty() && claimed != session.session_id {
431 return Err(CALLER_UNRESOLVED.into());
432 }
433 }
434 return Ok(Caller {
435 address: MailAddress::new(&machine, "claude-code", &session.session_id)
436 .map_err(|error| error.to_string())?,
437 name: format!("{}@{machine}", session.name),
438 });
439 }
440 if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
441 return Ok(Caller {
442 address: MailAddress::new(&machine, "codex", &session_id)
443 .map_err(|error| error.to_string())?,
444 name: format!("{}@{machine}", codex_name(&session_id)),
445 });
446 }
447 }
448 Err(CALLER_UNRESOLVED.into())
449}
450
451pub fn process_ancestry() -> Vec<u32> {
453 let Ok(output) = std::process::Command::new("ps")
454 .args(["-axo", "pid=,ppid="])
455 .output()
456 else {
457 return vec![std::process::id()];
458 };
459 let parents: std::collections::HashMap<u32, u32> = String::from_utf8_lossy(&output.stdout)
460 .lines()
461 .filter_map(|line| {
462 let mut fields = line.split_whitespace();
463 Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
464 })
465 .collect();
466 let mut chain = vec![std::process::id()];
467 let mut current = std::process::id();
468 while let Some(&parent) = parents.get(¤t) {
469 if parent <= 1 || chain.contains(&parent) {
470 break;
471 }
472 chain.push(parent);
473 current = parent;
474 }
475 chain
476}
477
478pub fn codex_hooks_path() -> std::path::PathBuf {
480 std::env::var_os("CODEX_HOME")
481 .map(std::path::PathBuf::from)
482 .or_else(|| {
483 std::env::var_os("HOME").map(|home| std::path::PathBuf::from(home).join(".codex"))
484 })
485 .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
486 .join("hooks.json")
487}
488
489pub fn codex_user_hook_installed() -> bool {
492 std::fs::read_to_string(codex_hooks_path())
493 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
494}
495
496pub fn codex_project_hook_installed(cwd: &Path) -> bool {
499 cwd.ancestors().any(|directory| {
500 std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
501 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
502 })
503}
504
505#[derive(Debug, Clone, PartialEq, Eq)]
507pub enum Delivered {
508 Steered,
510 Started,
512 Native {
515 busy: bool,
517 },
518 Hooked,
520 Queued,
522 Stored,
524 Operator,
526}
527
528#[derive(Debug, Clone, PartialEq, Eq)]
530pub enum Refused {
531 CannotQueueNative,
534 TooLong(usize),
536}
537
538pub const MAX_RELAYED_BYTES: usize = 100_000;
541
542pub async fn deliver(
546 envelope: &Envelope,
547 to: &MailAddress,
548 door: &Door,
549 wake: bool,
550 notify_when_idle: bool,
551) -> Result<Result<Delivered, Refused>, String> {
552 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
553 let mut final_reply = false;
557 let delivered = match door {
558 Door::Runtime(record) => {
559 let mut sent = envelope.clone();
560 let answers = matches!(envelope.reply_via, ReplyVia::Command);
561 if answers {
562 sent.reply_via = ReplyVia::FinalMessage {
563 destination: envelope.from_name.clone(),
564 };
565 }
566 match deliver_to_runtime(record, sent.render(), wake).await? {
567 RuntimeDelivery::Steered => {
568 mailbox
569 .deliver_read(&sent)
570 .map_err(|error| error.to_string())?;
571 final_reply = answers;
572 Delivered::Steered
573 }
574 RuntimeDelivery::Started => {
575 mailbox
576 .deliver_read(&sent)
577 .map_err(|error| error.to_string())?;
578 final_reply = answers;
579 Delivered::Started
580 }
581 RuntimeDelivery::NotWoken => {
582 mailbox
583 .deliver(envelope)
584 .map_err(|error| error.to_string())?;
585 Delivered::Queued
586 }
587 }
588 }
589 Door::Native(session) => {
590 let busy = session.status != Some(ClaudePeerStatus::Idle);
591 if !wake && !busy {
592 return Ok(Err(Refused::CannotQueueNative));
593 }
594 let text = envelope.render();
597 if text.len() > MAX_RELAYED_BYTES {
598 return Ok(Err(Refused::TooLong(text.len())));
599 }
600 match send_through_relay(
601 &envelope.from,
602 &envelope.from_name,
603 &session.name,
604 text,
605 &envelope.id,
606 )
607 .await
608 {
609 RelayReceipt::Delivered { .. } => {
610 mailbox
611 .deliver_read(envelope)
612 .map_err(|error| error.to_string())?;
613 Delivered::Native { busy }
614 }
615 RelayReceipt::Failed { detail } => return Err(detail),
616 }
617 }
618 Door::Hook => {
619 mailbox
620 .deliver(envelope)
621 .map_err(|error| error.to_string())?;
622 Delivered::Hooked
623 }
624 Door::Stored => {
625 mailbox
626 .deliver(envelope)
627 .map_err(|error| error.to_string())?;
628 Delivered::Stored
629 }
630 Door::Operator => {
631 mailbox
632 .deliver(envelope)
633 .map_err(|error| error.to_string())?;
634 Delivered::Operator
635 }
636 };
637 let notice = notify_when_idle && !matches!(door, Door::Operator);
638 if notice || final_reply {
639 let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
640 subscription.notice = notice;
641 subscription.final_reply = final_reply;
642 mailbox
643 .subscribe_idle(&subscription)
644 .map_err(|error| error.to_string())?;
645 if let Ok(program) = supercode_program() {
647 crate::claude_relay::ensure_machine_daemon(&program)
648 .await
649 .ok();
650 }
651 }
652 Ok(Ok(delivered))
653}
654
655#[derive(Debug, Clone, Copy, PartialEq, Eq)]
657pub enum UserTurn {
658 Steered,
660 Started,
662 Typed,
664 Waiting,
668}
669
670impl UserTurn {
671 pub const fn as_str(self) -> &'static str {
673 match self {
674 Self::Steered => "steered",
675 Self::Started => "started",
676 Self::Typed => "typed",
677 Self::Waiting => "waiting",
678 }
679 }
680}
681
682pub fn daemon_pane(session: &LiveSession) -> Option<String> {
684 let name = session.tmux.as_deref()?.split(':').next()?;
685 let rest = name.strip_prefix("p_")?;
686 (!rest.is_empty()
687 && rest
688 .chars()
689 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
690 .then(|| name.to_string())
691}
692
693pub async fn deliver_user_turn(
699 homes: &HarnessHomes,
700 envelope: &Envelope,
701 to: &MailAddress,
702) -> Result<UserTurn, String> {
703 let session = LiveSessions::read(homes)
704 .sessions
705 .into_iter()
706 .find(|session| &session.address == to)
707 .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
708 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
709 if let Door::Runtime(record) = &session.door {
710 let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
711 mailbox
712 .deliver_read(envelope)
713 .map_err(|error| error.to_string())?;
714 return Ok(match delivered {
715 RuntimeDelivery::Steered => UserTurn::Steered,
716 _ => UserTurn::Started,
717 });
718 }
719 let Some(pane) = daemon_pane(&session) else {
720 let name = session.name.split('@').next().unwrap_or(&session.name);
721 return Err(format!(
722 "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
723 pane with `supercode open {name}`; nothing was sent."
724 ));
725 };
726 mailbox
727 .deliver(envelope)
728 .map_err(|error| error.to_string())?;
729 let typed = type_user_turns(&mailbox, &pane).await;
730 Ok(if typed.contains(&envelope.id) {
731 UserTurn::Typed
732 } else {
733 UserTurn::Waiting
734 })
735}
736
737pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
741 let mut typed = Vec::new();
742 for stored in mailbox.user_turns().unwrap_or_default() {
743 match submit_when_composer_empty(pane, &stored.envelope.body).await {
744 Ok(true) => {
745 mailbox.mark_read(&stored).ok();
746 typed.push(stored.envelope.id.clone());
747 }
748 Ok(false) => break,
749 Err(error) => {
750 eprintln!(
751 "supercode: the user's turn {} for {} waits: {error}",
752 stored.envelope.id,
753 mailbox.address()
754 );
755 break;
756 }
757 }
758 }
759 typed
760}
761
762async fn submit_when_composer_empty(pane: &str, text: &str) -> Result<bool, String> {
765 let entry = crate::teams_entry().map_err(|error| error.to_string())?;
766 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
767 .ok()
768 .filter(|value| !value.trim().is_empty())
769 .unwrap_or_else(|| "node".into());
770 let output = tokio::process::Command::new(node)
771 .arg(entry)
772 .args(["input", pane, text, "--when-composer-empty"])
773 .stdin(std::process::Stdio::null())
774 .output()
775 .await
776 .map_err(|error| error.to_string())?;
777 if !output.status.success() {
778 return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
779 &output.stderr,
780 )));
781 }
782 let answer: serde_json::Value =
783 serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
784 Ok(answer["delivered"] == true)
785}