use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::time::Duration;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use crate::runtime::generated_session_id;
#[cfg(feature = "adapter-api")]
use crate::runtime::{HostedHarnessConnection, HostedHarnessRuntime};
use crate::sdk::{
discover_session_page, load_session, load_session_with_fidelity, SdkCapabilities, SdkError,
SdkErrorCode, SdkEvent, SdkOperation, SdkRequest, SdkRuntimeEvent, SdkService,
};
use crate::watch::{bound_session_view, message_json, normalized_session_json};
use crate::Fidelity;
#[cfg(feature = "adapter-api")]
use crate::SupercodeHttpRuntimeBackend;
use crate::{
discover_live_runtime, harness_support_registry, AcpRuntimeBackend, ClaudeCodeRuntimeBackend,
CodexRuntimeBackend, DiscoveryQuery, HarnessCatalog, HarnessId, ImplementationKind,
LiveRuntimeEndpoint, LiveRuntimeSource, OpenCodeRuntimeBackend, PiRuntimeBackend, Role,
RuntimeAttachRequest, RuntimeBackend, RuntimeConnection, RuntimeInput, RuntimeLaunch,
RuntimeStartRequest, Session, SessionFollower, SessionFormat, SessionLocator, SessionSource,
};
use crate::{reduce, tokens};
#[cfg(feature = "adapter-api")]
use crate::{register_live_runtime, resolve_live_runtime, LiveRuntimeRegistration};
pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
pub struct HarnessSessionService {
catalog: HarnessCatalog,
followers: BTreeMap<String, SessionFollower>,
followed_sources: BTreeMap<String, FollowedSource>,
next_subscription: u64,
runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
terminal_launches: BTreeMap<String, StructuredLaunch>,
runtime_sequences: BTreeMap<String, u64>,
next_runtime: u64,
reduction_store_root: Option<PathBuf>,
}
impl Default for HarnessSessionService {
fn default() -> Self {
Self::new()
}
}
impl HarnessSessionService {
pub fn new() -> Self {
Self {
catalog: HarnessCatalog::new(),
followers: BTreeMap::new(),
followed_sources: BTreeMap::new(),
next_subscription: 1,
runtimes: BTreeMap::new(),
terminal_launches: BTreeMap::new(),
runtime_sequences: BTreeMap::new(),
next_runtime: 1,
reduction_store_root: None,
}
}
pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
self.reduction_store_root = Some(root.into());
self
}
#[cfg(feature = "adapter-api")]
pub fn handle(&mut self, request: Value) -> Value {
let id = request.get("id").cloned().unwrap_or(Value::Null);
if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
}
let Some(method) = request.get("method").and_then(Value::as_str) else {
return rpc_error(id, -32600, "request is missing `method`");
};
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
match self.call(method, params) {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
}
}
#[cfg(feature = "adapter-api")]
pub async fn handle_async(&mut self, request: Value) -> Value {
let method = request
.get("method")
.and_then(Value::as_str)
.unwrap_or_default();
if matches!(
method,
"harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
) {
let id = request.get("id").cloned().unwrap_or(Value::Null);
if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
}
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
return match self.inventory_call(method, params).await {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
};
}
if method == "harness.v1.sessions.message" {
let id = request.get("id").cloned().unwrap_or(Value::Null);
if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
}
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
return match self.message_call(params).await {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
};
}
if let Some(operation) = SdkOperation::from_method(method) {
let id = request.get("id").cloned().unwrap_or(Value::Null);
if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
}
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
return match self.execute(SdkRequest { operation, params }).await {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(error) => sdk_rpc_error(id, &error),
};
}
if !method.starts_with("harness.v1.runtimes.") {
return self.handle(request);
}
let id = request.get("id").cloned().unwrap_or(Value::Null);
if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
}
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
match self.runtime_call(method, params).await {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
}
}
#[cfg(feature = "adapter-api")]
pub fn poll(&mut self) -> Vec<Value> {
let mut notifications = Vec::new();
for (subscription, follower) in &mut self.followers {
match follower.poll() {
Ok(Some(event)) => notifications.push(json!({
"jsonrpc": "2.0",
"method": SESSION_EVENT_METHOD,
"params": {
"subscription": subscription,
"event": event.to_json(),
}
})),
Ok(None) => {}
Err(error) => notifications.push(json!({
"jsonrpc": "2.0",
"method": SESSION_EVENT_METHOD,
"params": {
"subscription": subscription,
"event": {
"type": "watch_error",
"recoverable": true,
"message": error.to_string(),
},
}
})),
}
}
notifications
}
#[cfg(feature = "adapter-api")]
pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
let registry = crate::LocalRuntimeRegistry::new();
let authorization = crate::RuntimeAuthorization::observer();
let mut notifications = Vec::new();
for (subscription, source) in &mut self.followed_sources {
let state = match registry
.source_state(&source.harness, &source.session_id, &authorization)
.await
{
Ok(Some(state)) => state,
Ok(None) => crate::RuntimeRegistryState::Persisted,
Err(_) => continue,
};
if source.reported.as_deref() == Some(state.as_str()) {
continue;
}
source.reported = Some(state.as_str().to_string());
notifications.push(json!({
"jsonrpc": "2.0",
"method": SESSION_EVENT_METHOD,
"params": {
"subscription": subscription,
"event": {"type": "runtime_state", "state": state.as_str()},
},
}));
}
notifications
}
#[cfg(feature = "adapter-api")]
pub async fn poll_runtimes(&mut self) -> Vec<Value> {
self.poll_sdk_events()
.await
.into_iter()
.map(|(connection, runtime_event)| {
json!({
"jsonrpc": "2.0",
"method": RUNTIME_EVENT_METHOD,
"params": {
"connection": connection,
"session_id": runtime_event.session_id,
"sequence": runtime_event.event.sequence,
"event": {
"kind": runtime_event.event.kind,
"payload": runtime_event.event.payload,
},
},
})
})
.collect()
}
async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
let mut events = Vec::new();
let mut closed = Vec::new();
for (connection, runtime) in &mut self.runtimes {
let session_id = runtime.handle().runtime_id.clone();
match tokio::time::timeout(Duration::from_millis(1), runtime.next_event()).await {
Ok(Ok(Some(event))) => {
let terminal = event.kind == "transport_closed";
let next_sequence = self
.runtime_sequences
.entry(session_id.clone())
.or_insert(0);
let sequence = event.sequence.unwrap_or_else(|| {
*next_sequence = next_sequence.saturating_add(1);
*next_sequence
});
*next_sequence = (*next_sequence).max(sequence);
events.push((
connection.clone(),
SdkRuntimeEvent {
session_id: session_id.clone(),
event: SdkEvent {
sequence,
kind: event.kind,
payload: event.payload,
},
},
));
if terminal {
closed.push(connection.clone());
}
}
Ok(Ok(None)) => {
let sequence = self
.runtime_sequences
.entry(session_id.clone())
.or_insert(0);
*sequence = sequence.saturating_add(1);
events.push((
connection.clone(),
SdkRuntimeEvent {
session_id,
event: SdkEvent {
sequence: *sequence,
kind: "transport_closed".into(),
payload: json!({"message": "Harness runtime transport closed."}),
},
},
));
closed.push(connection.clone());
}
Err(_) => {}
Ok(Err(error)) => {
let sequence = self
.runtime_sequences
.entry(session_id.clone())
.or_insert(0);
*sequence = sequence.saturating_add(1);
events.push((
connection.clone(),
SdkRuntimeEvent {
session_id,
event: SdkEvent {
sequence: *sequence,
kind: "transport_error".into(),
payload: json!({"message": error.to_string(), "terminal": true}),
},
},
));
closed.push(connection.clone());
}
}
}
for connection in closed {
if let Some(runtime) = self.runtimes.remove(&connection) {
self.runtime_sequences.remove(&runtime.handle().runtime_id);
}
self.terminal_launches.remove(&connection);
}
events
}
fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
match method {
"harness.v1.capabilities" => Ok(json!({
"version": HARNESS_SERVICE_VERSION,
"sdk": self.capabilities(),
"methods": [
"harness.v1.support.report",
"harness.v1.harnesses.list",
"harness.v1.harnesses.probe",
"harness.v1.sessions.discover",
"harness.v1.sessions.load",
"harness.v1.sessions.follow",
"harness.v1.sessions.unfollow",
"harness.v1.sessions.message",
"harness.v1.sessions.import",
"harness.v1.sessions.export",
"harness.v1.sessions.translate",
"harness.v1.sessions.reduce",
"harness.v1.sessions.branch",
"harness.v1.sessions.handoff",
"harness.v1.sessions.resume_instructions",
"harness.v1.runtimes.capabilities",
"harness.v1.runtimes.start",
"harness.v1.runtimes.resume",
"harness.v1.runtimes.attach_existing",
"harness.v1.runtimes.attach",
"harness.v1.runtimes.send_input",
"harness.v1.runtimes.interrupt",
"harness.v1.runtimes.steer",
"harness.v1.runtimes.respond",
"harness.v1.runtimes.terminal_instructions",
"harness.v1.runtimes.close",
],
"notifications": [SESSION_EVENT_METHOD, RUNTIME_EVENT_METHOD],
"harnesses": harness_support_registry()
.harnesses
.into_iter()
.map(|harness| harness.id)
.collect::<Vec<_>>(),
})),
"harness.v1.support.report" => serde_json::to_value(harness_support_registry())
.map_err(|error| ServiceError::Operation(error.to_string())),
"harness.v1.sessions.discover" => {
let query = decode::<DiscoveryQuery>(params)?;
let page = discover_session_page(&query).map_err(operation)?;
let peers = if page
.sessions
.iter()
.any(|session| session.locator.harness.as_str() == HarnessId::CLAUDE_CODE)
{
crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(
&query.homes,
))
} else {
Vec::new()
};
let sessions = page
.sessions
.into_iter()
.map(|session| {
let mut value = serde_json::to_value(&session)
.map_err(|error| ServiceError::Operation(error.to_string()))?;
if let Some(workspace) = &session.cwd {
let source = LiveRuntimeSource {
harness: session.locator.harness.as_str().to_string(),
session_id: session.locator.session_id.clone(),
workspace: workspace.clone(),
};
if let Some(endpoint) = discover_live_runtime(&source)
.map_err(|error| ServiceError::Operation(error.to_string()))?
{
value["live_endpoint"] = json!(endpoint.as_str());
}
}
if let Some(peer) = peers.iter().find(|peer| {
session.locator.harness.as_str() == HarnessId::CLAUDE_CODE
&& peer.session_id == session.locator.session_id
}) {
if value.get("live_endpoint").is_none() {
value["live_endpoint"] = json!(peer.endpoint().as_str());
}
if let Some(status) = peer.status {
value["live_status"] = json!(status.as_str());
}
}
Ok(value)
})
.collect::<std::result::Result<Vec<_>, ServiceError>>()?;
Ok(json!({"sessions": sessions, "next_cursor": page.next_cursor}))
}
"harness.v1.sessions.load" => {
let params = decode::<LoadSessionParams>(params)?;
if let Some(options) = ¶ms.options {
options.validate()?;
return load_session(¶ms.read.locator)
.map(|session| projected_session_result(&session, options))
.map_err(operation);
}
let mut session = if params.read.display_history() {
self.catalog
.load_display_view(
¶ms.read.locator,
params.read.read_fidelity(),
params.read.tail_messages().unwrap_or(500),
)
.map_err(crate::Error::from)
} else if params.read.include_subagents() {
load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
} else {
self.catalog
.load_parent_with_fidelity(
¶ms.read.locator,
params.read.read_fidelity(),
)
.map_err(crate::Error::from)
}
.map_err(operation)?;
params.read.bound_session(&mut session);
Ok(json!({"session": normalized_session_json(&session)}))
}
"harness.v1.sessions.follow" => {
let params = decode::<LocatorParams>(params)?;
let mut follower = self
.catalog
.follow_read_view(
¶ms.locator,
params.read_fidelity(),
params.include_subagents(),
params.tail_messages(),
params.max_message_chars(),
params.display_history(),
)
.map_err(operation)?;
let initial = follower
.poll()
.map_err(operation)?
.map(|event| event.to_json());
let subscription = format!("sub-{}", self.next_subscription);
self.next_subscription += 1;
self.followers.insert(subscription.clone(), follower);
self.followed_sources.insert(
subscription.clone(),
FollowedSource {
harness: params.locator.harness.as_str().to_string(),
session_id: params.locator.session_id.clone(),
reported: None,
},
);
Ok(json!({"subscription": subscription, "initial": initial}))
}
"harness.v1.sessions.unfollow" => {
let params = decode::<UnfollowParams>(params)?;
self.followed_sources.remove(¶ms.subscription);
Ok(json!({
"removed": self.followers.remove(¶ms.subscription).is_some()
}))
}
"harness.v1.sessions.import" => {
let params = decode::<ImportSessionParams>(params)?;
let session = Session::load_str(¶ms.content, params.source_harness.into())
.map_err(operation)?;
Ok(json!({"session": normalized_session_json(&session)}))
}
"harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
let params = decode::<ExportSessionParams>(params)?;
let session = load_session(¶ms.locator).map_err(operation)?;
let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
Ok(json!({"artifact": artifact}))
}
"harness.v1.sessions.reduce" => {
let params = decode::<ReduceSessionParams>(params)?;
self.reduce_session(params)
}
"harness.v1.sessions.branch" => {
let params = decode::<BranchSessionParams>(params)?;
let session = load_session(¶ms.locator).map_err(operation)?;
let storage = params.locator.storage.path().display().to_string();
let bootstrap_prompt = format!(
"Continue as a new branch from {} session {}. The frozen parent transcript is at {}. Read or load that parent for context, summarize the relevant state, then continue independently without mutating the parent session.",
params.locator.harness.as_str(), params.locator.session_id, storage
);
let artifact = params
.target_harness
.map(|target| session_artifact(¶ms.locator, &session, target))
.transpose()?;
Ok(json!({
"parent": params.locator,
"session": normalized_session_json(&session),
"bootstrap_prompt": bootstrap_prompt,
"artifact": artifact,
}))
}
"harness.v1.sessions.handoff" => {
let params = decode::<HandoffSessionParams>(params)?;
let session = load_session(¶ms.locator).map_err(operation)?;
let cwd = params
.cwd
.or_else(|| session.meta.cwd.clone())
.unwrap_or_else(|| PathBuf::from("."));
let artifact =
handoff_artifact(¶ms.locator, &session, params.target_harness, &cwd)?;
let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
ServiceError::Operation(
"handoff artifact omitted target session identity".into(),
)
})?;
let instructions =
handoff_instructions(params.target_harness, target_session_id, &cwd);
Ok(json!({
"artifact": artifact,
"launch": instructions.launch,
"materialize": instructions.materialize,
"requires_materialization": instructions.requires_materialization,
"note": instructions.note,
}))
}
"harness.v1.sessions.resume_instructions" => {
let params = decode::<ResumeInstructionsParams>(params)?;
let session = load_session(¶ms.locator).map_err(operation)?;
let cwd = params
.cwd
.or(session.meta.cwd)
.unwrap_or_else(|| PathBuf::from("."));
let launch = resume_launch(
params.locator.harness.as_str(),
¶ms.locator.session_id,
&cwd,
params.policy,
)?;
Ok(json!({"launch": launch}))
}
_ => Err(ServiceError::MethodNotFound),
}
}
fn reduce_session(
&self,
params: ReduceSessionParams,
) -> std::result::Result<Value, ServiceError> {
let session = load_session(¶ms.locator).map_err(operation)?;
if session.messages.is_empty() {
return Err(ServiceError::InvalidParams(
"cannot reduce an empty session".into(),
));
}
let keep_last = params.keep_last.clamp(1, 128);
let policy = reduce::ReductionPolicy {
clear_turns_older_than: Some(keep_last),
..Default::default()
};
let (view, log) =
reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
if log.reductions.is_empty() {
return Err(ServiceError::UnsupportedAction(format!(
"session `{}` is already too small for a meaningful reversible reduction",
params.locator.session_id
)));
}
let source_tokens = tokens::estimate_view_tokens(&session.messages);
let reduced_tokens = tokens::estimate_view_tokens(&view);
if reduced_tokens >= source_tokens {
return Err(ServiceError::UnsupportedAction(format!(
"session `{}` has no token-reducing reversible projection",
params.locator.session_id
)));
}
let store_root = self
.reduction_store_root
.clone()
.unwrap_or_else(default_reduction_store_root);
let store = crate::SessionStore::open(&store_root).map_err(operation)?;
let rescue_id = format!("rescue-{}", generated_session_id());
let imported = session
.imported_message_count
.unwrap_or(session.messages.len())
.min(session.messages.len());
let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
let view_jsonl = messages_jsonl(&view)?;
let title = format!(
"Reduced {} continuation from {}",
params.target_harness.id(),
params.locator.session_id
);
store
.save_sidecar(&rescue_id, &sidecar_jsonl)
.map_err(operation)?;
store
.save_reduction_log(&rescue_id, &log)
.map_err(operation)?;
store
.save(&rescue_id, &title, &view_jsonl)
.map_err(operation)?;
let source_bytes = serde_json::to_vec(&session.messages)
.map_err(|error| ServiceError::Operation(error.to_string()))?
.len() as u64;
let reduced_bytes = serde_json::to_vec(&view)
.map_err(|error| ServiceError::Operation(error.to_string()))?
.len() as u64;
store
.set_reduction_stats(
&rescue_id,
&title,
source_bytes,
reduced_bytes,
log.reductions.len() as u32,
)
.map_err(operation)?;
let reloaded_sidecar = store
.load_sidecar(&rescue_id)
.map_err(operation)?
.ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
let reloaded_log = store
.load_reduction_log(&rescue_id)
.map_err(operation)?
.ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
let (restamped_view, restamped_log) =
reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
return Err(ServiceError::Operation(
"persisted reduction view does not match its durable log and sidecar".into(),
));
}
if restamped_log != reloaded_log {
return Err(ServiceError::Operation(
"reapplying the durable reduction log changed its identity".into(),
));
}
let inverted =
reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
if inverted != session.messages {
return Err(ServiceError::Operation(
"reduction inversion did not restore the source messages byte-exactly".into(),
));
}
let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
let sidecar_path = store.sidecar_path(&rescue_id);
let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
let bootstrap_prompt = reduced_bootstrap_prompt(
¶ms.locator,
params.target_harness,
&view_jsonl,
&sidecar_path,
&reduction_log_path,
);
let mut reduced_session = session.clone();
reduced_session.meta.session_id = Some(rescue_id.clone());
reduced_session.messages = view;
Ok(json!({
"session": normalized_session_json(&reduced_session),
"bootstrap_prompt": bootstrap_prompt,
"receipt": {
"id": rescue_id,
"sidecar_id": rescue_id,
"source_harness": params.locator.harness,
"target_harness": params.target_harness.id(),
"source_tokens": source_tokens,
"reduced_tokens": reduced_tokens,
"ratio": ratio,
"source_bytes": source_bytes,
"reduced_bytes": reduced_bytes,
"reductions": reloaded_log.reductions.len(),
"sidecar_path": sidecar_path,
"reduction_log_path": reduction_log_path,
"verified": true,
"reversible": true,
}
}))
}
async fn runtime_call(
&mut self,
method: &str,
params: Value,
) -> std::result::Result<Value, ServiceError> {
match method {
"harness.v1.runtimes.capabilities" => {
let params = decode::<RuntimeBackendParams>(params)?;
let backend = runtime_backend(¶ms)?;
Ok(json!({
"harness": backend.harness(),
"capabilities": backend.capabilities(),
}))
}
"harness.v1.runtimes.start" => {
let params = decode::<RuntimeStartParams>(params)?;
let backend = runtime_backend(¶ms.backend)?;
let capabilities = backend.capabilities();
let workspace = params.cwd.clone();
let runtime = backend
.start(RuntimeStartRequest {
cwd: params.cwd,
launch: runtime_launch(¶ms.backend),
})
.await
.map_err(operation)?;
self.insert_hosted_runtime(runtime, capabilities, workspace)
.await
}
"harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
let params = decode::<RuntimeAttachParams>(params)?;
let backend = runtime_backend(¶ms.backend)?;
let capabilities = backend.capabilities();
let workspace = params.cwd.clone().unwrap_or_else(|| {
std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))
});
let runtime = backend
.attach(RuntimeAttachRequest {
runtime_id: params.runtime_id,
cwd: params.cwd,
launch: runtime_launch(¶ms.backend),
})
.await
.map_err(operation)?;
self.insert_hosted_runtime(runtime, capabilities, workspace)
.await
}
"harness.v1.runtimes.attach_existing" => {
let params = decode::<RuntimeAttachParams>(params)?;
let backend: Box<dyn RuntimeBackend> = match params
.backend
.base_url
.as_deref()
.and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
{
Some(endpoint) => {
#[cfg(not(feature = "adapter-api"))]
{
let _ = endpoint;
return Err(ServiceError::UnsupportedAction(
"live HTTP attachment adapter is not compiled".into(),
));
}
#[cfg(feature = "adapter-api")]
{
let workspace = params.cwd.clone().ok_or_else(|| {
ServiceError::InvalidParams(
"Supercode live attach requires the project cwd".into(),
)
})?;
let source = LiveRuntimeSource {
harness: params.backend.harness.as_str().to_string(),
session_id: params.runtime_id.clone(),
workspace,
};
let receipt = resolve_live_runtime(&endpoint, &source)
.map_err(|error| ServiceError::Operation(error.to_string()))?;
Box::new(SupercodeHttpRuntimeBackend::new(receipt))
}
}
None => runtime_backend(¶ms.backend)?,
};
if !backend.capabilities().attach_existing_process {
return Err(ServiceError::Operation(format!(
"{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
backend.harness().as_str()
)));
}
let runtime = backend
.attach_existing(RuntimeAttachRequest {
runtime_id: params.runtime_id,
cwd: params.cwd,
launch: runtime_launch(¶ms.backend),
})
.await
.map_err(operation)?;
self.insert_runtime(runtime)
}
"harness.v1.runtimes.send_input" => {
let params = decode::<RuntimeInputParams>(params)?;
let runtime = self.runtime_mut(¶ms.connection)?;
let turn_id = runtime
.send_input(RuntimeInput { text: params.text })
.await
.map_err(operation)?;
Ok(json!({"turn_id": turn_id}))
}
"harness.v1.runtimes.interrupt" => {
let params = decode::<RuntimeConnectionParams>(params)?;
self.runtime_mut(¶ms.connection)?
.interrupt()
.await
.map_err(operation)?;
Ok(json!({}))
}
"harness.v1.runtimes.steer" => Err(ServiceError::UnsupportedAction(
"steer is not supported by this harness-native runtime adapter".into(),
)),
"harness.v1.runtimes.respond" => {
let params = decode::<RuntimeRespondParams>(params)?;
self.runtime_mut(¶ms.connection)?
.respond(params.request_id, params.response)
.await
.map_err(operation)?;
Ok(json!({}))
}
"harness.v1.runtimes.terminal_instructions" => {
let params = decode::<RuntimeConnectionParams>(params)?;
let launch = self
.terminal_launches
.get(¶ms.connection)
.ok_or_else(|| {
ServiceError::Operation(
"this runtime is not hosted for terminal attachment".into(),
)
})?;
Ok(json!({"launch":launch}))
}
"harness.v1.runtimes.close" => {
let params = decode::<RuntimeConnectionParams>(params)?;
let Some(mut runtime) = self.runtimes.remove(¶ms.connection) else {
return Err(ServiceError::InvalidParams(format!(
"unknown runtime connection `{}`",
params.connection
)));
};
self.terminal_launches.remove(¶ms.connection);
self.runtime_sequences.remove(&runtime.handle().runtime_id);
runtime.close().await.map_err(operation)?;
Ok(json!({"closed": true}))
}
_ => Err(ServiceError::MethodNotFound),
}
}
#[cfg(feature = "adapter-api")]
async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
let params = decode::<MessageSessionParams>(params)?;
Ok(message_live_session(¶ms, &crate::claude_peer::ProcessCourierRunner).await)
}
fn insert_runtime(
&mut self,
runtime: Box<dyn RuntimeConnection>,
) -> std::result::Result<Value, ServiceError> {
let connection = format!("runtime-{}", self.next_runtime);
self.next_runtime += 1;
let handle = runtime.handle().clone();
self.runtime_sequences
.entry(handle.runtime_id.clone())
.or_insert(0);
self.runtimes.insert(connection.clone(), runtime);
Ok(json!({"connection": connection, "handle": handle}))
}
#[cfg(feature = "adapter-api")]
async fn insert_hosted_runtime(
&mut self,
runtime: Box<dyn RuntimeConnection>,
capabilities: crate::RuntimeCapabilities,
workspace: PathBuf,
) -> std::result::Result<Value, ServiceError> {
let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities);
let token: std::sync::Arc<str> = crate::server::generate_token().into();
let server = crate::server::run_frontend_http(
host.clone(),
host.frontend_sender(),
"127.0.0.1:0",
token.clone(),
)
.await
.map_err(|error| ServiceError::Operation(error.to_string()))?;
let source = LiveRuntimeSource {
harness: connection.handle().harness.as_str().to_string(),
session_id: connection.handle().runtime_id.clone(),
workspace: workspace.clone(),
};
let registration = register_live_runtime(
connection.handle().runtime_id.clone(),
source.clone(),
format!("http://{}", server.address()),
token.to_string(),
)
.map_err(|error| ServiceError::Operation(error.to_string()))?;
let endpoint = registration.endpoint().to_string();
let launch = StructuredLaunch {
cwd: workspace,
program: std::env::current_exe()
.ok()
.map(|path| path.to_string_lossy().into_owned())
.unwrap_or_else(|| "supercode".into()),
arguments: vec![
"harness".into(),
"attach".into(),
"--endpoint".into(),
endpoint,
"--harness".into(),
source.harness,
"--session".into(),
source.session_id,
],
env: BTreeMap::new(),
};
let lease = HostedRuntimeLease {
connection,
_host: host,
_registration: registration,
_server: server,
};
let opened = self.insert_runtime(Box::new(lease))?;
let connection_id = opened["connection"]
.as_str()
.expect("insert_runtime returns a connection id")
.to_string();
self.terminal_launches.insert(connection_id, launch);
Ok(opened)
}
#[cfg(not(feature = "adapter-api"))]
async fn insert_hosted_runtime(
&mut self,
runtime: Box<dyn RuntimeConnection>,
_capabilities: crate::RuntimeCapabilities,
_workspace: PathBuf,
) -> std::result::Result<Value, ServiceError> {
self.insert_runtime(runtime)
}
fn runtime_mut(
&mut self,
connection: &str,
) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
self.runtimes.get_mut(connection).ok_or_else(|| {
ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
})
}
async fn inventory_call(
&self,
method: &str,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let mut params = decode::<HarnessInventoryParams>(params)?;
if method == "harness.v1.harnesses.probe" {
let harness = params.harness.take().ok_or_else(|| {
ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
})?;
params.harnesses = vec![harness];
}
let selected = params
.harnesses
.iter()
.map(HarnessId::as_str)
.collect::<std::collections::BTreeSet<_>>();
let supported = harness_support_registry()
.harnesses
.into_iter()
.filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
.collect::<Vec<_>>();
if !params.harnesses.is_empty() && supported.len() != selected.len() {
let known = supported
.iter()
.map(|harness| harness.id.as_str())
.collect::<std::collections::BTreeSet<_>>();
let missing = params
.harnesses
.iter()
.filter(|id| !known.contains(id.as_str()))
.map(HarnessId::as_str)
.collect::<Vec<_>>();
return Err(ServiceError::InvalidParams(format!(
"unknown harness(es): {}",
missing.join(", ")
)));
}
let global_counts = params
.include_sessions
.then(|| self.session_counts(None, ¶ms.harnesses));
let workspace_counts = params.include_sessions.then(|| {
params
.workspace
.as_deref()
.map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
});
let probes = supported.into_iter().map(|descriptor| {
let global = global_counts
.as_ref()
.map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
let workspace = workspace_counts
.as_ref()
.and_then(Option::as_ref)
.map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
self.probe_harness(descriptor, ¶ms, global, workspace)
});
let harnesses = futures::future::join_all(probes).await;
serde_json::to_value(HarnessInventoryReport {
probe: params.probe,
workspace: params.workspace,
harnesses,
})
.map_err(|error| ServiceError::Operation(error.to_string()))
}
async fn probe_harness(
&self,
descriptor: crate::HarnessSupportDescriptor,
params: &HarnessInventoryParams,
global: Option<usize>,
workspace: Option<usize>,
) -> LocalHarness {
let launch = descriptor.runtime.default_launch.as_ref();
let executable = launch.and_then(|launch| find_executable(&launch.program));
let installed = executable.is_some();
let version = if params.skip_versions {
None
} else {
match executable.as_deref() {
Some(path) => executable_version(path).await,
None => None,
}
};
let configured = auth_evidence(descriptor.id.as_str());
let mut auth = if configured {
HarnessAuthState::Configured
} else {
HarnessAuthState::Unknown
};
let mut runtime = if installed {
HarnessRuntimeState::Degraded
} else {
HarnessRuntimeState::Unavailable
};
let mut reason = (!installed).then(|| {
format!(
"{} is supported but `{}` was not found on PATH",
descriptor.display_name,
launch
.map(|launch| launch.program.as_str())
.unwrap_or("executable")
)
});
let mut repair = (!installed).then(|| {
format!(
"Install {} and ensure `{}` is on PATH.",
descriptor.display_name,
launch
.map(|launch| launch.program.as_str())
.unwrap_or("its executable")
)
});
if installed && params.probe == HarnessProbeLevel::Handshake {
let backend_params = RuntimeBackendParams {
harness: descriptor.id.clone(),
protocol: None,
launch: None,
base_url: None,
policy: RuntimePolicy::Default,
};
match runtime_backend(&backend_params) {
Ok(backend) => {
let cwd = params
.workspace
.clone()
.or_else(|| std::env::current_dir().ok())
.unwrap_or_else(|| PathBuf::from("."));
let isolated = descriptor
.runtime
.default_launch
.clone()
.and_then(|launch| {
IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok()
});
let Some(isolated) = isolated else {
reason = Some(
"No-prompt runtime handshake could not create its isolated harness home."
.into(),
);
repair = Some(
"Check temporary-directory permissions, then run the handshake probe again."
.into(),
);
return LocalHarness {
id: descriptor.id,
display_name: descriptor.display_name,
supported: true,
installed,
executable: executable.map(|path| path.to_string_lossy().into_owned()),
version,
auth,
runtime,
protocol: descriptor.runtime.protocol,
capabilities: descriptor.runtime.capabilities.clone(),
effective_capabilities: descriptor.runtime.capabilities,
sessions: HarnessSessionCounts { global, workspace },
reason,
repair,
};
};
match tokio::time::timeout(
Duration::from_secs(30),
backend.start(RuntimeStartRequest {
cwd,
launch: Some(isolated.launch.clone()),
}),
)
.await
{
Ok(Ok(mut connection)) => {
match stabilize_handshake(connection.as_mut()).await {
Ok(()) => {
auth = HarnessAuthState::Ready;
runtime = HarnessRuntimeState::Ready;
reason = Some(
"No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
.into(),
);
repair = None;
}
Err(message) => {
auth = if looks_like_auth_error(&message) {
HarnessAuthState::Required
} else if configured {
HarnessAuthState::Configured
} else {
HarnessAuthState::Unknown
};
reason = Some(format!(
"No-prompt runtime handshake became unhealthy during startup: {message}"
));
repair = Some(if auth == HarnessAuthState::Required {
format!(
"Run `{}` interactively once and complete sign-in, then probe again.",
launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
)
} else {
"Run the harness directly to inspect its startup failure, then probe again."
.into()
});
}
}
let _ =
tokio::time::timeout(Duration::from_secs(3), connection.close())
.await;
}
Ok(Err(error)) => {
let message = truncate_text(&error.to_string(), 500);
auth = if looks_like_auth_error(&message) {
HarnessAuthState::Required
} else if configured {
HarnessAuthState::Configured
} else {
HarnessAuthState::Unknown
};
reason = Some(format!("No-prompt runtime handshake failed: {message}"));
repair = Some(if auth == HarnessAuthState::Required {
format!(
"Run `{}` interactively once and complete sign-in, then probe again.",
launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
)
} else {
"Check the harness installation and run the handshake probe again."
.into()
});
}
Err(_) => {
reason = Some(
"No-prompt runtime handshake timed out after 30 seconds.".into(),
);
repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
}
}
let _ = isolated.cleanup();
tokio::time::sleep(Duration::from_millis(250)).await;
if let Err(error) = isolated.cleanup() {
auth = if configured {
HarnessAuthState::Configured
} else {
HarnessAuthState::Unknown
};
runtime = HarnessRuntimeState::Degraded;
reason = Some(format!(
"No-prompt runtime handshake could not remove its isolated harness home: {error}"
));
repair = Some(
"Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
.into(),
);
}
}
Err(error) => {
reason = Some(error_message(error));
}
}
} else if installed && configured {
reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
} else if installed {
reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
repair =
Some(format!(
"Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
));
}
let effective_capabilities = if installed {
descriptor.runtime.capabilities.clone()
} else {
unavailable_capabilities()
};
LocalHarness {
id: descriptor.id,
display_name: descriptor.display_name,
supported: true,
installed,
executable: executable.map(|path| path.to_string_lossy().into_owned()),
version,
auth,
runtime,
protocol: descriptor.runtime.protocol,
capabilities: descriptor.runtime.capabilities,
effective_capabilities,
sessions: HarnessSessionCounts { global, workspace },
reason,
repair,
}
}
fn session_counts(
&self,
workspace: Option<&Path>,
harnesses: &[HarnessId],
) -> BTreeMap<String, usize> {
let mut counts = BTreeMap::new();
for session in self
.catalog
.discover(&DiscoveryQuery {
workspace: workspace.map(Path::to_path_buf),
harnesses: harnesses.to_vec(),
..DiscoveryQuery::default()
})
.unwrap_or_default()
{
*counts
.entry(session.locator.harness.as_str().to_string())
.or_insert(0) += 1;
}
counts
}
}
#[async_trait::async_trait]
impl SdkService for HarnessSessionService {
fn capabilities(&self) -> SdkCapabilities {
SdkCapabilities::default()
}
async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
if request.operation == SdkOperation::Events {
let events = self
.poll_sdk_events()
.await
.into_iter()
.map(|(_, event)| event)
.collect::<Vec<_>>();
return serde_json::to_value(events).map_err(|error| {
SdkError::new(
SdkErrorCode::Execution,
request.operation,
error.to_string(),
)
});
}
let method = request
.operation
.method()
.ok_or_else(|| SdkError::unsupported(request.operation))?;
let result = match request.operation {
SdkOperation::Discover | SdkOperation::Load | SdkOperation::Export => {
self.call(method, request.params)
}
SdkOperation::Start
| SdkOperation::Resume
| SdkOperation::Input
| SdkOperation::Interrupt
| SdkOperation::Steer
| SdkOperation::Respond
| SdkOperation::Close => self.runtime_call(method, request.params).await,
SdkOperation::Events => unreachable!("handled before method dispatch"),
};
result.map_err(|error| sdk_error(request.operation, error))
}
async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
Ok(self
.poll_sdk_events()
.await
.into_iter()
.map(|(_, event)| event)
.collect())
}
}
#[cfg(feature = "adapter-api")]
struct HostedRuntimeLease {
connection: HostedHarnessConnection,
_host: std::sync::Arc<HostedHarnessRuntime>,
_registration: LiveRuntimeRegistration,
_server: crate::server::FrontendHttpServer,
}
#[async_trait::async_trait]
#[cfg(feature = "adapter-api")]
impl RuntimeConnection for HostedRuntimeLease {
fn handle(&self) -> &crate::RuntimeHandle {
self.connection.handle()
}
async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
self.connection.send_input(input).await
}
async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
self.connection.next_event().await
}
async fn interrupt(&mut self) -> crate::Result<()> {
self.connection.interrupt().await
}
async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
self.connection.respond(request_id, response).await
}
async fn close(&mut self) -> crate::Result<()> {
self.connection.close().await
}
}
async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
loop {
let now = tokio::time::Instant::now();
if now >= deadline {
return Ok(());
}
match tokio::time::timeout(deadline - now, connection.next_event()).await {
Err(_) => return Ok(()),
Ok(Ok(Some(event))) => {
if let Some(message) = handshake_event_failure(&event) {
return Err(truncate_text(&message, 500));
}
}
Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
Ok(Err(error)) => return Err(error.to_string()),
}
}
}
fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
let detail = event
.payload
.get("message")
.or_else(|| event.payload.get("line"))
.and_then(Value::as_str)
.unwrap_or(event.kind.as_str());
match event.kind.as_str() {
"transport_closed" => Some("runtime transport closed during startup".into()),
"transport_error" => Some(format!("runtime transport error: {detail}")),
"malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
_ => None,
}
}
fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
let total_messages = session.messages.len();
let (offset, end) = projected_message_window(total_messages, options);
json!({
"session": projected_session_json(session, options),
"summary": projected_session_summary(session, options),
"window": {
"has_more": offset > 0 || end < total_messages,
"has_newer": end < total_messages,
"has_older": offset > 0,
"newer_items": normalized_item_count(&session.messages[end..]),
"offset": offset,
"older_items": normalized_item_count(&session.messages[..offset]),
"returned": end.saturating_sub(offset),
"total_messages": total_messages,
}
})
}
fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
messages
.iter()
.map(|message| {
let conversation = usize::from(
matches!(message.role, Role::Assistant | Role::User)
&& message_has_content(message),
);
let tool_result =
usize::from(message.role == Role::Tool && message_has_content(message));
conversation + tool_result + message.tool_calls().len()
})
.sum()
}
fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
let mut conversational = session.messages.iter().filter(|message| {
matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
});
let first_message = conversational.clone().next();
let last_message = conversational.next_back();
let mut assistant = session
.messages
.iter()
.filter(|message| message.role == Role::Assistant && message_has_content(message));
let first_assistant_message = assistant.clone().next();
let last_assistant_message = assistant.next_back();
let end_of_turn = session
.messages
.iter()
.rev()
.find(|message| message.role != Role::System)
.is_some_and(|message| {
message.role == Role::Assistant
&& message_has_content(message)
&& message.tool_calls().is_empty()
});
let project = |message: Option<&crate::ChatMessage>| {
message.map(|message| project_inline_media(message_json(message), options))
};
json!({
"end_of_turn": end_of_turn,
"first_assistant_message": project(first_assistant_message),
"first_message": project(first_message),
"last_assistant_message": project(last_assistant_message),
"last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
"last_message": project(last_message),
})
}
fn message_has_content(message: &crate::ChatMessage) -> bool {
message
.content
.as_deref()
.is_some_and(|content| !content.trim().is_empty())
|| message
.content_parts
.as_ref()
.is_some_and(|parts| !parts.is_empty())
}
fn message_text(message: &crate::ChatMessage) -> String {
if let Some(content) = &message.content {
return content.clone();
}
message
.content_parts
.as_ref()
.into_iter()
.flatten()
.filter_map(|part| part.get("text").and_then(Value::as_str))
.collect::<Vec<_>>()
.join("\n")
}
fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
let (offset, end) = projected_message_window(session.messages.len(), options);
let messages = session.messages[offset..end]
.iter()
.map(|message| project_inline_media(message_json(message), options))
.collect::<Vec<_>>();
let subagents = if options.include_subagents.unwrap_or(true) {
let subagent_options = SessionLoadOptions {
message_limit: None,
message_offset: None,
message_tail: None,
..options.clone()
};
session
.subagents
.iter()
.map(|subagent| projected_session_json(subagent, &subagent_options))
.collect::<Vec<_>>()
} else {
Vec::new()
};
json!({
"source": match session.meta.source {
SessionSource::ClaudeCode => "claude_code",
SessionSource::Codex => "codex",
SessionSource::Gemini => "gemini",
SessionSource::Goose => "goose",
SessionSource::Grok => "grok",
SessionSource::Native => "native",
SessionSource::OpenCode => "opencode",
SessionSource::Pi => "pi",
},
"session_id": session.meta.session_id,
"model": session.meta.model,
"cwd": session.meta.cwd,
"system_prompt": session.meta.system_prompt,
"agent_id": session.meta.agent_id,
"parent_tool_use_id": session.meta.parent_tool_use_id,
"lineage": session.meta.lineage,
"messages": messages,
"subagents": subagents,
"raw_record_count": session.raw.len(),
"parse_error_lines": session.parse_error_lines,
})
}
fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
if let Some(tail) = options.message_tail {
return (total.saturating_sub(tail), total);
}
let offset = options.message_offset.unwrap_or(0).min(total);
let end = options
.message_limit
.map(|limit| offset.saturating_add(limit).min(total))
.unwrap_or(total);
(offset, end)
}
fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
return message;
};
for part in parts {
let Some(url) = part
.get("image_url")
.and_then(|image| image.get("url"))
.and_then(Value::as_str)
else {
continue;
};
let Some(rest) = url.strip_prefix("data:") else {
continue;
};
let Some((media_type, encoded)) = rest.split_once(";base64,") else {
continue;
};
let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
let decoded_bytes = encoded.len().saturating_mul(3) / 4;
let decoded_bytes = decoded_bytes.saturating_sub(padding);
let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
|| options
.max_inline_media_bytes
.is_some_and(|limit| decoded_bytes > limit);
if should_elide {
*part = json!({
"type": "media_reference",
"media_type": media_type,
"encoding": "base64",
"encoded_bytes": encoded.len(),
"decoded_bytes": decoded_bytes,
"omitted": true,
});
}
}
message
}
#[derive(Deserialize)]
struct LocatorParams {
locator: SessionLocator,
#[serde(default)]
fidelity: Option<Fidelity>,
#[serde(default)]
view: Option<SessionReadView>,
}
#[derive(Deserialize)]
struct SessionReadView {
#[serde(default)]
tail_messages: Option<usize>,
#[serde(default)]
include_subagents: bool,
#[serde(default)]
display_history: bool,
#[serde(default)]
max_message_chars: Option<usize>,
}
impl LocatorParams {
fn read_fidelity(&self) -> Fidelity {
self.fidelity.unwrap_or(Fidelity::Semantic)
}
fn include_subagents(&self) -> bool {
self.view
.as_ref()
.map(|view| view.include_subagents)
.unwrap_or(true)
}
fn tail_messages(&self) -> Option<usize> {
self.view
.as_ref()
.and_then(|view| view.tail_messages)
.map(|limit| limit.clamp(1, 5_000))
}
fn display_history(&self) -> bool {
self.view.as_ref().is_some_and(|view| view.display_history)
}
fn max_message_chars(&self) -> Option<usize> {
self.view
.as_ref()
.and_then(|view| view.max_message_chars)
.map(|limit| limit.clamp(256, 64_000))
}
fn bound_session(&self, session: &mut Session) {
bound_session_view(session, self.tail_messages(), self.max_message_chars());
}
}
#[derive(Debug, Clone, Copy, Default, Deserialize)]
#[serde(rename_all = "snake_case")]
enum InlineMediaMode {
#[default]
Full,
Metadata,
}
#[derive(Debug, Clone, Default, Deserialize)]
#[serde(default)]
struct SessionLoadOptions {
include_subagents: Option<bool>,
inline_media: InlineMediaMode,
max_inline_media_bytes: Option<usize>,
message_limit: Option<usize>,
message_offset: Option<usize>,
message_tail: Option<usize>,
}
impl SessionLoadOptions {
fn validate(&self) -> std::result::Result<(), ServiceError> {
if self.message_tail.is_some()
&& (self.message_limit.is_some() || self.message_offset.is_some())
{
return Err(ServiceError::InvalidParams(
"sessions.load options.message_tail cannot be combined with message_limit or message_offset"
.into(),
));
}
Ok(())
}
}
#[derive(Deserialize)]
struct LoadSessionParams {
#[serde(flatten)]
read: LocatorParams,
#[serde(default)]
options: Option<SessionLoadOptions>,
}
#[derive(Deserialize)]
struct UnfollowParams {
subscription: String,
}
#[derive(Deserialize)]
struct MessageSessionParams {
locator: SessionLocator,
text: String,
#[serde(default)]
homes: crate::HarnessHomes,
}
#[cfg(feature = "adapter-api")]
async fn message_live_session(
params: &MessageSessionParams,
runner: &dyn crate::claude_peer::CourierRunner,
) -> Value {
if params.locator.harness.as_str() != HarnessId::CLAUDE_CODE {
return json!({
"delivered_to_bus": false,
"refusal": {
"reason": crate::claude_peer::ClaudePeerRefusal::HarnessUnsupported.as_str(),
"message": format!(
"`{}` does not publish a live-session registry; only claude-code sessions can be messaged in place",
params.locator.harness.as_str()
),
},
});
}
match crate::claude_peer::message_claude_peer(
¶ms.homes,
¶ms.locator.session_id,
¶ms.text,
runner,
)
.await
{
Ok(delivery) => json!({
"delivered_to_bus": true,
"target": {
"session_id": delivery.target.session_id,
"name": delivery.target.name,
"pid": delivery.target.pid,
"cwd": delivery.target.cwd,
"status": delivery.target.status.map(|status| status.as_str()),
},
"courier": {
"model": crate::claude_peer::COURIER_MODEL,
"report": delivery.courier_report,
},
}),
Err(refusal) => json!({
"delivered_to_bus": false,
"refusal": {"reason": refusal.reason.as_str(), "message": refusal.message},
}),
}
}
#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
struct FollowedSource {
harness: String,
session_id: String,
reported: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
#[serde(rename_all = "kebab-case")]
enum TransferFormat {
ClaudeCode,
Codex,
#[serde(rename = "opencode", alias = "open-code")]
OpenCode,
Pi,
Grok,
Gemini,
Goose,
}
impl TransferFormat {
fn id(self) -> &'static str {
match self {
Self::ClaudeCode => HarnessId::CLAUDE_CODE,
Self::Codex => HarnessId::CODEX,
Self::OpenCode => HarnessId::OPENCODE,
Self::Pi => HarnessId::PI,
Self::Grok => HarnessId::GROK,
Self::Gemini => HarnessId::GEMINI,
Self::Goose => HarnessId::GOOSE,
}
}
}
impl From<TransferFormat> for SessionFormat {
fn from(value: TransferFormat) -> Self {
match value {
TransferFormat::ClaudeCode => Self::ClaudeCode,
TransferFormat::Codex => Self::Codex,
TransferFormat::OpenCode => Self::OpenCode,
TransferFormat::Pi => Self::Pi,
TransferFormat::Grok => Self::Grok,
TransferFormat::Gemini => Self::Gemini,
TransferFormat::Goose => Self::Goose,
}
}
}
#[derive(Deserialize)]
struct ImportSessionParams {
source_harness: TransferFormat,
content: String,
}
#[derive(Deserialize)]
struct ExportSessionParams {
locator: SessionLocator,
target_harness: TransferFormat,
}
#[derive(Deserialize)]
struct ReduceSessionParams {
locator: SessionLocator,
target_harness: TransferFormat,
#[serde(default = "default_keep_last")]
keep_last: usize,
}
fn default_keep_last() -> usize {
6
}
#[derive(Deserialize)]
struct BranchSessionParams {
locator: SessionLocator,
#[serde(default)]
target_harness: Option<TransferFormat>,
}
#[derive(Deserialize)]
struct HandoffSessionParams {
locator: SessionLocator,
target_harness: TransferFormat,
#[serde(default)]
cwd: Option<PathBuf>,
}
#[derive(Debug, Clone, Copy, Default, Deserialize)]
#[serde(rename_all = "snake_case")]
enum ResumePolicy {
#[default]
Default,
Yolo,
}
#[derive(Deserialize)]
struct ResumeInstructionsParams {
locator: SessionLocator,
#[serde(default)]
cwd: Option<PathBuf>,
#[serde(default)]
policy: ResumePolicy,
}
#[derive(Serialize)]
struct SessionArtifact {
source_harness: HarnessId,
target_harness: &'static str,
session_id: Option<String>,
content: String,
suggested_filename: String,
files: Vec<SessionArtifactFile>,
fidelity: Fidelity,
residue: Vec<String>,
}
#[derive(Serialize)]
struct SessionArtifactFile {
path: String,
content: String,
role: ArtifactFileRole,
}
#[derive(Serialize)]
#[serde(rename_all = "snake_case")]
enum ArtifactFileRole {
Primary,
Subagent,
Bundle,
SourceRecovery,
}
#[derive(Serialize)]
struct StructuredLaunch {
cwd: PathBuf,
program: String,
arguments: Vec<String>,
env: BTreeMap<String, String>,
}
struct HandoffInstructions {
launch: StructuredLaunch,
materialize: Option<StructuredLaunch>,
requires_materialization: bool,
note: String,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
enum HarnessProbeLevel {
#[default]
Passive,
Handshake,
}
#[derive(Default, Deserialize)]
#[serde(default)]
struct HarnessInventoryParams {
harness: Option<HarnessId>,
harnesses: Vec<HarnessId>,
workspace: Option<PathBuf>,
probe: HarnessProbeLevel,
include_sessions: bool,
skip_versions: bool,
}
#[derive(Serialize)]
struct HarnessInventoryReport {
probe: HarnessProbeLevel,
workspace: Option<PathBuf>,
harnesses: Vec<LocalHarness>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
enum HarnessAuthState {
Ready,
Configured,
Required,
Unknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
enum HarnessRuntimeState {
Ready,
Degraded,
Unavailable,
}
#[derive(Serialize)]
struct HarnessSessionCounts {
global: Option<usize>,
workspace: Option<usize>,
}
#[derive(Serialize)]
struct LocalHarness {
id: HarnessId,
display_name: String,
supported: bool,
installed: bool,
executable: Option<String>,
version: Option<String>,
auth: HarnessAuthState,
runtime: HarnessRuntimeState,
protocol: String,
capabilities: crate::RuntimeCapabilities,
effective_capabilities: crate::RuntimeCapabilities,
sessions: HarnessSessionCounts,
reason: Option<String>,
repair: Option<String>,
}
#[derive(Clone, Deserialize)]
struct RuntimeBackendParams {
harness: HarnessId,
#[serde(default)]
protocol: Option<String>,
#[serde(default)]
launch: Option<RuntimeLaunch>,
#[serde(default)]
base_url: Option<String>,
#[serde(default)]
policy: RuntimePolicy,
}
#[derive(Debug, Clone, Copy, Default, Deserialize)]
#[serde(rename_all = "snake_case")]
enum RuntimePolicy {
#[default]
Default,
Yolo,
}
#[derive(Deserialize)]
struct RuntimeStartParams {
#[serde(flatten)]
backend: RuntimeBackendParams,
cwd: PathBuf,
}
#[derive(Deserialize)]
struct RuntimeAttachParams {
#[serde(flatten)]
backend: RuntimeBackendParams,
runtime_id: String,
#[serde(default)]
cwd: Option<PathBuf>,
}
#[derive(Deserialize)]
struct RuntimeConnectionParams {
connection: String,
}
#[derive(Deserialize)]
struct RuntimeInputParams {
connection: String,
text: String,
}
#[derive(Deserialize)]
struct RuntimeRespondParams {
connection: String,
request_id: Value,
response: Value,
}
fn default_reduction_store_root() -> PathBuf {
if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
return PathBuf::from(root).join("sessions");
}
if let Some(home) = std::env::var_os("HOME") {
return PathBuf::from(home).join(".supercode").join("sessions");
}
PathBuf::from(".supercode").join("sessions")
}
fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
let mut output = String::new();
for message in messages {
output.push_str(
&serde_json::to_string(message)
.map_err(|error| ServiceError::Operation(error.to_string()))?,
);
output.push('\n');
}
Ok(output)
}
fn parse_messages_jsonl(
content: &str,
) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
content
.lines()
.enumerate()
.filter(|(_, line)| !line.trim().is_empty())
.map(|(index, line)| {
serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
ServiceError::Operation(format!(
"reduced transcript line {} is invalid: {error}",
index + 1
))
})
})
.collect()
}
fn reduced_bootstrap_prompt(
source: &SessionLocator,
target: TransferFormat,
view_jsonl: &str,
sidecar_path: &Path,
reduction_log_path: &Path,
) -> String {
format!(
"Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
\n\
The bounded working transcript is below. Treat reduction markers as transparent placeholders, not missing work. If a detail behind a marker is needed, use ordinary file-reading/search tools against the full Supercode sidecar at `{sidecar}` and its reduction index at `{log}`. Do not guess hidden content. Both files were reloaded and verified before this continuation was issued.\n\
\n\
<supercode-reduced-session source-session=\"{source_id}\">\n\
{view_jsonl}\
</supercode-reduced-session>\n\
\n\
Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
source_harness = source.harness.as_str(),
target_harness = target.id(),
sidecar = sidecar_path.display(),
log = reduction_log_path.display(),
source_id = source.session_id,
)
}
fn session_artifact(
locator: &SessionLocator,
session: &Session,
target: TransferFormat,
) -> std::result::Result<SessionArtifact, ServiceError> {
session_artifact_with_id(locator, session, target, None)
}
fn session_artifact_with_id(
locator: &SessionLocator,
session: &Session,
target: TransferFormat,
target_session_id: Option<&str>,
) -> std::result::Result<SessionArtifact, ServiceError> {
let format: SessionFormat = target.into();
let diagonal = format.source() == session.meta.source;
let has_appended_turns = session
.imported_message_count
.is_some_and(|imported| imported < session.messages.len());
let content = if let Some(id) = target_session_id {
if diagonal && format != SessionFormat::OpenCode {
session
.to_jsonl_spliced(format, Some(id))
.map_err(operation)?
} else {
let mut rewritten = session.clone();
rewritten.meta.session_id = Some(id.to_string());
rewritten.to_jsonl(format).map_err(operation)?
}
} else if diagonal && session.raw_is_verbatim && !has_appended_turns {
session.raw_verbatim()
} else if diagonal {
session.to_jsonl_spliced(format, None).map_err(operation)?
} else {
session.to_jsonl(format).map_err(operation)?
};
let stem = sanitize_filename(
target_session_id
.or(session.meta.session_id.as_deref())
.unwrap_or(&locator.session_id),
);
let suggested_filename = if diagonal && target == TransferFormat::Grok {
"chat_history.jsonl".to_string()
} else if target == TransferFormat::Goose {
format!("{stem}.goose.json")
} else {
format!("{stem}.{}.jsonl", target.id())
};
let mut files = vec![SessionArtifactFile {
path: suggested_filename.clone(),
content: content.clone(),
role: ArtifactFileRole::Primary,
}];
if target == TransferFormat::ClaudeCode {
let bundle_stem = Path::new(&suggested_filename)
.file_stem()
.and_then(|stem| stem.to_str())
.unwrap_or(&stem);
let mut child_paths = BTreeSet::new();
for (index, subagent) in session.subagents.iter().enumerate() {
let agent_id = subagent
.meta
.agent_id
.as_deref()
.map(|id| id.strip_prefix("agent-").unwrap_or(id))
.map(sanitize_filename)
.filter(|id| !id.is_empty())
.unwrap_or_else(|| format!("subagent-{}", index + 1));
let child_has_appended_turns = subagent
.imported_message_count
.is_some_and(|imported| imported < subagent.messages.len());
let child_content = if target_session_id.is_none()
&& subagent.meta.source == SessionSource::ClaudeCode
&& subagent.raw_is_verbatim
&& !child_has_appended_turns
{
subagent.raw_verbatim()
} else if subagent.meta.source == SessionSource::ClaudeCode {
subagent
.to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
.map_err(operation)?
} else {
let mut child = subagent.clone();
if let Some(id) = target_session_id {
child.meta.session_id = Some(id.to_string());
}
child
.to_jsonl(SessionFormat::ClaudeCode)
.map_err(operation)?
};
let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
if !child_paths.insert(path.clone()) {
return Err(ServiceError::Operation(format!(
"Claude subagent ids collide at artifact path `{path}`"
)));
}
files.push(SessionArtifactFile {
path,
content: child_content,
role: ArtifactFileRole::Subagent,
});
}
}
if diagonal && target == TransferFormat::Grok {
append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
}
if !diagonal || !session.raw_is_verbatim {
files.push(SessionArtifactFile {
path: "recovery/source.supercode.jsonl".into(),
content: session.to_native_jsonl(),
role: ArtifactFileRole::SourceRecovery,
});
for (index, subagent) in session.subagents.iter().enumerate() {
let id = subagent
.meta
.agent_id
.as_deref()
.map(sanitize_filename)
.unwrap_or_else(|| format!("subagent-{}", index + 1));
files.push(SessionArtifactFile {
path: format!("recovery/subagents/{id}.supercode.jsonl"),
content: subagent.to_native_jsonl(),
role: ArtifactFileRole::SourceRecovery,
});
}
}
if !diagonal && session.meta.source == SessionSource::Grok {
append_grok_bundle_files(
locator,
"recovery/grok/",
ArtifactFileRole::SourceRecovery,
&mut files,
)?;
}
let (fidelity, residue) = if diagonal
&& target_session_id.is_none()
&& session.raw_is_verbatim
&& !has_appended_turns
{
(Fidelity::ByteLossless, Vec::new())
} else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
(
Fidelity::ValueLossless,
vec![if target_session_id.is_some() {
"target identity was rewritten, so the artifact intentionally differs from source bytes".into()
} else {
"source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
}],
)
} else {
(
Fidelity::Semantic,
vec!["target schema has no portable slot for every source-native record and metadata field".into()],
)
};
Ok(SessionArtifact {
source_harness: locator.harness.clone(),
target_harness: target.id(),
session_id: target_session_id
.map(str::to_string)
.or_else(|| session.meta.session_id.clone()),
content,
suggested_filename,
files,
fidelity,
residue,
})
}
fn append_grok_bundle_files(
locator: &SessionLocator,
prefix: &str,
role: ArtifactFileRole,
files: &mut Vec<SessionArtifactFile>,
) -> std::result::Result<(), ServiceError> {
let primary = locator.storage.path();
if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
return Err(ServiceError::Operation(format!(
"Grok bundle locator must name chat_history.jsonl, got {}",
primary.display()
)));
}
let parent = primary.parent().ok_or_else(|| {
ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
})?;
for name in ["summary.json", "updates.jsonl"] {
let path = parent.join(name);
let metadata = match std::fs::symlink_metadata(&path) {
Ok(metadata) => metadata,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
Err(error) => return Err(ServiceError::Operation(error.to_string())),
};
if metadata.file_type().is_symlink() || !metadata.is_file() {
return Err(ServiceError::Operation(format!(
"refusing non-regular Grok bundle member {}",
path.display()
)));
}
let content = std::fs::read_to_string(&path).map_err(|error| {
ServiceError::Operation(format!(
"Grok bundle member {} is not representable as UTF-8: {error}",
path.display()
))
})?;
files.push(SessionArtifactFile {
path: format!("{prefix}{name}"),
content,
role: match role {
ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
_ => ArtifactFileRole::SourceRecovery,
},
});
}
Ok(())
}
fn handoff_artifact(
locator: &SessionLocator,
session: &Session,
target: TransferFormat,
cwd: &Path,
) -> std::result::Result<SessionArtifact, ServiceError> {
if target != TransferFormat::Grok {
let target_session_id = target_session_id(target);
return session_artifact_with_id(locator, session, target, Some(&target_session_id));
}
let mut importable = session.clone();
importable.meta.session_id = Some(target_session_id(TransferFormat::ClaudeCode));
importable.meta.cwd = Some(if cwd.is_absolute() {
cwd.to_path_buf()
} else {
std::env::current_dir()
.map_err(|error| ServiceError::Operation(error.to_string()))?
.join(cwd)
});
let content = importable
.to_jsonl(SessionFormat::ClaudeCode)
.map_err(operation)?;
let stem = sanitize_filename(
importable
.meta
.session_id
.as_deref()
.unwrap_or(&locator.session_id),
);
let suggested_filename = format!("{stem}.grok-import.claude-code.jsonl");
Ok(SessionArtifact {
source_harness: locator.harness.clone(),
target_harness: TransferFormat::ClaudeCode.id(),
session_id: importable.meta.session_id.clone(),
content: content.clone(),
suggested_filename: suggested_filename.clone(),
files: vec![SessionArtifactFile {
path: suggested_filename,
content,
role: ArtifactFileRole::Primary,
}],
fidelity: Fidelity::Semantic,
residue: vec!["Grok's stock importer accepts a Claude Code transcript, not a complete Grok updates/session bundle".into()],
})
}
fn target_session_id(target: TransferFormat) -> String {
let uuid = generated_session_id();
match target {
TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
TransferFormat::ClaudeCode
| TransferFormat::Codex
| TransferFormat::Pi
| TransferFormat::Grok
| TransferFormat::Gemini
| TransferFormat::Goose => uuid,
}
}
fn sanitize_filename(value: &str) -> String {
let value = value
.chars()
.map(|character| {
if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
character
} else {
'-'
}
})
.collect::<String>();
let value = value.trim_matches('-');
if value.is_empty() {
"session".into()
} else {
value.chars().take(100).collect()
}
}
fn handoff_instructions(
target: TransferFormat,
session_id: &str,
cwd: &Path,
) -> HandoffInstructions {
let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
cwd: cwd.to_path_buf(),
program: program.into(),
arguments,
env: BTreeMap::new(),
};
match target {
TransferFormat::ClaudeCode => HandoffInstructions {
launch: launch("claude", vec!["--resume".into(), session_id.into()]),
materialize: None,
requires_materialization: true,
note: "Write the artifact into Claude Code's native project session store before running the resume launch; Claude Code has no general transcript-import command.".into(),
},
TransferFormat::Codex => HandoffInstructions {
launch: launch("codex", vec!["resume".into(), session_id.into()]),
materialize: None,
requires_materialization: true,
note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
},
TransferFormat::OpenCode => HandoffInstructions {
launch: launch("opencode", vec!["--session".into(), session_id.into()]),
materialize: Some(launch(
"opencode",
vec!["import".into(), "{artifact_path}".into()],
)),
requires_materialization: true,
note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
},
TransferFormat::Pi => HandoffInstructions {
launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
materialize: None,
requires_materialization: true,
note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
},
TransferFormat::Grok => HandoffInstructions {
launch: launch(
"grok",
vec![
"--resume".into(),
"{imported_session_id}".into(),
"--fork-session".into(),
],
),
materialize: Some(launch(
"grok",
vec!["import".into(), "--json".into(), "{artifact_path}".into()],
)),
requires_materialization: true,
note: "The artifact is Claude Code JSONL for Grok's official importer. Write it to a file, run the materialize command, read sessionId from its NDJSON outcome=imported record, replace {imported_session_id} in the launch arguments, then launch a writable fork of the imported session.".into(),
},
TransferFormat::Gemini => HandoffInstructions {
launch: launch(
"gemini",
vec!["--session-file".into(), "{artifact_path}".into()],
),
materialize: None,
requires_materialization: true,
note: "Write the Gemini JSONL artifact to a file and replace {artifact_path}; Gemini imports it into the current project's chat store before opening the continuation.".into(),
},
TransferFormat::Goose => HandoffInstructions {
launch: launch(
"goose",
vec![
"session".into(),
"--resume".into(),
"--session-id".into(),
"{imported_session_id}".into(),
],
),
materialize: Some(launch(
"goose",
vec!["session".into(), "import".into(), "{artifact_path}".into()],
)),
requires_materialization: true,
note: "Write the Goose JSON artifact to a file, run the materialize command, read the imported session id from its output, replace {imported_session_id}, then resume that native Goose session.".into(),
},
}
}
fn resume_launch(
harness: &str,
session_id: &str,
cwd: &Path,
policy: ResumePolicy,
) -> std::result::Result<StructuredLaunch, ServiceError> {
let mut arguments = Vec::new();
let program = match harness {
HarnessId::GROK => {
if matches!(policy, ResumePolicy::Yolo) {
arguments.extend([
"--sandbox".into(),
"workspace".into(),
"--always-approve".into(),
]);
}
arguments.extend(["--resume".into(), session_id.into()]);
"grok"
}
HarnessId::CODEX => {
if matches!(policy, ResumePolicy::Yolo) {
arguments.extend([
"--dangerously-bypass-approvals-and-sandbox".into(),
"--dangerously-bypass-hook-trust".into(),
]);
}
arguments.extend(["resume".into(), session_id.into()]);
"codex"
}
HarnessId::CLAUDE_CODE => {
if matches!(policy, ResumePolicy::Yolo) {
arguments.push("--dangerously-skip-permissions".into());
}
arguments.extend(["--resume".into(), session_id.into()]);
"claude"
}
HarnessId::GEMINI => {
if matches!(policy, ResumePolicy::Yolo) {
arguments.push("--yolo".into());
}
arguments.extend(["--resume".into(), session_id.into()]);
"gemini"
}
HarnessId::GOOSE => {
arguments.extend([
"session".into(),
"--resume".into(),
"--session-id".into(),
session_id.into(),
]);
"goose"
}
HarnessId::PI => {
if matches!(policy, ResumePolicy::Yolo) {
arguments.push("--approve".into());
}
arguments.extend(["--session".into(), session_id.into()]);
"pi"
}
HarnessId::OPENCODE => {
arguments.extend(["--session".into(), session_id.into()]);
"opencode"
}
HarnessId::SUPERCODE => {
if matches!(policy, ResumePolicy::Yolo) {
arguments.push("--dangerous".into());
}
arguments.extend(["resume".into(), session_id.into()]);
"supercode"
}
other => {
return Err(ServiceError::InvalidParams(format!(
"no structured resume launch is registered for harness `{other}`"
)))
}
};
Ok(StructuredLaunch {
cwd: cwd.to_path_buf(),
program: program.into(),
arguments,
env: BTreeMap::new(),
})
}
fn runtime_backend(
params: &RuntimeBackendParams,
) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
if params.protocol.as_deref() == Some("acp") {
let launch = params
.launch
.clone()
.or_else(|| {
harness_support_registry()
.harnesses
.into_iter()
.find(|harness| harness.id == params.harness)
.filter(|harness| {
harness.runtime.implementation == ImplementationKind::GenericProtocol
&& harness.runtime.protocol.starts_with("acp")
})
.and_then(|harness| harness.runtime.default_launch)
})
.ok_or_else(|| {
ServiceError::InvalidParams(
"an ACP runtime requires `launch` unless the harness has a registered default"
.into(),
)
})?;
let resume_session = harness_support_registry()
.harnesses
.into_iter()
.find(|harness| harness.id == params.harness)
.is_some_and(|harness| harness.runtime.capabilities.resume_session);
return Ok(Box::new(
AcpRuntimeBackend::new(params.harness.clone(), launch)
.with_resume_support(resume_session),
));
}
let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
HarnessId::PI => Box::new(PiRuntimeBackend::new()),
HarnessId::OPENCODE => match ¶ms.base_url {
Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
None => Box::new(OpenCodeRuntimeBackend::new()),
},
harness => {
let descriptor = harness_support_registry()
.harnesses
.into_iter()
.find(|descriptor| descriptor.id.as_str() == harness)
.filter(|descriptor| {
descriptor.runtime.implementation == ImplementationKind::GenericProtocol
&& descriptor.runtime.protocol.starts_with("acp")
});
let Some(descriptor) = descriptor else {
return Err(ServiceError::InvalidParams(format!(
"no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
)));
};
let resume = descriptor.runtime.capabilities.resume_session;
Box::new(
AcpRuntimeBackend::new(
descriptor.id,
descriptor
.runtime
.default_launch
.expect("generic ACP registry entry includes its launch"),
)
.with_resume_support(resume),
)
}
};
Ok(backend)
}
fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
if let Some(launch) = ¶ms.launch {
return Some(launch.clone());
}
if !matches!(params.policy, RuntimePolicy::Yolo) {
return None;
}
let launch = match params.harness.as_str() {
HarnessId::GROK => RuntimeLaunch {
program: "grok".into(),
arguments: vec![
"--sandbox".into(),
"workspace".into(),
"--always-approve".into(),
"agent".into(),
"--no-leader".into(),
"stdio".into(),
],
env: BTreeMap::from([("GROK_AGENT_DASHBOARD".into(), "0".into())]),
},
HarnessId::CODEX => RuntimeLaunch {
program: "codex".into(),
arguments: vec![
"--dangerously-bypass-approvals-and-sandbox".into(),
"--dangerously-bypass-hook-trust".into(),
"app-server".into(),
],
env: BTreeMap::new(),
},
HarnessId::CLAUDE_CODE => RuntimeLaunch {
program: "claude".into(),
arguments: vec![
"--dangerously-skip-permissions".into(),
"--print".into(),
"--input-format".into(),
"stream-json".into(),
"--output-format".into(),
"stream-json".into(),
"--verbose".into(),
],
env: BTreeMap::new(),
},
HarnessId::PI => RuntimeLaunch {
program: "pi".into(),
arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
env: BTreeMap::new(),
},
HarnessId::OPENCODE => RuntimeLaunch {
program: "opencode".into(),
arguments: vec!["serve".into()],
env: BTreeMap::new(),
},
HarnessId::GEMINI => RuntimeLaunch {
program: "gemini".into(),
arguments: vec!["--acp".into(), "--yolo".into()],
env: BTreeMap::new(),
},
HarnessId::GOOSE => RuntimeLaunch {
program: "goose".into(),
arguments: vec!["acp".into()],
env: BTreeMap::new(),
},
HarnessId::SUPERCODE => RuntimeLaunch {
program: "supercode".into(),
arguments: vec!["acp".into(), "--dangerous".into()],
env: BTreeMap::new(),
},
_ => return None,
};
Some(launch)
}
struct IsolatedProbeHome {
launch: RuntimeLaunch,
root: PathBuf,
}
impl IsolatedProbeHome {
fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
let root = std::env::temp_dir().join(format!(
"supercode-harness-probe-{harness}-{}",
generated_session_id()
));
std::fs::create_dir_all(&root)?;
set_private_dir_permissions(&root)?;
if let Some(source_home) = std::env::var_os("HOME").map(PathBuf::from) {
for relative in probe_auth_files(harness) {
copy_probe_file(&source_home, &root, relative)?;
}
}
configure_isolated_probe_auth(harness, &root)?;
let root_text = root.to_string_lossy().into_owned();
for (key, value) in [
("HOME", root_text.clone()),
(
"XDG_CACHE_HOME",
root.join(".cache").to_string_lossy().into_owned(),
),
(
"XDG_CONFIG_HOME",
root.join(".config").to_string_lossy().into_owned(),
),
(
"XDG_DATA_HOME",
root.join(".local/share").to_string_lossy().into_owned(),
),
] {
launch.env.insert(key.into(), value);
}
let scoped = match harness {
HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
_ => None,
};
if let Some((key, value)) = scoped {
launch
.env
.insert(key.into(), value.to_string_lossy().into_owned());
}
Ok(Self { launch, root })
}
fn cleanup(&self) -> std::io::Result<()> {
match std::fs::remove_dir_all(&self.root) {
Ok(()) => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(error),
}
}
}
impl Drop for IsolatedProbeHome {
fn drop(&mut self) {
let _ = self.cleanup();
}
}
fn probe_auth_files(harness: &str) -> &'static [&'static str] {
match harness {
HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
HarnessId::CODEX => &[".codex/auth.json"],
HarnessId::GEMINI => &[
".gemini/google_accounts.json",
".gemini/oauth_creds.json",
".gemini/settings.json",
],
HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
HarnessId::OPENCODE => &[
".config/opencode/auth.json",
".local/share/opencode/auth.json",
],
HarnessId::PI => &[".pi/agent/auth.json"],
HarnessId::SUPERCODE => &[
".config/supercode/config.toml",
".config/supercode/credentials.toml",
],
_ => &[],
}
}
fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
let source = source_home.join(relative);
if !source.is_file() {
return Ok(());
}
let destination = probe_home.join(relative);
if let Some(parent) = destination.parent() {
std::fs::create_dir_all(parent)?;
set_private_dir_permissions(parent)?;
}
std::fs::copy(source, &destination)?;
set_private_file_permissions(&destination)
}
fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
if harness != HarnessId::GEMINI {
return Ok(());
}
let oauth = probe_home.join(".gemini/oauth_creds.json");
if !oauth.is_file() {
return Ok(());
}
let settings_path = probe_home.join(".gemini/settings.json");
let mut settings = std::fs::read_to_string(&settings_path)
.ok()
.and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
.unwrap_or_else(|| json!({}));
settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
std::fs::write(
&settings_path,
serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
)?;
set_private_file_permissions(&settings_path)
}
#[cfg(unix)]
fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
}
#[cfg(not(unix))]
fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
Ok(())
}
#[cfg(unix)]
fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
}
#[cfg(not(unix))]
fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
Ok(())
}
fn find_executable(program: &str) -> Option<PathBuf> {
let candidate = PathBuf::from(program);
if candidate.components().count() > 1 {
return candidate.is_file().then_some(candidate);
}
let path = std::env::var_os("PATH")?;
for directory in std::env::split_paths(&path) {
let candidate = directory.join(program);
if candidate.is_file() {
return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
}
#[cfg(windows)]
{
for extension in ["exe", "cmd", "bat"] {
let candidate = directory.join(format!("{program}.{extension}"));
if candidate.is_file() {
return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
}
}
}
}
None
}
async fn executable_version(executable: &Path) -> Option<String> {
let mut command = tokio::process::Command::new(executable);
command
.arg("--version")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
let output = tokio::time::timeout(Duration::from_secs(3), command.output())
.await
.ok()?
.ok()?;
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
stdout
.lines()
.chain(stderr.lines())
.map(str::trim)
.find(|line| !line.is_empty())
.map(|line| truncate_text(line, 200))
}
fn auth_evidence(harness: &str) -> bool {
let env_names: &[&str] = match harness {
HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
HarnessId::CODEX => &["OPENAI_API_KEY"],
HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
_ => &[],
};
if env_names
.iter()
.any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
{
return true;
}
let Some(home) = std::env::var_os("HOME").map(PathBuf::from) else {
return false;
};
let files: Vec<PathBuf> = match harness {
HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
HarnessId::CODEX => vec![home.join(".codex/auth.json")],
HarnessId::OPENCODE => vec![
home.join(".local/share/opencode/auth.json"),
home.join(".config/opencode/auth.json"),
],
HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
HarnessId::GROK => vec![home.join(".grok/auth.json")],
HarnessId::GEMINI => vec![
home.join(".gemini/oauth_creds.json"),
home.join(".gemini/google_accounts.json"),
],
HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
_ => Vec::new(),
};
if files.into_iter().any(|path| {
std::fs::metadata(path)
.map(|metadata| metadata.is_file() && metadata.len() > 2)
.unwrap_or(false)
}) {
return true;
}
if harness == HarnessId::CLAUDE_CODE {
return std::fs::read_to_string(home.join(".claude.json"))
.map(|text| text.contains("\"oauthAccount\""))
.unwrap_or(false);
}
false
}
fn looks_like_auth_error(message: &str) -> bool {
let message = message.to_ascii_lowercase();
[
"auth",
"login",
"sign in",
"sign-in",
"credential",
"unauthorized",
"forbidden",
"token",
]
.iter()
.any(|needle| message.contains(needle))
}
fn unavailable_capabilities() -> crate::RuntimeCapabilities {
crate::RuntimeCapabilities {
start_session: false,
resume_session: false,
attach_existing_process: false,
send_input: false,
stream_events: false,
interrupt: false,
respond_to_requests: false,
}
}
fn truncate_text(text: &str, max_chars: usize) -> String {
let mut chars = text.chars();
let truncated = chars.by_ref().take(max_chars).collect::<String>();
if chars.next().is_some() {
format!("{truncated}…")
} else {
truncated
}
}
fn error_message(error: ServiceError) -> String {
match error {
ServiceError::InvalidParams(message)
| ServiceError::Operation(message)
| ServiceError::UnsupportedAction(message) => message,
ServiceError::MethodNotFound => "runtime adapter is not available".into(),
ServiceError::Sdk(error) => error.to_string(),
}
}
#[derive(Debug)]
enum ServiceError {
InvalidParams(String),
MethodNotFound,
UnsupportedAction(String),
Operation(String),
Sdk(SdkError),
}
fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
match error {
ServiceError::InvalidParams(message) => {
SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
}
ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
SdkError::unsupported(operation)
}
ServiceError::Operation(message) => {
let code = if message.contains("already in progress") {
SdkErrorCode::Busy
} else if message.contains("not supported by this runtime") {
SdkErrorCode::UnsupportedAction
} else if message.contains("unknown runtime connection") {
SdkErrorCode::NotFound
} else {
SdkErrorCode::Execution
};
SdkError::new(code, operation, message)
}
ServiceError::Sdk(error) => error,
}
}
fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
let error_code = error.code();
let code = match error_code {
SdkErrorCode::Unauthenticated => -32030,
SdkErrorCode::Unauthorized => -32031,
SdkErrorCode::ControllerRequired => -32032,
SdkErrorCode::LeaseExpired => -32033,
SdkErrorCode::InvalidArgument => -32602,
SdkErrorCode::NotFound => -32004,
SdkErrorCode::Busy => -32000,
SdkErrorCode::UnsupportedAction => -32020,
SdkErrorCode::Execution => -32002,
SdkErrorCode::Transport => -32003,
};
json!({
"jsonrpc": "2.0",
"id": id,
"error": {
"code": code,
"name": error_code,
"operation": error.operation(),
"message": error.to_string(),
},
})
}
fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
}
fn operation(error: impl Into<crate::Error>) -> ServiceError {
let error = error.into();
match error {
crate::Error::Sdk(error) => ServiceError::Sdk(error),
error => ServiceError::Operation(error.to_string()),
}
}
fn rpc_error(id: Value, code: i64, message: &str) -> Value {
json!({
"jsonrpc": "2.0",
"id": id,
"error": {"code": code, "message": message},
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{HarnessEvent, HarnessId, RuntimeEndpoint, RuntimeHandle, StorageLocator};
use async_trait::async_trait;
use std::io::Write;
use std::path::PathBuf;
use std::time::Instant;
struct EndingRuntime {
handle: RuntimeHandle,
event: Option<HarnessEvent>,
}
#[async_trait]
impl RuntimeConnection for EndingRuntime {
fn handle(&self) -> &RuntimeHandle {
&self.handle
}
async fn send_input(&mut self, _input: RuntimeInput) -> crate::Result<Option<String>> {
unreachable!("ending runtime does not accept input")
}
async fn next_event(&mut self) -> crate::Result<Option<HarnessEvent>> {
Ok(self.event.take())
}
async fn interrupt(&mut self) -> crate::Result<()> {
Ok(())
}
async fn respond(&mut self, _request_id: Value, _response: Value) -> crate::Result<()> {
Ok(())
}
async fn close(&mut self) -> crate::Result<()> {
Ok(())
}
}
fn ending_runtime(event: Option<HarnessEvent>) -> Box<dyn RuntimeConnection> {
Box::new(EndingRuntime {
handle: RuntimeHandle {
harness: HarnessId::from(HarnessId::CLAUDE_CODE),
runtime_id: "ending-session".into(),
endpoint: RuntimeEndpoint::LocalProcess {
pid: None,
command: vec!["ending-runtime".into()],
protocol: "test".into(),
},
},
event,
})
}
fn request(id: u64, method: &str, params: Value) -> Value {
json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})
}
fn pi_locator() -> SessionLocator {
SessionLocator {
harness: HarnessId::from(HarnessId::PI),
session_id: "1e6f2a3b-0000-4000-8000-000000000001".into(),
storage: StorageLocator::File {
path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/pi_session.jsonl"),
},
}
}
fn opencode_locator() -> SessionLocator {
let session_id = "ses_fixtureAAAAAAAAAAAAAAA1";
SessionLocator {
harness: HarnessId::from(HarnessId::OPENCODE),
session_id: session_id.into(),
storage: StorageLocator::Sqlite {
path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/opencode_fixture/opencode.db"),
selector: session_id.into(),
},
}
}
fn grok_locator() -> SessionLocator {
SessionLocator {
harness: HarnessId::from(HarnessId::GROK),
session_id: "73c09283-4b33-41fa-90f1-0bcb0f7be523".into(),
storage: StorageLocator::File {
path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/grok_session/chat_history.jsonl"),
},
}
}
#[test]
fn capabilities_are_explicit_and_versioned() {
let mut service = HarnessSessionService::new();
let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
assert_eq!(response["result"]["version"], HARNESS_SERVICE_VERSION);
assert_eq!(
response["result"]["sdk"]["schema_version"],
crate::SDK_SCHEMA_VERSION
);
assert_eq!(
response["result"]["sdk"]["operations"]
.as_array()
.unwrap()
.len(),
SdkOperation::ALL.len()
);
assert_eq!(response["result"]["harnesses"].as_array().unwrap().len(), 8);
assert!(response["result"]["harnesses"]
.as_array()
.unwrap()
.iter()
.any(|harness| harness == HarnessId::GROK));
assert!(response["result"]["harnesses"]
.as_array()
.unwrap()
.iter()
.any(|harness| harness == HarnessId::GOOSE));
}
#[test]
fn handshake_health_uses_protocol_liveness_not_stderr_severity() {
let noisy_stderr = crate::HarnessEvent {
sequence: None,
kind: "transport_stderr".into(),
payload: json!({"line": "ERROR optional worker AuthorizationRequired"}),
};
assert_eq!(handshake_event_failure(&noisy_stderr), None);
let closed = crate::HarnessEvent {
sequence: None,
kind: "transport_closed".into(),
payload: json!({}),
};
assert!(handshake_event_failure(&closed).is_some());
}
#[tokio::test]
async fn runtime_eof_is_notified_and_removed_for_raw_and_explicit_close() {
let mut service = HarnessSessionService::new();
service
.runtimes
.insert("raw-eof".into(), ending_runtime(None));
service.runtimes.insert(
"explicit-close".into(),
ending_runtime(Some(HarnessEvent {
sequence: None,
kind: "transport_closed".into(),
payload: json!({"message": "native transport exited"}),
})),
);
let notifications = service.poll_runtimes().await;
assert_eq!(notifications.len(), 2);
assert!(notifications
.iter()
.all(|notification| { notification["params"]["event"]["kind"] == "transport_closed" }));
assert!(notifications.iter().all(|notification| {
notification["params"]["session_id"] == "ending-session"
&& notification["params"]["connection"].is_string()
}));
let mut sequences = notifications
.iter()
.filter_map(|notification| notification["params"]["sequence"].as_u64())
.collect::<Vec<_>>();
sequences.sort_unstable();
assert_eq!(sequences, vec![1, 2]);
assert!(service.runtimes.is_empty());
}
#[tokio::test]
async fn sdk_facade_returns_named_unsupported_actions() {
let mut service = HarnessSessionService::new();
let error = service
.execute(SdkRequest {
operation: SdkOperation::Steer,
params: json!({"connection": "runtime-1", "text": "go left"}),
})
.await
.unwrap_err();
assert_eq!(error.code(), SdkErrorCode::UnsupportedAction);
assert_eq!(error.operation(), Some(SdkOperation::Steer));
let response = service
.handle_async(request(
7,
"harness.v1.runtimes.steer",
json!({"connection": "runtime-1", "text": "go left"}),
))
.await;
assert_eq!(response["error"]["name"], "unsupported_action");
assert_eq!(response["error"]["operation"], "steer");
}
#[test]
fn support_report_and_grok_default_binding_share_the_registry() {
let mut service = HarnessSessionService::new();
let response = service.handle(request(1, "harness.v1.support.report", json!({})));
assert_eq!(response["result"]["schema"], crate::SUPPORT_REGISTRY_SCHEMA);
let params = RuntimeBackendParams {
harness: HarnessId::from(HarnessId::GROK),
protocol: None,
launch: None,
base_url: None,
policy: RuntimePolicy::Default,
};
let backend = match runtime_backend(¶ms) {
Ok(backend) => backend,
Err(_) => panic!("Grok should bind through its registered ACP launch"),
};
assert_eq!(backend.harness().as_str(), HarnessId::GROK);
assert!(backend.capabilities().start_session);
let registered = harness_support_registry()
.harnesses
.into_iter()
.find(|harness| harness.id.as_str() == HarnessId::GROK)
.and_then(|harness| harness.runtime.default_launch)
.unwrap();
assert!(!registered
.arguments
.iter()
.any(|argument| argument == "--always-approve"));
assert!(runtime_launch(¶ms).is_none());
let yolo = RuntimeBackendParams {
policy: RuntimePolicy::Yolo,
..params
};
assert!(runtime_launch(&yolo)
.unwrap()
.arguments
.iter()
.any(|argument| argument == "--always-approve"));
let mismatched_protocol = RuntimeBackendParams {
harness: HarnessId::from(HarnessId::CLAUDE_CODE),
protocol: Some("acp".into()),
launch: None,
base_url: None,
policy: RuntimePolicy::Default,
};
assert!(runtime_backend(&mismatched_protocol).is_err());
}
#[test]
fn load_follow_and_unfollow_share_the_same_locator() {
let mut service = HarnessSessionService::new();
let locator = pi_locator();
let loaded = service.handle(request(
1,
"harness.v1.sessions.load",
json!({"locator": locator}),
));
assert_eq!(
loaded["result"]["session"]["session_id"],
locator.session_id
);
let followed = service.handle(request(
2,
"harness.v1.sessions.follow",
json!({"locator": locator}),
));
assert_eq!(followed["result"]["subscription"], "sub-1");
assert_eq!(followed["result"]["initial"]["type"], "session_snapshot");
assert!(service.poll().is_empty());
let unfollowed = service.handle(request(
3,
"harness.v1.sessions.unfollow",
json!({"subscription": "sub-1"}),
));
assert_eq!(unfollowed["result"]["removed"], true);
}
#[test]
fn bounded_read_view_excludes_subagents_and_keeps_only_the_tail() {
let temp = std::env::temp_dir().join(format!(
"supercode-bounded-view-{}-{}",
std::process::id(),
generated_session_id()
));
let path = temp.join("parent.jsonl");
let subagents = temp.join("parent/subagents");
std::fs::create_dir_all(&subagents).unwrap();
let long_last = "x".repeat(300);
let parent_records = [
json!({"type":"user","uuid":"u1","parentUuid":null,"message":{"role":"user","content":"first"}}),
json!({"type":"assistant","uuid":"a1","parentUuid":"u1","message":{"role":"assistant","content":[{"type":"text","text":"middle"}]}}),
json!({"type":"user","uuid":"u2","parentUuid":"a1","message":{"role":"user","content":long_last}}),
];
std::fs::write(
&path,
format!(
"{}\n",
parent_records
.iter()
.map(Value::to_string)
.collect::<Vec<_>>()
.join("\n")
),
)
.unwrap();
std::fs::write(
subagents.join("agent-child.jsonl"),
concat!(
r#"{"type":"user","uuid":"cu","parentUuid":null,"agentId":"child","message":{"role":"user","content":"child work"}}"#,
"\n",
),
)
.unwrap();
let locator = SessionLocator {
harness: HarnessId::from(HarnessId::CLAUDE_CODE),
session_id: "parent".into(),
storage: StorageLocator::File { path },
};
let mut service = HarnessSessionService::new();
let complete = service.handle(request(
1,
"harness.v1.sessions.load",
json!({"locator": locator}),
));
assert_eq!(
complete["result"]["session"]["subagents"]
.as_array()
.unwrap()
.len(),
1
);
let bounded = service.handle(request(
2,
"harness.v1.sessions.load",
json!({
"locator": locator,
"view": {
"tail_messages": 1,
"max_message_chars": 256,
"include_subagents": false
},
}),
));
let session = &bounded["result"]["session"];
assert!(session["subagents"].as_array().unwrap().is_empty());
assert_eq!(session["messages"].as_array().unwrap().len(), 1);
assert_eq!(
session["messages"][0]["content"],
format!("{}\n…", "x".repeat(256))
);
let followed = service.handle(request(
3,
"harness.v1.sessions.follow",
json!({
"locator": locator,
"view": {
"tail_messages": 1,
"max_message_chars": 256,
"include_subagents": false
},
}),
));
let initial = &followed["result"]["initial"]["session"];
assert!(initial["subagents"].as_array().unwrap().is_empty());
assert_eq!(initial["messages"].as_array().unwrap().len(), 1);
let _ = std::fs::remove_dir_all(&temp);
}
#[test]
fn forty_megabyte_display_load_is_bounded_and_prompt() {
let temp = std::env::temp_dir().join(format!(
"supercode-large-display-view-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(&temp).unwrap();
let path = temp.join("rollout.jsonl");
let mut file = std::io::BufWriter::new(std::fs::File::create(&path).unwrap());
writeln!(
file,
r#"{{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{{"id":"large-display","cwd":"/tmp"}}}}"#
)
.unwrap();
let padding = "x".repeat(80 * 1024);
for index in 0..512 {
let marker = if index == 0 {
"OLDEST-SHOULD-NOT-LOAD"
} else if index == 511 {
"LATEST-MUST-LOAD"
} else {
"bulk"
};
writeln!(
file,
"{}",
json!({
"timestamp": "2026-01-01T00:00:01Z",
"type": "response_item",
"payload": {
"type": "message",
"role": "assistant",
"content": [{"type": "output_text", "text": format!("{marker}:{padding}")}],
},
})
)
.unwrap();
}
file.flush().unwrap();
drop(file);
assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
let locator = SessionLocator {
harness: HarnessId::from(HarnessId::CODEX),
session_id: "large-display".into(),
storage: StorageLocator::File { path },
};
let started = Instant::now();
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.sessions.load",
json!({
"locator": locator,
"view": {
"tail_messages": 500,
"max_message_chars": 1024,
"include_subagents": false,
"display_history": true,
},
}),
));
let elapsed = started.elapsed();
let wire = response.to_string();
eprintln!(
"bounded 40 MiB display load: {elapsed:?}, {} response bytes",
wire.len()
);
assert!(response.get("error").is_none(), "{response:#}");
assert!(wire.contains("LATEST-MUST-LOAD"));
assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
assert!(
wire.len() < 2 * 1024 * 1024,
"bounded wire was {} bytes",
wire.len()
);
assert!(
elapsed.as_secs_f64() < 3.0,
"bounded 40 MiB load took {elapsed:?}"
);
let _ = std::fs::remove_dir_all(&temp);
}
#[test]
fn forty_megabyte_goose_store_display_load_reads_only_the_tail() {
let temp = std::env::temp_dir().join(format!(
"supercode-large-goose-view-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(&temp).unwrap();
let path = temp.join("sessions.db");
let connection = rusqlite::Connection::open(&path).unwrap();
connection
.execute_batch(
"CREATE TABLE sessions (
id TEXT PRIMARY KEY, name TEXT NOT NULL, working_dir TEXT NOT NULL,
created_at TEXT NOT NULL, updated_at TEXT NOT NULL,
session_type TEXT NOT NULL, extension_data TEXT,
goose_mode TEXT NOT NULL, provider_name TEXT, model_config_json TEXT,
archived_at TEXT
);
CREATE TABLE messages (
id INTEGER PRIMARY KEY, session_id TEXT NOT NULL, message_id TEXT,
role TEXT NOT NULL, content_json TEXT NOT NULL,
created_timestamp INTEGER NOT NULL, metadata_json TEXT
);",
)
.unwrap();
connection
.execute(
"INSERT INTO sessions VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL)",
rusqlite::params![
"goose-large",
"Large Goose session",
"/tmp",
"2026-01-01 00:00:00",
"2026-01-01 00:00:02",
"user",
"{}",
"auto",
"anthropic",
r#"{"model_name":"claude-sonnet"}"#,
],
)
.unwrap();
let old_content = serde_json::to_string(&vec![json!({
"type": "text",
"text": format!("OLDEST-SHOULD-NOT-LOAD:{}", "x".repeat(40 * 1024 * 1024)),
})])
.unwrap();
connection
.execute(
"INSERT INTO messages VALUES (1, ?1, 'old', 'user', ?2, 1, '{}')",
rusqlite::params!["goose-large", old_content],
)
.unwrap();
connection
.execute(
"INSERT INTO messages VALUES (2, ?1, 'new', 'assistant', ?2, 2, '{}')",
rusqlite::params![
"goose-large",
r#"[{"type":"text","text":"LATEST-MUST-LOAD"}]"#
],
)
.unwrap();
drop(connection);
assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
let locator = SessionLocator {
harness: HarnessId::from(HarnessId::GOOSE),
session_id: "goose-large".into(),
storage: StorageLocator::Sqlite {
path,
selector: "goose-large".into(),
},
};
let started = Instant::now();
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.sessions.load",
json!({
"locator": locator,
"view": {
"tail_messages": 1,
"max_message_chars": 1024,
"include_subagents": false,
"display_history": true,
},
}),
));
let elapsed = started.elapsed();
let wire = response.to_string();
eprintln!(
"bounded 40 MiB Goose display load: {elapsed:?}, {} response bytes",
wire.len()
);
assert!(response.get("error").is_none(), "{response:#}");
assert!(wire.contains("LATEST-MUST-LOAD"));
assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
assert!(
wire.len() < 64 * 1024,
"bounded wire was {} bytes",
wire.len()
);
assert!(
elapsed.as_secs_f64() < 1.0,
"bounded Goose load took {elapsed:?}"
);
let _ = std::fs::remove_dir_all(&temp);
}
#[test]
fn display_view_keeps_codex_assistant_history_across_compaction() {
let temp = std::env::temp_dir().join(format!(
"supercode-codex-display-view-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(&temp).unwrap();
let path = temp.join("rollout.jsonl");
std::fs::write(
&path,
concat!(
r#"{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{"id":"codex-display","cwd":"/tmp"}}"#,
"\n",
r#"{"timestamp":"2026-01-01T00:00:01Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"old prompt"}]}}"#,
"\n",
r#"{"timestamp":"2026-01-01T00:00:02Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"old answer"}]}}"#,
"\n",
r#"{"timestamp":"2026-01-01T00:00:03Z","type":"compacted","payload":{"replacement_history":[{"type":"message","role":"user","content":[{"type":"input_text","text":"old prompt"}]},{"type":"compaction","encrypted_content":"opaque"}]}}"#,
"\n",
r#"{"timestamp":"2026-01-01T00:00:04Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"new prompt"}]}}"#,
"\n",
r#"{"timestamp":"2026-01-01T00:00:05Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"new answer"}]}}"#,
"\n",
),
)
.unwrap();
let locator = SessionLocator {
harness: HarnessId::from(HarnessId::CODEX),
session_id: "codex-display".into(),
storage: StorageLocator::File { path },
};
let mut service = HarnessSessionService::new();
let continuation = service.handle(request(
1,
"harness.v1.sessions.load",
json!({"locator": locator}),
));
let continuation_text = continuation["result"]["session"]["messages"].to_string();
assert!(!continuation_text.contains("old answer"));
let display = service.handle(request(
2,
"harness.v1.sessions.load",
json!({
"locator": locator,
"view": {
"tail_messages": 10,
"include_subagents": false,
"display_history": true,
},
}),
));
let display_text = display["result"]["session"]["messages"].to_string();
assert!(display_text.contains("old prompt"));
assert!(display_text.contains("old answer"));
assert!(display_text.contains("new prompt"));
assert!(display_text.contains("new answer"));
let _ = std::fs::remove_dir_all(&temp);
}
#[test]
fn load_supports_bounded_windows_and_media_metadata() {
let mut service = HarnessSessionService::new();
let locator = pi_locator();
let bounded = service.handle(request(
1,
"harness.v1.sessions.load",
json!({
"locator": locator,
"options": {
"include_subagents": false,
"message_limit": 2,
"message_offset": 1
}
}),
));
assert_eq!(bounded["result"]["window"]["offset"], 1);
assert_eq!(bounded["result"]["window"]["returned"], 2);
assert!(bounded["result"]["summary"]["first_message"].is_object());
assert!(bounded["result"]["summary"]["last_message"].is_object());
assert_eq!(
bounded["result"]["session"]["messages"]
.as_array()
.unwrap()
.len(),
2
);
assert!(bounded["result"]["session"]["subagents"]
.as_array()
.unwrap()
.is_empty());
let tail = service.handle(request(
2,
"harness.v1.sessions.load",
json!({"locator": locator, "options": {"message_tail": 1}}),
));
assert_eq!(tail["result"]["window"]["returned"], 1);
assert_eq!(tail["result"]["window"]["has_more"], true);
assert_eq!(tail["result"]["window"]["has_older"], true);
assert!(tail["result"]["window"]["older_items"].as_u64().unwrap() > 0);
assert!(tail["result"]["summary"]["first_message"].is_object());
let metadata_only = service.handle(request(
3,
"harness.v1.sessions.load",
json!({"locator": locator, "options": {"inline_media": "metadata"}}),
));
assert!(metadata_only["result"]["session"]
.to_string()
.contains("media_reference"));
assert!(!metadata_only["result"]["session"]
.to_string()
.contains("data:image/"));
}
#[test]
fn import_translate_branch_and_handoff_use_typed_artifacts() {
let mut service = HarnessSessionService::new();
let locator = pi_locator();
let translated = service.handle(request(
1,
"harness.v1.sessions.translate",
json!({"locator": locator, "target_harness": "grok"}),
));
assert_eq!(translated["result"]["artifact"]["source_harness"], "pi");
assert_eq!(translated["result"]["artifact"]["target_harness"], "grok");
assert!(translated["result"]["artifact"]["content"]
.as_str()
.is_some_and(|content| !content.is_empty()));
for target in ["opencode", "open-code"] {
let opencode = service.handle(request(
6,
"harness.v1.sessions.translate",
json!({"locator": locator, "target_harness": target}),
));
assert_eq!(opencode["result"]["artifact"]["target_harness"], "opencode");
}
let goose = service.handle(request(
7,
"harness.v1.sessions.translate",
json!({"locator": locator, "target_harness": "goose"}),
));
assert_eq!(goose["result"]["artifact"]["target_harness"], "goose");
assert!(serde_json::from_str::<Value>(
goose["result"]["artifact"]["content"].as_str().unwrap()
)
.unwrap()["conversation"]
.is_array());
let imported = service.handle(request(
2,
"harness.v1.sessions.import",
json!({
"source_harness": "grok",
"content": translated["result"]["artifact"]["content"],
}),
));
assert_eq!(imported["result"]["session"]["source"], "grok");
let branched = service.handle(request(
3,
"harness.v1.sessions.branch",
json!({"locator": locator, "target_harness": "codex"}),
));
assert_eq!(branched["result"]["parent"]["harness"], "pi");
assert!(branched["result"]["bootstrap_prompt"]
.as_str()
.unwrap()
.contains("frozen parent transcript"));
assert_eq!(branched["result"]["artifact"]["target_harness"], "codex");
let handoff = service.handle(request(
4,
"harness.v1.sessions.handoff",
json!({"locator": locator, "target_harness": "pi", "cwd": "/tmp/project"}),
));
assert_eq!(handoff["result"]["launch"]["program"], "pi");
assert_eq!(handoff["result"]["launch"]["cwd"], "/tmp/project");
assert_eq!(handoff["result"]["requires_materialization"], true);
let goose_handoff = service.handle(request(
8,
"harness.v1.sessions.handoff",
json!({"locator": locator, "target_harness": "goose", "cwd": "/tmp/project"}),
));
assert_eq!(goose_handoff["result"]["launch"]["program"], "goose");
assert_eq!(
goose_handoff["result"]["materialize"]["arguments"],
json!(["session", "import", "{artifact_path}"])
);
let resumed = service.handle(request(
5,
"harness.v1.sessions.resume_instructions",
json!({"locator": locator, "cwd": "/tmp/project", "policy": "yolo"}),
));
assert_eq!(resumed["result"]["launch"]["program"], "pi");
assert_eq!(resumed["result"]["launch"]["arguments"][0], "--approve");
}
#[test]
fn reduce_persists_and_reloads_a_byte_exact_reversible_bundle() {
let temp = std::env::temp_dir().join(format!(
"supercode-service-reduce-{}-{}",
std::process::id(),
generated_session_id()
));
let source_path = temp.join("source.jsonl");
let store_root = temp.join("store");
std::fs::create_dir_all(&temp).unwrap();
let mut records = vec![json!({
"timestamp": "2026-01-01T00:00:00Z",
"type": "session_meta",
"payload": {"id": "codex-reduce", "cwd": "/tmp/project"},
})];
for turn in 0..16 {
records.push(json!({
"timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 1),
"type": "response_item",
"payload": {
"type": "message",
"role": "user",
"content": [{
"type": "input_text",
"text": format!("request {turn}: {}", "context ".repeat(80)),
}],
},
}));
records.push(json!({
"timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 2),
"type": "response_item",
"payload": {
"type": "message",
"role": "assistant",
"content": [{
"type": "output_text",
"text": format!("answer {turn}: {}", "implementation detail ".repeat(80)),
}],
},
}));
}
let source = format!(
"{}\n",
records
.iter()
.map(Value::to_string)
.collect::<Vec<_>>()
.join("\n")
);
std::fs::write(&source_path, &source).unwrap();
let locator = SessionLocator {
harness: HarnessId::from(HarnessId::CODEX),
session_id: "codex-reduce".into(),
storage: StorageLocator::File {
path: source_path.clone(),
},
};
let original = load_session(&locator).unwrap();
let mut service =
HarnessSessionService::new().with_reduction_store_root(store_root.clone());
let response = service.handle(request(
1,
"harness.v1.sessions.reduce",
json!({
"locator": locator,
"target_harness": "claude-code",
"keep_last": 4,
}),
));
assert!(response.get("error").is_none(), "{response:#}");
let receipt = &response["result"]["receipt"];
assert_eq!(receipt["source_harness"], "codex");
assert_eq!(receipt["target_harness"], "claude-code");
assert_eq!(receipt["verified"], true);
assert_eq!(receipt["reversible"], true);
assert!(receipt["reductions"].as_u64().unwrap() > 0);
assert!(
receipt["source_tokens"].as_u64().unwrap()
> receipt["reduced_tokens"].as_u64().unwrap()
);
assert!(receipt["ratio"].as_f64().unwrap() > 1.0);
assert!(response["result"]["bootstrap_prompt"]
.as_str()
.unwrap()
.contains("Do not guess hidden content"));
let rescue_id = receipt["id"].as_str().unwrap();
let store = crate::SessionStore::open(&store_root).unwrap();
let sidecar =
Session::from_sidecar_str(&store.load_sidecar(rescue_id).unwrap().unwrap()).unwrap();
let log = store.load_reduction_log(rescue_id).unwrap().unwrap();
let persisted_view = parse_messages_jsonl(&store.load(rescue_id).unwrap()).unwrap();
let policy = reduce::ReductionPolicy {
clear_turns_older_than: Some(4),
..Default::default()
};
let (restamped_view, reapplied_log) =
reduce::project_messages(&sidecar.messages, &policy, &log);
assert_eq!(
messages_jsonl(&persisted_view).unwrap(),
messages_jsonl(&restamped_view).unwrap()
);
assert_eq!(reapplied_log, log);
reduce::verify_log(&log, &sidecar).unwrap();
assert_eq!(
reduce::invert(&restamped_view, &log, &sidecar).unwrap(),
original.messages
);
assert_eq!(std::fs::read_to_string(&source_path).unwrap(), source);
std::fs::remove_dir_all(temp).ok();
}
#[test]
fn read_surfaces_view_a_severed_claude_graph_while_transfer_still_refuses_it() {
let temp = std::env::temp_dir().join(format!(
"supercode-severed-view-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(&temp).unwrap();
let path = temp.join("severed.jsonl");
std::fs::write(
&path,
concat!(
r#"{"type":"user","uuid":"orphan-u","parentUuid":null,"message":{"role":"user","content":"stranded prompt"}}"#,
"\n",
r#"{"type":"assistant","uuid":"live-a","parentUuid":"pruned","message":{"id":"m","role":"assistant","content":[{"type":"text","text":"live answer"}]}}"#,
"\n",
),
)
.unwrap();
let locator = SessionLocator {
harness: HarnessId::from(HarnessId::CLAUDE_CODE),
session_id: "severed".into(),
storage: StorageLocator::File { path },
};
let mut service = HarnessSessionService::new();
let viewed = service.handle(request(
1,
"harness.v1.sessions.load",
json!({"locator": locator}),
));
let session = &viewed["result"]["session"];
assert_eq!(session["fidelity"], "semantic");
assert_eq!(session["messages"].as_array().unwrap().len(), 2);
assert!(session["residue"].as_array().unwrap().iter().any(|entry| {
entry
.as_str()
.is_some_and(|entry| entry.contains("live-a") && entry.contains("pruned"))
}));
let strict = service.handle(request(
2,
"harness.v1.sessions.load",
json!({"locator": locator, "fidelity": "byte_lossless"}),
));
assert!(strict["error"]["message"]
.as_str()
.unwrap()
.contains("cannot reconstruct lossless Claude continuation"));
let translated = service.handle(request(
3,
"harness.v1.sessions.translate",
json!({"locator": locator, "target_harness": "codex"}),
));
assert!(translated["error"]["message"]
.as_str()
.unwrap()
.contains("cannot reconstruct lossless Claude continuation"));
let resumed = service.handle(request(
4,
"harness.v1.sessions.resume_instructions",
json!({"locator": locator}),
));
assert!(resumed["error"]["message"]
.as_str()
.unwrap()
.contains("cannot reconstruct lossless Claude continuation"));
let _ = std::fs::remove_dir_all(&temp);
}
#[test]
fn structured_resume_launches_cover_gemini_goose_and_supercode() {
let gemini = resume_launch(
HarnessId::GEMINI,
"gemini-session",
Path::new("/tmp/project"),
ResumePolicy::Yolo,
)
.unwrap_or_else(|_| panic!("Gemini resume launch must be registered"));
assert_eq!(gemini.program, "gemini");
assert_eq!(gemini.arguments, ["--yolo", "--resume", "gemini-session"]);
let goose = resume_launch(
HarnessId::GOOSE,
"goose-session",
Path::new("/tmp/project"),
ResumePolicy::Yolo,
)
.unwrap_or_else(|_| panic!("Goose resume launch must be registered"));
assert_eq!(goose.program, "goose");
assert_eq!(
goose.arguments,
["session", "--resume", "--session-id", "goose-session"]
);
let supercode = resume_launch(
HarnessId::SUPERCODE,
"supercode-session",
Path::new("/tmp/project"),
ResumePolicy::Yolo,
)
.unwrap_or_else(|_| panic!("Supercode resume launch must be registered"));
assert_eq!(supercode.program, "supercode");
assert_eq!(
supercode.arguments,
["--dangerous", "resume", "supercode-session"]
);
}
#[test]
fn diagonal_artifacts_preserve_claude_subagents_and_grok_bundle_members() {
let temp = std::env::temp_dir().join(format!(
"supercode-harness-artifact-{}-{}",
std::process::id(),
generated_session_id()
));
let main_path = temp.join("parent.jsonl");
let subagent_path = temp.join("parent/subagents/agent-child.jsonl");
std::fs::create_dir_all(subagent_path.parent().unwrap()).unwrap();
let fixture = std::fs::read_to_string(
PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/claude_code_session.jsonl"),
)
.unwrap();
let parent = fixture.trim_end_matches('\n');
let child = fixture.trim_end_matches('\n');
std::fs::write(&main_path, parent).unwrap();
std::fs::write(&subagent_path, child).unwrap();
let locator = SessionLocator {
harness: HarnessId::from(HarnessId::CLAUDE_CODE),
session_id: "213bb148-51ea-453f-9206-f8b4b1168547".into(),
storage: StorageLocator::File {
path: main_path.clone(),
},
};
let mut service = HarnessSessionService::new();
let claude = service.handle(request(
1,
"harness.v1.sessions.translate",
json!({"locator": locator, "target_harness": "claude-code"}),
));
let artifact = &claude["result"]["artifact"];
assert_eq!(artifact["fidelity"], "byte_lossless");
assert_eq!(artifact["content"], parent);
let files = artifact["files"].as_array().unwrap();
assert!(files.iter().any(|file| {
file["role"] == "subagent"
&& file["path"]
.as_str()
.is_some_and(|path| path.ends_with("/subagents/agent-child.jsonl"))
&& file["content"] == child
}));
assert!(!artifact["content"].as_str().unwrap().ends_with('\n'));
let grok = service.handle(request(
2,
"harness.v1.sessions.translate",
json!({"locator": grok_locator(), "target_harness": "grok"}),
));
let files = grok["result"]["artifact"]["files"].as_array().unwrap();
for name in ["summary.json", "updates.jsonl"] {
let expected = std::fs::read_to_string(
PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/fixtures/grok_session")
.join(name),
)
.unwrap();
assert!(files.iter().any(|file| {
file["path"] == name && file["role"] == "bundle" && file["content"] == expected
}));
}
std::fs::remove_dir_all(temp).ok();
}
#[test]
fn every_non_grok_handoff_mints_and_uses_a_fresh_target_identity() {
let mut service = HarnessSessionService::new();
let source = pi_locator();
for (target, format) in [
("claude-code", SessionFormat::ClaudeCode),
("codex", SessionFormat::Codex),
("opencode", SessionFormat::OpenCode),
("pi", SessionFormat::Pi),
] {
let result = service.handle(request(
1,
"harness.v1.sessions.handoff",
json!({"locator": source, "target_harness": target, "cwd": "/tmp/project"}),
));
let artifact = &result["result"]["artifact"];
let target_id = artifact["session_id"].as_str().unwrap();
assert_ne!(target_id, source.session_id, "{target}");
let parsed = Session::load_str(artifact["content"].as_str().unwrap(), format).unwrap();
assert_eq!(
parsed.meta.session_id.as_deref(),
Some(target_id),
"{target}"
);
if target != "pi" {
assert!(result["result"]["launch"]["arguments"]
.as_array()
.unwrap()
.iter()
.any(|argument| argument == target_id));
}
if target == "opencode" {
assert!(target_id.starts_with("ses_"));
fn assert_session_ids(value: &Value, target_id: &str) {
match value {
Value::Object(fields) => {
if let Some(session_id) = fields.get("sessionID") {
assert_eq!(session_id, target_id);
}
for child in fields.values() {
assert_session_ids(child, target_id);
}
}
Value::Array(values) => {
for child in values {
assert_session_ids(child, target_id);
}
}
_ => {}
}
}
let document: Value =
serde_json::from_str(artifact["content"].as_str().unwrap()).unwrap();
assert_session_ids(&document, target_id);
}
}
let first = service.handle(request(
2,
"harness.v1.sessions.handoff",
json!({"locator": source, "target_harness": "codex"}),
));
let second = service.handle(request(
3,
"harness.v1.sessions.handoff",
json!({"locator": source, "target_harness": "codex"}),
));
assert_ne!(
first["result"]["artifact"]["session_id"],
second["result"]["artifact"]["session_id"]
);
}
#[test]
fn grok_handoff_uses_the_official_importer_contract() {
let mut service = HarnessSessionService::new();
let source = opencode_locator();
let response = service.handle(request(
1,
"harness.v1.sessions.handoff",
json!({
"locator": source,
"target_harness": "grok",
"cwd": "/tmp/grok-handoff-project",
}),
));
let result = &response["result"];
assert_eq!(result["artifact"]["target_harness"], "claude-code");
assert!(result["artifact"]["suggested_filename"]
.as_str()
.unwrap()
.ends_with(".grok-import.claude-code.jsonl"));
let artifact = Session::load_str(
result["artifact"]["content"].as_str().unwrap(),
SessionFormat::ClaudeCode,
)
.unwrap();
assert_eq!(
artifact.meta.cwd.as_deref(),
Some(Path::new("/tmp/grok-handoff-project"))
);
let target_session_id = artifact.meta.session_id.as_deref().unwrap();
assert_eq!(target_session_id.len(), 36);
assert_eq!(target_session_id.as_bytes()[14], b'4');
assert_ne!(target_session_id, opencode_locator().session_id);
assert_eq!(
result["artifact"]["session_id"],
artifact.meta.session_id.as_deref().unwrap()
);
assert_eq!(
result["materialize"]["arguments"],
json!(["import", "--json", "{artifact_path}"])
);
assert_eq!(
result["launch"]["arguments"],
json!(["--resume", "{imported_session_id}", "--fork-session"])
);
assert!(result["note"]
.as_str()
.unwrap()
.contains("outcome=imported"));
assert!(!result["launch"]["arguments"]
.as_array()
.unwrap()
.iter()
.any(|argument| argument == &opencode_locator().session_id));
}
#[tokio::test]
async fn inventory_rejects_unknown_harnesses_and_runtime_attach_is_honest() {
let mut service = HarnessSessionService::new();
let inventory = service
.handle_async(request(
1,
"harness.v1.harnesses.list",
json!({"harnesses": ["missing"]}),
))
.await;
assert_eq!(inventory["error"]["code"], -32602);
let attached = service
.handle_async(request(
2,
"harness.v1.runtimes.attach_existing",
json!({"harness": "codex", "runtime_id": "thread-1"}),
))
.await;
assert_eq!(attached["error"]["code"], -32000);
assert!(attached["error"]["message"]
.as_str()
.unwrap()
.contains("runtimes.resume"));
}
#[test]
fn invalid_params_and_unknown_methods_use_json_rpc_errors() {
let mut service = HarnessSessionService::new();
let invalid = service.handle(request(1, "harness.v1.sessions.load", json!({})));
assert_eq!(invalid["error"]["code"], -32602);
let unknown = service.handle(request(2, "harness.v1.unknown", json!({})));
assert_eq!(unknown["error"]["code"], -32601);
}
#[cfg(unix)]
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn async_service_drives_a_generic_acp_runtime() {
let _environment_guard = crate::live_runtime::test_environment_lock();
let script = r#"
i=0
while IFS= read -r line; do
i=$((i + 1))
case "$i" in
1) printf '%s\n' '{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentCapabilities":{},"authMethods":[]}}' ;;
2) printf '%s\n' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"svc_acp"}}' ;;
3)
printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"ok"}}}}'
printf '%s\n' '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
;;
4)
printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"from terminal"}}}}'
printf '%s\n' '{"jsonrpc":"2.0","id":4,"result":{"stopReason":"end_turn"}}'
;;
esac
done
"#;
let mut service = HarnessSessionService::new();
let started = service
.handle_async(request(
1,
"harness.v1.runtimes.start",
json!({
"harness": "codex",
"protocol": "acp",
"cwd": std::env::current_dir().unwrap(),
"launch": {"program": "/bin/sh", "arguments": ["-c", script], "env": {}},
}),
))
.await;
assert_eq!(started["result"]["connection"], "runtime-1");
assert_eq!(started["result"]["handle"]["runtime_id"], "svc_acp");
let terminal = service
.handle_async(request(
9,
"harness.v1.runtimes.terminal_instructions",
json!({"connection":"runtime-1"}),
))
.await;
let arguments = terminal["result"]["launch"]["arguments"]
.as_array()
.expect("hosted runtime should return terminal arguments");
let endpoint_index = arguments
.iter()
.position(|value| value == "--endpoint")
.expect("terminal command should use an opaque endpoint");
let endpoint = LiveRuntimeEndpoint::parse(
arguments[endpoint_index + 1]
.as_str()
.expect("endpoint argument should be text"),
)
.unwrap();
assert!(!terminal.to_string().contains("Bearer"));
let workspace = std::env::current_dir().unwrap();
let receipt = resolve_live_runtime(
&endpoint,
&LiveRuntimeSource {
harness: "codex".into(),
session_id: "svc_acp".into(),
workspace,
},
)
.unwrap();
let remote = crate::HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
.await
.unwrap();
let mut attachment = crate::FrontendRuntime::attach(remote.as_ref(), 100)
.await
.unwrap();
let sent = service
.handle_async(request(
2,
"harness.v1.runtimes.send_input",
json!({"connection": "runtime-1", "text": "hi"}),
))
.await;
assert_eq!(sent["result"]["turn_id"], "3");
let mut events = Vec::new();
for _ in 0..20 {
events.extend(service.poll_runtimes().await);
if events.len() >= 2 {
break;
}
tokio::time::sleep(Duration::from_millis(2)).await;
}
assert!(events
.iter()
.any(|event| { event["params"]["event"]["kind"] == "session/update" }));
assert!(events.iter().any(|event| {
event["params"]["event"]["kind"] == "supercode/acp_request_completed"
}));
let saw_editor_reply = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "text_delta" && event.payload["text"] == "ok" {
break;
}
}
})
.await;
assert!(
saw_editor_reply.is_ok(),
"terminal should observe the editor-driven turn"
);
crate::FrontendRuntime::submit(remote.as_ref(), "DRIVE FROM TERMINAL".into())
.await
.unwrap();
let saw_terminal_reply = tokio::time::timeout(Duration::from_secs(2), async {
loop {
let event = attachment.next_event().await.unwrap();
if event.kind == "text_delta" && event.payload["text"] == "from terminal" {
break;
}
}
})
.await;
assert!(
saw_terminal_reply.is_ok(),
"terminal should drive the same runtime"
);
let closed = service
.handle_async(request(
3,
"harness.v1.runtimes.close",
json!({"connection": "runtime-1"}),
))
.await;
assert_eq!(closed["result"]["closed"], true);
}
}