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 pub transcript: Option<std::path::PathBuf>,
111}
112
113impl LiveSession {
114 pub fn last_message_at_ms(&self, homes: &HarnessHomes) -> Option<u64> {
119 let harness = self.address.harness.as_str();
120 let path = match &self.transcript {
121 Some(path) => path.clone(),
122 None if harness == "claude-code" => {
123 let file = format!("{}.jsonl", self.address.session_id);
124 std::fs::read_dir(&homes.claude_code)
125 .ok()?
126 .flatten()
127 .map(|project| project.path().join(&file))
128 .find(|path| path.is_file())?
129 }
130 None => return None,
131 };
132 last_message_at_ms(&path, harness)
133 }
134}
135
136impl LiveSession {
137 pub fn pending_request(&self, homes: &HarnessHomes) -> Option<serde_json::Value> {
141 if self.status != "waiting" || self.address.harness != "claude-code" {
142 return None;
143 }
144 let file = format!("{}.jsonl", self.address.session_id);
145 let path = std::fs::read_dir(&homes.claude_code)
146 .ok()?
147 .flatten()
148 .map(|project| project.path().join(&file))
149 .find(|path| path.is_file())?;
150 pending_request(&path)
151 }
152}
153
154pub fn pending_request(path: &Path) -> Option<serde_json::Value> {
156 use std::io::{Read, Seek, SeekFrom};
157 let mut file = std::fs::File::open(path).ok()?;
158 let length = file.metadata().ok()?.len();
159 let start = length.saturating_sub(512 * 1024);
160 file.seek(SeekFrom::Start(start)).ok()?;
161 let mut bytes = Vec::new();
162 file.read_to_end(&mut bytes).ok()?;
163 let text = String::from_utf8_lossy(&bytes);
164 let mut answered = std::collections::HashSet::new();
165 for line in text.lines().rev() {
166 let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
167 continue;
168 };
169 if record["isSidechain"] == true {
170 continue;
171 }
172 let Some(content) = record
173 .pointer("/message/content")
174 .and_then(serde_json::Value::as_array)
175 else {
176 continue;
177 };
178 match record["type"].as_str() {
179 Some("user") => {
180 for item in content.iter().filter(|item| item["type"] == "tool_result") {
181 if let Some(id) = item["tool_use_id"].as_str() {
182 answered.insert(id.to_string());
183 }
184 }
185 }
186 Some("assistant") => {
187 let Some(call) = content.iter().rev().find(|item| item["type"] == "tool_use")
188 else {
189 continue;
190 };
191 if call["id"].as_str().is_some_and(|id| answered.contains(id)) {
192 return None;
193 }
194 let tool = call["name"].as_str().unwrap_or_default();
195 if tool == "AskUserQuestion" {
196 let questions = call["input"]["questions"]
197 .as_array()
198 .map(|questions| {
199 questions
200 .iter()
201 .map(|question| {
202 serde_json::json!({
203 "question": question["question"],
204 "header": question["header"],
205 "options": question["options"].as_array().map(|options| options.iter().map(|option| option["label"].clone()).collect::<Vec<_>>()).unwrap_or_default(),
206 })
207 })
208 .collect::<Vec<_>>()
209 })
210 .unwrap_or_default();
211 return Some(serde_json::json!({"tool": tool, "questions": questions}));
212 }
213 let input = call["input"].to_string();
214 let input: String = input.chars().take(500).collect();
215 return Some(serde_json::json!({"tool": tool, "input": input}));
216 }
217 _ => {}
218 }
219 }
220 None
221}
222
223fn last_message_at_ms(path: &Path, harness: &str) -> Option<u64> {
225 use std::io::{Read, Seek, SeekFrom};
226 for window in [256 * 1024_u64, 8 * 1024 * 1024] {
228 let mut file = std::fs::File::open(path).ok()?;
229 let length = file.metadata().ok()?.len();
230 let start = length.saturating_sub(window);
231 file.seek(SeekFrom::Start(start)).ok()?;
232 let mut bytes = Vec::new();
233 file.read_to_end(&mut bytes).ok()?;
234 let text = String::from_utf8_lossy(&bytes);
235 let mut lines = text.lines().rev().collect::<Vec<_>>();
236 if start > 0 {
237 lines.pop(); }
239 for line in lines {
240 let Ok(record) = serde_json::from_str::<serde_json::Value>(line) else {
241 continue;
242 };
243 let message = match harness {
244 "codex" => record["type"] == "response_item",
245 _ => {
246 matches!(record["type"].as_str(), Some("user" | "assistant"))
247 && record["isSidechain"] != true
248 }
249 };
250 if !message {
251 continue;
252 }
253 if let Some(at) = record["timestamp"]
254 .as_str()
255 .and_then(supercode_interchange::sidecar::rfc3339_to_ms)
256 .and_then(|at| u64::try_from(at).ok())
257 {
258 return Some(at);
259 }
260 }
261 if start == 0 {
262 return None;
263 }
264 }
265 None
266}
267
268#[derive(Debug, Default)]
272pub struct LiveSessions {
273 sessions: Vec<LiveSession>,
274}
275
276#[derive(Debug, Clone, PartialEq, Eq)]
278pub enum Unresolved {
279 Stale(String),
281 Unknown(String),
283}
284
285impl LiveSessions {
286 pub fn read(homes: &HarnessHomes) -> Self {
290 let machine = local_machine_name();
291 let registry = read_registry(®istry_dir(homes));
292 let registered = |record: &crate::live_runtime::LiveRuntimeRecord| {
295 registry.iter().find(|session| {
296 session.session_id == record.source.session_id
297 || session.session_id == record.runtime_session_id
298 })
299 };
300 let mut sessions: Vec<LiveSession> = crate::runtime_mail::controlled_runtimes()
301 .into_iter()
302 .filter_map(|record| {
303 let registered = registered(&record);
304 if registered.is_some_and(|session| session.name.starts_with(RELAY_NAME_PREFIX)) {
305 return None;
306 }
307 let address =
308 MailAddress::new(&machine, &record.source.harness, &record.source.session_id)
309 .ok()?;
310 let short: String = record.source.session_id.chars().take(8).collect();
311 let name = match registered {
312 Some(session) if !session.name.is_empty() => session.name.clone(),
313 _ => format!("{}-{short}", record.source.harness),
314 };
315 Some(LiveSession {
316 name: format!("{name}@{machine}"),
317 address,
318 status: "hosted".into(),
319 pid: Some(record.pid),
320 cwd: Some(record.source.workspace.clone()),
321 tmux: None,
322 transcript: None,
323 door: Door::Runtime(Box::new(record)),
324 })
325 })
326 .collect();
327 let controlled = |address: &MailAddress, sessions: &[LiveSession]| {
328 sessions.iter().any(|session| &session.address == address)
329 };
330 for session in registry {
331 if session.name.starts_with(RELAY_NAME_PREFIX) {
332 continue;
333 }
334 let Ok(address) = MailAddress::new(&machine, "claude-code", &session.session_id) else {
335 continue;
336 };
337 if controlled(&address, &sessions) {
338 continue;
339 }
340 sessions.push(LiveSession {
341 address,
342 name: format!("{}@{machine}", session.name),
343 status: session
344 .status
345 .as_ref()
346 .map(|status| status.as_str().to_string())
347 .unwrap_or_else(|| "unknown".into()),
348 pid: Some(session.pid),
349 cwd: session.cwd.clone(),
350 tmux: session.tmux.clone(),
351 transcript: None,
352 door: Door::Native(Box::new(session)),
353 });
354 }
355 let user_hook = codex_user_hook_installed();
356 for (path, status) in crate::codex_peer::live_rollouts(&homes.codex) {
357 let Some((session_id, None)) = crate::codex_peer::rollout_session(&path) else {
358 continue;
359 };
360 let Ok(address) = MailAddress::new(&machine, "codex", &session_id) else {
361 continue;
362 };
363 if controlled(&address, &sessions) {
364 continue;
365 }
366 let cwd = crate::codex_peer::rollout_cwd(&path);
367 let hooked = user_hook || cwd.as_deref().is_some_and(codex_project_hook_installed);
368 sessions.push(LiveSession {
369 name: format!("{}@{machine}", codex_name(&session_id)),
370 address,
371 status: status.as_str().to_string(),
372 pid: None,
373 cwd,
374 tmux: None,
375 transcript: Some(path),
376 door: if hooked { Door::Hook } else { Door::Stored },
377 });
378 }
379 Self { sessions }
380 }
381
382 pub fn all(&self) -> &[LiveSession] {
384 &self.sessions
385 }
386
387 pub fn door(&self, harness: &str, session_id: &str) -> Option<&'static str> {
390 self.sessions
391 .iter()
392 .find(|session| {
393 session.address.harness == harness && session.address.session_id == session_id
394 })
395 .map(|session| session.door.name())
396 }
397
398 pub fn resolve(&self, to: &str) -> Result<&LiveSession, Unresolved> {
401 let machine = local_machine_name();
402 if let Ok(address) = MailAddress::parse(to) {
403 return self
404 .sessions
405 .iter()
406 .find(|session| session.address == address)
407 .ok_or_else(|| {
408 Unresolved::Stale(format!(
409 "{to} is no longer running. Nothing was sent. Run supercode message list \
410 for the live sessions."
411 ))
412 });
413 }
414 let wanted = match to.split_once('@') {
415 Some((name, at)) if at == machine => name.to_string(),
416 Some((_, at)) => {
417 return Err(Unresolved::Unknown(format!(
418 "Not sent: {to} is on machine {at}, not this one. Nothing was sent."
419 )))
420 }
421 None => to.to_string(),
422 };
423 let short = |session: &LiveSession| {
424 session
425 .name
426 .split('@')
427 .next()
428 .unwrap_or_default()
429 .to_string()
430 };
431 let matching: Vec<&LiveSession> = self
432 .sessions
433 .iter()
434 .filter(|session| short(session) == wanted)
435 .collect();
436 if let [only] = matching.as_slice() {
437 return Ok(only);
438 }
439 let hint = if matching.len() > 1 {
440 format!(
441 " {} sessions are named {wanted}; use its address.",
442 matching.len()
443 )
444 } else {
445 let near: Vec<String> = self
446 .sessions
447 .iter()
448 .filter(|session| {
449 let name = short(session);
450 name.contains(&wanted)
451 || wanted.contains(&name)
452 || name
453 .chars()
454 .zip(wanted.chars())
455 .take_while(|(a, b)| a == b)
456 .count()
457 >= 4
458 })
459 .take(3)
460 .map(|session| {
461 format!(
462 "{} ({}, {})",
463 session.name, session.address.harness, session.status
464 )
465 })
466 .collect();
467 if near.is_empty() {
468 String::new()
469 } else {
470 format!(" Did you mean: {}?", near.join(", "))
471 }
472 };
473 Err(Unresolved::Unknown(format!(
474 "No session named \"{wanted}\" is reachable.{hint} Run supercode message list. Nothing \
475 was sent."
476 )))
477 }
478}
479
480pub fn has_message_tools(pid: u32) -> bool {
484 let Ok(output) = std::process::Command::new("ps")
485 .args(["-A", "-o", "ppid=,command="])
486 .output()
487 else {
488 return false;
489 };
490 String::from_utf8_lossy(&output.stdout).lines().any(|line| {
491 let line = line.trim_start();
492 let Some((ppid, command)) = line.split_once(' ') else {
493 return false;
494 };
495 ppid.parse::<u32>() == Ok(pid) && command.trim_end().ends_with(" message mcp")
496 })
497}
498
499pub fn codex_name(session_id: &str) -> String {
501 format!("codex-{}", session_id.chars().take(8).collect::<String>())
502}
503
504#[derive(Debug, Clone, PartialEq, Eq)]
506pub struct Caller {
507 pub address: MailAddress,
509 pub name: String,
511}
512
513pub const CALLER_UNRESOLVED: &str = "Can't tell which session is running this command, so \
515 replies would have nowhere to go. Nothing was sent. Run it from your agent session's own \
516 shell tool.";
517
518pub fn resolve_caller(homes: &HarnessHomes, pids: &[u32]) -> Result<Caller, String> {
523 let machine = local_machine_name();
524 let registry = read_registry(®istry_dir(homes));
525 let hosted = crate::runtime_mail::controlled_runtimes();
526 for &pid in pids {
527 if let Some(record) = hosted.iter().find(|record| record.pid == pid) {
528 let short: String = record.source.session_id.chars().take(8).collect();
529 return Ok(Caller {
530 address: MailAddress::new(
531 &machine,
532 &record.source.harness,
533 &record.source.session_id,
534 )
535 .map_err(|error| error.to_string())?,
536 name: format!("{}-{short}@{machine}", record.source.harness),
537 });
538 }
539 if let Some(session) = registry.iter().find(|session| session.pid == pid) {
540 if let Ok(claimed) = std::env::var("CLAUDE_CODE_SESSION_ID") {
541 if !claimed.is_empty() && claimed != session.session_id {
542 return Err(format!(
543 "{CALLER_UNRESOLVED} (CLAUDE_CODE_SESSION_ID names {claimed}, but the Claude \
544 process {pid} above this command is session {})",
545 session.session_id
546 ));
547 }
548 }
549 return Ok(Caller {
550 address: MailAddress::new(&machine, "claude-code", &session.session_id)
551 .map_err(|error| error.to_string())?,
552 name: format!("{}@{machine}", session.name),
553 });
554 }
555 if let Some(thread) = std::env::var("CODEX_THREAD_ID")
559 .ok()
560 .filter(|thread| !thread.is_empty())
561 .filter(|thread| crate::codex_peer::holds_session(pid, thread))
562 {
563 return Ok(Caller {
564 address: MailAddress::new(&machine, "codex", &thread)
565 .map_err(|error| error.to_string())?,
566 name: format!("{}@{machine}", codex_name(&thread)),
567 });
568 }
569 if let Some((session_id, _)) = crate::codex_peer::session_of_process(pid) {
570 return Ok(Caller {
571 address: MailAddress::new(&machine, "codex", &session_id)
572 .map_err(|error| error.to_string())?,
573 name: format!("{}@{machine}", codex_name(&session_id)),
574 });
575 }
576 }
577 #[cfg(windows)]
578 if let Some(session) = msys_cut_claim(®istry, pids) {
579 return Ok(Caller {
580 address: MailAddress::new(&machine, "claude-code", &session.session_id)
581 .map_err(|error| error.to_string())?,
582 name: format!("{}@{machine}", session.name),
583 });
584 }
585 Err(format!(
586 "{CALLER_UNRESOLVED} (looked for this command's processes {pids:?} among {} Claude \
587 sessions in {} and {} hosted runtimes)",
588 registry.len(),
589 registry_dir(homes).display(),
590 hosted.len()
591 ))
592}
593
594#[cfg(windows)]
601fn msys_cut_claim<'a>(
602 registry: &'a [ClaudePeerSession],
603 pids: &[u32],
604) -> Option<&'a ClaudePeerSession> {
605 let table = process_table();
606 let [.., shell, cut] = pids else {
607 return None;
608 };
609 if table.contains_key(cut) {
610 return None;
611 }
612 let name = table.get(shell)?.1.to_ascii_lowercase();
613 if !matches!(name.as_str(), "sh.exe" | "bash.exe" | "dash.exe") {
614 return None;
615 }
616 let claimed = std::env::var("CLAUDE_CODE_SESSION_ID").ok()?;
617 registry
618 .iter()
619 .find(|session| !claimed.is_empty() && session.session_id == claimed)
620}
621
622pub fn process_ancestry() -> Vec<u32> {
624 ancestry_of(std::process::id())
625}
626
627pub fn ancestry_of(pid: u32) -> Vec<u32> {
629 let parents = parent_pids();
630 let mut chain = vec![pid];
631 let mut current = pid;
632 while let Some(&parent) = parents.get(¤t) {
633 if parent <= 1 || chain.contains(&parent) {
634 break;
635 }
636 chain.push(parent);
637 current = parent;
638 }
639 chain
640}
641
642#[cfg(not(windows))]
644fn parent_pids() -> std::collections::HashMap<u32, u32> {
645 let Ok(output) = std::process::Command::new("ps")
646 .args(["-axo", "pid=,ppid="])
647 .output()
648 else {
649 return Default::default();
650 };
651 String::from_utf8_lossy(&output.stdout)
652 .lines()
653 .filter_map(|line| {
654 let mut fields = line.split_whitespace();
655 Some((fields.next()?.parse().ok()?, fields.next()?.parse().ok()?))
656 })
657 .collect()
658}
659
660#[cfg(windows)]
662fn parent_pids() -> std::collections::HashMap<u32, u32> {
663 process_table()
664 .into_iter()
665 .map(|(pid, (parent, _))| (pid, parent))
666 .collect()
667}
668
669#[cfg(windows)]
671fn process_table() -> std::collections::HashMap<u32, (u32, String)> {
672 use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
673 use windows_sys::Win32::System::Diagnostics::ToolHelp::{
674 CreateToolhelp32Snapshot, Process32FirstW, Process32NextW, PROCESSENTRY32W,
675 TH32CS_SNAPPROCESS,
676 };
677
678 let mut parents = std::collections::HashMap::new();
679 let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPPROCESS, 0) };
680 if snapshot == INVALID_HANDLE_VALUE {
681 return parents;
682 }
683 let mut entry: PROCESSENTRY32W = unsafe { std::mem::zeroed() };
684 entry.dwSize = std::mem::size_of::<PROCESSENTRY32W>() as u32;
685 let mut has_entry = unsafe { Process32FirstW(snapshot, &mut entry) } != 0;
686 while has_entry {
687 let length = entry
688 .szExeFile
689 .iter()
690 .position(|&unit| unit == 0)
691 .unwrap_or(entry.szExeFile.len());
692 let name = String::from_utf16_lossy(&entry.szExeFile[..length]);
693 parents.insert(entry.th32ProcessID, (entry.th32ParentProcessID, name));
694 has_entry = unsafe { Process32NextW(snapshot, &mut entry) } != 0;
695 }
696 unsafe {
697 CloseHandle(snapshot);
698 }
699 parents
700}
701
702pub fn codex_hooks_path() -> std::path::PathBuf {
704 std::env::var_os("CODEX_HOME")
705 .map(std::path::PathBuf::from)
706 .or_else(|| {
707 supercode_interchange::user_home()
708 .map(std::path::PathBuf::into_os_string)
709 .map(|home| std::path::PathBuf::from(home).join(".codex"))
710 })
711 .unwrap_or_else(|| std::path::PathBuf::from(".codex"))
712 .join("hooks.json")
713}
714
715pub fn codex_user_hook_installed() -> bool {
718 std::fs::read_to_string(codex_hooks_path())
719 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
720}
721
722pub fn codex_project_hook_installed(cwd: &Path) -> bool {
725 cwd.ancestors().any(|directory| {
726 std::fs::read_to_string(directory.join(".codex").join("hooks.json"))
727 .is_ok_and(|text| text.contains(CODEX_HOOK_ARGUMENTS))
728 })
729}
730
731#[derive(Debug, Clone, PartialEq, Eq)]
733pub enum Delivered {
734 Steered,
736 Started,
738 Native {
741 busy: bool,
743 },
744 Hooked,
746 Queued,
748 Stored,
750 Operator,
752}
753
754#[derive(Debug, Clone, PartialEq, Eq)]
756pub enum Refused {
757 CannotQueueNative,
760 TooLong(usize),
762}
763
764pub const MAX_RELAYED_BYTES: usize = 100_000;
767
768pub async fn deliver(
772 envelope: &Envelope,
773 to: &MailAddress,
774 door: &Door,
775 wake: bool,
776 notify_when_idle: bool,
777) -> Result<Result<Delivered, Refused>, String> {
778 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
779 let mut final_reply = false;
783 let delivered = match door {
784 Door::Runtime(record) => {
785 let mut sent = envelope.clone();
786 let answers = matches!(envelope.reply_via, ReplyVia::Command);
787 if answers {
788 sent.reply_via = ReplyVia::FinalMessage {
789 destination: envelope.from_name.clone(),
790 };
791 }
792 match deliver_to_runtime(record, sent.render(), wake).await? {
793 RuntimeDelivery::Steered => {
794 mailbox
795 .deliver_read(&sent)
796 .map_err(|error| error.to_string())?;
797 final_reply = answers;
798 Delivered::Steered
799 }
800 RuntimeDelivery::Started => {
801 mailbox
802 .deliver_read(&sent)
803 .map_err(|error| error.to_string())?;
804 final_reply = answers;
805 Delivered::Started
806 }
807 RuntimeDelivery::NotWoken => {
808 mailbox
809 .deliver(envelope)
810 .map_err(|error| error.to_string())?;
811 Delivered::Queued
812 }
813 }
814 }
815 Door::Native(session) => {
816 let busy = session.status != Some(ClaudePeerStatus::Idle);
817 if !wake && !busy {
818 return Ok(Err(Refused::CannotQueueNative));
819 }
820 let mut envelope = envelope.clone();
824 if envelope.reply_via == ReplyVia::Command && has_message_tools(session.pid) {
825 envelope.reply_via = ReplyVia::Tool;
826 }
827 let envelope = &envelope;
828 let text = envelope.render();
829 if text.len() > MAX_RELAYED_BYTES {
830 return Ok(Err(Refused::TooLong(text.len())));
831 }
832 match send_through_relay(
833 &envelope.from,
834 &envelope.from_name,
835 &session.name,
836 text,
837 &envelope.id,
838 )
839 .await
840 {
841 RelayReceipt::Delivered { .. } => {
842 mailbox
843 .deliver_read(envelope)
844 .map_err(|error| error.to_string())?;
845 Delivered::Native { busy }
846 }
847 RelayReceipt::Failed { detail } => return Err(detail),
848 }
849 }
850 Door::Hook => {
851 mailbox
852 .deliver(envelope)
853 .map_err(|error| error.to_string())?;
854 Delivered::Hooked
855 }
856 Door::Stored => {
857 mailbox
858 .deliver(envelope)
859 .map_err(|error| error.to_string())?;
860 Delivered::Stored
861 }
862 Door::Operator => {
863 mailbox
864 .deliver(envelope)
865 .map_err(|error| error.to_string())?;
866 Delivered::Operator
867 }
868 };
869 let notice = notify_when_idle && !matches!(door, Door::Operator);
870 if notice || final_reply {
871 let mut subscription = IdleSubscription::new(envelope.id.clone(), envelope.from.clone());
872 subscription.notice = notice;
873 subscription.final_reply = final_reply;
874 mailbox
875 .subscribe_idle(&subscription)
876 .map_err(|error| error.to_string())?;
877 if let Ok(program) = supercode_program() {
879 crate::claude_relay::ensure_machine_daemon(&program)
880 .await
881 .ok();
882 }
883 }
884 Ok(Ok(delivered))
885}
886
887#[derive(Debug, Clone, Copy, PartialEq, Eq)]
889pub enum UserTurn {
890 Steered,
892 Started,
894 Typed,
896 Waiting,
900}
901
902impl UserTurn {
903 pub const fn as_str(self) -> &'static str {
905 match self {
906 Self::Steered => "steered",
907 Self::Started => "started",
908 Self::Typed => "typed",
909 Self::Waiting => "waiting",
910 }
911 }
912}
913
914pub fn daemon_pane(session: &LiveSession) -> Option<String> {
916 let name = session.tmux.as_deref()?.split(':').next()?;
917 let rest = name.strip_prefix("p_")?;
918 (!rest.is_empty()
919 && rest
920 .chars()
921 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-'))
922 .then(|| name.to_string())
923}
924
925pub async fn deliver_user_turn(
931 homes: &HarnessHomes,
932 envelope: &Envelope,
933 to: &MailAddress,
934) -> Result<UserTurn, String> {
935 let session = LiveSessions::read(homes)
936 .sessions
937 .into_iter()
938 .find(|session| &session.address == to)
939 .ok_or_else(|| format!("{to} is not running; nothing was sent"))?;
940 let mailbox = Mailbox::open(&mail_root(), to).map_err(|error| error.to_string())?;
941 if let Door::Runtime(record) = &session.door {
942 let delivered = deliver_to_runtime(record, envelope.body.clone(), true).await?;
943 mailbox
944 .deliver_read(envelope)
945 .map_err(|error| error.to_string())?;
946 return Ok(match delivered {
947 RuntimeDelivery::Steered => UserTurn::Steered,
948 _ => UserTurn::Started,
949 });
950 }
951 let Some(pane) = daemon_pane(&session) else {
952 let name = session.name.split('@').next().unwrap_or(&session.name);
953 return Err(format!(
954 "{name} runs outside supercode, where nothing can speak as its user. Open it in a \
955 pane with `supercode open {name}`; nothing was sent."
956 ));
957 };
958 mailbox
959 .deliver(envelope)
960 .map_err(|error| error.to_string())?;
961 let typed = type_user_turns(&mailbox, &pane).await;
962 Ok(if typed.contains(&envelope.id) {
963 UserTurn::Typed
964 } else {
965 UserTurn::Waiting
966 })
967}
968
969pub async fn type_user_turns(mailbox: &Mailbox, pane: &str) -> Vec<String> {
973 let mut typed = Vec::new();
974 for waiting in mailbox.user_turns().unwrap_or_default() {
975 let Ok(Some(stored)) = mailbox.claim_user_turn(&waiting) else {
979 break;
980 };
981 match submit_when_composer_empty(pane, &stored.envelope.body).await {
982 Ok(true) => {
983 mailbox.acknowledge(&stored).ok();
984 typed.push(stored.envelope.id.clone());
985 }
986 Ok(false) => {
987 mailbox.release(&stored).ok();
988 break;
989 }
990 Err(error) => {
991 mailbox.release(&stored).ok();
992 eprintln!(
993 "supercode: the user's turn {} for {} waits: {error}",
994 stored.envelope.id,
995 mailbox.address()
996 );
997 break;
998 }
999 }
1000 }
1001 typed
1002}
1003
1004async fn submit_when_composer_empty(pane: &str, text: &str) -> Result<bool, String> {
1007 let entry = crate::teams_entry().map_err(|error| error.to_string())?;
1008 let node = std::env::var(crate::orchestrator_door::NODE_BIN_ENV)
1009 .ok()
1010 .filter(|value| !value.trim().is_empty())
1011 .unwrap_or_else(|| "node".into());
1012 let output = tokio::process::Command::new(node)
1013 .arg(entry)
1014 .args(["input", pane, text, "--when-composer-empty"])
1015 .stdin(std::process::Stdio::null())
1016 .output()
1017 .await
1018 .map_err(|error| error.to_string())?;
1019 if !output.status.success() {
1020 return Err(crate::mailbox::error_line(&String::from_utf8_lossy(
1021 &output.stderr,
1022 )));
1023 }
1024 let answer: serde_json::Value =
1025 serde_json::from_slice(&output.stdout).map_err(|error| error.to_string())?;
1026 Ok(answer["delivered"] == true)
1027}