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.
44pub async fn runtime_turn_state(record: &LiveRuntimeRecord) -> Result<FrontendTurnState, String> {
45    let remote = connect(record).await?;
46    FrontendRuntime::describe(remote.as_ref())
47        .await
48        .map(|descriptor| descriptor.turn_state)
49        .map_err(|error| error.to_string())
50}
51
52/// Each runtime's turn state, asked of every runtime at once with a short bound, from any context
53/// (a thread of its own runs the asking): `None` where one does not answer in time.
54pub fn runtime_turn_states(records: &[LiveRuntimeRecord]) -> Vec<Option<FrontendTurnState>> {
55    if records.is_empty() {
56        return Vec::new();
57    }
58    let count = records.len();
59    let records = records.to_vec();
60    std::thread::spawn(move || {
61        let runtime = tokio::runtime::Builder::new_current_thread()
62            .enable_all()
63            .build()
64            .ok()?;
65        Some(runtime.block_on(async {
66            futures::future::join_all(records.iter().map(|record| async move {
67                tokio::time::timeout(
68                    std::time::Duration::from_secs(2),
69                    runtime_turn_state(record),
70                )
71                .await
72                .ok()
73                .and_then(Result::ok)
74            }))
75            .await
76        }))
77    })
78    .join()
79    .ok()
80    .flatten()
81    .unwrap_or_else(|| vec![None; count])
82}
83
84/// Deliver `text` into a controlled session: steer a running turn, or start
85/// one when `wake` (a `--queue` send leaves an idle session alone).
86pub async fn deliver_to_runtime(
87    record: &LiveRuntimeRecord,
88    text: String,
89    wake: bool,
90) -> Result<RuntimeDelivery, String> {
91    let remote = connect(record).await?;
92    let descriptor = FrontendRuntime::describe(remote.as_ref())
93        .await
94        .map_err(|error| error.to_string())?;
95    if descriptor.session_id != record.runtime_session_id {
96        return Err("the live runtime's identity did not match its receipt".into());
97    }
98    match descriptor.turn_state {
99        FrontendTurnState::Busy if descriptor.actions.steer => {
100            FrontendRuntime::steer(remote.as_ref(), text)
101                .await
102                .map_err(|error| error.to_string())?;
103            Ok(RuntimeDelivery::Steered)
104        }
105        FrontendTurnState::Busy => Err(
106            "a turn is running and this runtime does not accept steering; the message was not \
107             delivered into it"
108                .into(),
109        ),
110        FrontendTurnState::Idle if wake => {
111            FrontendRuntime::send_input_with_images(remote.clone(), text, Vec::new())
112                .await
113                .map_err(|error| error.to_string())?;
114            Ok(RuntimeDelivery::Started)
115        }
116        FrontendTurnState::Idle => Ok(RuntimeDelivery::NotWoken),
117    }
118}
119
120/// The runtime's answer to the message `message_id`: the text of its last
121/// assistant message after the user message that carried that id, read from
122/// the session's own transcript (a hosted runtime's attach history can be
123/// empty; its harness's transcript never is). `None` until it has answered.
124pub fn answer_to(
125    homes: &crate::HarnessHomes,
126    record: &LiveRuntimeRecord,
127    message_id: &str,
128) -> Option<String> {
129    let query = crate::DiscoveryQuery {
130        harnesses: vec![crate::HarnessId::new(record.source.harness.clone())],
131        homes: homes.clone(),
132        workspace: Some(record.source.workspace.clone()),
133        ..Default::default()
134    };
135    let descriptor = crate::discover_sessions(&query)
136        .ok()?
137        .into_iter()
138        .find(|descriptor| {
139            descriptor.locator.session_id == record.source.session_id
140                || descriptor.locator.session_id == record.runtime_session_id
141        })?;
142    let session =
143        crate::sdk::load_session_with_fidelity(&descriptor.locator, crate::Fidelity::Semantic)
144            .ok()?;
145    let text = |message: &crate::ChatMessage| message.content.clone().unwrap_or_default();
146    let asked = session.messages.iter().rposition(|message| {
147        message.role == crate::Role::User && text(message).contains(message_id)
148    })?;
149    session.messages[asked + 1..]
150        .iter()
151        .rev()
152        .find(|message| message.role == crate::Role::Assistant && !text(message).trim().is_empty())
153        .map(text)
154}
155
156async fn connect(
157    record: &LiveRuntimeRecord,
158) -> Result<std::sync::Arc<HttpFrontendRuntime>, String> {
159    let receipt = resolve_live_runtime(&record.endpoint, &record.source)
160        .map_err(|error| error.to_string())?;
161    HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
162        .await
163        .map_err(|error| error.to_string())
164}