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    /// Settings passed with `--settings`, carrying the hooks.
173    pub settings: PathBuf,
174    /// The hosted runtime this relay runs as.
175    pub record: PathBuf,
176    /// Newest message id sent to each Claude session name, for threading.
177    pub sent: PathBuf,
178    /// Serializes sends through this relay.
179    pub lock: PathBuf,
180}
181
182impl RelayPaths {
183    /// Paths of the relay named `relay_name` under `root`.
184    pub fn new(root: &Path, relay_name: &str) -> Self {
185        let hash = blake3::hash(relay_name.as_bytes()).to_hex();
186        Self::in_directory(root.join("relays").join(&hash[..24]))
187    }
188
189    /// Paths of the relay whose directory is `directory`.
190    pub fn in_directory(directory: PathBuf) -> Self {
191        Self {
192            queue: directory.join("queue.json"),
193            receipt: directory.join("receipt.json"),
194            settings: directory.join("settings.json"),
195            record: directory.join("relay.json"),
196            sent: directory.join("sent.json"),
197            lock: directory.join("send.lock"),
198            directory,
199        }
200    }
201
202    /// Newest message id sent to each Claude session name.
203    pub fn last_sent(&self) -> HashMap<String, String> {
204        std::fs::read(&self.sent)
205            .ok()
206            .and_then(|bytes| serde_json::from_slice(&bytes).ok())
207            .unwrap_or_default()
208    }
209}
210
211/// Everything needed to start one relay.
212#[derive(Debug, Clone)]
213pub struct RelaySpec {
214    /// Session the relay speaks for.
215    pub represented: MailAddress,
216    /// Name the represented session is known by.
217    pub represented_name: String,
218    /// Registry name the relay registers under.
219    pub name: String,
220    /// Relay files.
221    pub paths: RelayPaths,
222    /// Program run by the hooks (this supercode binary).
223    pub program: PathBuf,
224}
225
226impl RelaySpec {
227    /// The relay that speaks for `represented`.
228    pub fn for_sender(represented: &MailAddress, represented_name: &str) -> std::io::Result<Self> {
229        let base = represented_name
230            .split('@')
231            .next()
232            .unwrap_or(represented_name);
233        let name = relay_name(base, &represented.machine);
234        Ok(Self {
235            represented: represented.clone(),
236            represented_name: represented_name.to_string(),
237            paths: RelayPaths::new(&mail_root(), &name),
238            name,
239            program: supercode_program()?,
240        })
241    }
242}
243
244/// Arguments of the Claude Code runtime a relay runs as.
245///
246/// `--setting-sources project` in an empty directory keeps user hooks and
247/// settings out; `--settings` then adds only the relay's hooks. The API is
248/// the relay endpoint, through the environment ([`relay_environment`]).
249pub fn relay_arguments(spec: &RelaySpec) -> Vec<String> {
250    vec![
251        "--print".into(),
252        "--input-format".into(),
253        "stream-json".into(),
254        "--output-format".into(),
255        "stream-json".into(),
256        "--verbose".into(),
257        "--permission-prompt-tool".into(),
258        "stdio".into(),
259        "--model".into(),
260        RELAY_MODEL.into(),
261        "--name".into(),
262        spec.name.clone(),
263        "--permission-mode".into(),
264        "bypassPermissions".into(),
265        "--setting-sources".into(),
266        "project".into(),
267        "--settings".into(),
268        spec.paths.settings.to_string_lossy().into_owned(),
269        "--tools".into(),
270        RELAY_TOOLS.into(),
271        "--no-session-persistence".into(),
272    ]
273}
274
275/// Environment of a relay: its API is the relay endpoint at `endpoint`, and
276/// its API key names its directory there. The key also keeps account
277/// connectors (MCP servers) out of the relay.
278pub fn relay_environment(spec: &RelaySpec, endpoint: &str) -> BTreeMap<String, String> {
279    let key = spec
280        .paths
281        .directory
282        .file_name()
283        .map(|name| name.to_string_lossy().into_owned())
284        .unwrap_or_default();
285    BTreeMap::from([
286        ("ANTHROPIC_BASE_URL".to_string(), endpoint.to_string()),
287        ("ANTHROPIC_API_KEY".to_string(), key),
288        (
289            "CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC".to_string(),
290            "1".to_string(),
291        ),
292    ])
293}
294
295/// Settings document installing a relay's hooks.
296pub fn relay_settings(spec: &RelaySpec) -> Value {
297    let program = shell_quote(&spec.program.to_string_lossy());
298    let directory = shell_quote(&spec.paths.directory.to_string_lossy());
299    let represented = shell_quote(&spec.represented.to_string());
300    let command = |verb: &str| format!("{program} message {verb} {directory}");
301    json!({
302        "hooks": {
303            // Every tool, not only the two named in --tools.
304            "PreToolUse": [{
305                "matcher": ".*",
306                "hooks": [{"type": "command", "command": command("gate")}],
307            }],
308            "PostToolUse": [{
309                "matcher": "SendMessage",
310                "hooks": [{"type": "command", "command": command("relay-receipt")}],
311            }],
312            // A SendMessage that fails fires this, with Claude Code's own error.
313            "PostToolUseFailure": [{
314                "matcher": "SendMessage",
315                "hooks": [{"type": "command", "command": command("relay-receipt")}],
316            }],
317            "Stop": [{
318                "hooks": [{"type": "command", "command": command("relay-receipt")}],
319            }],
320            // A turn an API error ends (a usage limit, an overload) fires this, not Stop.
321            "StopFailure": [{
322                "hooks": [{"type": "command", "command": command("relay-receipt")}],
323            }],
324            "UserPromptSubmit": [{
325                "hooks": [{
326                    "type": "command",
327                    "command": format!("{program} message relay-inbound {represented} {directory}"),
328                }],
329            }],
330        }
331    })
332}
333
334fn shell_quote(value: &str) -> String {
335    format!("'{}'", value.replace('\'', "'\\''"))
336}
337
338/// Whether a prompt Claude handed the relay is inbound mail (a peer envelope
339/// or a notice) rather than its host's own send turn.
340pub fn is_inbound_prompt(prompt: &str) -> bool {
341    let trimmed = prompt.trim_start();
342    trimmed.starts_with(NATIVE_PEER_OPENING)
343        || trimmed.starts_with(NATIVE_IDLE_NOTICE)
344        || trimmed.starts_with(NATIVE_DELIVERY_NOTICE)
345}
346
347/// The UserPromptSubmit answer that keeps a prompt from the model.
348pub fn inbound_block() -> Value {
349    json!({
350        "decision": "block",
351        "reason": "filed in the mailbox of the session this relay speaks for",
352    })
353}
354
355/// The host turn asking a relay to make one send.
356pub fn send_turn(send: &QueuedSend) -> String {
357    format!(
358        "Send one message with SendMessage.\n\
359         to: {}\n\
360         message: the exact text between the markers, without the markers\n\
361         ---BEGIN MESSAGE---\n{}\n---END MESSAGE---",
362        send.to, send.message
363    )
364}
365
366/// A native peer envelope as Claude delivered it to a relay.
367#[derive(Debug, Clone, PartialEq, Eq)]
368pub struct NativeEnvelope {
369    /// Claude's `from` (the sender's socket, `uds:/tmp/cc-socks/<pid>.sock`).
370    pub from: String,
371    /// Claude's `from-name`.
372    pub from_name: Option<String>,
373    /// Message body, unescaped as Claude delivered it.
374    pub body: String,
375}
376
377/// One inbound prompt, as Claude handed it to a relay's hook.
378#[derive(Debug, Clone, PartialEq, Eq)]
379pub enum RelayEvent {
380    /// A peer message arrived.
381    Peer(NativeEnvelope),
382    /// A `[Cross-session idle notice]`.
383    IdleNotice(String),
384    /// A `[Cross-session delivery notice]`.
385    DeliveryNotice(String),
386}
387
388/// The live Claude session behind a native `uds:` sender address.
389pub fn resolve_native_sender(
390    registry: &[ClaudePeerSession],
391    native_from: &str,
392) -> Option<ClaudePeerSession> {
393    let socket = native_from.strip_prefix("uds:").unwrap_or(native_from);
394    registry
395        .iter()
396        .find(|session| session.socket_path.to_string_lossy() == socket)
397        .cloned()
398}
399
400/// Outcome of one relayed send, as the sender reads it.
401#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
402#[serde(rename_all = "snake_case", tag = "outcome")]
403pub enum RelayReceipt {
404    /// Claude reported the message queued in the receiver's inbox.
405    Delivered {
406        /// Claude's own result text.
407        detail: String,
408        /// Claude's message id, when it reported one.
409        native_msg_id: Option<String>,
410    },
411    /// The send did not happen.
412    Failed {
413        /// Why, in words the sending agent can act on.
414        detail: String,
415    },
416}
417
418/// Read a receipt the relay's hooks wrote: Claude's own `SendMessage`
419/// result, a gate denial, or a turn that ended without sending.
420pub fn read_receipt(value: &Value) -> RelayReceipt {
421    if let Some(reason) = value.get("denied").and_then(Value::as_str) {
422        return RelayReceipt::Failed {
423            detail: format!("the relay's gate refused the send: {reason}"),
424        };
425    }
426    if let Some(error) = value.get("send_failed").and_then(Value::as_str) {
427        return RelayReceipt::Failed {
428            detail: format!("Claude Code refused the relay's send: {}", error.trim()),
429        };
430    }
431    if value.get("turn_ended").and_then(Value::as_bool) == Some(true) {
432        let said = value
433            .get("said")
434            .and_then(Value::as_str)
435            .filter(|said| !said.trim().is_empty())
436            .map(|said| format!("; it said: {}", said.trim()))
437            .unwrap_or_default();
438        return RelayReceipt::Failed {
439            detail: format!(
440                "the Claude relay ended its turn without sending; nothing was sent{said}"
441            ),
442        };
443    }
444    let response = value.get("tool_response").cloned().unwrap_or(Value::Null);
445    let parsed = match &response {
446        Value::String(text) => serde_json::from_str::<Value>(text).unwrap_or(response.clone()),
447        Value::Array(blocks) => blocks
448            .iter()
449            .find_map(|block| block.get("text").and_then(Value::as_str))
450            .and_then(|text| serde_json::from_str::<Value>(text).ok())
451            .unwrap_or(response.clone()),
452        _ => response.clone(),
453    };
454    if parsed.get("success").and_then(Value::as_bool) != Some(true) {
455        return RelayReceipt::Failed {
456            detail: format!("Claude did not report the send as successful: {parsed}"),
457        };
458    }
459    RelayReceipt::Delivered {
460        detail: parsed
461            .get("message")
462            .and_then(Value::as_str)
463            .unwrap_or_default()
464            .to_string(),
465        native_msg_id: parsed
466            .get("msg_id")
467            .and_then(Value::as_str)
468            .map(str::to_string),
469    }
470}
471
472/// Send `message` to the Claude session named `to`, from `sender` through
473/// its relay: start the relay as a hosted runtime when it is not running,
474/// queue the send, hand it the send turn through the runtime's own door, and
475/// wait for Claude's receipt.
476pub async fn send_through_relay(
477    sender: &MailAddress,
478    sender_name: &str,
479    to: &str,
480    message: String,
481    message_id: &str,
482) -> RelayReceipt {
483    let failed = |detail: String| RelayReceipt::Failed { detail };
484    let spec = match RelaySpec::for_sender(sender, sender_name) {
485        Ok(spec) => spec,
486        Err(error) => return failed(error.to_string()),
487    };
488    if let Err(error) = std::fs::create_dir_all(&spec.paths.directory) {
489        return failed(error.to_string());
490    }
491    // One send at a time through one relay: the queue file is the gate.
492    let lock = match tokio::task::spawn_blocking({
493        let path = spec.paths.lock.clone();
494        move || SendLock::acquire(&path)
495    })
496    .await
497    {
498        Ok(Ok(lock)) => lock,
499        Ok(Err(error)) => return failed(format!("could not take the relay's send lock: {error}")),
500        Err(error) => return failed(error.to_string()),
501    };
502    let runtime = match ensure_relay_runtime(&spec).await {
503        Ok(runtime) => runtime,
504        Err(detail) => return failed(detail),
505    };
506    let queued = QueuedSend {
507        to: to.to_string(),
508        message,
509    };
510    std::fs::remove_file(&spec.paths.receipt).ok();
511    if let Err(error) = std::fs::write(
512        &spec.paths.queue,
513        serde_json::to_vec(&queued).unwrap_or_default(),
514    ) {
515        return failed(error.to_string());
516    }
517    let delivered =
518        crate::runtime_mail::deliver_to_runtime(&runtime, send_turn(&queued), true).await;
519    let receipt = match delivered {
520        Err(detail) => failed(format!("could not reach the Claude relay: {detail}")),
521        Ok(_) => wait_for_receipt(&spec.paths.receipt).await,
522    };
523    std::fs::remove_file(&spec.paths.queue).ok();
524    if matches!(receipt, RelayReceipt::Delivered { .. }) {
525        let mut sent = spec.paths.last_sent();
526        sent.insert(to.to_string(), message_id.to_string());
527        std::fs::write(
528            &spec.paths.sent,
529            serde_json::to_vec(&sent).unwrap_or_default(),
530        )
531        .ok();
532    }
533    drop(lock);
534    receipt
535}
536
537async fn wait_for_receipt(path: &Path) -> RelayReceipt {
538    let started = Instant::now();
539    while started.elapsed() < SEND_TIMEOUT {
540        if let Some(value) = std::fs::read(path)
541            .ok()
542            .and_then(|bytes| serde_json::from_slice::<Value>(&bytes).ok())
543        {
544            return read_receipt(&value);
545        }
546        tokio::time::sleep(Duration::from_millis(200)).await;
547    }
548    RelayReceipt::Failed {
549        detail: format!(
550            "the Claude relay did not confirm the send within {} seconds; it may still arrive",
551            SEND_TIMEOUT.as_secs()
552        ),
553    }
554}
555
556/// The relay's hosted runtime, started in the machine daemon when it is not
557/// running.
558async fn ensure_relay_runtime(
559    spec: &RelaySpec,
560) -> Result<crate::live_runtime::LiveRuntimeRecord, String> {
561    #[derive(Serialize, Deserialize)]
562    struct Record {
563        runtime_id: String,
564        #[serde(default)]
565        endpoint: String,
566    }
567    let endpoint = relay_endpoint(&spec.program).await?;
568    // Written on every send, not only at start: Claude Code reloads a changed settings file, so
569    // a running relay takes the hooks this supercode installs.
570    std::fs::write(
571        &spec.paths.settings,
572        serde_json::to_vec_pretty(&relay_settings(spec)).unwrap_or_default(),
573    )
574    .map_err(|error| error.to_string())?;
575    if let Some(record) = std::fs::read(&spec.paths.record)
576        .ok()
577        .and_then(|bytes| serde_json::from_slice::<Record>(&bytes).ok())
578    {
579        if let Some(runtime) =
580            crate::runtime_mail::controlled_runtime("claude-code", &record.runtime_id)
581        {
582            if record.endpoint == endpoint {
583                return Ok(runtime);
584            }
585            // Started against an endpoint no longer serving, its requests
586            // would go nowhere: it ends, and a new one starts.
587            end_relay_process(&spec.name);
588        }
589    }
590    let params = json!({
591        "harness": "claude-code",
592        "launch": {
593            "program": "claude",
594            "arguments": relay_arguments(spec),
595            "env": relay_environment(spec, &endpoint),
596        },
597        "cwd": spec.paths.directory,
598    });
599    let result = machine_rpc(&spec.program, "runtimes.start", &params).await?;
600    let runtime_id = result
601        .pointer("/handle/runtime_id")
602        .and_then(Value::as_str)
603        .ok_or_else(|| format!("the machine daemon did not start the relay: {result}"))?
604        .to_string();
605    std::fs::write(
606        &spec.paths.record,
607        serde_json::to_vec(&Record {
608            runtime_id: runtime_id.clone(),
609            endpoint,
610        })
611        .unwrap_or_default(),
612    )
613    .map_err(|error| error.to_string())?;
614    crate::runtime_mail::controlled_runtime("claude-code", &runtime_id)
615        .ok_or_else(|| "the relay started but registered no live runtime".to_string())
616}
617
618/// End the Claude process registered under the relay name `name`.
619fn end_relay_process(name: &str) {
620    let registry = registry_dir(&HarnessHomes::default());
621    for session in read_registry(&registry) {
622        #[cfg(unix)]
623        if session.name == name {
624            // SAFETY: `kill` only sends a signal to the relay's own process.
625            unsafe {
626                libc::kill(session.pid as libc::pid_t, libc::SIGTERM);
627            }
628        }
629    }
630}
631
632/// The relay endpoint's base URL. The machine daemon's watcher serves it;
633/// the daemon is started when it is not running.
634async fn relay_endpoint(program: &Path) -> Result<String, String> {
635    if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
636        return Ok(url);
637    }
638    ensure_machine_daemon(program).await?;
639    let deadline = Instant::now() + Duration::from_secs(10);
640    while Instant::now() < deadline {
641        if let Some(url) = crate::relay_endpoint::relay_endpoint_url() {
642            return Ok(url);
643        }
644        tokio::time::sleep(Duration::from_millis(200)).await;
645    }
646    Err(
647        "the relay endpoint is not answering; the machine daemon's `supercode message watch` \
648         serves it (see mail/machine-daemon.log)"
649            .into(),
650    )
651}
652
653/// One `harness.v1` call into this machine's daemon (`supercode teams rpc`),
654/// starting the daemon when it is not running.
655pub async fn machine_rpc(program: &Path, method: &str, params: &Value) -> Result<Value, String> {
656    let call = || {
657        let mut command = tokio::process::Command::new(program);
658        command
659            .args(["teams", "rpc", method, &params.to_string()])
660            .stdin(std::process::Stdio::null())
661            .stdout(std::process::Stdio::piped())
662            .stderr(std::process::Stdio::piped());
663        command.output()
664    };
665    let output = call().await.map_err(|error| error.to_string())?;
666    if output.status.success() {
667        return serde_json::from_slice(&output.stdout)
668            .map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
669    }
670    // Not running: start it (this OS user's own daemon, no Teams context
671    // needed) and ask once more.
672    ensure_machine_daemon(program).await?;
673    let output = call().await.map_err(|error| error.to_string())?;
674    if output.status.success() {
675        return serde_json::from_slice(&output.stdout)
676            .map_err(|error| format!("unreadable answer from the machine daemon: {error}"));
677    }
678    Err(error_line(&String::from_utf8_lossy(&output.stderr)))
679}
680
681/// Start this OS user's machine daemon when it is not answering. It hosts
682/// relays (in its `harness serve`) and the idle watcher.
683pub async fn ensure_machine_daemon(program: &Path) -> Result<(), String> {
684    let describe = || {
685        tokio::process::Command::new(program)
686            .args(["teams", "describe"])
687            .stdin(std::process::Stdio::null())
688            .stdout(std::process::Stdio::null())
689            .stderr(std::process::Stdio::null())
690            .status()
691    };
692    if describe().await.is_ok_and(|status| status.success()) {
693        return Ok(());
694    }
695    // A service manager runs this user's daemon (a connector, or the machine
696    // unit): it is between restarts, and a daemon started here would take the
697    // socket and lock the service's daemon out. Wait for the service instead.
698    if crate::teams::service_owns_daemon() {
699        for _ in 0..100 {
700            tokio::time::sleep(Duration::from_millis(200)).await;
701            if describe().await.is_ok_and(|status| status.success()) {
702                return Ok(());
703            }
704        }
705        return Err(
706            "this machine's daemon is run by its service manager and did not answer within 20 seconds; \
707             see `supercode teams status`"
708                .into(),
709        );
710    }
711    let root = mail_root();
712    std::fs::create_dir_all(&root).map_err(|error| error.to_string())?;
713    let log = std::fs::OpenOptions::new()
714        .create(true)
715        .append(true)
716        .open(root.join("machine-daemon.log"))
717        .map_err(|error| error.to_string())?;
718    let mut command = std::process::Command::new(program);
719    command
720        .args(["teams", "machine", "start", "--supercode"])
721        .arg(program)
722        .stdin(std::process::Stdio::null())
723        .stdout(log.try_clone().map_err(|error| error.to_string())?)
724        .stderr(log);
725    #[cfg(unix)]
726    {
727        use std::os::unix::process::CommandExt;
728        // Its own process group, so the caller's exit does not take the
729        // daemon (and every relay) with it.
730        command.process_group(0);
731    }
732    command
733        .spawn()
734        .map_err(|error| format!("could not start the machine daemon: {error}"))?;
735    for _ in 0..50 {
736        tokio::time::sleep(Duration::from_millis(200)).await;
737        if describe().await.is_ok_and(|status| status.success()) {
738            return Ok(());
739        }
740    }
741    Err(format!(
742        "the machine daemon did not start within 10 seconds; see {}",
743        root.join("machine-daemon.log").display()
744    ))
745}
746
747/// The one line of a failed command's stderr an agent can act on.
748fn error_line(stderr: &str) -> String {
749    let lines: Vec<&str> = stderr
750        .lines()
751        .map(str::trim)
752        .filter(|line| !line.is_empty())
753        .collect();
754    lines
755        .iter()
756        .find(|line| line.starts_with("Error") || line.starts_with("error"))
757        .or(lines.last())
758        .copied()
759        .unwrap_or("the machine daemon failed without saying why")
760        .chars()
761        .take(300)
762        .collect()
763}
764
765/// An exclusive lock on a relay's sends, released on drop.
766struct SendLock {
767    _file: std::fs::File,
768}
769
770impl SendLock {
771    fn acquire(path: &Path) -> std::io::Result<Self> {
772        let file = std::fs::OpenOptions::new()
773            .create(true)
774            .truncate(false)
775            .write(true)
776            .open(path)?;
777        #[cfg(unix)]
778        {
779            use std::os::unix::io::AsRawFd;
780            // SAFETY: flock on a descriptor this function owns; the lock is
781            // released when the file is dropped.
782            if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX) } != 0 {
783                return Err(std::io::Error::last_os_error());
784            }
785        }
786        Ok(Self { _file: file })
787    }
788}
789
790/// File one inbound prompt in `represented`'s mailbox. `last_sent` maps a
791/// Claude session name to the newest message this relay sent it, for
792pub fn file_inbound_prompt(
793    homes: &HarnessHomes,
794    represented: &MailAddress,
795    prompt: &str,
796    last_sent: Option<&HashMap<String, String>>,
797) {
798    match parse_inbound_text(prompt) {
799        Some(RelayEvent::Peer(native)) => {
800            let registry = read_registry(&registry_dir(homes));
801            let sender = resolve_native_sender(&registry, &native.from);
802            let in_reply_to = sender
803                .as_ref()
804                .and_then(|session| last_sent.and_then(|sent| sent.get(&session.name).cloned()));
805            file_inbound(represented, &native, sender.as_ref(), in_reply_to);
806        }
807        Some(RelayEvent::IdleNotice(text) | RelayEvent::DeliveryNotice(text)) => {
808            file_notice(represented, &text);
809        }
810        _ => {}
811    }
812}
813/// Recognize an inbound peer envelope or notice in the text Claude handed a
814/// relay.
815pub fn parse_inbound_text(text: &str) -> Option<RelayEvent> {
816    let trimmed = text.trim_start();
817    if trimmed.starts_with(NATIVE_IDLE_NOTICE) {
818        return Some(RelayEvent::IdleNotice(trimmed.to_string()));
819    }
820    if trimmed.starts_with(NATIVE_DELIVERY_NOTICE) {
821        return Some(RelayEvent::DeliveryNotice(trimmed.to_string()));
822    }
823    if !trimmed.starts_with(NATIVE_PEER_PREAMBLE) && !trimmed.starts_with(NATIVE_PEER_OPENING) {
824        return None;
825    }
826    let start = text.find(NATIVE_PEER_OPENING)?;
827    let header_end = start + text[start..].find(">\n")?;
828    let header = &text[start + NATIVE_PEER_OPENING.len()..header_end];
829    let body_start = header_end + 2;
830    // The body is not escaped by Claude; the closing tag is the LAST one, and
831    // the fixed paragraph after it never contains the tag.
832    let body_end = text.rfind(&format!("\n{NATIVE_PEER_CLOSING}"))?;
833    if body_end < body_start {
834        return None;
835    }
836    let attributes = parse_attributes(header);
837    Some(RelayEvent::Peer(NativeEnvelope {
838        from: attributes.get("from")?.clone(),
839        from_name: attributes.get("from-name").cloned(),
840        body: text[body_start..body_end].to_string(),
841    }))
842}
843fn parse_attributes(header: &str) -> BTreeMap<String, String> {
844    let mut attributes = BTreeMap::new();
845    let mut rest = header;
846    while let Some(equals) = rest.find("=\"") {
847        let key = rest[..equals].trim().to_string();
848        let value_start = equals + 2;
849        let Some(length) = rest[value_start..].find('"') else {
850            break;
851        };
852        attributes.insert(key, rest[value_start..value_start + length].to_string());
853        rest = &rest[value_start + length + 1..];
854    }
855    attributes
856}
857/// Deliver mail to the session a relay speaks for. When it lives on another
858/// machine and Teams cannot take the message there, it is kept in this
859/// machine's copy of that mailbox (never dropped) and the failure is written
860fn deliver_home(represented: &MailAddress, envelope: &Envelope) {
861    if let Err(error) = crate::mailbox::deliver_to(represented, envelope) {
862        eprintln!(
863            "could not deliver {} to {represented}: {error}; kept in this machine's mailbox for it",
864            envelope.id
865        );
866        if let Ok(mailbox) = crate::mailbox::Mailbox::open(&mail_root(), represented) {
867            mailbox.deliver(envelope).ok();
868        }
869        return;
870    }
871    // Filed; now hand it to its session by the session's own door, as any
872    // sender's message is (a Claude session gets it at its next tool round).
873    // That can take as long as a relay send, which this hook must not wait
874    // on, so a detached `message push` does it.
875    if represented.machine == crate::mailbox::local_machine_name() {
876        if let Ok(program) = supercode_program() {
877            std::process::Command::new(program)
878                .args(["message", "push", &represented.to_string(), &envelope.id])
879                .stdin(std::process::Stdio::null())
880                .stdout(std::process::Stdio::null())
881                .stderr(std::process::Stdio::null())
882                .spawn()
883                .ok();
884        }
885    }
886}
887fn file_inbound(
888    represented: &MailAddress,
889    native: &NativeEnvelope,
890    sender: Option<&ClaudePeerSession>,
891    in_reply_to: Option<String>,
892) {
893    // The replying Claude session runs beside the relay, on this machine —
894    // not on the machine of the session the relay speaks for.
895    let machine = crate::mailbox::local_machine_name();
896    let (from, from_name) = match sender {
897        Some(session) => (
898            MailAddress::new(&machine, "claude-code", &session.session_id),
899            format!("{}@{machine}", session.name),
900        ),
901        None => (
902            MailAddress::new(&machine, "claude-code", "unknown"),
903            native
904                .from_name
905                .clone()
906                .unwrap_or_else(|| "an unknown Claude session".into()),
907        ),
908    };
909    let Ok(from) = from else { return };
910    let Ok(mut envelope) = Envelope::new(
911        from,
912        from_name,
913        MailKind::Peer,
914        ReplyVia::Command,
915        native.body.clone(),
916    ) else {
917        return;
918    };
919    envelope.native_from = Some(native.from.clone());
920    if let Some(in_reply_to) = in_reply_to {
921        envelope.in_reply_to = Some(in_reply_to);
922        envelope.in_reply_to_inferred = true;
923    }
924    // The session a relay speaks for may be on another machine; its mail goes
925    // home through Teams.
926    deliver_home(represented, &envelope);
927}
928fn file_notice(represented: &MailAddress, text: &str) {
929    let Ok(from) = MailAddress::new(
930        crate::mailbox::local_machine_name(),
931        "claude-code",
932        "notice",
933    ) else {
934        return;
935    };
936    let Ok(envelope) = Envelope::new(
937        from,
938        "Claude Code",
939        MailKind::Notice,
940        ReplyVia::None,
941        text.to_string(),
942    ) else {
943        return;
944    };
945    // The session a relay speaks for may be on another machine; its mail goes
946    // home through Teams.
947    deliver_home(represented, &envelope);
948}
949/// The `supercode` program that the relay hooks and the machine daemon run:
950/// this process when it is `supercode` (the CLI, `supercode harness serve`),
951pub fn supercode_program() -> std::io::Result<PathBuf> {
952    let current = std::env::current_exe()?;
953    if current.file_stem().and_then(|stem| stem.to_str()) == Some("supercode") {
954        return Ok(current);
955    }
956    std::env::var_os("PATH")
957        .iter()
958        .flat_map(std::env::split_paths)
959        .map(|directory| directory.join("supercode"))
960        .find(|candidate| candidate.is_file())
961        .ok_or_else(|| {
962            std::io::Error::new(
963                std::io::ErrorKind::NotFound,
964                "the supercode program is not on PATH; Claude relays and the machine daemon need it",
965            )
966        })
967}