use std::collections::{HashMap, VecDeque};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex as StdMutex};
use async_trait::async_trait;
use serde_json::{json, Value};
use tokio::sync::{broadcast, mpsc, oneshot};
use super::{HarnessEvent, RuntimeCapabilities, RuntimeConnection, RuntimeHandle, RuntimeInput};
use crate::frontend::{
FrontendActions, FrontendAttachment, FrontendConnectionState, FrontendDisplayCapabilities,
FrontendEvent, FrontendRuntime, FrontendRuntimeDescriptor, FrontendRuntimeError,
FrontendTurnState, FRONTEND_REPLAY_CAPACITY, FRONTEND_RUNTIME_SCHEMA_VERSION,
};
use crate::server::RuntimeSubmitError;
use crate::{Error, FrontendResponse, Result};
enum HostCommand {
Submit {
input: RuntimeInput,
reply: oneshot::Sender<Result<Option<String>>>,
},
Interrupt {
reply: oneshot::Sender<Result<()>>,
},
Steer {
text: String,
reply: oneshot::Sender<Result<()>>,
},
Respond {
request_id: Value,
response: Value,
reply: oneshot::Sender<Result<()>>,
},
Shutdown {
reply: oneshot::Sender<Result<()>>,
},
}
struct ProjectionState {
next_sequence: u64,
replay: VecDeque<FrontendEvent>,
}
pub struct HostedHarnessRuntime {
handle: RuntimeHandle,
capabilities: RuntimeCapabilities,
commands: mpsc::Sender<HostCommand>,
raw_events: broadcast::Sender<HarnessEvent>,
frontend_events: broadcast::Sender<FrontendEvent>,
projection: StdMutex<ProjectionState>,
pending_requests: StdMutex<HashMap<u64, Value>>,
busy: AtomicBool,
closed: AtomicBool,
}
impl HostedHarnessRuntime {
pub fn spawn(
runtime: Box<dyn RuntimeConnection>,
capabilities: RuntimeCapabilities,
) -> (Arc<Self>, HostedHarnessConnection) {
let handle = runtime.handle().clone();
let (commands, command_rx) = mpsc::channel(32);
let (raw_events, raw_rx) = broadcast::channel(1024);
let (frontend_events, _) = broadcast::channel(1024);
let host = Arc::new(Self {
handle: handle.clone(),
capabilities,
commands,
raw_events,
frontend_events,
projection: StdMutex::new(ProjectionState {
next_sequence: 1,
replay: VecDeque::new(),
}),
pending_requests: StdMutex::new(HashMap::new()),
busy: AtomicBool::new(false),
closed: AtomicBool::new(false),
});
tokio::spawn(run_native_runtime(
runtime,
Arc::downgrade(&host),
command_rx,
));
let connection = HostedHarnessConnection {
host: host.clone(),
handle,
events: raw_rx,
closed: false,
};
(host, connection)
}
pub fn frontend_sender(&self) -> broadcast::Sender<FrontendEvent> {
self.frontend_events.clone()
}
pub async fn shutdown(&self) -> Result<()> {
if self.closed.load(Ordering::SeqCst) {
return Ok(());
}
let (reply, response) = oneshot::channel();
self.commands
.send(HostCommand::Shutdown { reply })
.await
.map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
response
.await
.map_err(|_| Error::Other("hosted harness runtime stopped before shutdown".into()))?
}
fn claim_submit(&self) -> Result<()> {
if self.closed.load(Ordering::SeqCst) {
return Err(Error::Other("hosted harness runtime is closed".into()));
}
if self
.busy
.compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
.is_err()
{
return Err(Error::Other("a harness turn is already in progress".into()));
}
Ok(())
}
async fn submit_native_claimed(&self, input: RuntimeInput) -> Result<Option<String>> {
self.publish(json!({"type":"user_message", "text":input.text}));
self.publish(json!({"type":"turn_started"}));
let (reply, response) = oneshot::channel();
if self
.commands
.send(HostCommand::Submit { input, reply })
.await
.is_err()
{
self.busy.store(false, Ordering::SeqCst);
self.publish(
json!({"type":"turn_failed", "message":"Hosted harness runtime is closed."}),
);
self.publish(json!({"type":"turn_completed"}));
self.mark_closed("Harness runtime command channel closed.");
return Err(Error::Other("hosted harness runtime is closed".into()));
}
match response.await {
Ok(Ok(turn)) => Ok(turn),
Ok(Err(error)) => {
self.busy.store(false, Ordering::SeqCst);
self.publish(json!({"type":"turn_failed", "message":error.to_string()}));
self.publish(json!({"type":"turn_completed"}));
Err(error)
}
Err(_) => {
self.busy.store(false, Ordering::SeqCst);
self.publish(json!({"type":"turn_failed", "message":"Hosted harness runtime stopped before accepting input."}));
self.publish(json!({"type":"turn_completed"}));
self.mark_closed("Harness runtime stopped before accepting input.");
Err(Error::Other(
"hosted harness runtime stopped before accepting input".into(),
))
}
}
}
async fn submit_native(&self, text: String) -> Result<Option<String>> {
self.claim_submit()?;
self.submit_native_claimed(RuntimeInput {
text,
image_urls: Vec::new(),
})
.await
}
async fn interrupt_native(&self) -> Result<()> {
let (reply, response) = oneshot::channel();
self.commands
.send(HostCommand::Interrupt { reply })
.await
.map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
response
.await
.map_err(|_| Error::Other("hosted harness runtime stopped before interrupt".into()))?
}
async fn steer_native(&self, text: String) -> Result<()> {
let (reply, response) = oneshot::channel();
self.commands
.send(HostCommand::Steer { text, reply })
.await
.map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
response
.await
.map_err(|_| Error::Other("hosted harness runtime stopped before steering".into()))?
}
async fn respond_native(&self, request_id: Value, response: Value) -> Result<()> {
self.respond_native_as(request_id, response, None).await
}
async fn respond_native_as(
&self,
request_id: Value,
response: Value,
canonical: Option<Value>,
) -> Result<()> {
let (reply, completed) = oneshot::channel();
self.commands
.send(HostCommand::Respond {
request_id: request_id.clone(),
response: response.clone(),
reply,
})
.await
.map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
completed
.await
.map_err(|_| Error::Other("hosted harness runtime stopped before response".into()))??;
let resolved = {
let mut pending = self
.pending_requests
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let found = pending
.iter()
.find(|(_, native)| *native == &request_id)
.map(|(id, _)| *id);
found.and_then(|id| pending.remove(&id).map(|_| id))
};
if let Some(id) = resolved {
self.publish(json!({
"type": "request_resolved",
"request_id": id,
"response": canonical.unwrap_or(response),
}));
}
Ok(())
}
fn native_request_id(&self, request_id: u64) -> Option<Value> {
self.pending_requests
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&request_id)
.cloned()
}
fn publish(&self, payload: Value) {
let event = {
let mut projection = self
.projection
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let event = FrontendEvent::new(projection.next_sequence, payload);
projection.next_sequence = projection.next_sequence.saturating_add(1);
projection.replay.push_back(event.clone());
while projection.replay.len() > FRONTEND_REPLAY_CAPACITY {
projection.replay.pop_front();
}
event
};
let _ = self.frontend_events.send(event);
}
fn accept_native_event(&self, event: HarnessEvent) {
let _ = self.raw_events.send(event.clone());
for payload in project_native_event(self.handle.harness.as_str(), &event) {
let terminal = matches!(
payload.get("type").and_then(Value::as_str),
Some("turn_succeeded" | "turn_interrupted" | "turn_failed")
);
if terminal {
self.busy.store(false, Ordering::SeqCst);
}
if payload.get("type").and_then(Value::as_str) == Some("request") {
if let (Some(id), Some(native)) = (
payload.pointer("/request/id").and_then(Value::as_u64),
payload
.pointer("/request/payload/native_request_id")
.cloned(),
) {
self.pending_requests
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(id, native);
}
}
self.publish(payload);
}
}
fn mark_closed(&self, message: impl Into<String>) {
if self.closed.swap(true, Ordering::SeqCst) {
return;
}
let message = message.into();
self.busy.store(false, Ordering::SeqCst);
let _ = self.raw_events.send(HarnessEvent {
sequence: None,
kind: "transport_closed".into(),
payload: json!({"message":message, "terminal":true}),
});
self.publish(json!({"type":"runtime_disconnected", "message":message}));
}
fn descriptor(&self) -> FrontendRuntimeDescriptor {
FrontendRuntimeDescriptor {
schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
session_id: self.handle.runtime_id.clone(),
source_harness: Some(self.handle.harness.as_str().to_string()),
emulation_profile: None,
active_modules: Vec::new(),
commands: Vec::new(),
operations: Vec::new(),
actions: FrontendActions {
submit: self.capabilities.send_input,
interrupt: self.capabilities.interrupt,
steer: self.capabilities.steer,
respond: !hosted_answerable_decisions(self.handle.harness.as_str()).is_empty(),
detach: true,
close: false,
},
display: FrontendDisplayCapabilities {
event_kinds: vec![
"user_message".into(),
"turn_started".into(),
"turn_succeeded".into(),
"turn_interrupted".into(),
"turn_failed".into(),
"text_delta".into(),
"reasoning".into(),
"tool_call_started".into(),
"tool_call_completed".into(),
"request".into(),
"native_event".into(),
"runtime_disconnected".into(),
],
opaque_fallback: true,
},
model: self.handle.harness.as_str().to_string(),
turn_state: if self.busy.load(Ordering::SeqCst) {
FrontendTurnState::Busy
} else {
FrontendTurnState::Idle
},
connection_state: if self.closed.load(Ordering::SeqCst) {
FrontendConnectionState::ShuttingDown
} else {
FrontendConnectionState::Connected
},
extensions: Default::default(),
}
}
}
#[async_trait]
impl FrontendRuntime for HostedHarnessRuntime {
async fn describe(
&self,
) -> std::result::Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
Ok(self.descriptor())
}
async fn attach(
&self,
_history_limit: usize,
) -> std::result::Result<FrontendAttachment, FrontendRuntimeError> {
let live = self.frontend_events.subscribe();
let projection = self
.projection
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let replay = projection.replay.clone();
if let Some(first) = replay.front() {
if first.sequence > 1 {
return Err(FrontendRuntimeError::ReplayGap(first.sequence - 1));
}
}
Ok(FrontendAttachment::new(
self.descriptor(),
Vec::new(),
0,
replay,
live,
None,
))
}
async fn send_input(
self: Arc<Self>,
prompt: String,
) -> std::result::Result<(), FrontendRuntimeError> {
self.claim_submit().map_err(hosted_submit_error)?;
tokio::spawn(async move {
let _ = self
.submit_native_claimed(RuntimeInput {
text: prompt,
image_urls: Vec::new(),
})
.await;
});
Ok(())
}
async fn send_input_with_images(
self: Arc<Self>,
prompt: String,
image_urls: Vec<String>,
) -> std::result::Result<(), FrontendRuntimeError> {
self.claim_submit().map_err(hosted_submit_error)?;
tokio::spawn(async move {
let _ = self
.submit_native_claimed(RuntimeInput {
text: prompt,
image_urls,
})
.await;
});
Ok(())
}
async fn submit(&self, prompt: String) -> std::result::Result<String, FrontendRuntimeError> {
self.submit_native(prompt)
.await
.map(|turn| turn.unwrap_or_default())
.map_err(hosted_submit_error)
}
async fn interrupt(&self) -> std::result::Result<bool, FrontendRuntimeError> {
if !self.busy.load(Ordering::SeqCst) {
return Ok(false);
}
self.interrupt_native()
.await
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
Ok(true)
}
async fn steer(&self, prompt: String) -> std::result::Result<(), FrontendRuntimeError> {
if !self.busy.load(Ordering::SeqCst) || !self.capabilities.steer {
return Err(FrontendRuntimeError::UnsupportedAction("steer"));
}
self.steer_native(prompt)
.await
.map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
}
async fn respond(
&self,
response: FrontendResponse,
) -> std::result::Result<(), FrontendRuntimeError> {
let harness = self.handle.harness.as_str();
let FrontendResponse::Approval {
request_id,
decision,
} = &response
else {
return Err(FrontendRuntimeError::UnsupportedAction(
"respond: a hosted harness answers approval requests only",
));
};
let Some(native_id) = self.native_request_id(*request_id) else {
return Err(FrontendRuntimeError::UnknownRequest(*request_id));
};
let body = hosted_permission_reply(harness, *decision)
.map_err(FrontendRuntimeError::InvalidResponse)?;
let canonical = serde_json::to_value(&response).unwrap_or(Value::Null);
self.respond_native_as(native_id, body, Some(canonical))
.await
.map_err(|error| FrontendRuntimeError::Execution {
operation: crate::SdkOperation::Respond,
message: error.to_string(),
})
}
}
fn hosted_submit_error(error: Error) -> FrontendRuntimeError {
if error.to_string().contains("already in progress") {
FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)
} else {
FrontendRuntimeError::Transport(error.to_string())
}
}
pub struct HostedHarnessConnection {
host: Arc<HostedHarnessRuntime>,
handle: RuntimeHandle,
events: broadcast::Receiver<HarnessEvent>,
closed: bool,
}
#[async_trait]
impl RuntimeConnection for HostedHarnessConnection {
fn handle(&self) -> &RuntimeHandle {
&self.handle
}
async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
self.host.claim_submit()?;
self.host.submit_native_claimed(input).await
}
async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
match self.events.recv().await {
Ok(event) => Ok(Some(event)),
Err(broadcast::error::RecvError::Lagged(count)) => Err(Error::Other(format!(
"hosted harness event stream lost {count} event(s)"
))),
Err(broadcast::error::RecvError::Closed) => Ok(None),
}
}
async fn interrupt(&mut self) -> Result<()> {
self.host.interrupt_native().await
}
async fn steer(&mut self, text: String) -> Result<()> {
self.host.steer_native(text).await
}
async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
self.host.respond_native(request_id, response).await
}
async fn close(&mut self) -> Result<()> {
if self.closed {
return Ok(());
}
self.host.shutdown().await?;
self.closed = true;
Ok(())
}
}
async fn run_native_runtime(
mut runtime: Box<dyn RuntimeConnection>,
host: std::sync::Weak<HostedHarnessRuntime>,
mut commands: mpsc::Receiver<HostCommand>,
) {
loop {
tokio::select! {
command = commands.recv() => {
let Some(command) = command else {
let _ = runtime.close().await;
return;
};
match command {
HostCommand::Submit { input, reply } => {
let _ = reply.send(runtime.send_input(input).await);
}
HostCommand::Interrupt { reply } => {
let _ = reply.send(runtime.interrupt().await);
}
HostCommand::Steer { text, reply } => {
let _ = reply.send(runtime.steer(text).await);
}
HostCommand::Respond { request_id, response, reply } => {
let _ = reply.send(runtime.respond(request_id, response).await);
}
HostCommand::Shutdown { reply } => {
let result = runtime.close().await;
let closed = result.is_ok();
let _ = reply.send(result);
if closed {
if let Some(host) = host.upgrade() {
host.mark_closed("Harness runtime closed.");
}
return;
}
}
}
}
event = runtime.next_event() => {
let Some(host) = host.upgrade() else {
let _ = runtime.close().await;
return;
};
match event {
Ok(Some(event)) => host.accept_native_event(event),
Ok(None) => {
host.mark_closed("Harness runtime transport closed.");
return;
}
Err(error) => {
host.mark_closed(error.to_string());
return;
}
}
}
}
}
}
pub(crate) fn hosted_answerable_decisions(harness: &str) -> &'static [&'static str] {
match harness {
"claude-code" => &["allow", "deny"],
_ => &[],
}
}
fn hosted_permission_reply(
harness: &str,
decision: crate::FrontendApprovalDecision,
) -> std::result::Result<Value, String> {
use crate::FrontendApprovalDecision as Decision;
let accepted = hosted_answerable_decisions(harness);
let wire = match decision {
Decision::Allow => "allow",
Decision::AllowForSession => "allow_for_session",
Decision::Deny => "deny",
};
if !accepted.contains(&wire) {
return Err(format!(
"`{wire}` is not a decision a hosted {harness} runtime can carry; this harness accepts \
{}",
if accepted.is_empty() {
"no portable decision — its requests are observe-only".to_string()
} else {
accepted
.iter()
.map(|decision| format!("`{decision}`"))
.collect::<Vec<_>>()
.join(" or ")
},
));
}
match harness {
"claude-code" => Ok(match decision {
Decision::Allow => json!({"behavior": "allow"}),
Decision::Deny => json!({
"behavior": "deny",
"message": "denied through the supercode frontend",
}),
Decision::AllowForSession => unreachable!("refused above"),
}),
_ => Err(format!(
"no hosted permission reply is defined for {harness}"
)),
}
}
fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
let key = event.kind.to_ascii_lowercase().replace('-', "_");
let payload = &event.payload;
if key == "transport_closed" {
return vec![
json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
];
}
if key == "transport_error" || key == "error" {
return vec![
json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
];
}
if key == "session/update" {
let update = payload
.pointer("/params/update")
.or_else(|| payload.get("update"))
.unwrap_or(payload);
let update_kind = update
.get("sessionUpdate")
.or_else(|| update.get("type"))
.and_then(Value::as_str)
.unwrap_or_default()
.to_ascii_lowercase();
if update_kind == "agent_message_chunk" {
return vec![
json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
];
}
if update_kind == "agent_thought_chunk" {
return vec![
json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
];
}
if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
return vec![project_tool(update, payload)];
}
}
if key == "supercode/acp_request_completed" {
let failure = payload
.pointer("/params/error")
.or_else(|| payload.get("error"));
return completion(failure.and_then(extract_text));
}
match harness {
"codex" => project_codex(&key, payload),
"claude-code" => project_claude(&key, payload),
"pi" => project_pi(&key, payload),
"opencode" => project_opencode(&key, payload),
_ => project_generic(&key, payload),
}
}
fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
if key == "turn/started" {
return vec![native_payload(key, payload)];
}
if key == "turn/completed" {
let status = payload
.pointer("/params/turn/status")
.or_else(|| payload.pointer("/turn/status"))
.and_then(Value::as_str)
.unwrap_or("completed")
.to_ascii_lowercase();
return completion(
(status.contains("fail") || status.contains("error") || status.contains("cancel"))
.then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
);
}
if key.ends_with("/delta") {
let text = payload
.pointer("/params/delta")
.or_else(|| payload.get("delta"))
.and_then(extract_text)
.or_else(|| extract_text(payload))
.unwrap_or_default();
return vec![
json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
];
}
if key.contains("commandexecution")
|| key.contains("mcptool")
|| key.contains("filechange")
|| key.contains("tool")
{
let source = payload
.pointer("/params/item")
.or_else(|| payload.get("item"))
.unwrap_or(payload);
return vec![project_tool(source, payload)];
}
vec![native_payload(key, payload)]
}
fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
if key == "result" {
let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
|| payload
.get("subtype")
.and_then(Value::as_str)
.is_some_and(|subtype| subtype != "success");
return completion(
failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
);
}
if key == "stream_event" {
let stream = payload
.get("event")
.or_else(|| payload.get("stream_event"))
.unwrap_or(payload);
if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
return vec![native_payload(key, payload)];
}
}
if key == "assistant" {
return project_claude_blocks(payload, true);
}
if key == "user" {
return project_claude_blocks(payload, false);
}
if key == "control_request" {
return project_claude_control_request(payload);
}
if key == "system" {
if payload.get("subtype").and_then(Value::as_str) == Some("init") {
return Vec::new();
}
}
if key == "rate_limit_event" {
let status = payload
.pointer("/rate_limit_info/status")
.and_then(Value::as_str)
.unwrap_or_default();
if status.starts_with("allowed") {
return Vec::new();
}
}
project_generic(key, payload)
}
fn claude_blocks(payload: &Value) -> Option<&Vec<Value>> {
payload
.pointer("/message/content")
.or_else(|| payload.get("content"))
.and_then(Value::as_array)
}
fn project_claude_blocks(payload: &Value, assistant: bool) -> Vec<Value> {
let Some(blocks) = claude_blocks(payload) else {
let text = extract_text(payload).unwrap_or_default();
if text.is_empty() {
return vec![native_payload(
if assistant { "assistant" } else { "user" },
payload,
)];
}
return vec![
json!({"type":if assistant { "text_delta" } else { "user_message" }, "text":text, "raw":payload}),
];
};
let mut projected = Vec::new();
for block in blocks {
match block
.get("type")
.and_then(Value::as_str)
.unwrap_or_default()
{
"text" => {
let text = block
.get("text")
.and_then(Value::as_str)
.unwrap_or_default();
if !text.is_empty() {
projected.push(json!({
"type": if assistant { "text_delta" } else { "user_message" },
"text": text,
"raw": block,
}));
}
}
"thinking" | "redacted_thinking" => {
let text = extract_text(block).unwrap_or_default();
if !text.is_empty() {
projected.push(json!({"type":"reasoning", "text":text, "raw":block}));
}
}
"tool_use" => projected.push(json!({
"type": "tool_call_started",
"id": block.get("id").cloned().unwrap_or(Value::Null),
"name": block.get("name").cloned().unwrap_or(Value::Null),
"arguments": block.get("input").map(Value::to_string).unwrap_or_default(),
"raw": block,
})),
"tool_result" => projected.push(json!({
"type": "tool_call_completed",
"id": block.get("tool_use_id").cloned().unwrap_or(Value::Null),
"name": Value::Null,
"output": extract_text(block.get("content").unwrap_or(block)).unwrap_or_default(),
"is_error": block.get("is_error").and_then(Value::as_bool).unwrap_or(false),
"raw": block,
})),
_ => projected.push(native_payload(
if assistant { "assistant" } else { "user" },
block,
)),
}
}
projected
}
fn project_claude_control_request(payload: &Value) -> Vec<Value> {
let request = payload.get("request").unwrap_or(payload);
if request.get("subtype").and_then(Value::as_str) != Some("can_use_tool") {
return vec![native_payload("control_request", payload)];
}
let native_id = payload
.get("request_id")
.and_then(Value::as_str)
.unwrap_or_default();
vec![json!({
"type": "request",
"request": {
"id": claude_request_id(native_id),
"kind": "approval",
"payload": {
"tool": request.get("tool_name").or_else(|| request.get("toolName")).cloned().unwrap_or(Value::Null),
"arguments": request.get("input").cloned().unwrap_or(Value::Null),
"native_request_id": native_id,
"decisions": hosted_answerable_decisions("claude-code"),
},
},
})]
}
fn claude_request_id(native_id: &str) -> u64 {
let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
for byte in native_id.as_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
}
hash & ((1_u64 << 53) - 1)
}
fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
match key {
"agent_start" => vec![native_payload(key, payload)],
"agent_end" => completion(payload.get("error").and_then(extract_text)),
"message_update" => {
let update = payload
.get("assistantMessageEvent")
.or_else(|| payload.get("event"))
.unwrap_or(payload);
let kind = update
.get("type")
.and_then(Value::as_str)
.unwrap_or_default();
match kind {
"text_delta" => vec![
json!({"type":"text_delta", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
],
"thinking_delta" => vec![
json!({"type":"reasoning", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
],
_ => vec![native_payload(key, payload)],
}
}
_ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
_ => vec![native_payload(key, payload)],
}
}
fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
if key == "session.idle" {
return completion(None);
}
if key == "session.status" {
let status = payload
.pointer("/properties/status/type")
.or_else(|| payload.pointer("/status/type"))
.and_then(Value::as_str)
.unwrap_or_default();
if status == "busy" {
return vec![native_payload(key, payload)];
}
if status == "idle" {
return completion(None);
}
}
if key == "message.part.updated" {
let part = payload
.pointer("/properties/part")
.or_else(|| payload.get("part"))
.unwrap_or(payload);
if let Some(delta) = payload
.pointer("/properties/delta")
.or_else(|| payload.get("delta"))
.and_then(Value::as_str)
{
let reasoning = part
.get("type")
.and_then(Value::as_str)
.is_some_and(|kind| kind.contains("reasoning"));
return vec![
json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
];
}
if part
.get("type")
.and_then(Value::as_str)
.is_some_and(|kind| kind.contains("tool"))
{
return vec![project_tool(part, payload)];
}
}
if key == "session.error" {
return completion(Some(
extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
));
}
vec![native_payload(key, payload)]
}
fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
match key {
"turn_started" | "turn/started" | "agent_start" => {
vec![native_payload(key, payload)]
}
"turn_completed" | "turn/completed" | "agent_end" => {
completion(payload.get("error").and_then(extract_text))
}
"output_delta" | "content_delta" => vec![
json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
],
"reasoning_delta" => vec![
json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
],
"tool" => vec![project_tool(payload, payload)],
_ => vec![native_payload(key, payload)],
}
}
fn completion(error: Option<String>) -> Vec<Value> {
match error {
Some(message) => vec![
json!({"type":"turn_failed", "message":message}),
json!({"type":"turn_completed"}),
],
None => vec![
json!({"type":"turn_succeeded"}),
json!({"type":"turn_completed"}),
],
}
}
fn project_tool(source: &Value, raw: &Value) -> Value {
let status = source
.get("status")
.or_else(|| source.get("state"))
.or_else(|| source.get("sessionUpdate"))
.and_then(Value::as_str)
.unwrap_or_default()
.to_ascii_lowercase();
let completed = status.contains("complete")
|| status.contains("result")
|| status.contains("success")
|| status.contains("error")
|| status.contains("fail");
let arguments = source
.get("arguments")
.or_else(|| source.get("input"))
.or_else(|| source.get("rawInput"))
.cloned()
.unwrap_or(Value::Null);
json!({
"type": if completed { "tool_call_completed" } else { "tool_call_started" },
"id": source.get("toolCallId").or_else(|| source.get("tool_call_id")).or_else(|| source.get("callId")).or_else(|| source.get("id")).cloned().unwrap_or(Value::Null),
"name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
"arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
"output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
"is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
"raw": raw,
})
}
fn native_payload(kind: &str, payload: &Value) -> Value {
json!({"type":"native_event", "kind":kind, "raw":payload})
}
fn extract_text(value: &Value) -> Option<String> {
match value {
Value::String(text) => Some(text.clone()),
Value::Array(values) => {
let text = values
.iter()
.filter_map(extract_text)
.collect::<Vec<_>>()
.join("\n");
(!text.is_empty()).then_some(text)
}
Value::Object(object) => {
for key in ["text", "delta", "content", "message", "result", "error"] {
if let Some(text) = object.get(key).and_then(Value::as_str) {
return Some(text.to_string());
}
}
for key in [
"delta",
"content",
"message",
"error",
"data",
"part",
"params",
"properties",
"update",
"event",
] {
if let Some(text) = object.get(key).and_then(extract_text) {
return Some(text);
}
}
None
}
_ => None,
}
}
#[cfg(test)]
mod tests {
#[test]
fn pi_message_update_projects_only_the_deltas() {
use serde_json::json;
let delta = json!({"type":"message_update","assistantMessageEvent":{"type":"text_delta","contentIndex":0,"delta":"ack: hi"}});
let end = json!({"type":"message_update","assistantMessageEvent":{"type":"text_end","contentIndex":0,"content":"ack: hi"}});
let start = json!({"type":"message_update","assistantMessageEvent":{"type":"text_start","contentIndex":0}});
let projected = super::project_pi("message_update", &delta);
assert_eq!(projected[0]["type"], "text_delta");
assert_eq!(projected[0]["text"], "ack: hi");
for repeat in [&end, &start] {
let projected = super::project_pi("message_update", repeat);
assert_eq!(
projected[0]["type"], "native_event",
"{repeat} carries no new text"
);
}
let thinking = json!({"type":"message_update","assistantMessageEvent":{"type":"thinking_delta","contentIndex":0,"delta":"hm"}});
let projected = super::project_pi("message_update", &thinking);
assert_eq!(projected[0]["type"], "reasoning");
assert_eq!(projected[0]["text"], "hm");
}
use super::*;
use crate::{HarnessId, RuntimeEndpoint};
struct ControlledRuntime {
handle: RuntimeHandle,
events: mpsc::UnboundedReceiver<HarnessEvent>,
close_failures: usize,
}
#[async_trait]
impl RuntimeConnection for ControlledRuntime {
fn handle(&self) -> &RuntimeHandle {
&self.handle
}
async fn send_input(&mut self, _input: RuntimeInput) -> Result<Option<String>> {
Ok(Some("turn-1".into()))
}
async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
Ok(self.events.recv().await)
}
async fn interrupt(&mut self) -> Result<()> {
Ok(())
}
async fn respond(&mut self, _request_id: Value, _response: Value) -> Result<()> {
Ok(())
}
async fn close(&mut self) -> Result<()> {
if self.close_failures > 0 {
self.close_failures -= 1;
return Err(Error::Other("cleanup temporarily unavailable".into()));
}
Ok(())
}
}
fn controlled_runtime() -> (
Box<dyn RuntimeConnection>,
mpsc::UnboundedSender<HarnessEvent>,
) {
controlled_runtime_with_close_failures(0)
}
fn controlled_runtime_with_close_failures(
close_failures: usize,
) -> (
Box<dyn RuntimeConnection>,
mpsc::UnboundedSender<HarnessEvent>,
) {
let (events, event_rx) = mpsc::unbounded_channel();
(
Box::new(ControlledRuntime {
handle: RuntimeHandle {
harness: HarnessId::from(HarnessId::PI),
runtime_id: "shared-runtime".into(),
endpoint: RuntimeEndpoint::LocalProcess {
pid: None,
command: vec!["controlled-runtime".into()],
protocol: "test".into(),
},
},
events: event_rx,
close_failures,
}),
events,
)
}
#[tokio::test]
async fn failed_runtime_close_does_not_stop_the_host_or_discard_its_owner() {
let (runtime, _events) = controlled_runtime_with_close_failures(1);
let (host, mut owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
assert!(owner.close().await.is_err());
assert!(!owner.closed);
assert!(!host.closed.load(Ordering::SeqCst));
owner
.send_input(RuntimeInput {
text: "still usable".into(),
image_urls: Vec::new(),
})
.await
.unwrap();
owner.close().await.unwrap();
assert!(owner.closed);
assert!(host.closed.load(Ordering::SeqCst));
}
fn capabilities() -> RuntimeCapabilities {
RuntimeCapabilities {
start_session: true,
resume_session: true,
attach_existing_process: false,
send_input: true,
stream_events: true,
interrupt: true,
steer: false,
respond_to_requests: false,
}
}
#[tokio::test]
async fn native_eof_closes_the_raw_owner_connection() {
let (runtime, events) = controlled_runtime();
let (_host, mut connection) = HostedHarnessRuntime::spawn(runtime, capabilities());
drop(events);
let event =
tokio::time::timeout(std::time::Duration::from_secs(1), connection.next_event())
.await
.expect("raw owner should not hang after native EOF")
.unwrap()
.expect("EOF is projected as an explicit terminal event");
assert_eq!(event.kind, "transport_closed");
assert_eq!(event.payload["terminal"], true);
}
#[tokio::test]
async fn adapters_without_an_operation_route_preserve_the_requested_id() {
let (runtime, _events) = controlled_runtime();
let (host, _owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
let operation_id = "prompt:not-advertised".to_string();
let error = FrontendRuntime::invoke(
host.as_ref(),
crate::FrontendOperationInvocation::Prompt {
operation_id: operation_id.clone(),
arguments: String::new(),
},
)
.await
.unwrap_err();
assert!(
matches!(error, FrontendRuntimeError::UnsupportedOperation(id) if id == operation_id)
);
}
#[tokio::test]
async fn interrupt_stays_busy_until_the_native_terminal_event() {
let (runtime, events) = controlled_runtime();
let (host, mut owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
let mut terminal = FrontendRuntime::attach(host.as_ref(), 100).await.unwrap();
assert_eq!(
FrontendRuntime::submit(host.as_ref(), "hello".into())
.await
.unwrap(),
"turn-1"
);
assert_eq!(terminal.next_event().await.unwrap().kind, "user_message");
assert_eq!(terminal.next_event().await.unwrap().kind, "turn_started");
assert!(FrontendRuntime::interrupt(host.as_ref()).await.unwrap());
assert_eq!(
FrontendRuntime::describe(host.as_ref())
.await
.unwrap()
.turn_state,
FrontendTurnState::Busy
);
assert!(
tokio::time::timeout(std::time::Duration::from_millis(20), terminal.next_event())
.await
.is_err(),
"interrupt acceptance must not manufacture turn completion"
);
events
.send(HarnessEvent {
sequence: None,
kind: "agent_end".into(),
payload: json!({}),
})
.unwrap();
assert_eq!(terminal.next_event().await.unwrap().kind, "turn_succeeded");
assert_eq!(terminal.next_event().await.unwrap().kind, "turn_completed");
assert_eq!(
FrontendRuntime::describe(host.as_ref())
.await
.unwrap()
.turn_state,
FrontendTurnState::Idle
);
owner.close().await.unwrap();
}
#[test]
fn native_start_events_do_not_duplicate_the_hosted_turn_boundary() {
let event = HarnessEvent {
sequence: None,
kind: "turn/started".into(),
payload: json!({"method":"turn/started"}),
};
let projected = project_native_event(HarnessId::CODEX, &event);
assert_eq!(projected.len(), 1);
assert_eq!(projected[0]["type"], "native_event");
}
}