Skip to main content

supercode_harness/
claude_relay.rs

1//! Claude relays: the degraded tier's door into a Claude Code session that
2//! supercode does not control.
3//!
4//! A Claude session answers a message by copying its `from` into
5//! `SendMessage`. For that answer to come back, the message must come FROM a
6//! live Claude peer. A relay is that peer: a Claude Code runtime that
7//! supercode hosts like any other (`runtimes.start` in the machine daemon's
8//! `harness serve`, registered in the live-runtime registry), named
9//! `sc-<name>-on-<machine>`, sending on one session's behalf. Nothing here
10//! runs its own process supervisor or writes Claude's private socket frame.
11//!
12//! Three rules earn their place here:
13//!
14//! 1. **No model runs.** A relay's API is supercode's own deterministic
15//!    endpoint ([`crate::relay_endpoint`]), which answers a send turn with a
16//!    `SendMessage` of the queued pair. Every send is queued first, and a
17//!    PreToolUse hook (`supercode message gate`) still denies any tool call
18//!    but a `SendMessage` of exactly that pair.
19//! 2. **Receipts come from Claude, through hooks.** A PostToolUse hook writes
20//!    Claude's own `SendMessage` result beside the queue; a denied call and a
21//!    turn that ends without sending write theirs too. The sender waits on
22//!    that file, not on the model's words.
23//! 3. **Inbound mail never reaches the model.** Claude hands every inbound
24//!    peer envelope and idle or delivery notice to the relay's
25//!    `UserPromptSubmit` hook byte-exact; the hook files it in the represented
26//!    session's mailbox and blocks the prompt, so no model turn runs for it
27//!    (measured: `duration_api_ms: 0`).
28//!
29//! A relay's busy/idle is its own; its name carries the `sc-` prefix and
30//! every turn ends with [`RELAY_STATUS_LINE`].
31
32use 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
43/// Model name a relay's requests carry. The relay endpoint answers them
44/// without any model.
45pub const RELAY_MODEL: &str = "haiku";
46
47/// Prefix marking a registry name as a relay rather than a session.
48pub const RELAY_NAME_PREFIX: &str = "sc-";
49
50/// Text every relay turn ends with. Claude quotes a session's last line in
51/// its idle notices, so this is where the correction reaches the agent.
52pub const RELAY_STATUS_LINE: &str = "relay, not the session's status";
53
54/// How long a send waits for Claude's own receipt.
55pub const SEND_TIMEOUT: Duration = Duration::from_secs(90);
56
57/// Tools a relay may touch.
58const RELAY_TOOLS: &str = "ListAgents,SendMessage";
59
60/// Line Claude opens a delivered peer message with in a turn. The hook sees
61/// the envelope itself; one of the two starts is required, so text that
62/// merely quotes an envelope further in is never read as mail.
63const NATIVE_PEER_PREAMBLE: &str = "Another Claude session sent a message:";
64/// Tags of Claude's peer envelope.
65const 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
70/// Registry name of the relay that represents `name` on `machine`.
71///
72/// Not `name@machine`: Claude accepts `@` in `--name`, but its `SendMessage`
73/// refuses a `to` containing `@` ("must be a bare teammate name").
74pub 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/// The one send a relay is allowed to make next.
80#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
81pub struct QueuedSend {
82    /// Registry name of the receiving Claude session.
83    pub to: String,
84    /// Exact message text.
85    pub message: String,
86}
87
88/// Decide one PreToolUse hook call. `None` allows it; `Some(reason)` denies.
89///
90/// Only a `SendMessage` whose `to` and `message` equal the queued pair is
91/// allowed; `ListAgents` is read-only and always allowed; anything else, or a
92/// send with nothing queued, is denied.
93pub 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
151/// The hook's JSON answer for a denial.
152pub fn gate_denial(reason: &str) -> Value {
153    json!({
154        "hookSpecificOutput": {
155            "hookEventName": "PreToolUse",
156            "permissionDecision": "deny",
157            "permissionDecisionReason": reason,
158        }
159    })
160}
161
162/// Files of one relay, all in its own directory (also its working directory,
163/// which holds no project settings).
164#[derive(Debug, Clone)]
165pub struct RelayPaths {
166    /// The relay's directory.
167    pub directory: PathBuf,
168    /// The single queued send the gate compares against.
169    pub queue: PathBuf,
170    /// Claude's answer to the queued send, written by the relay's hooks.
171    pub receipt: PathBuf,
172    /// The prompt that began the relay's current turn, written as the turn starts: a turn's end
173    /// answers the queued send only when that send began it.
174    pub turn: PathBuf,
175    /// Settings passed with `--settings`, carrying the hooks.
176    pub settings: PathBuf,
177    /// The hosted runtime this relay runs as.
178    pub record: PathBuf,
179    /// Newest message id sent to each Claude session name, for threading.
180    pub sent: PathBuf,
181    /// Serializes sends through this relay.
182    pub lock: PathBuf,
183}
184
185impl RelayPaths {
186    /// Paths of the relay named `relay_name` under `root`.
187    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    /// Paths of the relay whose directory is `directory`.
193    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    /// Newest message id sent to each Claude session name.
207    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/// Everything needed to start one relay.
216#[derive(Debug, Clone)]
217pub struct RelaySpec {
218    /// Session the relay speaks for.
219    pub represented: MailAddress,
220    /// Name the represented session is known by.
221    pub represented_name: String,
222    /// Registry name the relay registers under.
223    pub name: String,
224    /// Relay files.
225    pub paths: RelayPaths,
226    /// Program run by the hooks (this supercode binary).
227    pub program: PathBuf,
228}
229
230impl RelaySpec {
231    /// The relay that speaks for `represented`.
232    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
248/// Arguments of the Claude Code runtime a relay runs as.
249///
250/// `--setting-sources project` in an empty directory keeps user hooks and
251/// settings out; `--settings` then adds only the relay's hooks. The API is
252/// the relay endpoint, through the environment ([`relay_environment`]).
253pub 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
279/// Environment of a relay: its API is the relay endpoint at `endpoint`, and
280/// its API key names its directory there. The key also keeps account
281/// connectors (MCP servers) out of the relay.
282pub 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
300/// Settings document installing a relay's hooks.
301pub 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            // Every tool, not only the two named in --tools.
309            "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            // A SendMessage that fails fires this, with Claude Code's own error.
318            "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            // A turn an API error ends (a usage limit, an overload) fires this, not Stop.
326            "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
343/// Whether a prompt Claude handed the relay is inbound mail (a peer envelope
344/// or a notice) rather than its host's own send turn.
345pub 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
352/// The UserPromptSubmit answer that keeps a prompt from the model.
353pub 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
360/// The host turn asking a relay to make one send.
361pub 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/// A native peer envelope as Claude delivered it to a relay.
372#[derive(Debug, Clone, PartialEq, Eq)]
373pub struct NativeEnvelope {
374    /// Claude's `from` (the sender's socket, `uds:/tmp/cc-socks/<pid>.sock`).
375    pub from: String,
376    /// Claude's `from-name`.
377    pub from_name: Option<String>,
378    /// Message body, unescaped as Claude delivered it.
379    pub body: String,
380}
381
382/// One inbound prompt, as Claude handed it to a relay's hook.
383#[derive(Debug, Clone, PartialEq, Eq)]
384pub enum RelayEvent {
385    /// A peer message arrived.
386    Peer(NativeEnvelope),
387    /// A `[Cross-session idle notice]`.
388    IdleNotice(String),
389    /// A `[Cross-session delivery notice]`.
390    DeliveryNotice(String),
391}
392
393/// The live Claude session behind a native `uds:` sender address.
394pub 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/// Outcome of one relayed send, as the sender reads it.
406#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
407#[serde(rename_all = "snake_case", tag = "outcome")]
408pub enum RelayReceipt {
409    /// Claude reported the message queued in the receiver's inbox.
410    Delivered {
411        /// Claude's own result text.
412        detail: String,
413        /// Claude's message id, when it reported one.
414        native_msg_id: Option<String>,
415    },
416    /// The send did not happen.
417    Failed {
418        /// Why, in words the sending agent can act on.
419        detail: String,
420    },
421}
422
423/// Read a receipt the relay's hooks wrote: Claude's own `SendMessage`
424/// result, a gate denial, or a turn that ended without sending.
425pub 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
477/// Send `message` to the Claude session named `to`, from `sender` through
478/// its relay: start the relay as a hosted runtime when it is not running,
479/// queue the send, hand it the send turn through the runtime's own door, and
480/// wait for Claude's receipt.
481pub async fn send_through_relay(
482    sender: &MailAddress,
483    sender_name: &str,
484    to: &str,
485    message: String,
486    message_id: &str,
487) -> RelayReceipt {
488    let failed = |detail: String| RelayReceipt::Failed { detail };
489    let spec = match RelaySpec::for_sender(sender, sender_name) {
490        Ok(spec) => spec,
491        Err(error) => return failed(error.to_string()),
492    };
493    if let Err(error) = std::fs::create_dir_all(&spec.paths.directory) {
494        return failed(error.to_string());
495    }
496    // One send at a time through one relay: the queue file is the gate.
497    let lock = match tokio::task::spawn_blocking({
498        let path = spec.paths.lock.clone();
499        move || SendLock::acquire(&path)
500    })
501    .await
502    {
503        Ok(Ok(lock)) => lock,
504        Ok(Err(error)) => return failed(format!("could not take the relay's send lock: {error}")),
505        Err(error) => return failed(error.to_string()),
506    };
507    let mut runtime = match ensure_relay_runtime(&spec).await {
508        Ok(runtime) => runtime,
509        Err(detail) => return failed(detail),
510    };
511    let queued = QueuedSend {
512        to: to.to_string(),
513        message,
514    };
515    std::fs::remove_file(&spec.paths.receipt).ok();
516    if let Err(error) = std::fs::write(
517        &spec.paths.queue,
518        serde_json::to_vec(&queued).unwrap_or_default(),
519    ) {
520        return failed(error.to_string());
521    }
522    let mut attempts = 0;
523    let receipt = loop {
524        let delivered =
525            crate::runtime_mail::deliver_to_runtime(&runtime, send_turn(&queued), true).await;
526        let receipt = match delivered {
527            Err(detail) => failed(format!("could not reach the Claude relay: {detail}")),
528            Ok(_) => wait_for_receipt(&spec.paths.receipt, relay_pid(&spec, &runtime)).await,
529        };
530        if matches!(receipt, RelayReceipt::Delivered { .. })
531            || relay_pid(&spec, &runtime).is_some_and(crate::claude_peer::process_is_live)
532            || attempts == 1
533        {
534            break receipt;
535        }
536        // Retain the same queue under the send lock. Only a proven-dead child
537        // permits a retry; an alive but slow relay may still deliver.
538        if relay_pid(&spec, &runtime).is_none() {
539            break receipt;
540        }
541        attempts += 1;
542        runtime = match ensure_relay_runtime(&spec).await {
543            Ok(runtime) => runtime,
544            Err(detail) => break failed(detail),
545        };
546        std::fs::remove_file(&spec.paths.receipt).ok();
547    };
548    std::fs::remove_file(&spec.paths.queue).ok();
549    if matches!(receipt, RelayReceipt::Delivered { .. }) {
550        let mut sent = spec.paths.last_sent();
551        sent.insert(to.to_string(), message_id.to_string());
552        std::fs::write(
553            &spec.paths.sent,
554            serde_json::to_vec(&sent).unwrap_or_default(),
555        )
556        .ok();
557    }
558    drop(lock);
559    receipt
560}
561
562async fn wait_for_receipt(path: &Path, pid: Option<u32>) -> RelayReceipt {
563    let started = Instant::now();
564    while started.elapsed() < SEND_TIMEOUT {
565        if let Some(value) = std::fs::read(path)
566            .ok()
567            .and_then(|bytes| serde_json::from_slice::<Value>(&bytes).ok())
568        {
569            return read_receipt(&value);
570        }
571        if pid.is_some_and(|pid| !crate::claude_peer::process_is_live(pid)) {
572            return RelayReceipt::Failed {
573                detail: "the Claude relay process exited before confirming the send".into(),
574            };
575        }
576        tokio::time::sleep(Duration::from_millis(200)).await;
577    }
578    RelayReceipt::Failed {
579        detail: format!(
580            "the Claude relay did not confirm the send within {} seconds; it may still arrive",
581            SEND_TIMEOUT.as_secs()
582        ),
583    }
584}
585
586#[derive(Serialize, Deserialize)]
587struct RelayRuntimeRecord {
588    runtime_id: String,
589    #[serde(default)]
590    pid: Option<u32>,
591    #[serde(default)]
592    endpoint: String,
593    /// The machine daemon's connection to the runtime: the door that closes it.
594    #[serde(default)]
595    connection: String,
596}
597/// The child's PID, including legacy relay records whose host PID was all the
598/// runtime receipt kept. Claude's native registry names the actual speaker.
599fn relay_pid(spec: &RelaySpec, runtime: &crate::live_runtime::LiveRuntimeRecord) -> Option<u32> {
600    runtime
601        .metadata
602        .child_pid
603        .or_else(|| {
604            std::fs::read(&spec.paths.record)
605                .ok()
606                .and_then(|bytes| serde_json::from_slice::<RelayRuntimeRecord>(&bytes).ok())
607                .filter(|record| record.runtime_id == runtime.runtime_session_id)
608                .and_then(|record| record.pid)
609        })
610        .or_else(|| {
611            read_registry(&registry_dir(&HarnessHomes::default()))
612                .into_iter()
613                .find(|session| {
614                    session.session_id == runtime.runtime_session_id && session.name == spec.name
615                })
616                .map(|session| session.pid)
617        })
618}
619
620/// The relay's hosted runtime, started in the machine daemon when it is not
621/// running.
622async fn ensure_relay_runtime(
623    spec: &RelaySpec,
624) -> Result<crate::live_runtime::LiveRuntimeRecord, String> {
625    let endpoint = relay_endpoint(&spec.program).await?;
626    // Written on every send, not only at start: Claude Code reloads a changed settings file, so
627    // a running relay takes the hooks this supercode installs.
628    std::fs::write(
629        &spec.paths.settings,
630        serde_json::to_vec_pretty(&relay_settings(spec)).unwrap_or_default(),
631    )
632    .map_err(|error| error.to_string())?;
633    if let Some(record) = std::fs::read(&spec.paths.record)
634        .ok()
635        .and_then(|bytes| serde_json::from_slice::<RelayRuntimeRecord>(&bytes).ok())
636    {
637        let runtime = crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id);
638        if let Some(runtime) = &runtime {
639            if relay_pid(spec, runtime).is_some_and(crate::claude_peer::process_is_live)
640                && (record.endpoint == endpoint
641                    || crate::relay_endpoint::endpoint_answers(&record.endpoint))
642            {
643                return Ok(runtime.clone());
644            }
645        }
646        // Close through the owner's door even after the receipt disappeared.
647        // The runtime id guard cannot close a reused daemon connection.
648        let result = machine_rpc(
649            &spec.program,
650            "runtimes.close",
651            &json!({ "connection": record.connection, "runtime_id": record.runtime_id }),
652        )
653        .await;
654        if let Err(error) = result {
655            if crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id).is_some()
656            {
657                return Err(format!("could not retire the old relay: {error}"));
658            }
659        }
660    }
661    let params = json!({
662        "harness": "claude-code",
663        "launch": {
664            "program": "claude",
665            "arguments": relay_arguments(spec),
666            "env": relay_environment(spec, &endpoint),
667        },
668        "cwd": spec.paths.directory,
669    });
670    let result = machine_rpc(&spec.program, "runtimes.start", &params).await?;
671    let runtime_id = result
672        .pointer("/handle/runtime_id")
673        .and_then(Value::as_str)
674        .ok_or_else(|| format!("the machine daemon did not start the relay: {result}"))?
675        .to_string();
676    let connection = result
677        .get("connection")
678        .and_then(Value::as_str)
679        .unwrap_or_default()
680        .to_string();
681    let pid = result
682        .pointer("/handle/endpoint/pid")
683        .and_then(Value::as_u64)
684        .map(|pid| pid as u32);
685    std::fs::write(
686        &spec.paths.record,
687        serde_json::to_vec(&RelayRuntimeRecord {
688            pid,
689            runtime_id: runtime_id.clone(),
690            endpoint,
691            connection,
692        })
693        .unwrap_or_default(),
694    )
695    .map_err(|error| error.to_string())?;
696    crate::runtime_mail::controlled_runtime("claude-code", &runtime_id)
697        .ok_or_else(|| "the relay started but registered no live runtime".to_string())
698}
699
700/// The relays of sessions no longer running, retired: each relay's runtime (a `claude --print` under
701/// `supercode harness serve`, some hundreds of MB) is closed through the machine daemon that
702/// started it (`runtimes.close`, so the daemon's own table forgets it: a relay process ended behind
703/// its back left sends failing until the daemon restarted), and its folder removed. A relay speaks
704/// for one session and receives the replies to what that session sent; once the session has ended,
705/// nothing reaches it. A relay for a board or an operator, which is not a session, is kept; so is a
706/// live relay recorded before its connection was kept, which no door can close cleanly.
707/// Answers the represented addresses whose relays were retired.
708pub async fn retire_relays_of_ended_sessions(homes: &HarnessHomes) -> Vec<String> {
709    #[derive(Deserialize)]
710    struct Record {
711        runtime_id: String,
712        #[serde(default)]
713        connection: String,
714    }
715    let Ok(entries) = std::fs::read_dir(mail_root().join("relays")) else {
716        return Vec::new();
717    };
718    let running: std::collections::HashSet<String> = crate::mail_route::LiveSessions::read(homes)
719        .all()
720        .iter()
721        .map(|session| session.address.to_string())
722        .collect();
723    let Ok(program) = supercode_program() else {
724        return Vec::new();
725    };
726    let mut retired = Vec::new();
727    for entry in entries.flatten() {
728        let paths = RelayPaths::in_directory(entry.path());
729        // the address the relay speaks for, as its inbound hook names it
730        let Some(represented) = std::fs::read_to_string(&paths.settings)
731            .ok()
732            .and_then(|text| {
733                let from = text.find("relay-inbound '")? + "relay-inbound '".len();
734                let len = text[from..].find('\'')?;
735                MailAddress::parse(&text[from..from + len]).ok()
736            })
737        else {
738            continue;
739        };
740        if matches!(represented.harness.as_str(), "board" | "operator")
741            || running.contains(&represented.to_string())
742        {
743            continue;
744        }
745        let record = std::fs::read(&paths.record)
746            .ok()
747            .and_then(|bytes| serde_json::from_slice::<Record>(&bytes).ok());
748        let alive = record.as_ref().is_some_and(|record| {
749            crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id).is_some()
750        });
751        if alive {
752            let (connection, runtime_id) = record
753                .map(|record| (record.connection, record.runtime_id))
754                .unwrap_or_default();
755            if connection.is_empty()
756                || machine_rpc(
757                    &program,
758                    "runtimes.close",
759                    &json!({ "connection": connection, "runtime_id": runtime_id }),
760                )
761                .await
762                .is_err()
763            {
764                continue;
765            }
766        }
767        if std::fs::remove_dir_all(&paths.directory).is_ok() {
768            retired.push(represented.to_string());
769        }
770    }
771    retired
772}
773
774/// The relay endpoint's base URL. The machine daemon's watcher serves it;
775/// the daemon is started when it is not running.
776async fn relay_endpoint(program: &Path) -> Result<String, String> {
777    if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
778        return Ok(url);
779    }
780    ensure_machine_daemon(program).await?;
781    let deadline = Instant::now() + Duration::from_secs(10);
782    while Instant::now() < deadline {
783        if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
784            return Ok(url);
785        }
786        tokio::time::sleep(Duration::from_millis(200)).await;
787    }
788    Err(
789        "the relay endpoint is not answering; the machine daemon's `supercode message watch` \
790         serves it (see mail/machine-daemon.log)"
791            .into(),
792    )
793}
794
795/// One `harness.v1` call into this machine's daemon (`supercode teams rpc`),
796/// starting the daemon when it is not running.
797pub async fn machine_rpc(program: &Path, method: &str, params: &Value) -> Result<Value, String> {
798    let started = std::time::Instant::now();
799    let answer = machine_rpc_now(program, method, params).await;
800    crate::slow_log::note_step(
801        &format!("machine rpc {method}"),
802        started.elapsed().as_millis(),
803    );
804    answer
805}
806
807async fn machine_rpc_now(program: &Path, method: &str, params: &Value) -> Result<Value, String> {
808    let call = || {
809        let mut command = tokio::process::Command::new(program);
810        command
811            .args(["teams", "rpc", method, &params.to_string()])
812            .stdin(std::process::Stdio::null())
813            .stdout(std::process::Stdio::piped())
814            .stderr(std::process::Stdio::piped());
815        command.output()
816    };
817    let output = call().await.map_err(|error| error.to_string())?;
818    if output.status.success() {
819        return serde_json::from_slice(&output.stdout)
820            .map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
821    }
822    // Not running: start it (this OS user's own daemon, no Teams context
823    // needed) and ask once more.
824    ensure_machine_daemon(program).await?;
825    let output = call().await.map_err(|error| error.to_string())?;
826    if output.status.success() {
827        return serde_json::from_slice(&output.stdout)
828            .map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
829    }
830    Err(error_line(&String::from_utf8_lossy(&output.stderr)))
831}
832
833/// Start this OS user's machine daemon when it is not answering. It hosts
834/// relays (in its `harness serve`) and the idle watcher.
835pub async fn ensure_machine_daemon(program: &Path) -> Result<(), String> {
836    let describe = || {
837        tokio::process::Command::new(program)
838            .args(["teams", "describe"])
839            .stdin(std::process::Stdio::null())
840            .stdout(std::process::Stdio::null())
841            .stderr(std::process::Stdio::null())
842            .status()
843    };
844    if describe().await.is_ok_and(|status| status.success()) {
845        return Ok(());
846    }
847    // A service manager runs this user's daemon (a connector, or the machine
848    // unit): it is between restarts, and a daemon started here would take the
849    // socket and lock the service's daemon out. Wait for the service instead.
850    if crate::teams::service_owns_daemon() {
851        for _ in 0..100 {
852            tokio::time::sleep(Duration::from_millis(200)).await;
853            if describe().await.is_ok_and(|status| status.success()) {
854                return Ok(());
855            }
856        }
857        return Err(
858            "this machine's daemon is run by its service manager and did not answer within 20 seconds; \
859             see `supercode teams status`"
860                .into(),
861        );
862    }
863    let root = mail_root();
864    std::fs::create_dir_all(&root).map_err(|error| error.to_string())?;
865    let log = std::fs::OpenOptions::new()
866        .create(true)
867        .append(true)
868        .open(root.join("machine-daemon.log"))
869        .map_err(|error| error.to_string())?;
870    let mut command = std::process::Command::new(program);
871    command
872        .args(["teams", "machine", "start", "--supercode"])
873        .arg(program)
874        .stdin(std::process::Stdio::null())
875        .stdout(log.try_clone().map_err(|error| error.to_string())?)
876        .stderr(log);
877    #[cfg(unix)]
878    {
879        use std::os::unix::process::CommandExt;
880        // Its own process group, so the caller's exit does not take the
881        // daemon (and every relay) with it.
882        command.process_group(0);
883    }
884    command
885        .spawn()
886        .map_err(|error| format!("could not start the machine daemon: {error}"))?;
887    for _ in 0..50 {
888        tokio::time::sleep(Duration::from_millis(200)).await;
889        if describe().await.is_ok_and(|status| status.success()) {
890            return Ok(());
891        }
892    }
893    Err(format!(
894        "the machine daemon did not start within 10 seconds; see {}",
895        root.join("machine-daemon.log").display()
896    ))
897}
898
899/// The one line of a failed command's stderr an agent can act on.
900fn error_line(stderr: &str) -> String {
901    let lines: Vec<&str> = stderr
902        .lines()
903        .map(str::trim)
904        .filter(|line| !line.is_empty())
905        .collect();
906    lines
907        .iter()
908        .find(|line| line.starts_with("Error") || line.starts_with("error"))
909        .or(lines.last())
910        .copied()
911        .unwrap_or("the machine daemon failed without saying why")
912        .chars()
913        .take(300)
914        .collect()
915}
916
917/// An exclusive lock on a relay's sends, released on drop.
918struct SendLock {
919    _file: std::fs::File,
920}
921
922impl SendLock {
923    fn acquire(path: &Path) -> std::io::Result<Self> {
924        let file = std::fs::OpenOptions::new()
925            .create(true)
926            .truncate(false)
927            .write(true)
928            .open(path)?;
929        #[cfg(unix)]
930        {
931            use std::os::unix::io::AsRawFd;
932            // SAFETY: flock on a descriptor this function owns; the lock is
933            // released when the file is dropped.
934            if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) } != 0 {
935                return Err(std::io::Error::last_os_error());
936            }
937        }
938        Ok(Self { _file: file })
939    }
940}
941
942/// File one inbound prompt in `represented`'s mailbox. `last_sent` maps a
943/// Claude session name to the newest message this relay sent it, for
944pub fn file_inbound_prompt(
945    homes: &HarnessHomes,
946    represented: &MailAddress,
947    prompt: &str,
948    last_sent: Option<&HashMap<String, String>>,
949) {
950    match parse_inbound_text(prompt) {
951        Some(RelayEvent::Peer(native)) => {
952            let registry = read_registry(&registry_dir(homes));
953            let sender = resolve_native_sender(&registry, &native.from);
954            let in_reply_to = sender
955                .as_ref()
956                .and_then(|session| last_sent.and_then(|sent| sent.get(&session.name).cloned()));
957            file_inbound(represented, &native, sender.as_ref(), in_reply_to);
958        }
959        Some(RelayEvent::IdleNotice(text) | RelayEvent::DeliveryNotice(text)) => {
960            file_notice(represented, &text);
961        }
962        _ => {}
963    }
964}
965/// Recognize an inbound peer envelope or notice in the text Claude handed a
966/// relay.
967pub fn parse_inbound_text(text: &str) -> Option<RelayEvent> {
968    let trimmed = text.trim_start();
969    if trimmed.starts_with(NATIVE_IDLE_NOTICE) {
970        return Some(RelayEvent::IdleNotice(trimmed.to_string()));
971    }
972    if trimmed.starts_with(NATIVE_DELIVERY_NOTICE) {
973        return Some(RelayEvent::DeliveryNotice(trimmed.to_string()));
974    }
975    if !trimmed.starts_with(NATIVE_PEER_PREAMBLE) && !trimmed.starts_with(NATIVE_PEER_OPENING) {
976        return None;
977    }
978    let start = text.find(NATIVE_PEER_OPENING)?;
979    let header_end = start + text[start..].find(">\n")?;
980    let header = &text[start + NATIVE_PEER_OPENING.len()..header_end];
981    let body_start = header_end + 2;
982    // The body is not escaped by Claude; the closing tag is the LAST one, and
983    // the fixed paragraph after it never contains the tag.
984    let body_end = text.rfind(&format!("\n{NATIVE_PEER_CLOSING}"))?;
985    if body_end < body_start {
986        return None;
987    }
988    let attributes = parse_attributes(header);
989    Some(RelayEvent::Peer(NativeEnvelope {
990        from: attributes.get("from")?.clone(),
991        from_name: attributes.get("from-name").cloned(),
992        body: text[body_start..body_end].to_string(),
993    }))
994}
995fn parse_attributes(header: &str) -> BTreeMap<String, String> {
996    let mut attributes = BTreeMap::new();
997    let mut rest = header;
998    while let Some(equals) = rest.find("=\"") {
999        let key = rest[..equals].trim().to_string();
1000        let value_start = equals + 2;
1001        let Some(length) = rest[value_start..].find('"') else {
1002            break;
1003        };
1004        attributes.insert(key, rest[value_start..value_start + length].to_string());
1005        rest = &rest[value_start + length + 1..];
1006    }
1007    attributes
1008}
1009/// Deliver mail to the session a relay speaks for. When it lives on another
1010/// machine and Teams cannot take the message there, it is kept in this
1011/// machine's copy of that mailbox (never dropped) and the failure is written
1012fn deliver_home(represented: &MailAddress, envelope: &Envelope) {
1013    if let Err(error) = crate::mailbox::deliver_to(represented, envelope) {
1014        eprintln!(
1015            "could not deliver {} to {represented}: {error}; kept in this machine's mailbox for it",
1016            envelope.id
1017        );
1018        if let Ok(mailbox) = crate::mailbox::Mailbox::open(&mail_root(), represented) {
1019            mailbox.deliver(envelope).ok();
1020        }
1021        return;
1022    }
1023    // Filed; now hand it to its session by the session's own door, as any
1024    // sender's message is (a Claude session gets it at its next tool round).
1025    // That can take as long as a relay send, which this hook must not wait
1026    // on, so a detached `message push` does it.
1027    if represented.machine == crate::mailbox::local_machine_name() {
1028        if let Ok(program) = supercode_program() {
1029            std::process::Command::new(program)
1030                .args(["message", "push", &represented.to_string(), &envelope.id])
1031                .stdin(std::process::Stdio::null())
1032                .stdout(std::process::Stdio::null())
1033                .stderr(std::process::Stdio::null())
1034                .spawn()
1035                .ok();
1036        }
1037    }
1038}
1039fn file_inbound(
1040    represented: &MailAddress,
1041    native: &NativeEnvelope,
1042    sender: Option<&ClaudePeerSession>,
1043    in_reply_to: Option<String>,
1044) {
1045    // The replying Claude session runs beside the relay, on this machine —
1046    // not on the machine of the session the relay speaks for.
1047    let machine = crate::mailbox::local_machine_name();
1048    let (from, from_name) = match sender {
1049        Some(session) => (
1050            MailAddress::new(&machine, "claude-code", &session.session_id),
1051            format!("{}@{machine}", session.name),
1052        ),
1053        None => (
1054            MailAddress::new(&machine, "claude-code", "unknown"),
1055            native
1056                .from_name
1057                .clone()
1058                .unwrap_or_else(|| "an unknown Claude session".into()),
1059        ),
1060    };
1061    let Ok(from) = from else { return };
1062    let Ok(mut envelope) = Envelope::new(
1063        from,
1064        from_name,
1065        MailKind::Peer,
1066        ReplyVia::Command,
1067        native.body.clone(),
1068    ) else {
1069        return;
1070    };
1071    envelope.native_from = Some(native.from.clone());
1072    if let Some(in_reply_to) = in_reply_to {
1073        envelope.thread = crate::mailbox::thread_of_reply(Some(&in_reply_to));
1074        envelope.in_reply_to = Some(in_reply_to);
1075        envelope.in_reply_to_inferred = true;
1076    }
1077    // The session a relay speaks for may be on another machine; its mail goes
1078    // home through Teams.
1079    deliver_home(represented, &envelope);
1080}
1081fn file_notice(represented: &MailAddress, text: &str) {
1082    let Ok(from) = MailAddress::new(
1083        crate::mailbox::local_machine_name(),
1084        "claude-code",
1085        "notice",
1086    ) else {
1087        return;
1088    };
1089    let Ok(envelope) = Envelope::new(
1090        from,
1091        "Claude Code",
1092        MailKind::Notice,
1093        ReplyVia::None,
1094        text.to_string(),
1095    ) else {
1096        return;
1097    };
1098    // The session a relay speaks for may be on another machine; its mail goes
1099    // home through Teams.
1100    deliver_home(represented, &envelope);
1101}
1102/// The `supercode` program that the relay hooks and the machine daemon run:
1103/// this process when it is `supercode` (the CLI, `supercode harness serve`),
1104pub fn supercode_program() -> std::io::Result<PathBuf> {
1105    let current = std::env::current_exe()?;
1106    if current.file_stem().and_then(|stem| stem.to_str()) == Some("supercode") {
1107        return Ok(current);
1108    }
1109    std::env::var_os("PATH")
1110        .iter()
1111        .flat_map(std::env::split_paths)
1112        .map(|directory| directory.join("supercode"))
1113        .find(|candidate| candidate.is_file())
1114        .ok_or_else(|| {
1115            std::io::Error::new(
1116                std::io::ErrorKind::NotFound,
1117                "the supercode program is not on PATH; Claude relays and the machine daemon need it",
1118            )
1119        })
1120}