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 {
49 pane: Option<String>,
51 idle: bool,
53 },
54 Stored,
56 Operator,
58}
59
60impl Door {
61 pub fn name(&self) -> &'static str {
63 match self {
64 Self::Runtime(_) => "runtime",
65 Self::Native(_) => "native",
66 Self::Hook { .. } => "hook",
67 Self::Stored => "stored",
68 Self::Operator => "operator",
69 }
70 }
71}
72
73#[derive(Debug, Clone, PartialEq, Eq)]
75pub enum NoDoor {
76 OtherMachine(String),
78 NotRunning,
80}
81
82pub fn door_for(homes: &HarnessHomes, to: &MailAddress) -> Result<Door, NoDoor> {
84 if to.machine != local_machine_name() {
85 return Err(NoDoor::OtherMachine(to.machine.clone()));
86 }
87 if to.harness == "operator" || to.harness == "board" {
88 return Ok(Door::Operator);
89 }
90 LiveSessions::read(homes)
91 .sessions
92 .into_iter()
93 .find(|session| &session.address == to)
94 .map(|session| session.door)
95 .ok_or(NoDoor::NotRunning)
96}
97
98#[derive(Debug, Clone, PartialEq, Eq)]
100pub struct LiveSession {
101 pub address: MailAddress,
103 pub name: String,
105 pub status: String,
107 pub door: Door,
109 pub pid: Option<u32>,
112 pub cwd: Option<std::path::PathBuf>,
114 pub tmux: Option<String>,
116 pub transcript: Option<std::path::PathBuf>,
118}
119
120impl LiveSession {
121 pub fn last_message_at_ms(&self, homes: &HarnessHomes) -> Option<u64> {
126 let harness = self.address.harness.as_str();
127 let path = match &self.transcript {
128 Some(path) => path.clone(),
129 None if harness == "claude-code" => {
130 let file = format!("{}.jsonl", self.address.session_id);
131 std::fs::read_dir(&homes.claude_code)
132 .ok()?
133 .flatten()
134 .map(|project| project.path().join(&file))
135 .find(|path| path.is_file())?
136 }
137 None => return None,
138 };
139 last_message_at_ms(&path, harness)
140 }
141}
142
143impl LiveSession {
144 pub fn pending_request(&self, homes: &HarnessHomes) -> Option<serde_json::Value> {
148 if self.status != "waiting" {
149 return None;
150 }
151 if self.address.harness == "codex" {
153 let prompt = daemon_pane(self).and_then(|pane| pane_prompt(&pane))?;
154 return Some(serde_json::json!({"prompt": prompt, "source": "screen"}));
155 }
156 if self.address.harness != "claude-code" {
157 return None;
158 }
159 let file = format!("{}.jsonl", self.address.session_id);
160 let from_transcript = std::fs::read_dir(&homes.claude_code)
161 .ok()
162 .and_then(|projects| {
163 projects
164 .flatten()
165 .map(|project| project.path().join(&file))
166 .find(|path| path.is_file())
167 })
168 .and_then(|path| pending_request(&path));
169 from_transcript.or_else(|| {
172 let prompt = daemon_pane(self).and_then(|pane| pane_prompt(&pane))?;
173 Some(serde_json::json!({"prompt": prompt, "source": "screen"}))
174 })
175 }
176}
177
178pub fn pending_request(path: &Path) -> Option<serde_json::Value> {
180 use std::io::{Read, Seek, SeekFrom};
181 let mut file = std::fs::File::open(path).ok()?;
182 let length = file.metadata().ok()?.len();
183 let start = length.saturating_sub(512 * 1024);
184 file.seek(SeekFrom::Start(start)).ok()?;
185 let mut bytes = Vec::new();
186 file.read_to_end(&mut bytes).ok()?;
187 let text = String::from_utf8_lossy(&bytes);
188 let mut answered = std::collections::HashSet::new();
189 for line in text.lines().rev() {
190 let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
191 continue;
192 };
193 if record["isSidechain"] == true {
194 continue;
195 }
196 let Some(content) = record
197 .pointer("/message/content")
198 .and_then(serde_json::Value::as_array)
199 else {
200 continue;
201 };
202 match record["type"].as_str() {
203 Some("user") => {
204 for item in content.iter().filter(|item| item["type"] == "tool_result") {
205 if let Some(id) = item["tool_use_id"].as_str() {
206 answered.insert(id.to_string());
207 }
208 }
209 }
210 Some("assistant") => {
211 let Some(call) = content.iter().rev().find(|item| item["type"] == "tool_use")
212 else {
213 continue;
214 };
215 if call["id"].as_str().is_some_and(|id| answered.contains(id)) {
216 return None;
217 }
218 let tool = call["name"].as_str().unwrap_or_default();
219 if tool == "AskUserQuestion" {
220 let questions = call["input"]["questions"]
221 .as_array()
222 .map(|questions| {
223 questions
224 .iter()
225 .map(|question| {
226 serde_json::json!({
227 "question": question["question"],
228 "header": question["header"],
229 "options": question["options"].as_array().map(|options| options.iter().map(|option| option["label"].clone()).collect::<Vec<_>>()).unwrap_or_default(),
230 })
231 })
232 .collect::<Vec<_>>()
233 })
234 .unwrap_or_default();
235 return Some(serde_json::json!({"tool": tool, "questions": questions}));
236 }
237 let input = call["input"].to_string();
238 let input: String = input.chars().take(500).collect();
239 return Some(serde_json::json!({"tool": tool, "input": input}));
240 }
241 _ => {}
242 }
243 }
244 None
245}
246
247fn last_message_at_ms(path: &Path, harness: &str) -> Option<u64> {
249 use std::io::{Read, Seek, SeekFrom};
250 for window in [256 * 1024_u64, 8 * 1024 * 1024] {
252 let mut file = std::fs::File::open(path).ok()?;
253 let length = file.metadata().ok()?.len();
254 let start = length.saturating_sub(window);
255 file.seek(SeekFrom::Start(start)).ok()?;
256 let mut bytes = Vec::new();
257 file.read_to_end(&mut bytes).ok()?;
258 let text = String::from_utf8_lossy(&bytes);
259 let mut lines = text.lines().rev().collect::<Vec<_>>();
260 if start > 0 {
261 lines.pop(); }
263 for line in lines {
264 let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
265 continue;
266 };
267 let message = match harness {
268 "codex" => record["type"] == "response_item",
269 _ => {
270 matches!(record["type"].as_str(), Some("user" | "assistant"))
271 && record["isSidechain"] != true
272 }
273 };
274 if !message {
275 continue;
276 }
277 if let Some(at) = record["timestamp"]
278 .as_str()
279 .and_then(supercode_interchange::sidecar::rfc3339_to_ms)
280 .and_then(|at| u64::try_from(at).ok())
281 {
282 return Some(at);
283 }
284 }
285 if start == 0 {
286 return None;
287 }
288 }
289 None
290}
291
292#[derive(Debug, Default)]
296pub struct LiveSessions {
297 sessions: Vec<LiveSession>,
298}
299
300#[derive(Debug, Clone, PartialEq, Eq)]
302pub enum Unresolved {
303 Stale(String),
305 Unknown(String),
307}
308
309impl LiveSessions {
310 pub fn read(homes: &HarnessHomes) -> Self {
314 let machine = local_machine_name();
315 let registry = read_registry(®istry_dir(homes));
316 let registered = |record: &crate::live_runtime::LiveRuntimeRecord| {
319 registry.iter().find(|session| {
320 session.session_id == record.source.session_id
321 || session.session_id == record.runtime_session_id
322 })
323 };
324 let records = crate::runtime_mail::controlled_runtimes();
325 let turns = crate::runtime_mail::runtime_turn_states(&records);
327 let mut sessions: Vec<LiveSession> = records
328 .into_iter()
329 .zip(turns)
330 .filter_map(|(record, turn)| {
331 let registered = registered(&record);
332 if registered.is_some_and(|session| session.name.starts_with(RELAY_NAME_PREFIX)) {
333 return None;
334 }
335 let address =
336 MailAddress::new(&machine, &record.source.harness, &record.source.session_id)
337 .ok()?;
338 let short: String = record.source.session_id.chars().take(8).collect();
339 let name = match registered {
340 Some(session) if !session.name.is_empty() => session.name.clone(),
341 _ => format!("{}-{short}", record.source.harness),
342 };
343 Some(LiveSession {
344 name: format!("{name}@{machine}"),
345 address,
346 status: match turn {
348 Some(crate::frontend::FrontendTurnState::Busy) => "busy".into(),
349 Some(crate::frontend::FrontendTurnState::Idle) => "idle".into(),
350 None => "hosted".into(),
351 },
352 pid: Some(record.pid),
353 cwd: Some(record.source.workspace.clone()),
354 tmux: None,
355 transcript: None,
356 door: Door::Runtime(Box::new(record)),
357 })
358 })
359 .collect();
360 let controlled = |address: &MailAddress, sessions: &[LiveSession]| {
361 sessions.iter().any(|session| &session.address == address)
362 };
363 for session in registry {
364 if session.name.starts_with(RELAY_NAME_PREFIX) {
365 continue;
366 }
367 let Ok(address) = MailAddress::new(&machine, "claude-code", &session.session_id) else {
368 continue;
369 };
370 if controlled(&address, &sessions) {
371 continue;
372 }
373 sessions.push(LiveSession {
374 address,
375 name: format!("{}@{machine}", session.name),
376 status: session
377 .status
378 .as_ref()
379 .map(|status| status.as_str().to_string())
380 .unwrap_or_else(|| "unknown".into()),
381 pid: Some(session.pid),
382 cwd: session.cwd.clone(),
383 tmux: session.tmux.clone(),
384 transcript: None,
385 door: Door::Native(Box::new(session)),
386 });
387 }
388 let user_hook = codex_user_hook_installed();
389 let rollouts = crate::codex_peer::live_rollouts(&homes.codex);
390 let panes = crate::codex_peer::session_panes();
391 let live_panes = (!panes.is_empty()).then(live_daemon_panes).flatten();
394 for (path, status) in rollouts {
395 let Some((session_id, None)) = crate::codex_peer::rollout_session(&path) else {
396 continue;
397 };
398 let Ok(address) = MailAddress::new(&machine, "codex", &session_id) else {
399 continue;
400 };
401 if controlled(&address, &sessions) {
402 continue;
403 }
404 let cwd = crate::codex_peer::rollout_cwd(&path);
405 let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
406 let pane = panes
407 .get(&session_id)
408 .filter(|pane| {
409 live_panes
410 .as_ref()
411 .is_none_or(|live| live.contains_key(*pane))
412 })
413 .cloned();
414 let status = match status {
417 crate::codex_peer::CodexPeerStatus::Busy
418 | crate::codex_peer::CodexPeerStatus::Running
419 if pane.as_ref().is_some_and(|p| {
420 live_panes.as_ref().and_then(|live| live.get(p)) == Some(&true)
421 }) =>
422 {
423 "idle".to_string()
424 }
425
426 crate::codex_peer::CodexPeerStatus::Busy
427 | crate::codex_peer::CodexPeerStatus::Running
428 if pane.as_deref().and_then(pane_prompt).is_some() =>
429 {
430 "waiting".to_string()
431 }
432 status => status.as_str().to_string(),
433 };
434 let door = if hooked {
435 Door::Hook {
436 pane: pane.clone(),
437 idle: status == "idle",
438 }
439 } else {
440 Door::Stored
441 };
442 sessions.push(LiveSession {
443 name: format!("{}@{machine}", codex_name(&session_id)),
444 address,
445 status,
446 pid: None,
447 cwd,
448 tmux: pane,
450 transcript: Some(path),
451 door,
452 });
453 }
454 for (session_id, pane) in &panes {
457 let Ok(address) = MailAddress::new(&machine, "codex", session_id) else {
458 continue;
459 };
460 if controlled(&address, &sessions)
461 || !live_panes
462 .as_ref()
463 .is_some_and(|live| live.contains_key(pane))
464 {
465 continue;
466 }
467 let Some(path) = crate::codex_peer::rollout_of_session(&homes.codex, session_id) else {
468 continue;
469 };
470 let cwd = crate::codex_peer::rollout_cwd(&path);
471 let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
472 sessions.push(LiveSession {
473 name: format!("{}@{machine}", codex_name(session_id)),
474 address,
475 status: "idle".to_string(),
476 pid: None,
477 cwd,
478 tmux: Some(pane.clone()),
479 transcript: Some(path),
480 door: if hooked {
481 Door::Hook {
482 pane: Some(pane.clone()),
483 idle: true,
484 }
485 } else {
486 Door::Stored
487 },
488 });
489 }
490 Self { sessions }
491 }
492
493 pub fn all(&self) -> &[LiveSession] {
495 &self.sessions
496 }
497
498 pub fn door(&self, harness: &str, session_id: &str) -> Option<&'static str> {
501 self.sessions
502 .iter()
503 .find(|session| {
504 session.address.harness == harness && session.address.session_id == session_id
505 })
506 .map(|session| session.door.name())
507 }
508
509 pub fn resolve(&self, to: &str) -> Result<&LiveSession, Unresolved> {
512 let machine = local_machine_name();
513 if let Ok(address) = MailAddress::parse(to) {
514 return self
515 .sessions
516 .iter()
517 .find(|session| session.address == address)
518 .ok_or_else(|| {
519 Unresolved::Stale(format!(
520 "{to} is no longer running. Nothing was sent. Run supercode message list \
521 for the live sessions."
522 ))
523 });
524 }
525 let wanted = match to.split_once('@') {
526 Some((name, at)) if at == machine => name.to_string(),
527 Some((_, at)) => {
528 return Err(Unresolved::Unknown(format!(
529 "Not sent: {to} is on machine {at}, not this one. Nothing was sent."
530 )))
531 }
532 None => to.to_string(),
533 };
534 let short = |session: &LiveSession| {
535 session
536 .name
537 .split('@')
538 .next()
539 .unwrap_or_default()
540 .to_string()
541 };
542 let matching: Vec<&LiveSession> = self
543 .sessions
544 .iter()
545 .filter(|session| short(session) == wanted)
546 .collect();
547 if let [only] = matching.as_slice() {
548 return Ok(only);
549 }
550 let hint = if matching.len() > 1 {
551 format!(
552 " {} sessions are named {wanted}; use its address.",
553 matching.len()
554 )
555 } else {
556 let near: Vec<String> = self
557 .sessions
558 .iter()
559 .filter(|session| {
560 let name = short(session);
561 name.contains(&wanted)
562 || wanted.contains(&name)
563 || name
564 .chars()
565 .zip(wanted.chars())
566 .take_while(|(a, b)| a == b)
567 .count()
568 >= 4
569 })
570 .take(3)
571 .map(|session| {
572 format!(
573 "{} ({}, {})",
574 session.name, session.address.harness, session.status
575 )
576 })
577 .collect();
578 if near.is_empty() {
579 String::new()
580 } else {
581 format!(" Did you mean: {}?", near.join(", "))
582 }
583 };
584 Err(Unresolved::Unknown(format!(
585 "No session named \"{wanted}\" is reachable.{hint} Run supercode message list. Nothing \
586 was sent."
587 )))
588 }
589}
590
591pub fn has_message_tools(pid: u32) -> bool {
595 let Ok(output) = std::process::Command::new("ps")
596 .args(["-A", "-o", "ppid=,command="])
597 .output()
598 else {
599 return false;
600 };
601 String::from_utf8_lossy(&output.stdout).lines().any(|line| {
602 let line = line.trim_start();
603 let Some((ppid, command)) = line.split_once(' ') else {
604 return false;
605 };
606 ppid.parse::<u32>() == Ok(pid) && command.trim_end().ends_with(" message mcp")
607 })
608}
609
610pub fn codex_name(session_id: &str) -> String {
612 format!("codex-{}", session_id.chars().take(8).collect::<String>())
613}
614
615#[derive(Debug, Clone, PartialEq, Eq)]
617pub struct Caller {
618 pub address: MailAddress,
620 pub name: String,
622}
623
624pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
626 replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
627 shell tool.";
628
629pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
634 let machine = local_machine_name();
635 let registry = read_registry(®istry_dir(homes));
636 let hosted = crate::runtime_mail::controlled_runtimes();
637 for &pid in pids {
638 if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
639 let short: String = record.source.session_id.chars().take(8).collect();
640 return Ok(Caller {
641 address: MailAddress::new(
642 &machine,
643 &record.source.harness,
644 &record.source.session_id,
645 )
646 .map_err(|error| error.to_string())?,
647 name: format!("{}-{short}@{machine}", record.source.harness),
648 });
649 }
650 if let Some(session) = registry.iter().find(|session| session.pid == pid) {
651 if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
652 if !claimed.is_empty() && claimed != session.session_id {
653 return Err(format!(
654 "{CALLER_UNRESOLVED} (CLAUDE_CODE_SESSION_ID names {claimed}, but the Claude \
655 process {pid} above this command is session {})",
656 session.session_id
657 ));
658 }
659 }
660 return Ok(Caller {
661 address: MailAddress::new(&machine, "claude-code", &session.session_id)
662 .map_err(|error| error.to_string())?,
663 name: format!("{}@{machine}", session.name),
664 });
665 }
666 if let Some(thread) = std::env::var("CODEX_THREAD_ID")
670 .ok()
671 .filter(|thread| !thread.is_empty())
672 .filter(|thread| crate::codex_peer::holds_session(pid, thread))
673 {
674 return Ok(Caller {
675 address: MailAddress::new(&machine, "codex", &thread)
676 .map_err(|error| error.to_string())?,
677 name: format!("{}@{machine}", codex_name(&thread)),
678 });
679 }
680 if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
681 return Ok(Caller {
682 address: MailAddress::new(&machine, "codex", &session_id)
683 .map_err(|error| error.to_string())?,
684 name: format!("{}@{machine}", codex_name(&session_id)),
685 });
686 }
687 }
688 #[cfg(windows)]
689 if let Some(session) = msys_cut_claim(®istry, pids) {
690 return Ok(Caller {
691 address: MailAddress::new(&machine, "claude-code", &session.session_id)
692 .map_err(|error| error.to_string())?,
693 name: format!("{}@{machine}", session.name),
694 });
695 }
696 Err(format!(
697 "{CALLER_UNRESOLVED} (looked for this command's processes {pids:?} among {} Claude \
698 sessions in {} and {} hosted runtimes)",
699 registry.len(),
700 registry_dir(homes).display(),
701 hosted.len()
702 ))
703}
704
705#[cfg(windows)]
712fn msys_cut_claim<'a>(
713 registry: &'a [ClaudePeerSession],
714 pids: &[u32],
715) -> Option<&'a ClaudePeerSession> {
716 let table = process_table();
717 let [.., shell, cut] = pids else {
718 return None;
719 };
720 if table.contains_key(cut) {
721 return None;
722 }
723 let name = table.get(shell)?.1.to_ascii_lowercase();
724 if !matches!(name.as_str(), "sh.exe" | "bash.exe" | "dash.exe") {
725 return None;
726 }
727 let claimed = std::env::var("CLAUDE_CODE_SESSION_ID").ok()?;
728 registry
729 .iter()
730 .find(|session| !claimed.is_empty() && session.session_id == claimed)
731}
732
733pub fn process_ancestry() -> Vec<u32> {
735 ancestry_of(std::process::id())
736}
737
738pub fn ancestry_of(pid: u32) -> Vec<u32> {
740 let parents = parent_pids();
741 let mut chain = vec![pid];
742 let mut current = pid;
743 while let Some(&parent) = parents.get(¤t) {
744 if parent <= 1 || chain.contains(&parent) {
745 break;
746 }
747 chain.push(parent);
748 current = parent;
749 }
750 chain
751}
752
753#[cfg(not(windows))]
755fn parent_pids() -> std::collections::HashMap<u32, u32> {
756 let Ok(output) = std::process::Command::new("ps")
757 .args(["-axo", "pid=,ppid="])
758 .output()
759 else {
760 return Default::default();
761 };
762 String::from_utf8_lossy(&output.stdout)
763 .lines()
764 .filter_map(|line| {
765 let mut fields = line.split_whitespace();
766 Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
767 })
768 .collect()
769}
770
771#[cfg(windows)]
773fn parent_pids() -> std::collections::HashMap<u32, u32> {
774 process_table()
775 .into_iter()
776 .map(|(pid, (parent, _))| (pid, parent))
777 .collect()
778}
779
780#[cfg(windows)]
782fn process_table() -> std::collections::HashMap<u32, (u32, String)> {
783 use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
784 use windows_sys::Win32::System::Diagnostics::ToolHelp::{
785 CreateToolhelp32Snapshot, Process32FirstW, Process32NextW, PROCESSENTRY32W,
786 TH32CS_SNAPPROCESS,
787 };
788
789 let mut parents = std::collections::HashMap::new();
790 let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPPROCESS, 0) };
791 if snapshot == INVALID_HANDLE_VALUE {
792 return parents;
793 }
794 let mut entry: PROCESSENTRY32W = unsafe { std::mem::zeroed() };
795 entry.dwSize = std::mem::size_of::<PROCESSENTRY32W>() as u32;
796 let mut has_entry = unsafe { Process32FirstW(snapshot, &mut entry) } != 0;
797 while has_entry {
798 let length = entry
799 .szExeFile
800 .iter()
801 .position(|&unit| unit == 0)
802 .unwrap_or(entry.szExeFile.len());
803 let name = String::from_utf16_lossy(&entry.szExeFile[..length]);
804 parents.insert(entry.th32ProcessID, (entry.th32ParentProcessID, name));
805 has_entry = unsafe { Process32NextW(snapshot, &mut entry) } != 0;
806 }
807 unsafe {
808 CloseHandle(snapshot);
809 }
810 parents
811}
812
813pub fn codex_hooks_path() -> std::path::PathBuf {
815 std::env::var_os("CODEX_HOME")
816 .map(std::path::PathBuf::from)
817 .or_else(|| {
818 supercode_interchange::user_home()
819 .map(std::path::PathBuf::into_os_string)
820 .map(|home| std::path::PathBuf::from(home).join(".codex"))
821 })
822 .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
823 .join("hooks.json")
824}
825
826pub fn codex_user_hook_installed() -> bool {
829 std::fs::read_to_string(codex_hooks_path())
830 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
831}
832
833pub fn codex_project_hook_installed(cwd: &Path) -> bool {
836 cwd.ancestors().any(|directory| {
837 std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
838 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
839 })
840}
841
842#[derive(Debug, Clone, PartialEq, Eq)]
844pub enum Delivered {
845 Steered,
847 Started,
849 Native {
852 busy: bool,
854 },
855 Hooked,
857 HookWoken,
860 Queued,
862 Stored,
864 Operator,
866}
867
868#[derive(Debug, Clone, PartialEq, Eq)]
870pub enum Refused {
871 CannotQueueNative,
874 TooLong(usize),
876}
877
878pub const MAX_RELAYED_BYTES: usize = 100_000;
881
882pub async fn deliver(
886 envelope: &Envelope,
887 to: &MailAddress,
888 door: &Door,
889 wake: bool,
890 notify_when_idle: bool,
891) -> Result<Result<Delivered, Refused>, String> {
892 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
893 let mut final_reply = false;
897 let delivered = match door {
898 Door::Runtime(record) => {
899 let mut sent = envelope.clone();
900 let answers = matches!(envelope.reply_via, ReplyVia::Command);
901 if answers {
902 sent.reply_via = ReplyVia::FinalMessage {
903 destination: envelope.from_name.clone(),
904 };
905 }
906 match deliver_to_runtime(record, sent.render(), wake).await? {
907 RuntimeDelivery::Steered => {
908 mailbox
909 .deliver_read(&sent)
910 .map_err(|error| error.to_string())?;
911 final_reply = answers;
912 Delivered::Steered
913 }
914 RuntimeDelivery::Started => {
915 mailbox
916 .deliver_read(&sent)
917 .map_err(|error| error.to_string())?;
918 final_reply = answers;
919 Delivered::Started
920 }
921 RuntimeDelivery::NotWoken => {
922 mailbox
923 .deliver(envelope)
924 .map_err(|error| error.to_string())?;
925 Delivered::Queued
926 }
927 }
928 }
929 Door::Native(session) => {
930 let busy = session.status != Some(ClaudePeerStatus::Idle);
931 if !wake && !busy {
932 return Ok(Err(Refused::CannotQueueNative));
933 }
934 let mut envelope = envelope.clone();
938 if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
939 envelope.reply_via = ReplyVia::Tool;
940 }
941 let envelope = &envelope;
942 let text = envelope.render();
943 if text.len() > MAX_RELAYED_BYTES {
944 return Ok(Err(Refused::TooLong(text.len())));
945 }
946 match send_through_relay(
947 &envelope.from,
948 &envelope.from_name,
949 &session.name,
950 text,
951 &envelope.id,
952 )
953 .await
954 {
955 RelayReceipt::Delivered { .. } => {
956 mailbox
957 .deliver_read(envelope)
958 .map_err(|error| error.to_string())?;
959 Delivered::Native { busy }
960 }
961 RelayReceipt::Failed { detail } => return Err(detail),
962 }
963 }
964 Door::Hook { .. } => {
965 mailbox
966 .deliver(envelope)
967 .map_err(|error| error.to_string())?;
968 if wake {
969 mailbox
970 .request_wake(&envelope.id)
971 .map_err(|error| error.to_string())?;
972 }
973 Delivered::Hooked
974 }
975 Door::Stored => {
976 mailbox
977 .deliver(envelope)
978 .map_err(|error| error.to_string())?;
979 Delivered::Stored
980 }
981 Door::Operator => {
982 mailbox
983 .deliver(envelope)
984 .map_err(|error| error.to_string())?;
985 Delivered::Operator
986 }
987 };
988 let notice = notify_when_idle && !matches!(door, Door::Operator);
989 if notice || final_reply {
990 let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
991 subscription.notice = notice;
992 subscription.final_reply = final_reply;
993 mailbox
994 .subscribe_idle(&subscription)
995 .map_err(|error| error.to_string())?;
996 if let Ok(program) = supercode_program() {
998 crate::claude_relay::ensure_machine_daemon(&program)
999 .await
1000 .ok();
1001 }
1002 }
1003 Ok(Ok(delivered))
1004}
1005
1006#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1008pub enum UserTurn {
1009 Steered,
1011 Started,
1013 Typed,
1015 Waiting,
1019}
1020
1021impl UserTurn {
1022 pub const fn as_str(self) -> &'static str {
1024 match self {
1025 Self::Steered => "steered",
1026 Self::Started => "started",
1027 Self::Typed => "typed",
1028 Self::Waiting => "waiting",
1029 }
1030 }
1031}
1032
1033pub fn pane_prompt(pane: &str) -> Option<Vec<String>> {
1038 let screen = pane_screen(pane)?;
1039 let lines: Vec<&str> = screen.lines().map(str::trim_end).collect();
1040 let is_option = |line: &str| {
1041 let line = line.trim_start();
1042 let rest = line
1044 .strip_prefix('❯')
1045 .or_else(|| line.strip_prefix('›'))
1046 .or_else(|| line.strip_prefix('>'))
1047 .unwrap_or(line)
1048 .trim_start();
1049 let digits = rest.chars().take_while(char::is_ascii_digit).count();
1050 digits > 0 && matches!(rest[digits..].chars().next(), Some('.' | ')'))
1051 };
1052 let selected = lines.iter().position(|line| {
1053 let line = line.trim_start();
1054 (line.starts_with('❯') || line.starts_with('›') || line.starts_with('>')) && is_option(line)
1055 })?;
1056 let block_start = |end: usize| {
1058 lines[..end]
1059 .iter()
1060 .rposition(|line| line.trim().is_empty())
1061 .map_or(0, |blank| blank + 1)
1062 };
1063 let options = block_start(selected);
1064 let above = lines[..options]
1065 .iter()
1066 .rposition(|line| !line.trim().is_empty())
1067 .map(|last| block_start(last));
1068 let start = above.unwrap_or(options);
1069 let end = lines[selected..]
1070 .iter()
1071 .position(|line| line.trim().is_empty())
1072 .map_or(lines.len(), |blank| selected + blank);
1073 Some(
1074 lines[start..end]
1075 .iter()
1076 .map(|line| line.trim().to_string())
1077 .filter(|line| !line.is_empty())
1078 .take(20)
1079 .collect(),
1080 )
1081}
1082
1083fn live_daemon_panes() -> Option<std::collections::HashMap<String, bool>> {
1085 let entry = crate::teams_entry().ok()?;
1086 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1087 .ok()
1088 .filter(|value| !value.trim().is_empty())
1089 .unwrap_or_else(|| "node".into());
1090 let output = std::process::Command::new(node)
1091 .arg(entry)
1092 .args(["panes", "ls", "--presence-only"])
1093 .stdin(std::process::Stdio::null())
1094 .output()
1095 .ok()?;
1096 if !output.status.success() {
1097 return None;
1098 }
1099 let rows: Vec<serde_json::Value> = serde_json::from_slice(&output.stdout).ok()?;
1100 rows.iter()
1101 .map(|row| {
1102 Some((
1103 row["pane"].as_str()?.to_owned(),
1104 row["idleComposer"].as_bool()?,
1105 ))
1106 })
1107 .collect()
1108}
1109
1110fn pane_screen(pane: &str) -> Option<String> {
1113 #[cfg(unix)]
1114 {
1115 let output = std::process::Command::new("tmux")
1116 .args(["capture-pane", "-p", "-t", &format!("{pane}:")])
1117 .output()
1118 .ok()?;
1119 output
1120 .status
1121 .success()
1122 .then(|| String::from_utf8_lossy(&output.stdout).into_owned())
1123 }
1124 #[cfg(not(unix))]
1125 {
1126 let entry = crate::teams_entry().ok()?;
1127 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1128 .ok()
1129 .filter(|value| !value.trim().is_empty())
1130 .unwrap_or_else(|| "node".into());
1131 let output = std::process::Command::new(node)
1132 .arg(entry)
1133 .args(["panes", "capture", pane, "--lines", "60"])
1134 .stdin(std::process::Stdio::null())
1135 .output()
1136 .ok()?;
1137 output
1138 .status
1139 .success()
1140 .then(|| String::from_utf8_lossy(&output.stdout).into_owned())
1141 }
1142}
1143
1144pub fn daemon_pane(session: &LiveSession) -> Option<String> {
1146 let name = session.tmux.as_deref()?.split(':').next()?;
1147 let rest = name.strip_prefix("p_")?;
1148 (!rest.is_empty()
1149 && rest
1150 .chars()
1151 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
1152 .then(|| name.to_string())
1153}
1154
1155pub async fn deliver_user_turn(
1161 homes: &HarnessHomes,
1162 envelope: &Envelope,
1163 to: &MailAddress,
1164) -> Result<UserTurn, String> {
1165 let session = LiveSessions::read(homes)
1166 .sessions
1167 .into_iter()
1168 .find(|session| &session.address == to)
1169 .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
1170 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
1171 if let Door::Runtime(record) = &session.door {
1172 let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
1173 mailbox
1174 .deliver_read(envelope)
1175 .map_err(|error| error.to_string())?;
1176 return Ok(match delivered {
1177 RuntimeDelivery::Steered => UserTurn::Steered,
1178 _ => UserTurn::Started,
1179 });
1180 }
1181 let Some(pane) = daemon_pane(&session) else {
1182 let name = session.name.split('@').next().unwrap_or(&session.name);
1183 return Err(format!(
1184 "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
1185 pane with `supercode open {name}`; nothing was sent."
1186 ));
1187 };
1188 mailbox
1189 .deliver(envelope)
1190 .map_err(|error| error.to_string())?;
1191 let typed = type_user_turns(&mailbox, &pane).await;
1192 Ok(if typed.contains(&envelope.id) {
1193 UserTurn::Typed
1194 } else {
1195 UserTurn::Waiting
1196 })
1197}
1198
1199pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
1203 let mut typed = Vec::new();
1204 for waiting in mailbox.user_turns().unwrap_or_default() {
1205 let Ok(Some(stored)) = mailbox.claim_user_turn(&waiting) else {
1209 break;
1210 };
1211 match submit_mail_batch(pane, &stored.envelope.body, &[stored.envelope.id.clone()]).await {
1212 Ok(true) => {
1213 mailbox.acknowledge(&stored).ok();
1214 typed.push(stored.envelope.id.clone());
1215 }
1216 Ok(false) => {
1217 mailbox.release(&stored).ok();
1218 break;
1219 }
1220 Err(error) => {
1221 mailbox.release(&stored).ok();
1222 eprintln!(
1223 "supercode: the user's turn {} for {} waits: {error}",
1224 stored.envelope.id,
1225 mailbox.address()
1226 );
1227 break;
1228 }
1229 }
1230 }
1231 typed
1232}
1233
1234pub(crate) async fn wake_hook_mailbox(
1237 mailbox: &Mailbox,
1238 pane: &str,
1239 ids: &[String],
1240) -> Result<bool, String> {
1241 let messages = ids
1242 .iter()
1243 .map(|id| mailbox.find(id))
1244 .collect::<Result<Vec<_>, _>>()
1245 .map_err(|error| error.to_string())?;
1246 let messages: Vec<_> = messages.into_iter().flatten().collect();
1247 if messages.is_empty() {
1248 return Ok(false);
1249 }
1250 let mut delivered = true;
1251 for message in messages {
1252 delivered &= submit_mail_batch(
1253 pane,
1254 &message.envelope.render(),
1255 &[message.envelope.id.clone()],
1256 )
1257 .await?;
1258 }
1259 Ok(delivered)
1260}
1261
1262async fn submit_mail_batch(pane: &str, text: &str, ids: &[String]) -> Result<bool, String> {
1263 let entry = crate::teams_entry().map_err(|error| error.to_string())?;
1264 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1265 .ok()
1266 .filter(|value| !value.trim().is_empty())
1267 .unwrap_or_else(|| "node".into());
1268 let output = tokio::process::Command::new(node)
1269 .arg(entry)
1270 .args([
1271 "input",
1272 pane,
1273 text,
1274 "--when-composer-empty",
1275 "--message-ids",
1276 &ids.join(","),
1277 ])
1278 .env_remove("SUPERCODE_CALLER")
1281 .stdin(std::process::Stdio::null())
1282 .output()
1283 .await
1284 .map_err(|error| error.to_string())?;
1285 if !output.status.success() {
1286 return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
1287 &output.stderr,
1288 )));
1289 }
1290 let answer: serde_json::Value =
1291 serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
1292 Ok(answer["delivered"] == true)
1293}