1use std::collections::{BTreeMap, HashMap};
33use std::path::{Path, PathBuf};
34use std::time::{Duration, Instant};
35
36use serde::{Deserialize, Serialize};
37use serde_json::{json, Value};
38
39use crate::claude_peer::{read_registry, registry_dir, ClaudePeerSession};
40use crate::mailbox::{mail_root, Envelope, MailAddress, MailKind, ReplyVia};
41use crate::HarnessHomes;
42
43pub const RELAY_MODEL: &str = "haiku";
46
47pub const RELAY_NAME_PREFIX: &str = "sc-";
49
50pub const RELAY_STATUS_LINE: &str = "relay, not the session's status";
53
54pub const SEND_TIMEOUT: Duration = Duration::from_secs(90);
56
57const RELAY_TOOLS: &str = "ListAgents,SendMessage";
59
60const NATIVE_PEER_PREAMBLE: &str = "Another Claude session sent a message:";
64const NATIVE_PEER_OPENING: &str = "<cross-session-message ";
66const NATIVE_PEER_CLOSING: &str = "</cross-session-message>";
67const NATIVE_IDLE_NOTICE: &str = "[Cross-session idle notice]";
68const NATIVE_DELIVERY_NOTICE: &str = "[Cross-session delivery notice]";
69
70pub fn relay_name(name: &str, machine: &str) -> String {
75 let base = name.strip_prefix(RELAY_NAME_PREFIX).unwrap_or(name);
76 format!("{RELAY_NAME_PREFIX}{base}-on-{machine}")
77}
78
79#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
81pub struct QueuedSend {
82 pub to: String,
84 pub message: String,
86}
87
88pub fn gate_decision(hook_input: &Value, queued: Option<&QueuedSend>) -> Option<String> {
94 let tool = hook_input
95 .get("tool_name")
96 .and_then(Value::as_str)
97 .unwrap_or_default();
98 if tool == "ListAgents" {
99 return None;
100 }
101 if tool != "SendMessage" {
102 return Some(format!("a relay may not use {tool}"));
103 }
104 let Some(queued) = queued else {
105 return Some(
106 "nothing is queued to send; mail for the host is not answered by the relay".into(),
107 );
108 };
109 let input = hook_input.get("tool_input").cloned().unwrap_or(Value::Null);
110 let to = input.get("to").and_then(Value::as_str).unwrap_or_default();
111 let message = input
112 .get("message")
113 .and_then(Value::as_str)
114 .unwrap_or_default();
115 if to != queued.to {
116 return Some(format!(
117 "not the queued outbound message: `to` is {to:?}, the queue says {:?}",
118 queued.to
119 ));
120 }
121 if message != queued.message {
122 let at = message
123 .char_indices()
124 .zip(queued.message.chars())
125 .find(|((_, sent), queued)| sent != queued)
126 .map(|((index, _), _)| index)
127 .unwrap_or_else(|| message.len().min(queued.message.len()));
128 return Some(format!(
129 "not the queued outbound message: `message` ({} bytes) differs from the queue ({} \
130 bytes) at byte {at}: sent {:?}, queued {:?}",
131 message.len(),
132 queued.message.len(),
133 message
134 .get(at..)
135 .unwrap_or_default()
136 .chars()
137 .take(40)
138 .collect::<String>(),
139 queued
140 .message
141 .get(at..)
142 .unwrap_or_default()
143 .chars()
144 .take(40)
145 .collect::<String>(),
146 ));
147 }
148 None
149}
150
151pub fn gate_denial(reason: &str) -> Value {
153 json!({
154 "hookSpecificOutput": {
155 "hookEventName": "PreToolUse",
156 "permissionDecision": "deny",
157 "permissionDecisionReason": reason,
158 }
159 })
160}
161
162#[derive(Debug, Clone)]
165pub struct RelayPaths {
166 pub directory: PathBuf,
168 pub queue: PathBuf,
170 pub receipt: PathBuf,
172 pub turn: PathBuf,
175 pub settings: PathBuf,
177 pub record: PathBuf,
179 pub sent: PathBuf,
181 pub lock: PathBuf,
183}
184
185impl RelayPaths {
186 pub fn new(root: &Path, relay_name: &str) -> Self {
188 let hash = blake3::hash(relay_name.as_bytes()).to_hex();
189 Self::in_directory(root.join("relays").join(&hash[..24]))
190 }
191
192 pub fn in_directory(directory: PathBuf) -> Self {
194 Self {
195 queue: directory.join("queue.json"),
196 receipt: directory.join("receipt.json"),
197 turn: directory.join("turn.txt"),
198 settings: directory.join("settings.json"),
199 record: directory.join("relay.json"),
200 sent: directory.join("sent.json"),
201 lock: directory.join("send.lock"),
202 directory,
203 }
204 }
205
206 pub fn last_sent(&self) -> HashMap<String, String> {
208 std::fs::read(&self.sent)
209 .ok()
210 .and_then(|bytes| serde_json::from_slice(&bytes).ok())
211 .unwrap_or_default()
212 }
213}
214
215#[derive(Debug, Clone)]
217pub struct RelaySpec {
218 pub represented: MailAddress,
220 pub represented_name: String,
222 pub name: String,
224 pub paths: RelayPaths,
226 pub program: PathBuf,
228}
229
230impl RelaySpec {
231 pub fn for_sender(represented: &MailAddress, represented_name: &str) -> std::io::Result<Self> {
233 let base = represented_name
234 .split('@')
235 .next()
236 .unwrap_or(represented_name);
237 let name = relay_name(base, &represented.machine);
238 Ok(Self {
239 represented: represented.clone(),
240 represented_name: represented_name.to_string(),
241 paths: RelayPaths::new(&mail_root(), &name),
242 name,
243 program: supercode_program()?,
244 })
245 }
246}
247
248pub fn relay_arguments(spec: &RelaySpec) -> Vec<String> {
254 vec![
255 "--print".into(),
256 "--input-format".into(),
257 "stream-json".into(),
258 "--output-format".into(),
259 "stream-json".into(),
260 "--verbose".into(),
261 "--permission-prompt-tool".into(),
262 "stdio".into(),
263 "--model".into(),
264 RELAY_MODEL.into(),
265 "--name".into(),
266 spec.name.clone(),
267 "--permission-mode".into(),
268 "bypassPermissions".into(),
269 "--setting-sources".into(),
270 "project".into(),
271 "--settings".into(),
272 spec.paths.settings.to_string_lossy().into_owned(),
273 "--tools".into(),
274 RELAY_TOOLS.into(),
275 "--no-session-persistence".into(),
276 ]
277}
278
279pub fn relay_environment(spec: &RelaySpec, endpoint: &str) -> BTreeMap<String, String> {
283 let key = spec
284 .paths
285 .directory
286 .file_name()
287 .map(|name| name.to_string_lossy().into_owned())
288 .unwrap_or_default();
289 BTreeMap::from([
290 ("SUPERCODE_CLAUDE_RELAY".to_string(), "1".to_string()),
291 ("ANTHROPIC_BASE_URL".to_string(), endpoint.to_string()),
292 ("ANTHROPIC_API_KEY".to_string(), key),
293 (
294 "CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC".to_string(),
295 "1".to_string(),
296 ),
297 ])
298}
299
300pub fn relay_settings(spec: &RelaySpec) -> Value {
302 let program = shell_quote(&spec.program.to_string_lossy());
303 let directory = shell_quote(&spec.paths.directory.to_string_lossy());
304 let represented = shell_quote(&spec.represented.to_string());
305 let command = |verb: &str| format!("{program} message {verb} {directory}");
306 json!({
307 "hooks": {
308 "PreToolUse": [{
310 "matcher": ".*",
311 "hooks": [{"type": "command", "command": command("gate")}],
312 }],
313 "PostToolUse": [{
314 "matcher": "SendMessage",
315 "hooks": [{"type": "command", "command": command("relay-receipt")}],
316 }],
317 "PostToolUseFailure": [{
319 "matcher": "SendMessage",
320 "hooks": [{"type": "command", "command": command("relay-receipt")}],
321 }],
322 "Stop": [{
323 "hooks": [{"type": "command", "command": command("relay-receipt")}],
324 }],
325 "StopFailure": [{
327 "hooks": [{"type": "command", "command": command("relay-receipt")}],
328 }],
329 "UserPromptSubmit": [{
330 "hooks": [{
331 "type": "command",
332 "command": format!("{program} message relay-inbound {represented} {directory}"),
333 }],
334 }],
335 }
336 })
337}
338
339fn shell_quote(value: &str) -> String {
340 format!("'{}'", value.replace('\'', "'\\''"))
341}
342
343pub fn is_inbound_prompt(prompt: &str) -> bool {
346 let trimmed = prompt.trim_start();
347 trimmed.starts_with(NATIVE_PEER_OPENING)
348 || trimmed.starts_with(NATIVE_IDLE_NOTICE)
349 || trimmed.starts_with(NATIVE_DELIVERY_NOTICE)
350}
351
352pub fn inbound_block() -> Value {
354 json!({
355 "decision": "block",
356 "reason": "filed in the mailbox of the session this relay speaks for",
357 })
358}
359
360pub fn send_turn(send: &QueuedSend) -> String {
362 format!(
363 "Send one message with SendMessage.\n\
364 to: {}\n\
365 message: the exact text between the markers, without the markers\n\
366 ---BEGIN MESSAGE---\n{}\n---END MESSAGE---",
367 send.to, send.message
368 )
369}
370
371#[derive(Debug, Clone, PartialEq, Eq)]
373pub struct NativeEnvelope {
374 pub from: String,
376 pub from_name: Option<String>,
378 pub body: String,
380}
381
382#[derive(Debug, Clone, PartialEq, Eq)]
384pub enum RelayEvent {
385 Peer(NativeEnvelope),
387 IdleNotice(String),
389 DeliveryNotice(String),
391}
392
393pub fn resolve_native_sender(
395 registry: &[ClaudePeerSession],
396 native_from: &str,
397) -> Option<ClaudePeerSession> {
398 let socket = native_from.strip_prefix("uds:").unwrap_or(native_from);
399 registry
400 .iter()
401 .find(|session| session.socket_path.to_string_lossy() == socket)
402 .cloned()
403}
404
405#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
407#[serde(rename_all = "snake_case", tag = "outcome")]
408pub enum RelayReceipt {
409 Delivered {
411 detail: String,
413 native_msg_id: Option<String>,
415 },
416 Failed {
418 detail: String,
420 },
421}
422
423pub fn read_receipt(value: &Value) -> RelayReceipt {
426 if let Some(reason) = value.get("denied").and_then(Value::as_str) {
427 return RelayReceipt::Failed {
428 detail: format!("the relay's gate refused the send: {reason}"),
429 };
430 }
431 if let Some(error) = value.get("send_failed").and_then(Value::as_str) {
432 return RelayReceipt::Failed {
433 detail: format!("Claude Code refused the relay's send: {}", error.trim()),
434 };
435 }
436 if value.get("turn_ended").and_then(Value::as_bool) == Some(true) {
437 let said = value
438 .get("said")
439 .and_then(Value::as_str)
440 .filter(|said| !said.trim().is_empty())
441 .map(|said| format!("; it said: {}", said.trim()))
442 .unwrap_or_default();
443 return RelayReceipt::Failed {
444 detail: format!(
445 "the Claude relay ended its turn without sending; nothing was sent{said}"
446 ),
447 };
448 }
449 let response = value.get("tool_response").cloned().unwrap_or(Value::Null);
450 let parsed = match &response {
451 Value::String(text) => serde_json::from_str::<Value>(text).unwrap_or(response.clone()),
452 Value::Array(blocks) => blocks
453 .iter()
454 .find_map(|block| block.get("text").and_then(Value::as_str))
455 .and_then(|text| serde_json::from_str::<Value>(text).ok())
456 .unwrap_or(response.clone()),
457 _ => response.clone(),
458 };
459 if parsed.get("success").and_then(Value::as_bool) != Some(true) {
460 return RelayReceipt::Failed {
461 detail: format!("Claude did not report the send as successful: {parsed}"),
462 };
463 }
464 RelayReceipt::Delivered {
465 detail: parsed
466 .get("message")
467 .and_then(Value::as_str)
468 .unwrap_or_default()
469 .to_string(),
470 native_msg_id: parsed
471 .get("msg_id")
472 .and_then(Value::as_str)
473 .map(str::to_string),
474 }
475}
476
477pub async fn send_through_relay(
487 sender: &MailAddress,
488 sender_name: &str,
489 receiver: &ClaudePeerSession,
490 message: String,
491 message_id: &str,
492) -> RelayReceipt {
493 let failed = |detail: String| RelayReceipt::Failed { detail };
494 let to = receiver_address(receiver);
495 let to = to.as_str();
496 let spec = match RelaySpec::for_sender(sender, sender_name) {
497 Ok(spec) => spec,
498 Err(error) => return failed(error.to_string()),
499 };
500 if let Err(error) = std::fs::create_dir_all(&spec.paths.directory) {
501 return failed(error.to_string());
502 }
503 let lock = match tokio::task::spawn_blocking({
505 let path = spec.paths.lock.clone();
506 move || SendLock::acquire(&path)
507 })
508 .await
509 {
510 Ok(Ok(lock)) => lock,
511 Ok(Err(error)) => return failed(format!("could not take the relay's send lock: {error}")),
512 Err(error) => return failed(error.to_string()),
513 };
514 let mut runtime = match ensure_relay_runtime(&spec).await {
515 Ok(runtime) => runtime,
516 Err(detail) => return failed(detail),
517 };
518 let queued = QueuedSend {
519 to: to.to_string(),
520 message,
521 };
522 std::fs::remove_file(&spec.paths.receipt).ok();
523 if let Err(error) = std::fs::write(
524 &spec.paths.queue,
525 serde_json::to_vec(&queued).unwrap_or_default(),
526 ) {
527 return failed(error.to_string());
528 }
529 let mut attempts = 0;
530 let receipt = loop {
531 let delivered =
532 crate::runtime_mail::deliver_to_runtime(&runtime, send_turn(&queued), true).await;
533 let receipt = match delivered {
534 Err(detail) => failed(format!("could not reach the Claude relay: {detail}")),
535 Ok(_) => wait_for_receipt(&spec.paths.receipt, relay_pid(&spec, &runtime)).await,
536 };
537 if matches!(receipt, RelayReceipt::Delivered { .. })
538 || relay_pid(&spec, &runtime).is_some_and(crate::claude_peer::process_is_live)
539 || attempts == 1
540 {
541 break receipt;
542 }
543 if relay_pid(&spec, &runtime).is_none() {
546 break receipt;
547 }
548 attempts += 1;
549 runtime = match ensure_relay_runtime(&spec).await {
550 Ok(runtime) => runtime,
551 Err(detail) => break failed(detail),
552 };
553 std::fs::remove_file(&spec.paths.receipt).ok();
554 };
555 std::fs::remove_file(&spec.paths.queue).ok();
556 if matches!(receipt, RelayReceipt::Delivered { .. }) {
557 let mut sent = spec.paths.last_sent();
558 sent.insert(receiver.session_id.clone(), message_id.to_string());
559 std::fs::write(
560 &spec.paths.sent,
561 serde_json::to_vec(&sent).unwrap_or_default(),
562 )
563 .ok();
564 }
565 drop(lock);
566 receipt
567}
568
569fn receiver_address(receiver: &ClaudePeerSession) -> String {
572 if receiver.socket_path.as_os_str().is_empty() {
573 receiver.name.clone()
574 } else {
575 format!("uds:{}", receiver.socket_path.display())
576 }
577}
578
579async fn wait_for_receipt(path: &Path, pid: Option<u32>) -> RelayReceipt {
580 let started = Instant::now();
581 while started.elapsed() < SEND_TIMEOUT {
582 if let Some(value) = std::fs::read(path)
583 .ok()
584 .and_then(|bytes| serde_json::from_slice::<Value>(&bytes).ok())
585 {
586 return read_receipt(&value);
587 }
588 if pid.is_some_and(|pid| !crate::claude_peer::process_is_live(pid)) {
589 return RelayReceipt::Failed {
590 detail: "the Claude relay process exited before confirming the send".into(),
591 };
592 }
593 tokio::time::sleep(Duration::from_millis(200)).await;
594 }
595 RelayReceipt::Failed {
596 detail: format!(
597 "the Claude relay did not confirm the send within {} seconds; it may still arrive",
598 SEND_TIMEOUT.as_secs()
599 ),
600 }
601}
602
603#[derive(Serialize, Deserialize)]
604struct RelayRuntimeRecord {
605 runtime_id: String,
606 #[serde(default)]
607 pid: Option<u32>,
608 #[serde(default)]
609 endpoint: String,
610 #[serde(default)]
612 connection: String,
613}
614fn relay_pid(spec: &RelaySpec, runtime: &crate::live_runtime::LiveRuntimeRecord) -> Option<u32> {
617 runtime
618 .metadata
619 .child_pid
620 .or_else(|| {
621 std::fs::read(&spec.paths.record)
622 .ok()
623 .and_then(|bytes| serde_json::from_slice::<RelayRuntimeRecord>(&bytes).ok())
624 .filter(|record| record.runtime_id == runtime.runtime_session_id)
625 .and_then(|record| record.pid)
626 })
627 .or_else(|| {
628 read_registry(®istry_dir(&HarnessHomes::default()))
629 .into_iter()
630 .find(|session| {
631 session.session_id == runtime.runtime_session_id && session.name == spec.name
632 })
633 .map(|session| session.pid)
634 })
635}
636
637async fn ensure_relay_runtime(
640 spec: &RelaySpec,
641) -> Result<crate::live_runtime::LiveRuntimeRecord, String> {
642 let endpoint = relay_endpoint(&spec.program).await?;
643 std::fs::write(
646 &spec.paths.settings,
647 serde_json::to_vec_pretty(&relay_settings(spec)).unwrap_or_default(),
648 )
649 .map_err(|error| error.to_string())?;
650 if let Some(record) = std::fs::read(&spec.paths.record)
651 .ok()
652 .and_then(|bytes| serde_json::from_slice::<RelayRuntimeRecord>(&bytes).ok())
653 {
654 let runtime = crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id);
655 if let Some(runtime) = &runtime {
656 if relay_pid(spec, runtime).is_some_and(crate::claude_peer::process_is_live)
657 && (record.endpoint == endpoint
658 || crate::relay_endpoint::endpoint_answers(&record.endpoint))
659 {
660 return Ok(runtime.clone());
661 }
662 }
663 let result = machine_rpc(
666 &spec.program,
667 "runtimes.close",
668 &json!({ "connection": record.connection, "runtime_id": record.runtime_id }),
669 )
670 .await;
671 if let Err(error) = result {
672 if crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id).is_some()
673 {
674 return Err(format!("could not retire the old relay: {error}"));
675 }
676 }
677 }
678 let params = json!({
679 "harness": "claude-code",
680 "launch": {
681 "program": "claude",
682 "arguments": relay_arguments(spec),
683 "env": relay_environment(spec, &endpoint),
684 },
685 "cwd": spec.paths.directory,
686 });
687 let result = machine_rpc(&spec.program, "runtimes.start", ¶ms).await?;
688 let runtime_id = result
689 .pointer("/handle/runtime_id")
690 .and_then(Value::as_str)
691 .ok_or_else(|| format!("the machine daemon did not start the relay: {result}"))?
692 .to_string();
693 let connection = result
694 .get("connection")
695 .and_then(Value::as_str)
696 .unwrap_or_default()
697 .to_string();
698 let pid = result
699 .pointer("/handle/endpoint/pid")
700 .and_then(Value::as_u64)
701 .map(|pid| pid as u32);
702 std::fs::write(
703 &spec.paths.record,
704 serde_json::to_vec(&RelayRuntimeRecord {
705 pid,
706 runtime_id: runtime_id.clone(),
707 endpoint,
708 connection,
709 })
710 .unwrap_or_default(),
711 )
712 .map_err(|error| error.to_string())?;
713 crate::runtime_mail::controlled_runtime("claude-code", &runtime_id)
714 .ok_or_else(|| "the relay started but registered no live runtime".to_string())
715}
716
717pub async fn retire_relays_of_ended_sessions(homes: &HarnessHomes) -> Vec<String> {
726 #[derive(Deserialize)]
727 struct Record {
728 runtime_id: String,
729 #[serde(default)]
730 connection: String,
731 }
732 let Ok(entries) = std::fs::read_dir(mail_root().join("relays")) else {
733 return Vec::new();
734 };
735 let running: std::collections::HashSet<String> = crate::mail_route::LiveSessions::read(homes)
736 .all()
737 .iter()
738 .map(|session| session.address.to_string())
739 .collect();
740 let Ok(program) = supercode_program() else {
741 return Vec::new();
742 };
743 let mut retired = Vec::new();
744 for entry in entries.flatten() {
745 let paths = RelayPaths::in_directory(entry.path());
746 let Some(represented) = std::fs::read_to_string(&paths.settings)
748 .ok()
749 .and_then(|text| {
750 let from = text.find("relay-inbound '")? + "relay-inbound '".len();
751 let len = text[from..].find('\'')?;
752 MailAddress::parse(&text[from..from + len]).ok()
753 })
754 else {
755 continue;
756 };
757 if matches!(represented.harness.as_str(), "board" | "operator")
758 || running.contains(&represented.to_string())
759 {
760 continue;
761 }
762 let record = std::fs::read(&paths.record)
763 .ok()
764 .and_then(|bytes| serde_json::from_slice::<Record>(&bytes).ok());
765 let alive = record.as_ref().is_some_and(|record| {
766 crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id).is_some()
767 });
768 if alive {
769 let (connection, runtime_id) = record
770 .map(|record| (record.connection, record.runtime_id))
771 .unwrap_or_default();
772 if connection.is_empty()
773 || machine_rpc(
774 &program,
775 "runtimes.close",
776 &json!({ "connection": connection, "runtime_id": runtime_id }),
777 )
778 .await
779 .is_err()
780 {
781 continue;
782 }
783 }
784 if std::fs::remove_dir_all(&paths.directory).is_ok() {
785 retired.push(represented.to_string());
786 }
787 }
788 retired
789}
790
791async fn relay_endpoint(program: &Path) -> Result<String, String> {
794 if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
795 return Ok(url);
796 }
797 ensure_machine_daemon(program).await?;
798 let deadline = Instant::now() + Duration::from_secs(10);
799 while Instant::now() < deadline {
800 if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
801 return Ok(url);
802 }
803 tokio::time::sleep(Duration::from_millis(200)).await;
804 }
805 Err(
806 "the relay endpoint is not answering; the machine daemon's `supercode message watch` \
807 serves it (see mail/machine-daemon.log)"
808 .into(),
809 )
810}
811
812pub async fn machine_rpc(program: &Path, method: &str, params: &Value) -> Result<Value, String> {
815 let started = std::time::Instant::now();
816 let answer = machine_rpc_now(program, method, params).await;
817 crate::slow_log::note_step(
818 &format!("machine rpc {method}"),
819 started.elapsed().as_millis(),
820 );
821 answer
822}
823
824async fn machine_rpc_now(program: &Path, method: &str, params: &Value) -> Result<Value, String> {
825 let call = || {
826 let mut command = tokio::process::Command::new(program);
827 command
828 .args(["teams", "rpc", method, ¶ms.to_string()])
829 .stdin(std::process::Stdio::null())
830 .stdout(std::process::Stdio::piped())
831 .stderr(std::process::Stdio::piped());
832 command.output()
833 };
834 let output = call().await.map_err(|error| error.to_string())?;
835 if output.status.success() {
836 return serde_json::from_slice(&output.stdout)
837 .map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
838 }
839 ensure_machine_daemon(program).await?;
842 let output = call().await.map_err(|error| error.to_string())?;
843 if output.status.success() {
844 return serde_json::from_slice(&output.stdout)
845 .map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
846 }
847 Err(error_line(&String::from_utf8_lossy(&output.stderr)))
848}
849
850pub async fn ensure_machine_daemon(program: &Path) -> Result<(), String> {
853 let describe = || {
854 tokio::process::Command::new(program)
855 .args(["teams", "describe"])
856 .stdin(std::process::Stdio::null())
857 .stdout(std::process::Stdio::null())
858 .stderr(std::process::Stdio::null())
859 .status()
860 };
861 if describe().await.is_ok_and(|status| status.success()) {
862 return Ok(());
863 }
864 if crate::teams::service_owns_daemon() {
868 for _ in 0..100 {
869 tokio::time::sleep(Duration::from_millis(200)).await;
870 if describe().await.is_ok_and(|status| status.success()) {
871 return Ok(());
872 }
873 }
874 return Err(
875 "this machine's daemon is run by its service manager and did not answer within 20 seconds; \
876 see `supercode teams status`"
877 .into(),
878 );
879 }
880 let root = mail_root();
881 std::fs::create_dir_all(&root).map_err(|error| error.to_string())?;
882 let log = std::fs::OpenOptions::new()
883 .create(true)
884 .append(true)
885 .open(root.join("machine-daemon.log"))
886 .map_err(|error| error.to_string())?;
887 let mut command = std::process::Command::new(program);
888 command
889 .args(["teams", "machine", "start", "--supercode"])
890 .arg(program)
891 .stdin(std::process::Stdio::null())
892 .stdout(log.try_clone().map_err(|error| error.to_string())?)
893 .stderr(log);
894 #[cfg(unix)]
895 {
896 use std::os::unix::process::CommandExt;
897 command.process_group(0);
900 }
901 command
902 .spawn()
903 .map_err(|error| format!("could not start the machine daemon: {error}"))?;
904 for _ in 0..50 {
905 tokio::time::sleep(Duration::from_millis(200)).await;
906 if describe().await.is_ok_and(|status| status.success()) {
907 return Ok(());
908 }
909 }
910 Err(format!(
911 "the machine daemon did not start within 10 seconds; see {}",
912 root.join("machine-daemon.log").display()
913 ))
914}
915
916fn error_line(stderr: &str) -> String {
918 let lines: Vec<&str> = stderr
919 .lines()
920 .map(str::trim)
921 .filter(|line| !line.is_empty())
922 .collect();
923 lines
924 .iter()
925 .find(|line| line.starts_with("Error") || line.starts_with("error"))
926 .or(lines.last())
927 .copied()
928 .unwrap_or("the machine daemon failed without saying why")
929 .chars()
930 .take(300)
931 .collect()
932}
933
934struct SendLock {
936 _file: std::fs::File,
937}
938
939impl SendLock {
940 fn acquire(path: &Path) -> std::io::Result<Self> {
941 let file = std::fs::OpenOptions::new()
942 .create(true)
943 .truncate(false)
944 .write(true)
945 .open(path)?;
946 #[cfg(unix)]
947 {
948 use std::os::unix::io::AsRawFd;
949 if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) } != 0 {
952 return Err(std::io::Error::last_os_error());
953 }
954 }
955 Ok(Self { _file: file })
956 }
957}
958
959pub fn file_inbound_prompt(
962 homes: &HarnessHomes,
963 represented: &MailAddress,
964 prompt: &str,
965 last_sent: Option<&HashMap<String, String>>,
966) {
967 match parse_inbound_text(prompt) {
968 Some(RelayEvent::Peer(native)) => {
969 let registry = read_registry(®istry_dir(homes));
970 let sender = resolve_native_sender(®istry, &native.from);
971 let in_reply_to = sender
972 .as_ref()
973 .and_then(|session| last_sent.and_then(|sent| sent.get(&session.session_id).cloned()));
974 file_inbound(represented, &native, sender.as_ref(), in_reply_to);
975 }
976 Some(RelayEvent::IdleNotice(text) | RelayEvent::DeliveryNotice(text)) => {
977 file_notice(represented, &text);
978 }
979 _ => {}
980 }
981}
982pub fn parse_inbound_text(text: &str) -> Option<RelayEvent> {
985 let trimmed = text.trim_start();
986 if trimmed.starts_with(NATIVE_IDLE_NOTICE) {
987 return Some(RelayEvent::IdleNotice(trimmed.to_string()));
988 }
989 if trimmed.starts_with(NATIVE_DELIVERY_NOTICE) {
990 return Some(RelayEvent::DeliveryNotice(trimmed.to_string()));
991 }
992 if !trimmed.starts_with(NATIVE_PEER_PREAMBLE) && !trimmed.starts_with(NATIVE_PEER_OPENING) {
993 return None;
994 }
995 let start = text.find(NATIVE_PEER_OPENING)?;
996 let header_end = start + text[start..].find(">\n")?;
997 let header = &text[start + NATIVE_PEER_OPENING.len()..header_end];
998 let body_start = header_end + 2;
999 let body_end = text.rfind(&format!("\n{NATIVE_PEER_CLOSING}"))?;
1002 if body_end < body_start {
1003 return None;
1004 }
1005 let attributes = parse_attributes(header);
1006 Some(RelayEvent::Peer(NativeEnvelope {
1007 from: attributes.get("from")?.clone(),
1008 from_name: attributes.get("from-name").cloned(),
1009 body: text[body_start..body_end].to_string(),
1010 }))
1011}
1012fn parse_attributes(header: &str) -> BTreeMap<String, String> {
1013 let mut attributes = BTreeMap::new();
1014 let mut rest = header;
1015 while let Some(equals) = rest.find("=\"") {
1016 let key = rest[..equals].trim().to_string();
1017 let value_start = equals + 2;
1018 let Some(length) = rest[value_start..].find('"') else {
1019 break;
1020 };
1021 attributes.insert(key, rest[value_start..value_start + length].to_string());
1022 rest = &rest[value_start + length + 1..];
1023 }
1024 attributes
1025}
1026fn deliver_home(represented: &MailAddress, envelope: &Envelope) {
1030 if let Err(error) = crate::mailbox::deliver_to(represented, envelope) {
1031 eprintln!(
1032 "could not deliver {} to {represented}: {error}; kept in this machine's mailbox for it",
1033 envelope.id
1034 );
1035 if let Ok(mailbox) = crate::mailbox::Mailbox::open(&mail_root(), represented) {
1036 mailbox.deliver(envelope).ok();
1037 }
1038 return;
1039 }
1040 if represented.machine == crate::mailbox::local_machine_name() {
1045 if let Ok(program) = supercode_program() {
1046 std::process::Command::new(program)
1047 .args(["message", "push", &represented.to_string(), &envelope.id])
1048 .stdin(std::process::Stdio::null())
1049 .stdout(std::process::Stdio::null())
1050 .stderr(std::process::Stdio::null())
1051 .spawn()
1052 .ok();
1053 }
1054 }
1055}
1056fn file_inbound(
1057 represented: &MailAddress,
1058 native: &NativeEnvelope,
1059 sender: Option<&ClaudePeerSession>,
1060 in_reply_to: Option<String>,
1061) {
1062 let machine = crate::mailbox::local_machine_name();
1065 let (from, from_name) = match sender {
1066 Some(session) => (
1067 MailAddress::new(&machine, "claude-code", &session.session_id),
1068 format!("{}@{machine}", session.name),
1069 ),
1070 None => (
1071 MailAddress::new(&machine, "claude-code", "unknown"),
1072 native
1073 .from_name
1074 .clone()
1075 .unwrap_or_else(|| "an unknown Claude session".into()),
1076 ),
1077 };
1078 let Ok(from) = from else { return };
1079 let Ok(mut envelope) = Envelope::new(
1080 from,
1081 from_name,
1082 MailKind::Peer,
1083 ReplyVia::Command,
1084 native.body.clone(),
1085 ) else {
1086 return;
1087 };
1088 envelope.native_from = Some(native.from.clone());
1089 if let Some(in_reply_to) = in_reply_to {
1090 envelope.thread = crate::mailbox::thread_of_reply(Some(&in_reply_to));
1091 envelope.in_reply_to = Some(in_reply_to);
1092 envelope.in_reply_to_inferred = true;
1093 }
1094 deliver_home(represented, &envelope);
1097}
1098fn file_notice(represented: &MailAddress, text: &str) {
1099 let Ok(from) = MailAddress::new(
1100 crate::mailbox::local_machine_name(),
1101 "claude-code",
1102 "notice",
1103 ) else {
1104 return;
1105 };
1106 let Ok(envelope) = Envelope::new(
1107 from,
1108 "Claude Code",
1109 MailKind::Notice,
1110 ReplyVia::None,
1111 text.to_string(),
1112 ) else {
1113 return;
1114 };
1115 deliver_home(represented, &envelope);
1118}
1119pub fn supercode_program() -> std::io::Result<PathBuf> {
1122 let current = std::env::current_exe()?;
1123 if current.file_stem().and_then(|stem| stem.to_str()) == Some("supercode") {
1124 return Ok(current);
1125 }
1126 std::env::var_os("PATH")
1127 .iter()
1128 .flat_map(std::env::split_paths)
1129 .map(|directory| directory.join("supercode"))
1130 .find(|candidate| candidate.is_file())
1131 .ok_or_else(|| {
1132 std::io::Error::new(
1133 std::io::ErrorKind::NotFound,
1134 "the supercode program is not on PATH; Claude relays and the machine daemon need it",
1135 )
1136 })
1137}