supercode_harness/
runtime_mail.rs1use crate::frontend::{FrontendRuntime, FrontendTurnState, HttpFrontendRuntime};
15use crate::live_runtime::{list_live_runtimes, resolve_live_runtime, LiveRuntimeRecord};
16
17pub 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
26pub fn controlled_runtimes() -> Vec<LiveRuntimeRecord> {
28 list_live_runtimes().unwrap_or_default()
29}
30
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
33pub enum RuntimeDelivery {
34 Steered,
36 Started,
38 NotWoken,
41}
42
43pub 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
65pub 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
97pub 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
133pub 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}