Skip to main content

supercode_harness/
runtime_mail.rs

1//! The default delivery tier of the mailbox: sessions supercode controls.
2//!
3//! A session supercode hosts as a runtime (any harness: Codex over
4//! app-server, Claude Code over stream-json, pi, the ACP harnesses, OpenCode)
5//! registers a live-runtime receipt ([`crate::live_runtime`]). Mail to such a
6//! session goes through the runtime's own doors, the same ones every
7//! frontend uses: `steer` while a turn runs, `send_input` to start one. No
8//! hook, relay or per-harness code is involved; those are the degraded tier,
9//! for sessions supercode does not control.
10//!
11//! The text delivered is the rendered envelope: who sent it, that it is not
12//! the user, and how to answer.
13
14use crate::frontend::{FrontendRuntime, FrontendTurnState, HttpFrontendRuntime};
15use crate::live_runtime::{list_live_runtimes, resolve_live_runtime, LiveRuntimeRecord};
16
17/// The live runtime supercode hosts for `harness`'s session `session_id`, if
18/// it controls that session.
19pub fn controlled_runtime(harness: &str, session_id: &str) -> Option<LiveRuntimeRecord> {
20    list_live_runtimes().ok()?.into_iter().find(|record| {
21        record.source.harness == harness
22            && (record.source.session_id == session_id || record.runtime_session_id == session_id)
23    })
24}
25
26/// Every session supercode controls on this machine.
27pub fn controlled_runtimes() -> Vec<LiveRuntimeRecord> {
28    list_live_runtimes().unwrap_or_default()
29}
30
31/// How a message entered a controlled session.
32#[derive(Debug, Clone, Copy, PartialEq, Eq)]
33pub enum RuntimeDelivery {
34    /// A turn was running; the message was steered into it.
35    Steered,
36    /// The session was idle; the message started a turn.
37    Started,
38    /// The session was idle and the sender asked not to wake it; nothing was
39    /// delivered into the runtime.
40    NotWoken,
41}
42
43/// The runtime's turn state right now, asked as an observer over the client kept for its endpoint:
44/// a full `connect` opened an event stream and two describes per read, each on a new connection
45/// (t_f3f102db t_4f29b43e).
46pub async fn runtime_turn_state(record: &LiveRuntimeRecord) -> Result<FrontendTurnState, String> {
47    let receipt = resolve_live_runtime(&record.endpoint, &record.source)
48        .map_err(|error| error.to_string())?;
49    let client_id = crate::RuntimeClientId::parse(format!(
50        "mail-{}",
51        record
52            .endpoint
53            .as_str()
54            .rsplit('/')
55            .next()
56            .unwrap_or("probe")
57    ))
58    .map_err(|error| error.to_string())?;
59    HttpFrontendRuntime::probe_described_kept(receipt.base_url, receipt.token, client_id)
60        .await
61        .map(|(_, descriptor)| descriptor.turn_state)
62        .map_err(|error| error.to_string())
63}
64
65/// Each runtime's turn state, asked of every runtime at once with a short bound, from any context
66/// (a thread of its own runs the asking): `None` where one does not answer in time.
67pub fn runtime_turn_states(records: &[LiveRuntimeRecord]) -> Vec<Option<FrontendTurnState>> {
68    if records.is_empty() {
69        return Vec::new();
70    }
71    let count = records.len();
72    let records = records.to_vec();
73    std::thread::spawn(move || {
74        let runtime = tokio::runtime::Builder::new_current_thread()
75            .enable_all()
76            .build()
77            .ok()?;
78        Some(runtime.block_on(async {
79            futures::future::join_all(records.iter().map(|record| async move {
80                tokio::time::timeout(
81                    std::time::Duration::from_secs(2),
82                    runtime_turn_state(record),
83                )
84                .await
85                .ok()
86                .and_then(Result::ok)
87            }))
88            .await
89        }))
90    })
91    .join()
92    .ok()
93    .flatten()
94    .unwrap_or_else(|| vec![None; count])
95}
96
97/// Deliver `text` into a controlled session: steer a running turn, or start
98/// one when `wake` (a `--queue` send leaves an idle session alone).
99pub async fn deliver_to_runtime(
100    record: &LiveRuntimeRecord,
101    text: String,
102    wake: bool,
103) -> Result<RuntimeDelivery, String> {
104    let remote = connect(record).await?;
105    let descriptor = FrontendRuntime::describe(remote.as_ref())
106        .await
107        .map_err(|error| error.to_string())?;
108    if descriptor.session_id != record.runtime_session_id {
109        return Err("the live runtime's identity did not match its receipt".into());
110    }
111    match descriptor.turn_state {
112        FrontendTurnState::Busy if descriptor.actions.steer => {
113            FrontendRuntime::steer(remote.as_ref(), text)
114                .await
115                .map_err(|error| error.to_string())?;
116            Ok(RuntimeDelivery::Steered)
117        }
118        FrontendTurnState::Busy => Err(
119            "a turn is running and this runtime does not accept steering; the message was not \
120             delivered into it"
121                .into(),
122        ),
123        FrontendTurnState::Idle if wake => {
124            FrontendRuntime::send_input_with_images(remote.clone(), text, Vec::new())
125                .await
126                .map_err(|error| error.to_string())?;
127            Ok(RuntimeDelivery::Started)
128        }
129        FrontendTurnState::Idle => Ok(RuntimeDelivery::NotWoken),
130    }
131}
132
133/// The runtime's answer to the message `message_id`: the text of its last
134/// assistant message after the user message that carried that id, read from
135/// the session's own transcript (a hosted runtime's attach history can be
136/// empty; its harness's transcript never is). `None` until it has answered.
137pub fn answer_to(
138    homes: &crate::HarnessHomes,
139    record: &LiveRuntimeRecord,
140    message_id: &str,
141) -> Option<String> {
142    let query = crate::DiscoveryQuery {
143        harnesses: vec![crate::HarnessId::new(record.source.harness.clone())],
144        homes: homes.clone(),
145        workspace: Some(record.source.workspace.clone()),
146        ..Default::default()
147    };
148    let descriptor = crate::discover_sessions(&query)
149        .ok()?
150        .into_iter()
151        .find(|descriptor| {
152            descriptor.locator.session_id == record.source.session_id
153                || descriptor.locator.session_id == record.runtime_session_id
154        })?;
155    let session =
156        crate::sdk::load_session_with_fidelity(&descriptor.locator, crate::Fidelity::Semantic)
157            .ok()?;
158    let text = |message: &crate::ChatMessage| message.content.clone().unwrap_or_default();
159    let asked = session.messages.iter().rposition(|message| {
160        message.role == crate::Role::User && text(message).contains(message_id)
161    })?;
162    session.messages[asked + 1..]
163        .iter()
164        .rev()
165        .find(|message| message.role == crate::Role::Assistant && !text(message).trim().is_empty())
166        .map(text)
167}
168
169async fn connect(
170    record: &LiveRuntimeRecord,
171) -> Result<std::sync::Arc<HttpFrontendRuntime>, String> {
172    let receipt = resolve_live_runtime(&record.endpoint, &record.source)
173        .map_err(|error| error.to_string())?;
174    HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
175        .await
176        .map_err(|error| error.to_string())
177}