use std::collections::{BTreeMap, BTreeSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use tokio::sync::Notify;
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, HarnessHomes, HarnessId,
ImplementationKind, LiveRuntimeEndpoint, LiveRuntimeSource, OpenCodeRuntimeBackend,
PiRuntimeBackend, Role, RuntimeAttachRequest, RuntimeBackend, RuntimeConnection, RuntimeInput,
RuntimeLaunch, RuntimeStartRequest, Session, SessionDescriptor, 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_METHODS: &[&str] = &[
"harness.v1.support.report",
"harness.v1.harnesses.list",
"harness.v1.harnesses.probe",
"harness.v1.harnesses.settings",
"harness.v1.harnesses.configure",
"harness.v1.harnesses.auth.methods",
"harness.v1.harnesses.auth.begin",
"harness.v1.harnesses.auth.verify",
"harness.v1.sessions.discover",
"harness.v1.sessions.load",
"harness.v1.sessions.follow",
"harness.v1.sessions.unfollow",
"harness.v1.sessions.activity.subscribe",
"harness.v1.sessions.activity.unsubscribe",
"harness.v1.sessions.index.subscribe",
"harness.v1.sessions.index.unsubscribe",
"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.skills.list",
"harness.v1.skills.install",
"harness.v1.skills.remove",
"harness.v1.memory.show",
"harness.v1.memory.search",
"harness.v1.jobs.list",
"harness.v1.jobs.get",
"harness.v1.jobs.create",
"harness.v1.jobs.update",
"harness.v1.jobs.pause",
"harness.v1.jobs.resume",
"harness.v1.jobs.run",
"harness.v1.jobs.delete",
"harness.v1.sessions.new",
"harness.v1.sessions.reset",
"harness.v1.sessions.archive",
"harness.v1.sessions.delete",
"harness.v1.runs.list",
"harness.v1.runs.get",
"harness.v1.approvals.list",
"harness.v1.approvals.resolve",
"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",
"harness.v1.profiles.list",
"harness.v1.profiles.get",
"harness.v1.profiles.create",
"harness.v1.profiles.delete",
"harness.v1.channels.list",
"harness.v1.routes.list",
"harness.v1.triggers.list",
"harness.v1.channels.status",
"harness.v1.world.load",
"harness.v1.world.save",
"harness.v1.world.compile",
"harness.v1.world.decompile",
"harness.v1.world.import",
"harness.v1.world.export",
];
pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
pub const SESSION_ACTIVITY_EVENT_METHOD: &str = "harness.v1.sessions.activity_event";
pub const SESSION_INDEX_EVENT_METHOD: &str = "harness.v1.sessions.index_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>,
activity_subscriptions: BTreeMap<String, ActivitySubscription>,
index_subscriptions: BTreeMap<String, crate::session_index::SessionIndexSubscription>,
index_notifier: Arc<Notify>,
#[cfg(feature = "adapter-api")]
activity_monitor: crate::session_activity::SessionActivityMonitor,
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>,
approvals: crate::approvals::ApprovalRegistry,
subagent_approvals: Option<Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>>,
}
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(),
activity_subscriptions: BTreeMap::new(),
index_subscriptions: BTreeMap::new(),
index_notifier: Arc::new(Notify::new()),
#[cfg(feature = "adapter-api")]
activity_monitor: Default::default(),
next_subscription: 1,
runtimes: BTreeMap::new(),
terminal_launches: BTreeMap::new(),
runtime_sequences: BTreeMap::new(),
next_runtime: 1,
reduction_store_root: None,
approvals: crate::approvals::ApprovalRegistry::new(),
subagent_approvals: None,
}
}
pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
self.reduction_store_root = Some(root.into());
self
}
pub fn observe_subagent_approvals(
&mut self,
queue: Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>,
) {
self.subagent_approvals = Some(queue);
}
pub fn approvals(&self, query: &crate::approvals::ApprovalsQuery) -> Vec<crate::ApprovalRow> {
let now = crate::approvals::now_ms();
let mut rows = self.approvals.rows(now);
if let Some(queue) = self.subagent_approvals.as_ref() {
let queued = queue
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
rows.extend(crate::approvals::subagent_rows(&queued, now));
}
rows.retain(|row| query.matches(row));
rows.sort_by(|left, right| {
left.requested_at_ms
.cmp(&right.requested_at_ms)
.then_with(|| left.id.cmp(&right.id))
});
rows
}
async fn approvals_resolve(
&mut self,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let params = decode::<crate::approvals::ApprovalsResolveParams>(params)?;
if params.id.trim().is_empty() {
return Err(ServiceError::InvalidParams(
"approvals resolve requires the `id` of a listed approval row".into(),
));
}
let choice = match (params.decision, params.option_id.as_deref()) {
(Some(_), Some(_)) => {
return Err(ServiceError::InvalidParams(
"approvals resolve takes either `decision` or `option_id`, not both".into(),
))
}
(Some(decision), None) => crate::approvals::ApprovalChoice::Decision(decision),
(None, Some(option)) => crate::approvals::ApprovalChoice::Option(option.to_string()),
(None, None) => {
return Err(ServiceError::InvalidParams(format!(
"approvals resolve requires `decision` ({}) or an explicit `option_id`",
crate::approvals::ApprovalDecision::ALL
.map(|decision| decision.as_str())
.join(" | "),
)))
}
};
let resolution = self
.approvals
.resolution(¶ms.id, &choice)
.map_err(|error| ServiceError::InvalidParams(error.to_string()))?;
self.runtime_call(
"harness.v1.runtimes.respond",
json!({
"connection": resolution.connection,
"request_id": resolution.request_id,
"response": resolution.response,
}),
)
.await?;
Ok(json!({
"id": params.id,
"decision": params.decision.map(|decision| decision.as_str()),
"option_id": resolution.option_id,
"resolved": true,
}))
}
#[cfg(feature = "adapter-api")]
pub fn session_index_notifier(&self) -> Arc<Notify> {
Arc::clone(&self.index_notifier)
}
#[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 matches!(
method,
"harness.v1.harnesses.auth.methods"
| "harness.v1.harnesses.auth.begin"
| "harness.v1.harnesses.auth.verify"
) {
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.harness_authentication_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 matches!(
method,
"harness.v1.sessions.new"
| "harness.v1.sessions.reset"
| "harness.v1.sessions.archive"
| "harness.v1.sessions.delete"
) {
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!({}));
let verb = match method {
"harness.v1.sessions.new" => crate::SessionVerb::New,
"harness.v1.sessions.reset" => crate::SessionVerb::Reset,
"harness.v1.sessions.archive" => crate::SessionVerb::Archive,
_ => crate::SessionVerb::Delete,
};
return match self.mutate_session(verb, 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 matches!(
method,
"harness.v1.harnesses.settings" | "harness.v1.harnesses.configure"
) {
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.harness_settings_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),
};
}
if method == "harness.v1.sessions.activity.subscribe" {
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.subscribe_session_activity(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_session_activities(&mut self) -> Vec<Value> {
let subscriptions = self
.activity_subscriptions
.iter()
.map(|(id, subscription)| {
(
id.clone(),
subscription.locators.clone(),
subscription.homes.clone(),
)
})
.collect::<Vec<_>>();
let mut notifications = Vec::new();
for (subscription_id, locators, homes) in subscriptions {
let Ok(activities) = self.activity_monitor.resolve(&locators, &homes).await else {
continue;
};
let Some(subscription) = self.activity_subscriptions.get_mut(&subscription_id) else {
continue;
};
let mut changed = Vec::new();
for activity in activities {
let key = activity.key();
if subscription
.reported
.get(&key)
.is_some_and(|previous| previous.same_state(&activity))
{
continue;
}
subscription.reported.insert(key, activity.clone());
changed.push(activity);
}
if !changed.is_empty() {
notifications.push(json!({
"jsonrpc": "2.0",
"method": SESSION_ACTIVITY_EVENT_METHOD,
"params": {
"subscription": subscription_id,
"activities": changed,
},
}));
}
}
notifications
}
#[cfg(feature = "adapter-api")]
pub fn poll_session_indexes(&mut self) -> Vec<Value> {
let mut notifications = Vec::new();
for (subscription, index) in &mut self.index_subscriptions {
let homes = index.homes().clone();
match index.poll() {
Ok(Some(delta)) => match live_index_changes(delta.changes, &homes) {
Ok(changes) => notifications.push(json!({
"jsonrpc": "2.0",
"method": SESSION_INDEX_EVENT_METHOD,
"params": {
"subscription": subscription,
"revision": delta.revision,
"changes": changes,
},
})),
Err(error) => notifications.push(json!({
"jsonrpc": "2.0",
"method": SESSION_INDEX_EVENT_METHOD,
"params": {
"subscription": subscription,
"error": {"recoverable": true, "message": error_message(error)},
},
})),
},
Ok(None) => {}
Err(error) => notifications.push(json!({
"jsonrpc": "2.0",
"method": SESSION_INDEX_EVENT_METHOD,
"params": {
"subscription": subscription,
"error": {"recoverable": true, "message": error},
},
})),
}
}
notifications
}
#[cfg(feature = "adapter-api")]
async fn subscribe_session_activity(
&mut self,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let params = decode::<ActivitySubscribeParams>(params)?;
if params.locators.is_empty() {
return Err(ServiceError::InvalidParams(
"sessions.activity.subscribe requires at least one locator".into(),
));
}
if params.locators.len() > 2_048 {
return Err(ServiceError::InvalidParams(
"sessions.activity.subscribe accepts at most 2048 locators".into(),
));
}
let initial = self
.activity_monitor
.resolve(¶ms.locators, ¶ms.homes)
.await
.map_err(ServiceError::Sdk)?;
let subscription = format!("activity-sub-{}", self.next_subscription);
self.next_subscription += 1;
let reported = initial
.iter()
.cloned()
.map(|activity| (activity.key(), activity))
.collect();
self.activity_subscriptions.insert(
subscription.clone(),
ActivitySubscription {
locators: params.locators,
homes: params.homes,
reported,
},
);
Ok(json!({"subscription": subscription, "initial": initial}))
}
#[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();
let now_ms = crate::approvals::now_ms();
for (connection, runtime) in &mut self.runtimes {
let session_id = runtime.handle().runtime_id.clone();
let harness = runtime.handle().harness.clone();
match tokio::time::timeout(Duration::from_millis(1), runtime.next_event()).await {
Ok(Ok(Some(event))) => {
let terminal = event.kind == "transport_closed";
self.approvals
.observe(connection, &harness, &session_id, &event, now_ms);
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);
self.approvals.forget(&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_SERVICE_METHODS,
"notifications": [
SESSION_EVENT_METHOD,
SESSION_ACTIVITY_EVENT_METHOD,
SESSION_INDEX_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.profiles.list" | "harness.v1.profiles.get" => profiles_call(method, params),
"harness.v1.profiles.create" => {
mutate_profile(crate::profiles_control::ProfileVerb::Create, params)
}
"harness.v1.profiles.delete" => {
mutate_profile(crate::profiles_control::ProfileVerb::Delete, params)
}
"harness.v1.channels.list" | "harness.v1.channels.status" => {
channels_call(method, params)
}
"harness.v1.routes.list" => routes_call(params),
"harness.v1.triggers.list" => triggers_call(params),
"harness.v1.world.load" => {
let params = decode::<WorldLoadParams>(params)?;
let read =
crate::world_doors::load(¶ms.root, params.flavor).map_err(operation)?;
serde_json::to_value(read)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.world.save" => {
let params = decode::<WorldSaveParams>(params)?;
let saved = crate::world_doors::save(¶ms.root, params.world, params.vault)
.map_err(operation)?;
serde_json::to_value(saved)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.world.compile" => {
let params = decode::<WorldCompileParams>(params)?;
let read =
crate::world_doors::compile(params.from, ¶ms.home).map_err(operation)?;
serde_json::to_value(read)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.world.decompile" => {
let params = decode::<WorldDecompileParams>(params)?;
let report = crate::world_doors::decompile(
params.to,
params.world,
¶ms.source,
params.source_flavor,
¶ms.dest,
params.vault,
)
.map_err(operation)?;
serde_json::to_value(report)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.world.import" => {
let params = decode::<WorldImportParams>(params)?;
let imported = crate::world_doors::import(params.from, ¶ms.home, ¶ms.into)
.map_err(operation)?;
serde_json::to_value(imported)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.world.export" => {
let params = decode::<WorldExportParams>(params)?;
let report = crate::world_doors::export(params.to, ¶ms.root, ¶ms.dest)
.map_err(operation)?;
serde_json::to_value(report)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
"harness.v1.skills.list" => {
let query = decode::<crate::skills::SkillsQuery>(params)?;
if let Some(harness) = query.harness.as_deref() {
if !crate::skills::SKILL_HARNESSES.contains(&harness) {
return Err(ServiceError::UnsupportedAction(format!(
"`{harness}` has no skills root supercode reads"
)));
}
}
serde_json::to_value(crate::skills::list_skills(&query))
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.skills.install" => {
mutate_skill(crate::skills_control::SkillVerb::Install, params)
}
"harness.v1.skills.remove" => {
mutate_skill(crate::skills_control::SkillVerb::Remove, params)
}
"harness.v1.approvals.list" => {
let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
if let Some(harness) = query.harness.as_deref() {
if !crate::approvals::lists_approvals(harness) {
return Err(ServiceError::UnsupportedAction(format!(
"`{harness}` has no runtime door that carries an approval request"
)));
}
}
serde_json::to_value(self.approvals(&query))
.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 activities = crate::session_activity::resolve_stock_session_activities(
&page
.sessions
.iter()
.map(|session| session.locator.clone())
.collect::<Vec<_>>(),
&query.homes,
)
.into_iter()
.map(|activity| (activity.key(), activity))
.collect::<BTreeMap<_, _>>();
let sessions = page
.sessions
.into_iter()
.map(|session| {
let mut value = live_descriptor_value(&session, &peers)?;
let activity_key = (
session.locator.harness.as_str().to_string(),
session.locator.session_id.clone(),
);
if let Some(activity) = activities.get(&activity_key) {
value["activity"] = serde_json::to_value(activity)
.map_err(|error| ServiceError::Operation(error.to_string()))?;
if let Some(status) = legacy_live_status(activity) {
value["live_status"] = json!(status);
}
}
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.activity.unsubscribe" => {
let params = decode::<UnfollowParams>(params)?;
Ok(json!({
"removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
}))
}
"harness.v1.sessions.index.subscribe" => {
let query = decode::<DiscoveryQuery>(params)?;
crate::session_index::validate_query(&query)
.map_err(ServiceError::InvalidParams)?;
let homes = query.homes.clone();
let (index, initial) = crate::session_index::SessionIndexSubscription::open(
query,
Arc::clone(&self.index_notifier),
)
.map_err(ServiceError::Operation)?;
let peers = peers_for_descriptors(&initial, &homes);
let initial = initial
.iter()
.map(|descriptor| live_descriptor_value(descriptor, &peers))
.collect::<std::result::Result<Vec<_>, ServiceError>>()?;
let subscription = format!("index-sub-{}", self.next_subscription);
self.next_subscription += 1;
self.index_subscriptions.insert(subscription.clone(), index);
Ok(json!({
"subscription": subscription,
"revision": 1,
"initial": initial,
}))
}
"harness.v1.sessions.index.unsubscribe" => {
let params = decode::<UnfollowParams>(params)?;
Ok(json!({
"removed": self.index_subscriptions.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.jobs.list" => {
let query = decode::<crate::jobs::JobsQuery>(params)?;
if let Some(harness) = query.harness.as_deref() {
refuse_harness_without_jobs(harness, "jobs.list")?;
}
let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
serde_json::to_value(listing)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.jobs.get" => {
let params = decode::<JobsGetParams>(params)?;
refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
.map_err(operation)?
{
Some((job, source)) => Ok(json!({"job": job, "source": source})),
None => Err(ServiceError::Operation(format!(
"`{}` has no scheduled job `{}`",
params.harness, params.id
))),
}
}
"harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
"harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
"harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
"harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
"harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
"harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
"harness.v1.runs.list" => {
let query = decode::<crate::runs::RunsQuery>(params)?;
if let Some(harness) = query.harness.as_deref() {
refuse_harness_without_runs(harness, "runs.list")?;
}
let listing = crate::runs::list_runs(&query).map_err(operation)?;
serde_json::to_value(listing)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.runs.get" => {
let params = decode::<RunsGetParams>(params)?;
refuse_harness_without_runs(¶ms.harness, "runs.get")?;
match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
.map_err(operation)?
{
Some((run, source)) => Ok(json!({"run": run, "source": source})),
None => Err(ServiceError::Operation(format!(
"`{}` has no run `{}`",
params.harness, params.id
))),
}
}
"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),
mcp_servers: params.mcp_servers,
})
.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 image_urls = validate_runtime_image_urls(params.image_urls)?;
let runtime = self.runtime_mut(¶ms.connection)?;
let turn_id = runtime
.send_input(RuntimeInput {
text: params.text,
image_urls,
})
.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" => {
let params = decode::<RuntimeInputParams>(params)?;
if !params.image_urls.is_empty() {
return Err(ServiceError::InvalidParams(
"runtime steering accepts text only".into(),
));
}
let text = params.text.trim();
if text.is_empty() || text.chars().count() > 50_000 {
return Err(ServiceError::InvalidParams(
"runtime steering requires 1 to 50,000 text characters".into(),
));
}
self.runtime_mut(¶ms.connection)?
.steer(text.to_string())
.await
.map_err(operation)?;
Ok(json!({}))
}
"harness.v1.runtimes.respond" => {
let params = decode::<RuntimeRespondParams>(params)?;
let request_id = params.request_id.clone();
self.runtime_mut(¶ms.connection)?
.respond(params.request_id, params.response)
.await
.map_err(operation)?;
self.approvals.answered(¶ms.connection, &request_id);
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);
self.approvals.forget(¶ms.connection);
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)
}
#[cfg(feature = "adapter-api")]
fn harness_settings_call(
&self,
method: &str,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let homes = crate::HarnessHomes::default();
match method {
"harness.v1.harnesses.settings" => {
let params = decode::<HarnessSettingsParams>(params)?;
let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
.map_err(|error| ServiceError::Operation(error.to_string()))?;
serde_json::to_value(report)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.harnesses.configure" => {
let params = decode::<ConfigureHarnessParams>(params)?;
let report = crate::configure_harness_interop_settings(
&homes,
¶ms.harness,
¶ms.changes,
params.expected_revision.as_deref(),
)
.map_err(|error| ServiceError::Operation(error.to_string()))?;
serde_json::to_value(report)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
_ => Err(ServiceError::MethodNotFound),
}
}
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 mutate_session(
&mut self,
verb: crate::SessionVerb,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let mutation = decode::<crate::SessionMutation>(params)?;
let door = crate::sessions_control::door(&mutation.harness, verb)
.map_err(session_control_error)?;
let outcome = match door {
#[cfg(not(feature = "adapter-api"))]
crate::SessionDoor::Live(command) => {
return Err(ServiceError::Operation(format!(
"`{}` performs `sessions.{}` by typing `{command}` into a live driven \
session, which needs this build's `adapter-api` feature",
mutation.harness,
verb.as_str()
)));
}
#[cfg(feature = "adapter-api")]
crate::SessionDoor::Live(command) => {
let connection = mutation
.connection
.clone()
.filter(|value| !value.trim().is_empty())
.ok_or_else(|| {
ServiceError::InvalidParams(format!(
"`{}` performs `sessions.{}` by typing `{command}` into a live \
driven session: pass the `connection` of an open runtime \
(`harness.v1.runtimes.start`)",
mutation.harness,
verb.as_str()
))
})?;
let runtime = self.runtime_mut(&connection)?;
let session = mutation
.session
.clone()
.filter(|value| !value.trim().is_empty())
.unwrap_or_else(|| runtime.handle().runtime_id.clone());
runtime
.send_input(RuntimeInput {
text: command.to_string(),
image_urls: Vec::new(),
})
.await
.map_err(operation)?;
crate::sessions_control::live_outcome(verb, &mutation, command, session)
.map_err(session_control_error)?
}
_ => crate::sessions_control::mutate(verb, &mutation)
.await
.map_err(session_control_error)?,
};
serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
}
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()))
}
#[cfg(feature = "adapter-api")]
async fn harness_authentication_call(
&self,
method: &str,
params: Value,
) -> std::result::Result<Value, ServiceError> {
match method {
"harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
let params = decode::<HarnessAuthenticationParams>(params)?;
serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
"harness.v1.harnesses.auth.begin" => {
let params = decode::<BeginHarnessAuthenticationParams>(params)?;
let cwd = params
.cwd
.or_else(|| std::env::current_dir().ok())
.unwrap_or_else(|| PathBuf::from("."));
let plan = crate::harness_authentication_plan(
¶ms.harness,
params.environment,
params.method,
&cwd,
)
.map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
serde_json::to_value(plan)
.map_err(|error| ServiceError::Operation(error.to_string()))
}
_ => Err(ServiceError::MethodNotFound),
}
}
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 orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
.then(crate::orchestrator::daemon_entry)
.and_then(Result::ok);
let executable = match &orchestrator_entry {
Some(entry) => Some(entry.clone()),
None => launch.and_then(|launch| find_executable(&launch.program)),
};
let installed = executable.is_some();
let version = if params.skip_versions || orchestrator_entry.is_some() {
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 if matches!(
descriptor.id.as_str(),
HarnessId::CLAUDE_CODE | HarnessId::CODEX
) {
HarnessAuthState::Required
} else {
HarnessAuthState::Unknown
};
let mut runtime = if installed {
HarnessRuntimeState::Degraded
} else {
HarnessRuntimeState::Unavailable
};
let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
let mut reason = (!installed).then(|| {
if is_orchestrator {
format!(
"{} is supported but its daemon entry `{}` was not found",
descriptor.display_name,
crate::orchestrator::DAEMON_ENTRY
)
} else {
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(|| {
if is_orchestrator {
format!(
"Install the `supercode-orchestrator` package so `{}` resolves.",
crate::orchestrator::DAEMON_ENTRY
)
} else {
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(),
);
let running = probe_running_instance(descriptor.id.as_str());
return LocalHarness {
gateway: gateway_health(
descriptor.id.as_str(),
installed,
running.as_ref(),
version.as_deref(),
),
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 },
running,
reason,
repair,
};
};
match tokio::time::timeout(
Duration::from_secs(30),
backend.start(RuntimeStartRequest {
cwd,
launch: Some(isolated.launch.clone()),
mcp_servers: Vec::new(),
}),
)
.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 && auth == HarnessAuthState::Required {
reason =
Some("Executable found, but no native authentication evidence is present.".into());
repair = Some(format!(
"Run `supercode harness login {}` to use the harness-owned sign-in flow.",
descriptor.id.as_str()
));
} 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()
};
let running = probe_running_instance(descriptor.id.as_str());
LocalHarness {
gateway: gateway_health(
descriptor.id.as_str(),
installed,
running.as_ref(),
version.as_deref(),
),
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 },
running,
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(),
)
});
}
if self.runtimes.is_empty()
&& matches!(
request.operation,
SdkOperation::Input
| SdkOperation::Interrupt
| SdkOperation::Steer
| SdkOperation::Respond
| SdkOperation::Close
)
{
return Err(SdkError::unsupported(request.operation));
}
let method = request
.operation
.method()
.ok_or_else(|| SdkError::unsupported(request.operation))?;
let result = match request.operation {
SdkOperation::Discover
| SdkOperation::Load
| SdkOperation::Export
| SdkOperation::ProfilesList
| SdkOperation::ProfilesGet
| SdkOperation::ProfilesCreate
| SdkOperation::ProfilesDelete
| SdkOperation::SkillsList
| SdkOperation::SkillsInstall
| SdkOperation::SkillsRemove
| SdkOperation::ChannelsList
| SdkOperation::RoutesList
| SdkOperation::TriggersList
| SdkOperation::ChannelsStatus
| SdkOperation::MemoryShow
| SdkOperation::MemorySearch
| SdkOperation::JobsList
| SdkOperation::JobsGet
| SdkOperation::JobsCreate
| SdkOperation::JobsUpdate
| SdkOperation::JobsPause
| SdkOperation::JobsResume
| SdkOperation::JobsRun
| SdkOperation::JobsDelete
| SdkOperation::RunsList
| SdkOperation::RunsGet
| SdkOperation::ApprovalsList
| SdkOperation::WorldLoad
| SdkOperation::WorldSave
| SdkOperation::WorldCompile
| SdkOperation::WorldDecompile
| SdkOperation::WorldImport
| SdkOperation::WorldExport => self.call(method, request.params),
SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
SdkOperation::Start
| SdkOperation::Resume
| SdkOperation::Input
| SdkOperation::Interrupt
| SdkOperation::Steer
| SdkOperation::Respond
| SdkOperation::Close => self.runtime_call(method, request.params).await,
SdkOperation::SessionsNew => {
self.mutate_session(crate::SessionVerb::New, request.params)
.await
}
SdkOperation::SessionsReset => {
self.mutate_session(crate::SessionVerb::Reset, request.params)
.await
}
SdkOperation::SessionsArchive => {
self.mutate_session(crate::SessionVerb::Archive, request.params)
.await
}
SdkOperation::SessionsDelete => {
self.mutate_session(crate::SessionVerb::Delete, 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::OpenClaw => "openclaw",
SessionSource::Hermes => "hermes",
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 ActivitySubscribeParams {
locators: Vec<SessionLocator>,
#[serde(default)]
homes: crate::HarnessHomes,
}
#[derive(Deserialize)]
struct MessageSessionParams {
locator: SessionLocator,
text: String,
#[serde(default)]
homes: crate::HarnessHomes,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct HarnessSettingsParams {
harness: String,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct ConfigureHarnessParams {
harness: String,
#[serde(default)]
changes: Vec<crate::HarnessSettingChange>,
#[serde(default)]
expected_revision: Option<String>,
}
fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
Ok(report) => (
serde_json::to_value(report).unwrap_or(Value::Null),
Value::Null,
),
Err(error) => (
Value::Null,
Value::String(format!(
"Supercode could not inspect Claude Code inbound controls: {error}"
)),
),
}
}
#[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()
),
},
});
}
let (inbound_controls, inbound_controls_error) =
claude_inbound_controls_or_error(¶ms.homes);
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,
},
"inbound_controls": inbound_controls,
"inbound_controls_error": inbound_controls_error,
}),
Err(refusal) => json!({
"delivered_to_bus": false,
"refusal": {"reason": refusal.reason.as_str(), "message": refusal.message},
"inbound_controls": inbound_controls,
"inbound_controls_error": inbound_controls_error,
}),
}
}
#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
struct FollowedSource {
harness: String,
session_id: String,
reported: Option<String>,
}
#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
struct ActivitySubscription {
locators: Vec<SessionLocator>,
homes: crate::HarnessHomes,
reported: BTreeMap<(String, String), crate::SessionActivity>,
}
fn peers_for_descriptors(
descriptors: &[SessionDescriptor],
homes: &HarnessHomes,
) -> Vec<crate::claude_peer::ClaudePeerSession> {
if descriptors
.iter()
.any(|session| session.locator.harness.as_str() == HarnessId::CLAUDE_CODE)
{
crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
} else {
Vec::new()
}
}
fn live_descriptor_value(
session: &SessionDescriptor,
peers: &[crate::claude_peer::ClaudePeerSession],
) -> std::result::Result<Value, ServiceError> {
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 value.get("live_endpoint").is_none() {
if let Some(peer) = peers.iter().find(|peer| {
session.locator.harness.as_str() == HarnessId::CLAUDE_CODE
&& peer.session_id == session.locator.session_id
}) {
value["live_endpoint"] = json!(peer.endpoint().as_str());
}
}
Ok(value)
}
fn live_index_changes(
changes: Vec<crate::session_index::SessionIndexChange>,
homes: &HarnessHomes,
) -> std::result::Result<Vec<Value>, ServiceError> {
use crate::session_index::SessionIndexChange;
let has_claude = changes.iter().any(|change| match change {
SessionIndexChange::Added { descriptor } | SessionIndexChange::Updated { descriptor } => {
descriptor.locator.harness.as_str() == HarnessId::CLAUDE_CODE
}
SessionIndexChange::Removed { .. } => false,
});
let peers = if has_claude {
crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
} else {
Vec::new()
};
changes
.into_iter()
.map(|change| match change {
SessionIndexChange::Added { descriptor } => Ok(json!({
"kind": "added",
"descriptor": live_descriptor_value(&descriptor, &peers)?,
})),
SessionIndexChange::Updated { descriptor } => Ok(json!({
"kind": "updated",
"descriptor": live_descriptor_value(&descriptor, &peers)?,
})),
SessionIndexChange::Removed { key } => Ok(json!({
"kind": "removed",
"key": key,
})),
})
.collect()
}
fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
use crate::{SessionPresence, SessionTurnState};
match (activity.presence, activity.turn) {
(SessionPresence::Persisted, _) => None,
(SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
(SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
(SessionPresence::Running, SessionTurnState::Unknown)
if activity.evidence.native_state.is_none() =>
{
None
}
(SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
}
}
#[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(Deserialize)]
struct WorldLoadParams {
root: PathBuf,
#[serde(default)]
flavor: crate::world_doors::HomeFlavor,
}
#[derive(Deserialize)]
struct WorldSaveParams {
root: PathBuf,
world: crate::world::World,
#[serde(default)]
vault: BTreeMap<String, String>,
}
#[derive(Deserialize)]
struct WorldCompileParams {
from: crate::world_doors::WorldHarness,
home: PathBuf,
}
#[derive(Deserialize)]
struct WorldDecompileParams {
to: crate::world_doors::WorldHarness,
world: crate::world::World,
source: PathBuf,
#[serde(default)]
source_flavor: crate::world_doors::SourceFlavor,
dest: PathBuf,
#[serde(default)]
vault: BTreeMap<String, String>,
}
#[derive(Deserialize)]
struct WorldImportParams {
from: crate::world_doors::WorldHarness,
home: PathBuf,
into: PathBuf,
}
#[derive(Deserialize)]
struct WorldExportParams {
to: crate::world_doors::WorldHarness,
root: PathBuf,
dest: PathBuf,
}
#[derive(Deserialize)]
struct JobsGetParams {
harness: String,
id: String,
#[serde(default)]
homes: crate::HarnessHomes,
}
fn mutate_job(
verb: crate::jobs_control::JobVerb,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
}
fn mutate_skill(
verb: crate::skills_control::SkillVerb,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
if !crate::skills_control::supports_skill_control(&mutation.harness) {
return Err(ServiceError::UnsupportedAction(format!(
"`{}` has no skills root supercode reads; `skills.{}` is supported for: {}",
mutation.harness,
verb.as_str(),
crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
)));
}
let outcome =
crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
}
fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
match error {
crate::skills_control::SkillControlError::Unsupported(message) => {
ServiceError::UnsupportedAction(message)
}
crate::skills_control::SkillControlError::Invalid(message) => {
ServiceError::InvalidParams(message)
}
crate::skills_control::SkillControlError::Failed(message) => {
ServiceError::Operation(message)
}
}
}
fn mutate_profile(
verb: crate::profiles_control::ProfileVerb,
params: Value,
) -> std::result::Result<Value, ServiceError> {
let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
let outcome =
crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
}
fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
match error {
crate::profiles_control::ProfileControlError::Unsupported(message) => {
ServiceError::UnsupportedAction(message)
}
crate::profiles_control::ProfileControlError::Invalid(message) => {
ServiceError::InvalidParams(message)
}
crate::profiles_control::ProfileControlError::Failed(message) => {
ServiceError::Operation(message)
}
}
}
fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
match error {
crate::jobs_control::JobControlError::Unsupported(message) => {
ServiceError::UnsupportedAction(message)
}
crate::jobs_control::JobControlError::Invalid(message) => {
ServiceError::InvalidParams(message)
}
crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
}
}
fn session_control_error(error: crate::SessionControlError) -> ServiceError {
match error {
crate::SessionControlError::Unsupported(message) => {
ServiceError::UnsupportedAction(message)
}
crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
}
}
fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
if crate::jobs::supports_jobs(harness) {
return Ok(());
}
Err(ServiceError::UnsupportedAction(format!(
"`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
crate::jobs::JOB_HARNESSES.join(", ")
)))
}
#[derive(Deserialize)]
struct RunsGetParams {
harness: String,
id: String,
#[serde(default)]
homes: crate::HarnessHomes,
}
fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
if crate::runs::supports_runs(harness) {
return Ok(());
}
Err(ServiceError::UnsupportedAction(format!(
"`{harness}` keeps no run store; `{verb}` is supported for: {}",
crate::runs::RUN_HARNESSES.join(", ")
)))
}
#[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(Deserialize)]
struct HarnessAuthenticationParams {
harness: HarnessId,
}
#[derive(Deserialize)]
struct BeginHarnessAuthenticationParams {
harness: HarnessId,
#[serde(default = "local_browser_authentication_environment")]
environment: crate::HarnessAuthenticationEnvironment,
#[serde(default)]
method: Option<crate::HarnessAuthenticationMethodId>,
#[serde(default)]
cwd: Option<PathBuf>,
}
fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
crate::HarnessAuthenticationEnvironment::LocalBrowser
}
#[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(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum GatewayState {
Up,
Down,
Unknown,
}
#[derive(Debug, Clone, Serialize)]
pub struct GatewayHealth {
pub state: GatewayState,
#[serde(skip_serializing_if = "Option::is_none")]
pub endpoint: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
pub evidence: String,
pub checked_at_ms: u64,
}
fn openclaw_gateway_endpoint(home: &Path) -> String {
let config_path = home.join(".openclaw/openclaw.json");
let gateway = std::fs::read_to_string(&config_path)
.ok()
.and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
.and_then(|config| config.get("gateway").cloned());
if let Some(url) = gateway
.as_ref()
.and_then(|gateway| gateway.get("url"))
.and_then(serde_json::Value::as_str)
{
return url.to_string();
}
let port = gateway
.as_ref()
.and_then(|gateway| gateway.get("port"))
.and_then(serde_json::Value::as_u64)
.unwrap_or(18789);
format!("ws://127.0.0.1:{port}")
}
fn hermes_gateway_status() -> Option<(GatewayState, String)> {
let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
let output = std::process::Command::new(&program)
.args(["gateway", "status"])
.stdin(std::process::Stdio::null())
.output()
.ok()?;
let text = format!(
"{}{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
let verdict = text.lines().find_map(|line| {
let l = line.trim();
if l.contains("supervised by launchd (PID")
|| l.contains("supervised by systemd (PID")
|| l.contains("Gateway is running")
|| l.contains("process is running")
{
Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
} else if l.contains("not running") || l.contains("not installed") {
Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
} else {
None
}
});
verdict
}
fn gateway_health(
id: &str,
installed: bool,
running: Option<&RunningInstance>,
version: Option<&str>,
) -> GatewayHealth {
let checked_at_ms = now_epoch_ms();
let home = std::env::var_os("HOME").map(PathBuf::from);
match id {
HarnessId::HERMES | HarnessId::OPENCLAW => {
let endpoint = (id == HarnessId::OPENCLAW)
.then(|| home.as_deref().map(openclaw_gateway_endpoint))
.flatten();
let (state, evidence) = match running {
Some(instance) => (GatewayState::Up, instance.evidence.clone()),
None if !installed => (
GatewayState::Unknown,
format!("`{id}` is not installed; no gateway to probe"),
),
None if id == HarnessId::HERMES => match hermes_gateway_status() {
Some((state, evidence)) => (state, evidence),
None => (
GatewayState::Down,
"no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
),
},
None => (
GatewayState::Down,
format!(
"no TCP listener at {}",
endpoint.as_deref().unwrap_or("the gateway endpoint")
),
),
};
GatewayHealth {
state,
endpoint,
version: version.map(str::to_string),
evidence,
checked_at_ms,
}
}
HarnessId::ORCHESTRATOR => {
let root = crate::HarnessHomes::default().orchestrator;
let (state, evidence) = match crate::orchestrator::read_lease(&root) {
Some(lease) if crate::orchestrator::pid_is_live(lease.pid) => (
GatewayState::Up,
format!(
"`{}` names pid {} (started {}), which is live",
crate::orchestrator::lock_path(&root).display(),
lease.pid,
lease.started_at
),
),
Some(lease) => (
GatewayState::Down,
format!(
"stale lease `{}`: pid {} is gone",
crate::orchestrator::lock_path(&root).display(),
lease.pid
),
),
None => (
GatewayState::Down,
format!(
"no lease at `{}`; `supercode orchestrator start` writes one",
crate::orchestrator::lock_path(&root).display()
),
),
};
GatewayHealth {
state,
endpoint: None,
version: version.map(str::to_string),
evidence,
checked_at_ms,
}
}
_ => GatewayHealth {
state: GatewayState::Unknown,
endpoint: None,
version: version.map(str::to_string),
evidence: format!("`{id}` runs per session, not as a gateway"),
checked_at_ms,
},
}
}
#[derive(Debug, Clone, Serialize)]
struct RunningInstance {
method: RunningInstanceMethod,
evidence: String,
checked_at_ms: u64,
}
#[derive(Debug, Clone, Copy, Serialize)]
#[serde(rename_all = "snake_case")]
enum RunningInstanceMethod {
GatewayConnect,
StoreWalActivity,
}
fn now_epoch_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|elapsed| elapsed.as_millis() as u64)
.unwrap_or(0)
}
fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
let config_path = home.join(".openclaw/openclaw.json");
let text = std::fs::read_to_string(&config_path).ok();
let gateway = text
.as_deref()
.and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
.and_then(|config| config.get("gateway").cloned());
let address = gateway
.as_ref()
.and_then(|gateway| gateway.get("url"))
.and_then(serde_json::Value::as_str)
.and_then(|url| {
url.split("://").nth(1).map(|rest| {
rest.trim_end_matches('/')
.split('/')
.next()
.unwrap_or(rest)
.to_string()
})
})
.unwrap_or_else(|| {
let port = gateway
.as_ref()
.and_then(|gateway| gateway.get("port"))
.and_then(serde_json::Value::as_u64)
.unwrap_or(18789);
format!("127.0.0.1:{port}")
});
let reachable = std::net::TcpStream::connect_timeout(
&address.parse().ok()?,
std::time::Duration::from_millis(400),
)
.is_ok();
reachable.then(|| RunningInstance {
method: RunningInstanceMethod::GatewayConnect,
evidence: format!(
"gateway endpoint {address} accepted a TCP connect (from {})",
config_path.display()
),
checked_at_ms: now_epoch_ms(),
})
}
fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
let wal = home.join(".hermes/state.db-wal");
let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
let age_ms = std::time::SystemTime::now()
.duration_since(modified)
.map(|age| age.as_millis() as u64)
.unwrap_or(u64::MAX);
(age_ms <= max_wal_age_ms).then(|| RunningInstance {
method: RunningInstanceMethod::StoreWalActivity,
evidence: format!(
"{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
wal.display()
),
checked_at_ms: now_epoch_ms(),
})
}
fn probe_running_instance(id: &str) -> Option<RunningInstance> {
let home = std::env::var_os("HOME").map(PathBuf::from)?;
match id {
HarnessId::OPENCLAW => probe_openclaw_running(&home),
HarnessId::HERMES => probe_hermes_running(&home, 300_000),
_ => None,
}
}
#[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,
#[serde(skip_serializing_if = "Option::is_none")]
running: Option<RunningInstance>,
gateway: GatewayHealth,
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,
#[serde(default)]
mcp_servers: Vec<crate::McpServerLaunch>,
}
#[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,
#[serde(default)]
image_urls: Vec<String>,
}
const MAX_RUNTIME_IMAGES: usize = 4;
const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
if image_urls.len() > MAX_RUNTIME_IMAGES {
return Err(ServiceError::InvalidParams(format!(
"a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
)));
}
let mut total = 0usize;
for url in &image_urls {
if !(url.starts_with("data:image/")
|| url.starts_with("https://")
|| url.starts_with("http://"))
{
return Err(ServiceError::InvalidParams(
"runtime images must be image data URLs or HTTP(S) URLs".into(),
));
}
if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
return Err(ServiceError::InvalidParams(format!(
"one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
)));
}
total = total.saturating_add(url.len());
}
if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
return Err(ServiceError::InvalidParams(format!(
"runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
)));
}
Ok(image_urls)
}
#[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) {
if crate::support::self_sandbox_supported() {
arguments.extend(["--sandbox".into(), "workspace".into()]);
}
arguments.push("--always-approve".into());
}
arguments.extend(["--resume".into(), session_id.into()]);
"grok"
}
HarnessId::CODEX => {
let cwd_key = serde_json::to_string(cwd.to_string_lossy().as_ref())
.expect("a filesystem path always serializes as JSON text");
arguments.extend([
"-c".into(),
"check_for_update_on_startup=false".into(),
"-c".into(),
format!("projects.{cwd_key}.trust_level=\"trusted\""),
]);
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 openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
let digest = blake3::hash(address.as_bytes()).to_hex();
let path = std::env::temp_dir().join(format!(
"supercode-openclaw-gateway-token-{}",
&digest.as_str()[..16]
));
#[cfg(unix)]
{
use std::io::Write;
use std::os::unix::fs::OpenOptionsExt;
let mut file = std::fs::OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.mode(0o600)
.open(&path)?;
file.write_all(secret.as_bytes())?;
}
#[cfg(not(unix))]
std::fs::write(&path, secret)?;
Ok(path)
}
fn open_connect_descriptor(
descriptor: &crate::HarnessSupportDescriptor,
home: &Path,
) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
let Some(connect) = &descriptor.runtime.connect_launch else {
return Err(ServiceError::InvalidParams(format!(
"harness `{}` has no registered connect-mode launch",
descriptor.id.as_str()
)));
};
let resolved = connect
.resolve(home)
.map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
match (descriptor.id.as_str(), connect.protocol.as_str()) {
(HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
if let Some(token) = resolved.auth {
backend = backend.with_bearer(token);
}
Ok(Box::new(backend))
}
(HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
let mut env = BTreeMap::new();
let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
if let Some(token) = resolved.auth {
let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
.map_err(|error| {
ServiceError::UnsupportedAction(format!(
"could not stage the gateway credential for the bridge: {error}"
))
})?;
arguments.push("--token-file".into());
arguments.push(token_path.to_string_lossy().into_owned());
env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
}
let program = descriptor
.runtime
.default_launch
.as_ref()
.map(|launch| launch.program.clone())
.unwrap_or_else(|| "openclaw".into());
let launch = RuntimeLaunch {
program,
arguments,
env,
};
Ok(Box::new(
crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
.with_resume_support(descriptor.runtime.capabilities.resume_session),
))
}
_ => Err(ServiceError::UnsupportedAction(format!(
"connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
descriptor.id.as_str(),
connect.protocol
))),
}
}
fn registry_connect_descriptor(
params: &RuntimeBackendParams,
) -> Option<crate::HarnessSupportDescriptor> {
if params.launch.is_some() || params.base_url.is_some() {
return None;
}
harness_support_registry()
.harnesses
.into_iter()
.find(|descriptor| descriptor.id == params.harness)
.filter(|descriptor| descriptor.runtime.connect_launch.is_some())
}
fn service_home() -> std::result::Result<PathBuf, ServiceError> {
std::env::var_os("HOME").map(PathBuf::from).ok_or_else(|| {
ServiceError::UnsupportedAction(
"connect-mode launches need HOME to locate the harness config".into(),
)
})
}
fn runtime_backend(
params: &RuntimeBackendParams,
) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
if let Some(descriptor) = registry_connect_descriptor(params) {
return open_connect_descriptor(&descriptor, &service_home()?);
}
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: {
let mut arguments: Vec<String> = Vec::new();
if crate::support::self_sandbox_supported() {
arguments.extend(["--sandbox".into(), "workspace".into()]);
}
arguments.extend([
"--always-approve".into(),
"agent".into(),
"--no-leader".into(),
"stdio".into(),
]);
arguments
},
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::OPENCLAW => &[".openclaw/openclaw.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))
}
pub(crate) 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,
steer: 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()),
}
}
#[derive(Debug, Clone, Deserialize, Default)]
#[serde(default)]
struct MemoryRequest {
harness: Option<String>,
query: Option<String>,
profile: Option<String>,
session: Option<String>,
full: bool,
regex: bool,
cwd: Option<std::path::PathBuf>,
homes: crate::HarnessHomes,
}
fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
let request = decode::<MemoryRequest>(params)?;
let harness = request
.harness
.clone()
.ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
let to_service = |error: crate::memory::MemoryError| match error {
crate::memory::MemoryError::UnsupportedHarness { .. }
| crate::memory::MemoryError::SessionNotScoped { .. } => {
ServiceError::UnsupportedAction(error.to_string())
}
other => ServiceError::InvalidParams(other.to_string()),
};
match method {
"harness.v1.memory.show" => {
let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
harness,
profile: request.profile,
session: request.session,
full: request.full,
cwd: request.cwd,
homes: request.homes,
})
.map_err(to_service)?;
Ok(json!({
"schema": crate::memory::MEMORY_SCHEMA,
"documents": documents,
}))
}
"harness.v1.memory.search" => {
let query = request
.query
.ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
harness,
query,
profile: request.profile,
regex: request.regex,
cwd: request.cwd,
homes: request.homes,
})
.map_err(to_service)?;
Ok(json!({
"schema": crate::memory::MEMORY_SCHEMA,
"matches": matches,
}))
}
_ => Err(ServiceError::MethodNotFound),
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(default)]
struct ProfilesQuery {
harness: Option<String>,
name: Option<String>,
homes: crate::HarnessHomes,
}
impl Default for ProfilesQuery {
fn default() -> Self {
Self {
harness: None,
name: None,
homes: crate::HarnessHomes::default(),
}
}
}
fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
let query = decode::<ProfilesQuery>(params)?;
let to_service = |error: crate::profiles::ProfileError| match error {
crate::profiles::ProfileError::UnsupportedHarness { .. } => {
ServiceError::UnsupportedAction(error.to_string())
}
crate::profiles::ProfileError::NotFound { .. } => {
ServiceError::InvalidParams(error.to_string())
}
};
match method {
"harness.v1.profiles.list" => {
let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
.map_err(to_service)?;
Ok(json!({
"schema": crate::profiles::PROFILES_SCHEMA,
"profiles": profiles,
}))
}
"harness.v1.profiles.get" => {
let harness = query
.harness
.ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
let name = query
.name
.ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
let profile =
crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
Ok(json!({
"schema": crate::profiles::PROFILES_SCHEMA,
"profile": profile,
}))
}
_ => Err(ServiceError::MethodNotFound),
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(default)]
struct ChannelsQuery {
harness: Option<String>,
name: Option<String>,
homes: crate::HarnessHomes,
}
impl Default for ChannelsQuery {
fn default() -> Self {
Self {
harness: None,
name: None,
homes: crate::HarnessHomes::default(),
}
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(default)]
struct RoutesQuery {
harness: Option<String>,
profile: Option<String>,
homes: crate::HarnessHomes,
}
impl Default for RoutesQuery {
fn default() -> Self {
Self {
harness: None,
profile: None,
homes: crate::HarnessHomes::default(),
}
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(default)]
struct TriggersQuery {
harness: Option<String>,
homes: crate::HarnessHomes,
}
impl Default for TriggersQuery {
fn default() -> Self {
Self {
harness: None,
homes: crate::HarnessHomes::default(),
}
}
}
fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
let query = decode::<TriggersQuery>(params)?;
let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
.map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
Ok(json!({
"schema": crate::triggers::TRIGGERS_SCHEMA,
"triggers": triggers,
}))
}
fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
let query = decode::<RoutesQuery>(params)?;
let routes = crate::routes::list_routes(
&query.homes,
query.harness.as_deref(),
query.profile.as_deref(),
)
.map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
Ok(json!({
"schema": crate::routes::ROUTES_SCHEMA,
"routes": routes,
}))
}
fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
let query = decode::<ChannelsQuery>(params)?;
let to_service = |error: crate::channels::ChannelError| match error {
crate::channels::ChannelError::UnsupportedHarness { .. } => {
ServiceError::UnsupportedAction(error.to_string())
}
crate::channels::ChannelError::NotFound { .. } => {
ServiceError::InvalidParams(error.to_string())
}
};
match method {
"harness.v1.channels.list" => {
let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
.map_err(to_service)?;
Ok(json!({
"schema": crate::channels::CHANNELS_SCHEMA,
"channels": channels,
}))
}
"harness.v1.channels.status" => {
let harness = query
.harness
.ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
let name = query
.name
.ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
let channel = crate::channels::channel_status(&query.homes, &harness, &name)
.map_err(to_service)?;
Ok(json!({
"schema": crate::channels::CHANNELS_SCHEMA,
"channel": channel,
}))
}
_ => Err(ServiceError::MethodNotFound),
}
}
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;
#[test]
fn indexed_claude_descriptor_keeps_the_live_peer_address() {
let descriptor = SessionDescriptor {
locator: SessionLocator {
harness: HarnessId::new(HarnessId::CLAUDE_CODE),
session_id: "live-session".into(),
storage: StorageLocator::File {
path: PathBuf::from("/tmp/live-session.jsonl"),
},
},
cwd: Some(PathBuf::from("/project")),
title: None,
preview_candidates: Vec::new(),
latest_message_candidates: Vec::new(),
updated_at_ms: Some(1),
message_count: None,
model: None,
parent_session_id: None,
child_session_count: 0,
nouns: Default::default(),
};
let peer = crate::claude_peer::ClaudePeerSession {
pid: 42,
session_id: "live-session".into(),
cwd: Some(PathBuf::from("/project")),
name: "peer".into(),
socket_path: PathBuf::from("/tmp/peer.sock"),
status: Some(crate::claude_peer::ClaudePeerStatus::Busy),
updated_at_ms: Some(1),
version: Some("test".into()),
};
let value = live_descriptor_value(&descriptor, &[peer]).unwrap();
assert!(value["live_endpoint"]
.as_str()
.is_some_and(|endpoint| endpoint.starts_with("cc-peer:v1:42:peer:")));
}
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 hermes_store() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/hermes_home/state.db")
}
fn hermes_discovery(params: Value) -> Value {
let mut response =
HarnessSessionService::new().handle(request(1, "harness.v1.sessions.discover", params));
let store = hermes_store().display().to_string();
for session in response["result"]["sessions"]
.as_array_mut()
.expect("sessions array")
{
if session["locator"]["storage"]["path"] == json!(store) {
session["locator"]["storage"]["path"] = json!("<fixtures>/hermes_home/state.db");
}
session.as_object_mut().unwrap().remove("activity");
}
response["result"].take()
}
fn hermes_query() -> Value {
json!({
"harnesses": ["hermes"],
"homes": {"hermes": hermes_store()},
})
}
fn row<'a>(result: &'a Value, id: &str) -> &'a Value {
result["sessions"]
.as_array()
.expect("sessions array")
.iter()
.find(|session| session["locator"]["session_id"] == json!(id))
.unwrap_or_else(|| panic!("no discovered row for `{id}` in {result:#}"))
}
#[test]
fn orch6_discover_rows_carry_the_conversation_nouns() {
let result = hermes_discovery(hermes_query());
let dm = row(&result, "tg-dm-1");
assert_eq!(dm["trigger"], json!("channel"));
assert_eq!(dm["surface"]["platform"], json!("telegram"));
assert_eq!(dm["surface"]["kind"], json!("dm"));
assert_eq!(dm["surface"]["chat_id"], json!("123456"));
assert_eq!(dm["surface"]["participant_id"], json!("u1"));
assert_eq!(
dm["workspace"],
json!({"kind": "channel", "value": "telegram:123456"})
);
assert!(dm.get("profile").is_none(), "{dm:#}");
let fire = row(&result, "cron_job42_20260902_120000");
assert_eq!(fire["trigger"], json!("cron"));
assert_eq!(
fire["recurrence"],
json!({"job_id": "job42", "kind": "cron"})
);
assert_eq!(fire["workspace"]["kind"], json!("repo"));
let coder = row(&result, "tg-coder-1");
assert_eq!(coder["trigger"], json!("channel"));
assert_eq!(coder["profile"], json!("coder"));
assert_eq!(coder["surface"]["thread_id"], json!("55"));
assert_eq!(
coder["surface"]["key"],
json!("agent:coder:telegram:group:-100777:55")
);
assert_eq!(
coder["workspace"],
json!({"kind": "repo", "value": "/workspace/project"})
);
assert_eq!(
coder["cross_surface"],
json!({"state": "pending", "platform": "discord"})
);
let acp = row(&result, "cef97234-e8e8-428a-99ab-e8fff4e7e613");
assert_eq!(acp["trigger"], json!("human"));
assert!(acp.get("surface").is_none(), "{acp:#}");
assert_eq!(acp["workspace"], json!({"kind": "none"}));
}
#[test]
fn orch6_discover_filters_by_harness_and_profile() {
let mut params = hermes_query();
params["profile"] = json!("coder");
let result = hermes_discovery(params);
let ids: Vec<&str> = result["sessions"]
.as_array()
.expect("sessions array")
.iter()
.map(|session| session["locator"]["session_id"].as_str().unwrap())
.collect();
assert_eq!(ids, vec!["tg-coder-1"]);
let mut missing = hermes_query();
missing["profile"] = json!("nobody");
assert_eq!(hermes_discovery(missing)["sessions"], json!([]));
let elsewhere = json!({"harnesses": ["codex"], "homes": {"codex": hermes_store()}});
assert_eq!(hermes_discovery(elsewhere)["sessions"], json!([]));
}
#[test]
fn orch6_load_reports_the_same_nouns_as_discovery() {
let mut service = HarnessSessionService::new();
let loaded = service.handle(request(
1,
"harness.v1.sessions.load",
json!({"locator": {
"harness": "hermes",
"session_id": "tg-coder-1",
"storage": {"kind": "file", "path": hermes_store()},
}}),
));
let session = &loaded["result"]["session"];
let discovered = hermes_discovery(hermes_query());
let row = row(&discovered, "tg-coder-1");
for noun in [
"trigger",
"surface",
"profile",
"recurrence",
"cross_surface",
"workspace",
] {
assert_eq!(
session[noun],
row.get(noun).cloned().unwrap_or(Value::Null),
"`{noun}` disagrees between sessions.load and sessions.discover"
);
}
}
fn profile_fixture_homes() -> Value {
let fixtures = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
json!({
"hermes": fixtures.join("hermes_home/state.db"),
"openclaw": fixtures.join("openclaw_home"),
})
}
fn profile_row<'a>(response: &'a Value, harness: &str, name: &str) -> &'a Value {
response["result"]["profiles"]
.as_array()
.unwrap_or_else(|| panic!("no profiles array in {response}"))
.iter()
.find(|row| row["harness"] == harness && row["name"] == name)
.unwrap_or_else(|| panic!("no `{harness}` profile `{name}` in {response}"))
}
#[test]
fn profiles_list_reads_every_source_uniformly() {
let mut service = HarnessSessionService::new();
let response = service.handle(request(
1,
"harness.v1.profiles.list",
json!({"homes": profile_fixture_homes()}),
));
assert_eq!(
response["result"]["schema"],
crate::profiles::PROFILES_SCHEMA
);
let default = profile_row(&response, "hermes", "default");
assert_eq!(default["kind"], "hermes_profile");
assert_eq!(default["default"], true);
assert_eq!(default["routes"], 0);
assert_eq!(default["sessions"], 11);
assert_eq!(default["model"], "anthropic/claude-sonnet-4-5");
let coder = profile_row(&response, "hermes", "coder");
assert_eq!(coder["kind"], "hermes_profile");
assert_eq!(coder["default"], false);
assert_eq!(coder["routes"], 1, "gateway.profile_routes targets coder");
assert_eq!(coder["sessions"], 1, "state.db profile_name = 'coder'");
assert_eq!(coder["model"], "anthropic/claude-opus-4-8");
assert!(coder["home"]
.as_str()
.unwrap()
.ends_with("hermes_home/profiles/coder"));
let main = profile_row(&response, "openclaw", "main");
assert_eq!(main["kind"], "openclaw_agent");
assert_eq!(main["default"], true);
assert_eq!(main["routes"], 0);
assert_eq!(main["sessions"], 4);
assert_eq!(
main["model"],
Value::Null,
"`agents.defaults.model` is an install default, not this agent's pin"
);
let design = profile_row(&response, "openclaw", "design");
assert_eq!(design["default"], false);
assert_eq!(design["routes"], 1, "one binding names agentId `design`");
assert_eq!(design["sessions"], 0);
assert_eq!(design["model"], "anthropic/claude-opus-4-8");
let preset = profile_row(&response, "supercode", "supercode-default");
assert_eq!(preset["kind"], "preset");
assert_eq!(preset["default"], true);
assert_eq!(preset["home"], Value::Null);
assert_eq!(preset["routes"], Value::Null);
}
#[test]
fn profiles_list_reads_codex_profile_tables() {
let codex_home = std::env::temp_dir().join(format!(
"supercode-orch10-codex-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(codex_home.join("sessions")).unwrap();
std::fs::write(
codex_home.join("config.toml"),
"profile = \"review\"\n\n[profiles.review]\nmodel = \"gpt-5.1-codex\"\n\n[profiles.fast]\nmodel = \"gpt-5.1-codex-mini\"\n",
)
.unwrap();
let mut service = HarnessSessionService::new();
let response = service.handle(request(
1,
"harness.v1.profiles.list",
json!({"harness": "codex", "homes": {"codex": codex_home.join("sessions")}}),
));
let rows = response["result"]["profiles"].as_array().unwrap();
assert_eq!(rows.len(), 2, "{response}");
let review = profile_row(&response, "codex", "review");
assert_eq!(review["kind"], "codex_profile");
assert_eq!(review["default"], true);
assert_eq!(review["model"], "gpt-5.1-codex");
assert_eq!(review["home"], Value::Null);
assert_eq!(profile_row(&response, "codex", "fast")["default"], false);
let got = service.handle(request(
2,
"harness.v1.profiles.get",
json!({
"harness": "codex",
"name": "fast",
"homes": {"codex": codex_home.join("sessions")},
}),
));
assert_eq!(got["result"]["profile"]["model"], "gpt-5.1-codex-mini");
std::fs::remove_dir_all(&codex_home).ok();
}
#[test]
fn profiles_refuse_harnesses_without_the_concept() {
let mut service = HarnessSessionService::new();
let response = service.handle(request(
1,
"harness.v1.profiles.list",
json!({"harness": "claude-code"}),
));
assert_eq!(response["error"]["code"], -32020, "{response}");
let missing = service.handle(request(
2,
"harness.v1.profiles.get",
json!({
"harness": "hermes",
"name": "no-such-profile",
"homes": profile_fixture_homes(),
}),
));
assert_eq!(missing["error"]["code"], -32602, "{missing}");
}
#[test]
fn profiles_methods_are_advertised() {
let mut service = HarnessSessionService::new();
let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
let methods = response["result"]["methods"].as_array().unwrap();
for method in ["harness.v1.profiles.list", "harness.v1.profiles.get"] {
assert!(
methods.iter().any(|entry| entry == method),
"{method} is not advertised"
);
}
}
fn channel_row<'a>(response: &'a Value, harness: &str, name: &str) -> &'a Value {
response["result"]["channels"]
.as_array()
.unwrap_or_else(|| panic!("no channels array in {response}"))
.iter()
.find(|row| row["harness"] == harness && row["name"] == name)
.unwrap_or_else(|| panic!("no `{harness}` channel `{name}` in {response}"))
}
fn channels_list(harness: Option<&str>) -> Value {
let mut params = json!({"homes": profile_fixture_homes()});
if let Some(harness) = harness {
params["harness"] = json!(harness);
}
HarnessSessionService::new().handle(request(1, "harness.v1.channels.list", params))
}
#[test]
fn channels_list_reads_both_gateway_harnesses_uniformly() {
let response = channels_list(None);
assert_eq!(
response["result"]["schema"],
crate::channels::CHANNELS_SCHEMA
);
let telegram = channel_row(&response, "hermes", "telegram");
assert_eq!(telegram["kind"], "telegram");
assert_eq!(telegram["enabled"], true);
assert_eq!(telegram["configured"], true);
assert_eq!(telegram["sessions"], 2);
let api = channel_row(&response, "hermes", "api_server");
assert_eq!(api["configured"], true, "extra.key is a credential key");
assert_eq!(api["sessions"], 0);
let webhook = channel_row(&response, "hermes", "webhook");
assert_eq!(webhook["enabled"], false);
assert_eq!(webhook["configured"], true);
let linked = channel_row(&response, "openclaw", "slack/T0FIXTURE");
assert_eq!(linked["kind"], "slack");
assert_eq!(linked["account"], "T0FIXTURE");
assert_eq!(linked["enabled"], true);
assert_eq!(linked["configured"], true);
let unlinked = channel_row(&response, "openclaw", "slack/T1FIXTURE");
assert_eq!(unlinked["enabled"], false);
assert_eq!(
unlinked["configured"], false,
"an account with no credential key is not configured"
);
let telegram = channel_row(&response, "openclaw", "telegram");
assert_eq!(telegram["account"], "hermes-fixture-bot");
assert_eq!(telegram["configured"], true);
for row in response["result"]["channels"].as_array().unwrap() {
assert_eq!(row["status"], "unknown", "{row}");
}
}
#[test]
fn channels_rows_never_carry_a_fixture_secret() {
let secrets = [
"FAKE-TOKEN-DO-NOT-EMIT",
"FAKE-API-SERVER-KEY-DO-NOT-EMIT",
"FAKE-SLACK-BOT-TOKEN-DO-NOT-EMIT",
"FAKE-SLACK-APP-TOKEN-DO-NOT-EMIT",
"FAKE-TELEGRAM-TOKEN-DO-NOT-EMIT",
];
let fixtures = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
let raw = format!(
"{}{}",
std::fs::read_to_string(fixtures.join("hermes_home/config.yaml")).unwrap(),
std::fs::read_to_string(fixtures.join("openclaw_home/openclaw.json")).unwrap(),
);
for secret in secrets {
assert!(raw.contains(secret), "fixture no longer holds `{secret}`");
}
let emitted = serde_json::to_string(&channels_list(None)["result"]).unwrap();
for secret in secrets {
assert!(
!emitted.contains(secret),
"`{secret}` leaked into a channel row: {emitted}"
);
}
for row in channels_list(None)["result"]["channels"]
.as_array()
.unwrap()
{
for key in row.as_object().unwrap().keys() {
let key = key.to_ascii_lowercase();
assert!(
!["token", "key", "secret", "password", "credential"]
.iter()
.any(|marker| key.ends_with(marker)),
"`{key}` is a credential-shaped field on a channel row"
);
}
}
}
#[test]
fn channels_status_reads_one_row_by_name() {
let mut service = HarnessSessionService::new();
let got = service.handle(request(
1,
"harness.v1.channels.status",
json!({
"harness": "openclaw",
"name": "slack/T0FIXTURE",
"homes": profile_fixture_homes(),
}),
));
assert_eq!(got["result"]["channel"]["kind"], "slack");
assert_eq!(got["result"]["channel"]["account"], "T0FIXTURE");
assert_eq!(got["result"]["channel"]["status"], "unknown");
let missing = service.handle(request(
2,
"harness.v1.channels.status",
json!({
"harness": "openclaw",
"name": "no-such-channel",
"homes": profile_fixture_homes(),
}),
));
assert_eq!(missing["error"]["code"], -32602, "{missing}");
}
#[test]
fn channels_refuse_harnesses_without_the_concept() {
let response = channels_list(Some("claude-code"));
assert_eq!(response["error"]["code"], -32020, "{response}");
let codex = channels_list(Some("codex"));
assert_eq!(codex["error"]["code"], -32020, "{codex}");
}
#[test]
fn channels_list_filters_by_harness() {
let response = channels_list(Some("openclaw"));
let rows = response["result"]["channels"].as_array().unwrap();
assert!(!rows.is_empty(), "{response}");
assert!(
rows.iter().all(|row| row["harness"] == "openclaw"),
"harness filter leaked: {response}"
);
}
#[test]
fn channels_methods_are_advertised() {
let mut service = HarnessSessionService::new();
let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
let methods = response["result"]["methods"].as_array().unwrap();
for method in ["harness.v1.channels.list", "harness.v1.channels.status"] {
assert!(
methods.iter().any(|entry| entry == method),
"{method} is not advertised"
);
}
}
#[test]
fn orch6_story_fixture_matches_the_live_discovery_response() {
let path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("../../sdk/ui/stories/fixtures/hermes-discovery.json");
let mut result = hermes_discovery(hermes_query());
result.as_object_mut().unwrap().remove("next_cursor");
let rendered = format!("{}\n", serde_json::to_string_pretty(&result).unwrap());
if std::env::var_os("SUPERCODE_UPDATE_FIXTURES").is_some() {
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, &rendered).unwrap();
}
let committed = std::fs::read_to_string(&path).unwrap_or_default();
assert_eq!(
committed, rendered,
"sdk/ui/stories/fixtures/hermes-discovery.json is stale — \
re-run with SUPERCODE_UPDATE_FIXTURES=1"
);
}
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"),
},
}
}
fn fixture_homes() -> Value {
let fixtures = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
json!({
"claude_code": fixtures.join("__absent__"),
"codex": fixtures.join("__absent__"),
"opencode": fixtures.join("__absent__"),
"pi": fixtures.join("__absent__"),
"agents": fixtures.join("__absent__"),
"hermes": fixtures.join("hermes_home"),
"openclaw": fixtures.join("openclaw_home"),
})
}
fn skills_rows(params: Value) -> Vec<Value> {
let response =
HarnessSessionService::new().handle(request(1, "harness.v1.skills.list", params));
assert!(response.get("error").is_none(), "{response:#}");
response["result"].as_array().cloned().unwrap_or_default()
}
#[test]
fn skills_list_reads_the_hermes_and_openclaw_roots() {
let fixtures = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
let rows = skills_rows(json!({
"homes": fixture_homes(),
"cwd": fixtures.join("hermes_home"),
}));
let arxiv = rows
.iter()
.find(|row| row["name"] == json!("arxiv-search"))
.unwrap_or_else(|| panic!("no arxiv row in {rows:#?}"));
assert_eq!(arxiv["harness"], json!(HarnessId::HERMES));
assert_eq!(arxiv["scope"], json!("user"));
assert_eq!(arxiv["version"], json!("1.4.0"));
assert!(arxiv["location"]
.as_str()
.unwrap()
.ends_with("hermes_home/skills/research/arxiv"));
let bare = rows
.iter()
.find(|row| row["name"] == json!("bare-skill"))
.unwrap_or_else(|| panic!("no bare-skill row in {rows:#?}"));
assert_eq!(bare["enabled"], json!(null));
assert!(bare.get("description").is_none());
let demo = rows
.iter()
.find(|row| row["name"] == json!("clawhub-demo"))
.unwrap_or_else(|| panic!("no clawhub-demo row in {rows:#?}"));
assert_eq!(demo["harness"], json!(HarnessId::OPENCLAW));
assert_eq!(demo["scope"], json!("managed"));
assert_eq!(demo["enabled"], json!(false));
}
#[test]
fn skills_list_filters_by_harness_and_scope() {
let fixtures = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
let hermes = skills_rows(json!({
"homes": fixture_homes(),
"cwd": fixtures.join("hermes_home"),
"harness": HarnessId::HERMES,
}));
assert!(!hermes.is_empty());
assert!(hermes
.iter()
.all(|row| row["harness"] == json!(HarnessId::HERMES)));
let managed = skills_rows(json!({
"homes": fixture_homes(),
"cwd": fixtures.join("openclaw_home"),
"harness": HarnessId::OPENCLAW,
"scope": "managed",
}));
assert_eq!(managed.len(), 1, "{managed:#?}");
assert_eq!(managed[0]["name"], json!("clawhub-demo"));
let bundled = skills_rows(json!({
"homes": fixture_homes(),
"cwd": fixtures.join("openclaw_home"),
"harness": HarnessId::OPENCLAW,
"scope": "bundled",
}));
assert!(bundled.is_empty(), "{bundled:#?}");
}
#[test]
fn skills_list_refuses_an_unknown_harness() {
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.skills.list",
json!({"harness": "not-a-harness", "homes": fixture_homes()}),
));
assert_eq!(response["error"]["code"], json!(-32020), "{response:#}");
assert!(response["error"]["message"]
.as_str()
.unwrap()
.contains("not-a-harness"));
}
#[test]
fn skills_list_is_an_advertised_method_and_sdk_operation() {
assert!(HARNESS_SERVICE_METHODS.contains(&"harness.v1.skills.list"));
assert_eq!(
SdkOperation::from_method("harness.v1.skills.list"),
Some(SdkOperation::SkillsList)
);
}
#[test]
fn skills_install_and_remove_are_advertised_methods_and_sdk_operations() {
assert!(HARNESS_SERVICE_METHODS.contains(&"harness.v1.skills.install"));
assert!(HARNESS_SERVICE_METHODS.contains(&"harness.v1.skills.remove"));
assert_eq!(
SdkOperation::from_method("harness.v1.skills.install"),
Some(SdkOperation::SkillsInstall)
);
assert_eq!(
SdkOperation::from_method("harness.v1.skills.remove"),
Some(SdkOperation::SkillsRemove)
);
}
#[test]
fn skills_install_and_remove_drive_the_directory_door() {
let root = std::env::temp_dir().join(format!(
"supercode-orch22-rpc-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let source = root.join("probe-src");
std::fs::create_dir_all(&source).unwrap();
std::fs::write(
source.join("SKILL.md"),
"---\nname: orch22-rpc\ndescription: a probe\n---\nbody\n",
)
.unwrap();
let homes = json!({
"claude_code": root.join("claude_home"),
"codex": root.join("__absent__"),
"opencode": root.join("__absent__"),
"pi": root.join("__absent__"),
"hermes": root.join("__absent__"),
"openclaw": root.join("__absent__"),
"agents": root.join("__absent__"),
});
let mut service = HarnessSessionService::new();
let installed = service.handle(request(
1,
"harness.v1.skills.install",
json!({
"harness": HarnessId::CLAUDE_CODE,
"source": source,
"scope": "user",
"cwd": root,
"homes": homes,
}),
));
let result = &installed["result"];
assert_eq!(result["name"], json!("orch22-rpc"), "{installed:#}");
assert_eq!(result["verb"], json!("install"));
assert!(result["ran"]
.as_str()
.is_some_and(|ran| ran.starts_with("cp -R ")));
assert_eq!(result["skill"]["scope"], json!("user"));
let removed = service.handle(request(
2,
"harness.v1.skills.remove",
json!({
"harness": HarnessId::CLAUDE_CODE,
"name": "orch22-rpc",
"scope": "user",
"cwd": root,
"homes": homes,
}),
));
assert_eq!(removed["result"]["removed"], json!(true), "{removed:#}");
assert!(!root.join("claude_home/skills/orch22-rpc").exists());
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn skills_remove_refuses_openclaw_at_the_pin() {
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.skills.remove",
json!({"harness": HarnessId::OPENCLAW, "name": "clawhub-demo"}),
));
assert_eq!(response["error"]["code"], json!(-32020), "{response:#}");
assert!(response["error"]["message"]
.as_str()
.unwrap()
.contains("no `skills remove` verb"));
}
#[test]
fn skills_install_refuses_a_harness_without_a_skills_root() {
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.skills.install",
json!({"harness": "not-a-harness", "source": "/tmp/x"}),
));
assert_eq!(response["error"]["code"], json!(-32020), "{response:#}");
assert!(response["error"]["message"]
.as_str()
.unwrap()
.contains("not-a-harness"));
}
fn memory_homes() -> Value {
let fixtures = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures");
json!({
"claude_code": fixtures.join("__absent__"),
"codex": fixtures.join("__absent__"),
"opencode": fixtures.join("__absent__"),
"pi": fixtures.join("__absent__"),
"grok": fixtures.join("__absent__"),
"gemini": fixtures.join("__absent__"),
"goose": fixtures.join("__absent__"),
"supercode": fixtures.join("__absent__"),
"hermes": fixtures.join("hermes_home/state.db"),
"openclaw": fixtures.join("openclaw_home"),
})
}
fn memory_call_ok(method: &str, params: Value, key: &str) -> Vec<Value> {
let response = HarnessSessionService::new().handle(request(1, method, params));
assert!(response.get("error").is_none(), "{response:#}");
assert_eq!(response["result"]["schema"], json!("supercode.memory.v1"));
response["result"][key]
.as_array()
.cloned()
.unwrap_or_default()
}
fn memory_documents(params: Value) -> Vec<Value> {
memory_call_ok("harness.v1.memory.show", params, "documents")
}
fn memory_matches(params: Value) -> Vec<Value> {
memory_call_ok("harness.v1.memory.search", params, "matches")
}
fn find_document<'a>(rows: &'a [Value], profile: &str, name: &str) -> &'a Value {
rows.iter()
.find(|row| row["profile"] == profile && row["name"] == name)
.unwrap_or_else(|| panic!("no `{profile}` document `{name}` in {rows:#?}"))
}
#[test]
fn memory_show_reads_the_hermes_profile_homes() {
let rows = memory_documents(json!({"harness": "hermes", "homes": memory_homes()}));
let notes = find_document(&rows, "default", "MEMORY.md");
assert_eq!(notes["harness"], "hermes");
assert_eq!(notes["scope"], "user");
assert!(notes["size"].as_u64().unwrap() > 0);
assert!(notes["updated_at"].is_string(), "{notes:#?}");
assert!(notes.get("content").is_none(), "{notes:#?}");
assert_eq!(notes["truncated"], true);
assert_eq!(notes["preview"].as_array().unwrap().len(), 5);
let user = find_document(&rows, "default", "USER.md");
assert_eq!(user["scope"], "user");
assert!(user["preview"]
.as_array()
.unwrap()
.iter()
.any(|line| line.as_str().unwrap().contains("neovim")));
let topic = find_document(&rows, "default", "memories/2026-09-01-notes.md");
assert!(topic["path"]
.as_str()
.unwrap()
.ends_with("hermes_home/memories/2026-09-01-notes.md"));
let coder = find_document(&rows, "coder", "MEMORY.md");
assert_eq!(coder["scope"], "profile");
assert!(coder["path"]
.as_str()
.unwrap()
.ends_with("hermes_home/profiles/coder/MEMORY.md"));
}
#[test]
fn memory_show_returns_bodies_only_under_full_and_narrows_by_profile() {
let rows = memory_documents(json!({
"harness": "hermes",
"profile": "coder",
"full": true,
"homes": memory_homes(),
}));
assert!(
rows.iter().all(|row| row["profile"] == "coder"),
"{rows:#?}"
);
let coder = find_document(&rows, "coder", "MEMORY.md");
assert!(coder["content"]
.as_str()
.expect("full returns the body")
.contains("anthropic/claude-opus-4-8"));
}
#[test]
fn memory_show_reads_the_openclaw_agent_workspaces() {
let rows = memory_documents(json!({"harness": "openclaw", "homes": memory_homes()}));
let main = find_document(&rows, "main", "MEMORY.md");
assert_eq!(main["scope"], "agent");
assert!(main["path"]
.as_str()
.unwrap()
.ends_with("openclaw_home/workspace/MEMORY.md"));
let topic = find_document(&rows, "main", "memory/2026-09-01-standup.md");
assert!(topic["path"]
.as_str()
.unwrap()
.ends_with("openclaw_home/workspace/memory/2026-09-01-standup.md"));
let design = find_document(&rows, "design", "MEMORY.md");
assert!(design["path"]
.as_str()
.unwrap()
.ends_with("openclaw_home/workspace-design/MEMORY.md"));
}
#[test]
fn memory_show_reads_a_claude_code_project_auto_memory_directory() {
let scratch = std::env::temp_dir().join(format!(
"supercode-orch12-cc-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let project = scratch.join("repo");
std::fs::create_dir_all(project.join(".git")).unwrap();
let worktree = project.join("crates/harness");
std::fs::create_dir_all(&worktree).unwrap();
let slug: String = project
.to_string_lossy()
.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '-' })
.collect();
let projects = scratch.join("claude/projects");
let memory = projects.join(&slug).join("memory");
std::fs::create_dir_all(&memory).unwrap();
std::fs::write(
memory.join("MEMORY.md"),
"# index\n- [build box](build-box.md) — the pinned harnesses\n",
)
.unwrap();
std::fs::write(
memory.join("build-box.md"),
"hermes 0.21.0 and openclaw 2026.7.1-2 are the pins\n",
)
.unwrap();
let mut homes = memory_homes();
homes["claude_code"] = json!(projects);
let rows = memory_documents(json!({
"harness": "claude-code",
"cwd": worktree,
"homes": homes,
}));
let index = find_document(&rows, &slug, "MEMORY.md");
assert_eq!(index["harness"], "claude-code");
assert_eq!(index["scope"], "project");
let topic = find_document(&rows, &slug, "build-box.md");
assert!(topic["preview"]
.as_array()
.unwrap()
.iter()
.any(|line| line.as_str().unwrap().contains("2026.7.1-2")));
let hits = memory_matches(json!({
"harness": "claude-code",
"query": "pinned harnesses",
"cwd": worktree,
"homes": homes,
}));
assert_eq!(hits.len(), 1, "{hits:#?}");
assert_eq!(hits[0]["name"], "MEMORY.md");
assert_eq!(hits[0]["line"], 2);
let _ = std::fs::remove_dir_all(&scratch);
}
#[test]
fn memory_show_resolves_the_default_workspace_without_an_openclaw_config() {
let state = std::env::temp_dir().join(format!(
"supercode-orch12-oc-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(state.join("agents/main/agent")).unwrap();
std::fs::create_dir_all(state.join("workspace")).unwrap();
std::fs::write(
state.join("workspace/MEMORY.md"),
"the gateway websocket needs credentials\n",
)
.unwrap();
let mut homes = memory_homes();
homes["openclaw"] = json!(state);
let rows = memory_documents(json!({"harness": "openclaw", "homes": homes}));
assert_eq!(rows.len(), 1, "{rows:#?}");
let row = find_document(&rows, "main", "MEMORY.md");
assert_eq!(row["scope"], "agent");
assert!(row["path"]
.as_str()
.unwrap()
.ends_with("workspace/MEMORY.md"));
let _ = std::fs::remove_dir_all(&state);
}
#[test]
fn memory_search_reports_hits_by_line_and_misses_as_empty() {
let hit = memory_matches(json!({
"harness": "hermes",
"query": "NEOVIM",
"homes": memory_homes(),
}));
assert_eq!(hit.len(), 1, "{hit:#?}");
assert_eq!(hit[0]["harness"], "hermes");
assert_eq!(hit[0]["name"], "USER.md");
assert_eq!(hit[0]["scope"], "user");
assert_eq!(hit[0]["line"], 5);
assert!(hit[0]["excerpt"].as_str().unwrap().contains("neovim"));
let regex = memory_matches(json!({
"harness": "hermes",
"query": "neo(vim|vi)",
"regex": true,
"homes": memory_homes(),
}));
assert_eq!(regex.len(), 1, "{regex:#?}");
let miss = memory_matches(json!({
"harness": "hermes",
"query": "no-memory-line-says-this",
"homes": memory_homes(),
}));
assert!(miss.is_empty(), "{miss:#?}");
}
#[test]
fn memory_refuses_harnesses_without_a_store_and_misplaced_session_scoping() {
for method in ["harness.v1.memory.show", "harness.v1.memory.search"] {
let response = HarnessSessionService::new().handle(request(
1,
method,
json!({"harness": "codex", "query": "anything", "homes": memory_homes()}),
));
assert_eq!(response["error"]["code"], json!(-32020), "{response:#}");
assert!(response["error"]["message"]
.as_str()
.unwrap()
.contains("codex"));
}
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.memory.show",
json!({"harness": "hermes", "session": "abc", "homes": memory_homes()}),
));
assert_eq!(response["error"]["code"], json!(-32020), "{response:#}");
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.memory.show",
json!({"homes": memory_homes()}),
));
assert_eq!(response["error"]["code"], json!(-32602), "{response:#}");
}
#[test]
fn memory_methods_are_advertised_and_map_to_sdk_operations() {
assert!(HARNESS_SERVICE_METHODS.contains(&"harness.v1.memory.show"));
assert!(HARNESS_SERVICE_METHODS.contains(&"harness.v1.memory.search"));
assert_eq!(
SdkOperation::from_method("harness.v1.memory.show"),
Some(SdkOperation::MemoryShow)
);
assert_eq!(
SdkOperation::from_method("harness.v1.memory.search"),
Some(SdkOperation::MemorySearch)
);
}
struct RequestingRuntime {
handle: RuntimeHandle,
events: std::collections::VecDeque<HarnessEvent>,
answered: std::sync::Arc<std::sync::Mutex<Vec<Value>>>,
}
#[async_trait]
impl RuntimeConnection for RequestingRuntime {
fn handle(&self) -> &RuntimeHandle {
&self.handle
}
async fn send_input(&mut self, _input: RuntimeInput) -> crate::Result<Option<String>> {
unreachable!("this runtime only raises requests")
}
async fn next_event(&mut self) -> crate::Result<Option<HarnessEvent>> {
match self.events.pop_front() {
Some(event) => Ok(Some(event)),
None => std::future::pending().await,
}
}
async fn interrupt(&mut self) -> crate::Result<()> {
Ok(())
}
async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
self.answered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(json!({"request_id": request_id, "response": response}));
Ok(())
}
async fn close(&mut self) -> crate::Result<()> {
Ok(())
}
}
fn requesting_runtime(
harness: &str,
events: Vec<HarnessEvent>,
answered: std::sync::Arc<std::sync::Mutex<Vec<Value>>>,
) -> Box<dyn RuntimeConnection> {
requesting_runtime_named(harness, "hermes-live-session", events, answered)
}
fn requesting_runtime_named(
harness: &str,
runtime_id: &str,
events: Vec<HarnessEvent>,
answered: std::sync::Arc<std::sync::Mutex<Vec<Value>>>,
) -> Box<dyn RuntimeConnection> {
Box::new(RequestingRuntime {
handle: RuntimeHandle {
harness: HarnessId::from(harness),
runtime_id: runtime_id.into(),
endpoint: RuntimeEndpoint::LocalProcess {
pid: None,
command: vec!["hermes-acp".into()],
protocol: "acp".into(),
},
},
events: events.into(),
answered,
})
}
fn permission_event(id: u64, title: &str) -> HarnessEvent {
HarnessEvent {
sequence: None,
kind: "session/request_permission".into(),
payload: json!({
"jsonrpc": "2.0",
"id": id,
"method": "session/request_permission",
"params": {
"sessionId": "hermes-live-session",
"toolCall": {"toolCallId": "call-1", "title": title, "kind": "execute"},
"options": [
{"optionId": "allow_once", "name": "Allow once", "kind": "allow_once"},
{"optionId": "allow_for_session", "name": "Allow for session", "kind": "allow_always"},
{"optionId": "deny", "name": "Deny", "kind": "reject_once"},
],
},
}),
}
}
fn approvals(service: &mut HarnessSessionService, params: Value) -> Value {
let response = service.handle(request(1, "harness.v1.approvals.list", params));
assert!(response.get("error").is_none(), "{response:#}");
response["result"].clone()
}
#[tokio::test]
async fn a_claude_code_permission_request_lists_and_resolves_on_the_uniform_door() {
let answered = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let mut service = HarnessSessionService::new();
service.runtimes.insert(
"runtime-cc".into(),
requesting_runtime_named(
HarnessId::CLAUDE_CODE,
"claude-live-session",
vec![HarnessEvent {
sequence: None,
kind: "control_request".into(),
payload: json!({
"type": "control_request",
"request_id": "053f8a2d-3445-4011-a259-4261b31c7326",
"request": {
"subtype": "can_use_tool",
"tool_name": "Bash",
"display_name": "Bash",
"input": {"command": "touch probe-artifact.txt"},
"tool_use_id": "toolu_mock_1",
},
}),
}],
answered.clone(),
),
);
let notifications = service.poll_runtimes().await;
assert_eq!(notifications.len(), 1, "{notifications:#?}");
let rows = approvals(&mut service, json!({"harness": HarnessId::CLAUDE_CODE}));
assert_eq!(rows.as_array().map(Vec::len), Some(1), "{rows:#}");
let row = &rows[0];
assert_eq!(row["id"], "runtime-cc/053f8a2d-3445-4011-a259-4261b31c7326");
assert_eq!(row["harness"], HarnessId::CLAUDE_CODE);
assert_eq!(row["status"], "pending");
assert_eq!(row["subject"], "Bash touch probe-artifact.txt");
assert_eq!(row["runtime_id"], "claude-live-session");
assert_eq!(
row["options"]
.as_array()
.unwrap()
.iter()
.map(|option| option["id"].as_str().unwrap())
.collect::<Vec<_>>(),
vec!["allow", "deny"],
);
let response = resolve(
&mut service,
json!({"id": row["id"], "decision": "allow_once"}),
)
.await;
assert!(response.get("error").is_none(), "{response:#}");
assert_eq!(response["result"]["option_id"], "allow");
assert_eq!(
answered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_slice(),
&[json!({
"request_id": "053f8a2d-3445-4011-a259-4261b31c7326",
"response": {"behavior": "allow"},
})],
);
assert_eq!(
approvals(&mut service, json!({"harness": HarnessId::CLAUDE_CODE}))
.as_array()
.map(Vec::len),
Some(0),
);
}
#[tokio::test]
async fn a_live_permission_request_lists_until_it_is_answered() {
let answered = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let mut service = HarnessSessionService::new();
service.runtimes.insert(
"runtime-1".into(),
requesting_runtime(
HarnessId::HERMES,
vec![permission_event(7, "rm -rf build")],
answered.clone(),
),
);
let notifications = service.poll_runtimes().await;
assert_eq!(notifications.len(), 1, "{notifications:#?}");
let rows = approvals(&mut service, json!({}));
assert_eq!(rows.as_array().map(Vec::len), Some(1), "{rows:#}");
let row = &rows[0];
assert_eq!(row["id"], "runtime-1/7");
assert_eq!(row["harness"], HarnessId::HERMES);
assert_eq!(row["kind"], "live");
assert_eq!(row["status"], "pending");
assert_eq!(row["subject"], "rm -rf build");
assert_eq!(row["session_id"], "hermes-live-session");
assert_eq!(row["runtime_id"], "hermes-live-session");
assert!(row["requested_at_ms"].as_i64().is_some(), "{row:#}");
assert!(
row["age_ms"].as_i64().is_some_and(|age| age >= 0),
"{row:#}"
);
assert_eq!(
row["options"]
.as_array()
.unwrap()
.iter()
.map(|option| option["id"].as_str().unwrap())
.collect::<Vec<_>>(),
vec!["allow_once", "allow_for_session", "deny"],
);
assert_eq!(
approvals(&mut service, json!({"harness": HarnessId::HERMES}))
.as_array()
.map(Vec::len),
Some(1),
);
assert_eq!(
approvals(&mut service, json!({"session": "some-other-session"}))
.as_array()
.map(Vec::len),
Some(0),
);
let response = service
.handle_async(request(
2,
"harness.v1.runtimes.respond",
json!({
"connection": "runtime-1",
"request_id": 7,
"response": {"outcome": {"outcome": "selected", "optionId": "allow_once"}},
}),
))
.await;
assert!(response.get("error").is_none(), "{response:#}");
assert_eq!(
answered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_slice(),
&[json!({
"request_id": 7,
"response": {"outcome": {"outcome": "selected", "optionId": "allow_once"}},
})],
);
let rows = approvals(&mut service, json!({}));
assert_eq!(rows.as_array().map(Vec::len), Some(0), "{rows:#}");
}
#[test]
fn queued_subagent_approvals_list_through_the_same_door() {
let queue = std::sync::Arc::new(std::sync::Mutex::new(vec![
crate::subagents::QueuedApproval {
child_agent_id: "child-7".into(),
tool: "shell".into(),
subject: Some("cargo publish --dry-run".into()),
queued_at_ms: 1,
outcome: None,
},
crate::subagents::QueuedApproval {
child_agent_id: "child-8".into(),
tool: "write_file".into(),
subject: None,
queued_at_ms: 2,
outcome: Some(crate::subagents::QueuedApprovalOutcome::Denied),
},
]));
let mut service = HarnessSessionService::new();
service.observe_subagent_approvals(queue);
let rows = approvals(&mut service, json!({}));
assert_eq!(rows.as_array().map(Vec::len), Some(2), "{rows:#}");
assert_eq!(rows[0]["id"], "supercode/subagent/child-7/1/0");
assert_eq!(rows[0]["harness"], HarnessId::SUPERCODE);
assert_eq!(rows[0]["status"], "pending");
assert_eq!(rows[0]["subject"], "shell cargo publish --dry-run");
assert_eq!(rows[1]["status"], "denied");
assert!(rows[1]["options"].as_array().unwrap().is_empty());
let only = approvals(&mut service, json!({"session": "child-8"}));
assert_eq!(only.as_array().map(Vec::len), Some(1), "{only:#}");
assert_eq!(only[0]["id"], "supercode/subagent/child-8/2/1");
}
#[test]
fn approvals_list_refuses_a_harness_that_cannot_carry_a_request() {
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.approvals.list",
json!({"harness": "not-a-harness"}),
));
assert_eq!(response["error"]["code"], json!(-32020), "{response:#}");
assert!(response["error"]["message"]
.as_str()
.unwrap()
.contains("not-a-harness"));
for harness in [HarnessId::CLAUDE_CODE, HarnessId::CODEX] {
let response = HarnessSessionService::new().handle(request(
1,
"harness.v1.approvals.list",
json!({"harness": harness}),
));
assert!(response.get("error").is_none(), "{harness}: {response:#}");
}
}
#[test]
fn approvals_list_is_an_advertised_method_and_an_observed_tier() {
assert!(HARNESS_SERVICE_METHODS.contains(&"harness.v1.approvals.list"));
assert_eq!(
SdkOperation::from_method("harness.v1.approvals.list"),
Some(SdkOperation::ApprovalsList)
);
let registry = harness_support_registry();
for id in [
HarnessId::HERMES,
HarnessId::OPENCLAW,
HarnessId::CODEX,
HarnessId::CLAUDE_CODE,
] {
let concept = registry
.harnesses
.iter()
.find(|harness| harness.id.as_str() == id)
.unwrap()
.orchestration
.concepts
.iter()
.find(|concept| concept.concept == "pending_request")
.unwrap();
assert_eq!(concept.observed, crate::ImplementationKind::BuiltIn, "{id}");
assert!(concept
.methods
.iter()
.any(|method| method == "harness.v1.approvals.list"));
}
}
async fn resolve(service: &mut HarnessSessionService, params: Value) -> Value {
service
.handle_async(request(3, "harness.v1.approvals.resolve", params))
.await
}
#[tokio::test]
async fn a_listed_row_resolves_with_one_uniform_decision_and_then_is_gone() {
let answered = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let mut service = HarnessSessionService::new();
service.runtimes.insert(
"runtime-1".into(),
requesting_runtime(
HarnessId::HERMES,
vec![permission_event(7, "rm -rf build")],
answered.clone(),
),
);
service.poll_runtimes().await;
let rows = approvals(&mut service, json!({}));
assert_eq!(rows[0]["id"], "runtime-1/7");
let response = resolve(
&mut service,
json!({"id": "runtime-1/7", "decision": "allow_once"}),
)
.await;
assert!(response.get("error").is_none(), "{response:#}");
assert_eq!(
response["result"],
json!({
"id": "runtime-1/7",
"decision": "allow_once",
"option_id": "allow_once",
"resolved": true,
}),
);
assert_eq!(
answered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_slice(),
&[json!({
"request_id": 7,
"response": {"outcome": {"outcome": "selected", "optionId": "allow_once"}},
})],
);
assert_eq!(
approvals(&mut service, json!({})).as_array().map(Vec::len),
Some(0),
);
let response = resolve(
&mut service,
json!({"id": "runtime-1/7", "decision": "allow_once"}),
)
.await;
assert_eq!(response["error"]["code"], json!(-32602), "{response:#}");
}
#[tokio::test]
async fn deny_selects_the_requests_own_reject_option() {
let answered = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let mut service = HarnessSessionService::new();
service.runtimes.insert(
"runtime-1".into(),
requesting_runtime(
HarnessId::HERMES,
vec![permission_event(11, "git push --force")],
answered.clone(),
),
);
service.poll_runtimes().await;
let response = resolve(
&mut service,
json!({"id": "runtime-1/11", "decision": "deny"}),
)
.await;
assert!(response.get("error").is_none(), "{response:#}");
assert_eq!(response["result"]["option_id"], "deny");
assert_eq!(
answered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)[0]["response"],
json!({"outcome": {"outcome": "selected", "optionId": "deny"}}),
);
assert_eq!(
approvals(&mut service, json!({})).as_array().map(Vec::len),
Some(0),
);
}
#[tokio::test]
async fn a_decision_the_request_does_not_offer_is_refused_with_the_offered_ones() {
let answered = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let mut service = HarnessSessionService::new();
let mut event = permission_event(3, "rm -rf build");
event.payload["params"]["options"] = json!([
{"optionId": "allow_once", "name": "Allow edit", "kind": "allow_once"},
{"optionId": "deny", "name": "Deny", "kind": "reject_once"},
]);
service.runtimes.insert(
"runtime-1".into(),
requesting_runtime(HarnessId::HERMES, vec![event], answered.clone()),
);
service.poll_runtimes().await;
let response = resolve(
&mut service,
json!({"id": "runtime-1/3", "decision": "allow_always"}),
)
.await;
assert_eq!(response["error"]["code"], json!(-32602), "{response:#}");
let message = response["error"]["message"].as_str().unwrap();
assert!(message.contains("allow_always"), "{message}");
assert!(message.contains("allow_once, deny"), "{message}");
assert!(answered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_empty());
assert_eq!(
approvals(&mut service, json!({})).as_array().map(Vec::len),
Some(1),
);
}
#[tokio::test]
async fn a_queued_subagent_row_is_refused_by_name_rather_than_silently_answered() {
let queue = std::sync::Arc::new(std::sync::Mutex::new(vec![
crate::subagents::QueuedApproval {
child_agent_id: "child-7".into(),
tool: "shell".into(),
subject: Some("cargo publish --dry-run".into()),
queued_at_ms: 1,
outcome: None,
},
]));
let mut service = HarnessSessionService::new();
service.observe_subagent_approvals(queue.clone());
let row = approvals(&mut service, json!({}))[0]["id"]
.as_str()
.unwrap()
.to_string();
assert_eq!(row, "supercode/subagent/child-7/1/0");
let response = resolve(&mut service, json!({"id": row, "decision": "allow_once"})).await;
assert_eq!(response["error"]["code"], json!(-32602), "{response:#}");
let message = response["error"]["message"].as_str().unwrap();
assert!(message.contains("queued subagent record"), "{message}");
assert!(message.contains("request"), "{message}");
assert!(queue
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)[0]
.outcome
.is_none());
}
#[tokio::test]
async fn an_unknown_row_and_a_missing_decision_are_both_named() {
let mut service = HarnessSessionService::new();
let response = resolve(
&mut service,
json!({"id": "runtime-9/4", "decision": "deny"}),
)
.await;
assert_eq!(response["error"]["code"], json!(-32602), "{response:#}");
assert!(response["error"]["message"]
.as_str()
.unwrap()
.contains("runtime-9/4"));
let response = resolve(&mut service, json!({"id": "runtime-9/4"})).await;
let message = response["error"]["message"].as_str().unwrap();
assert!(
message.contains("allow_once | allow_always | deny"),
"{message}"
);
let response = resolve(
&mut service,
json!({"id": "runtime-9/4", "decision": "deny", "option_id": "deny"}),
)
.await;
assert!(response["error"]["message"]
.as_str()
.unwrap()
.contains("not both"));
}
#[test]
fn approvals_resolve_is_an_advertised_method_and_a_controlled_tier() {
assert!(HARNESS_SERVICE_METHODS.contains(&"harness.v1.approvals.resolve"));
assert_eq!(
SdkOperation::from_method("harness.v1.approvals.resolve"),
Some(SdkOperation::ApprovalsResolve)
);
assert_eq!(
SdkOperation::ApprovalsResolve.action_name(),
"approvals_resolve"
);
let registry = harness_support_registry();
for id in [
HarnessId::HERMES,
HarnessId::OPENCLAW,
HarnessId::CODEX,
HarnessId::CLAUDE_CODE,
] {
let concept = registry
.harnesses
.iter()
.find(|harness| harness.id.as_str() == id)
.unwrap()
.orchestration
.concepts
.iter()
.find(|concept| concept.concept == "pending_request")
.unwrap();
assert_eq!(
concept.controlled,
crate::ImplementationKind::BuiltIn,
"{id}"
);
assert!(
concept
.methods
.iter()
.any(|method| method == "harness.v1.approvals.resolve"),
"{id}"
);
}
}
#[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(),
11
);
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());
}
#[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 codex = resume_launch(
HarnessId::CODEX,
"codex-session",
Path::new("/tmp/project"),
ResumePolicy::Yolo,
)
.unwrap_or_else(|_| panic!("Codex resume launch must be registered"));
assert_eq!(codex.program, "codex");
assert_eq!(
codex.arguments,
[
"-c",
"check_for_update_on_startup=false",
"-c",
"projects.\"/tmp/project\".trust_level=\"trusted\"",
"--dangerously-bypass-approvals-and-sandbox",
"--dangerously-bypass-hook-trust",
"resume",
"codex-session",
]
);
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);
}
#[test]
fn running_instances_are_detected_from_mock_gateway_and_active_wal() {
let home = connect_scratch_home("uni7-running");
assert!(probe_hermes_running(&home, 300_000).is_none());
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
std::fs::create_dir_all(home.join(".openclaw")).unwrap();
std::fs::write(
home.join(".openclaw/openclaw.json"),
format!(r#"{{"gateway": {{"mode": "local", "port": {port}, "auth": {{"mode": "token", "token": "t"}}}}}}"#),
)
.unwrap();
let running = probe_openclaw_running(&home).expect("listening gateway must be detected");
assert!(matches!(
running.method,
RunningInstanceMethod::GatewayConnect
));
assert!(running.evidence.contains(&format!("127.0.0.1:{port}")));
drop(listener);
let mut closed_detected = probe_openclaw_running(&home).is_some();
for _ in 0..3 {
if !closed_detected {
break;
}
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
drop(listener);
std::fs::write(
home.join(".openclaw/openclaw.json"),
format!(r#"{{"gateway": {{"mode": "local", "port": {port}, "auth": {{"mode": "token", "token": "t"}}}}}}"#),
)
.unwrap();
closed_detected = probe_openclaw_running(&home).is_some();
}
assert!(
!closed_detected,
"a closed gateway must not read as running"
);
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
std::fs::write(
home.join(".openclaw/openclaw.json"),
format!(r#"{{"gateway": {{"url": "ws://127.0.0.1:{port}", "auth": {{"mode": "token", "token": "t"}}}}}}"#),
)
.unwrap();
assert!(probe_openclaw_running(&home).is_some());
drop(listener);
std::fs::create_dir_all(home.join(".hermes")).unwrap();
let wal = home.join(".hermes/state.db-wal");
std::fs::write(&wal, b"wal").unwrap();
let running = probe_hermes_running(&home, 300_000).expect("fresh WAL must be detected");
assert!(matches!(
running.method,
RunningInstanceMethod::StoreWalActivity
));
assert!(running.evidence.contains("state.db-wal"));
let stale = std::time::SystemTime::now() - std::time::Duration::from_secs(3_600);
std::fs::File::options()
.append(true)
.open(&wal)
.unwrap()
.set_modified(stale)
.unwrap();
assert!(
probe_hermes_running(&home, 300_000).is_none(),
"a stale WAL (crash leftover) must not read as running"
);
}
fn connect_scratch_home(tag: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"supercode-connect-service-{tag}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
dir
}
async fn mock_opencode_endpoint() -> (String, tokio::sync::mpsc::UnboundedReceiver<String>) {
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let (request_sender, request_receiver) = tokio::sync::mpsc::unbounded_channel();
tokio::spawn(async move {
loop {
let Ok((mut stream, _)) = listener.accept().await else {
break;
};
let request_sender = request_sender.clone();
tokio::spawn(async move {
let (reader, mut writer) = stream.split();
let mut reader = BufReader::new(reader);
let mut request_line = String::new();
if reader.read_line(&mut request_line).await.unwrap_or(0) == 0 {
return;
}
let request_line = request_line.trim_end().to_string();
let mut authorization = String::new();
let mut content_length = 0usize;
loop {
let mut line = String::new();
if reader.read_line(&mut line).await.unwrap_or(0) == 0 {
return;
}
let line = line.trim_end();
if line.is_empty() {
break;
}
let lower = line.to_ascii_lowercase();
if let Some(value) = lower.strip_prefix("authorization:") {
authorization = value.trim().to_string();
}
if let Some(value) = lower.strip_prefix("content-length:") {
content_length = value.trim().parse().unwrap_or(0);
}
}
if content_length > 0 {
let mut body = vec![0u8; content_length];
let _ = reader.read_exact(&mut body).await;
}
let _ = request_sender.send(format!("{request_line} :: {authorization}"));
if request_line.starts_with("GET /event") {
let _ = writer
.write_all(
b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\n\r\n",
)
.await;
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
return;
}
let body = if request_line.starts_with("POST /session") {
r#"{"id":"mock-session"}"#
} else {
r#"{"status":"ok"}"#
};
let response = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
let _ = writer.write_all(response.as_bytes()).await;
});
}
});
(format!("http://{address}"), request_receiver)
}
fn connect_descriptor(protocol: &str) -> crate::HarnessSupportDescriptor {
crate::HarnessSupportDescriptor {
orchestration: Default::default(),
id: HarnessId::from(HarnessId::OPENCODE),
display_name: "OpenCode".into(),
native: crate::NativeSupport {
discover: crate::ImplementationKind::Absent,
load: crate::ImplementationKind::Absent,
follow: crate::ImplementationKind::Absent,
import: crate::ImplementationKind::Absent,
export: crate::ImplementationKind::Absent,
},
runtime: crate::RuntimeSupport {
implementation: crate::ImplementationKind::BuiltIn,
protocol: protocol.into(),
default_launch: None,
connect_launch: Some(crate::RuntimeConnectLaunch {
config_path: "~/opencode-tui.json".into(),
address_pointer: "/server/url".into(),
port_pointer: None,
default_address: None,
auth_pointer: Some("/server/token".into()),
protocol: protocol.into(),
}),
capabilities: crate::RuntimeCapabilities {
start_session: true,
resume_session: true,
attach_existing_process: true,
send_input: true,
stream_events: true,
interrupt: true,
steer: false,
respond_to_requests: true,
},
},
}
}
#[tokio::test]
async fn connect_mode_descriptor_opens_a_running_endpoint_with_config_sourced_auth() {
let (base_url, mut requests) = mock_opencode_endpoint().await;
let home = connect_scratch_home("open");
std::fs::write(
home.join("opencode-tui.json"),
format!(r#"{{"server": {{"url": "{base_url}", "token": "connect-secret"}}}}"#),
)
.unwrap();
let descriptor = connect_descriptor("opencode-http-sse");
let backend = open_connect_descriptor(&descriptor, &home).unwrap();
assert!(backend.capabilities().attach_existing_process);
let connection = backend
.start(crate::RuntimeStartRequest {
cwd: home.clone(),
launch: None,
mcp_servers: Vec::new(),
})
.await
.unwrap();
let handle = connection.handle();
assert_eq!(handle.runtime_id, "mock-session");
match &handle.endpoint {
crate::RuntimeEndpoint::Http {
base_url: endpoint, ..
} => assert_eq!(endpoint, &base_url),
other => panic!("connect mode must join the running endpoint, got {other:?}"),
}
let mut seen = Vec::new();
while let Ok(line) = requests.try_recv() {
seen.push(line);
}
assert!(seen
.iter()
.any(|line| line.starts_with("GET /global/health")
&& line.contains("bearer connect-secret")));
assert!(seen.iter().any(
|line| line.starts_with("POST /session") && line.contains("bearer connect-secret")
));
}
#[tokio::test]
async fn openclaw_connect_mode_attaches_lists_and_resumes_via_a_mock_bridge() {
let home = connect_scratch_home("openclaw");
std::fs::create_dir_all(home.join(".openclaw")).unwrap();
std::fs::write(
home.join(".openclaw/openclaw.json"),
r#"{"gateway": {"remote": {"url": "ws://127.0.0.1:19789"}, "auth": {"mode": "token", "token": "mock-gateway-token"}}}"#,
)
.unwrap();
let script = home.join("openclaw");
std::fs::write(
&script,
r#"#!/bin/sh
# Fake `openclaw acp` bridge: verify the connect-mode contract, then speak ACP.
[ "$1" = "acp" ] || { echo "unexpected argv: $*" >&2; exit 9; }
[ "$2" = "--url" ] && [ "$3" = "ws://127.0.0.1:19789" ] || { echo "missing --url: $*" >&2; exit 9; }
[ "$4" = "--token-file" ] || { echo "missing --token-file: $*" >&2; exit 9; }
[ "$(cat "$5")" = "mock-gateway-token" ] || { echo "token file wrong" >&2; exit 9; }
while IFS= read -r line; do
case "$line" in
*'"initialize"'*)
printf '%s
' '{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentCapabilities":{"loadSession":true,"sessionCapabilities":{"list":{},"resume":{}}},"agentInfo":{"name":"openclaw-acp","version":"2026.7.1-2"},"authMethods":[]}}' ;;
*'"session/resume"'*)
printf '%s
' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"agent:main:main"}}' ;;
*'"session/new"'*)
printf '%s
' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"agent:main:fresh"}}' ;;
*'"session/prompt"'*)
printf '%s
' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"agent:main:main","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"joined"}}}}'
printf '%s
' '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}' ;;
esac
done
"#,
)
.unwrap();
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap();
let mut descriptor = crate::harness_support_registry()
.harnesses
.into_iter()
.find(|harness| harness.id.as_str() == HarnessId::OPENCLAW)
.expect("openclaw must be registered");
descriptor
.runtime
.connect_launch
.as_mut()
.unwrap()
.config_path = "~/.openclaw/openclaw.json".into();
descriptor.runtime.default_launch.as_mut().unwrap().program =
script.to_string_lossy().into_owned();
let backend = open_connect_descriptor(&descriptor, &home).unwrap();
assert!(backend.capabilities().resume_session);
let joined = backend
.attach(crate::RuntimeAttachRequest {
runtime_id: "agent:main:main".into(),
cwd: Some(home.clone()),
launch: None,
})
.await;
let mut connection = joined.expect("mock bridge attach must succeed");
assert_eq!(connection.handle().runtime_id, "agent:main:main");
let turn = connection
.send_input(crate::RuntimeInput {
text: "hello".into(),
image_urls: Vec::new(),
})
.await;
assert!(turn.is_ok(), "prompt through the mock bridge: {turn:?}");
connection.close().await.unwrap();
}
#[tokio::test]
async fn connect_mode_fails_closed_without_a_protocol_client_or_config() {
let home = connect_scratch_home("fail");
std::fs::write(
home.join("opencode-tui.json"),
r#"{"server": {"url": "http://127.0.0.1:1", "token": "connect-secret"}}"#,
)
.unwrap();
let gateway_only = connect_descriptor("acp-v1-jsonrpc");
let Err(error) = open_connect_descriptor(&gateway_only, &home) else {
panic!("an ACP connect endpoint has no gateway client yet");
};
let message = format!("{error:?}");
assert!(message.contains("acp-v1-jsonrpc"));
assert!(!message.contains("connect-secret"));
let unreadable = connect_descriptor("opencode-http-sse");
let missing_home = connect_scratch_home("missing");
let Err(error) = open_connect_descriptor(&unreadable, &missing_home) else {
panic!("an unreadable connect config must fail closed");
};
let message = format!("{error:?}");
assert!(message.contains("opencode-tui.json"));
assert!(!message.contains("connect-secret"));
}
fn jobs_fixture_root() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures")
}
fn jobs_fixture_homes() -> Value {
let root = jobs_fixture_root();
json!({
"claude_code": root.join("claude_jobs_home/projects"),
"hermes": root.join("hermes_home/state.db"),
"openclaw": root.join("openclaw_home"),
})
}
fn jobs_list(params: Value) -> Value {
let mut service = HarnessSessionService::new();
service.handle(request(1, "harness.v1.jobs.list", params))
}
fn job_row<'a>(result: &'a Value, id: &str) -> &'a Value {
result["jobs"]
.as_array()
.expect("jobs is an array")
.iter()
.find(|job| job["id"] == id)
.unwrap_or_else(|| panic!("no job `{id}` in {result}"))
}
#[test]
fn gateway_health_derives_from_running_probe_and_install_state() {
let running = RunningInstance {
method: RunningInstanceMethod::GatewayConnect,
evidence: "gateway endpoint 127.0.0.1:18789 accepted a TCP connect".into(),
checked_at_ms: 1,
};
let up = gateway_health(
HarnessId::OPENCLAW,
true,
Some(&running),
Some("2026.7.1-2"),
);
assert_eq!(up.state, GatewayState::Up);
assert!(up.endpoint.as_deref().unwrap().starts_with("ws://"));
assert_eq!(up.version.as_deref(), Some("2026.7.1-2"));
let dir = std::env::temp_dir().join(format!("supercode-orch17-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let fake = dir.join("hermes");
let write_fake = |body: &str| {
std::fs::write(&fake, format!("#!/bin/sh\n{body}\n")).unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&fake, std::fs::Permissions::from_mode(0o755)).unwrap();
}
};
write_fake("echo '✗ Gateway service is not installed'");
crate::harness_command::TEST_PROGRAM_OVERRIDE.with(|slot| {
*slot.borrow_mut() = Some((
HarnessId::HERMES.to_string(),
fake.to_string_lossy().into_owned(),
))
});
let down = gateway_health(HarnessId::HERMES, true, None, None);
assert_eq!(down.state, GatewayState::Down, "{down:?}");
assert!(down.endpoint.is_none());
assert!(down.evidence.contains("not installed"));
write_fake("echo 'Launchd plist: /x/ai.hermes.gateway.plist'; echo '✓ Gateway is supervised by launchd (PID 4242)'");
let idle_but_up = gateway_health(HarnessId::HERMES, true, None, Some("0.21.0"));
assert_eq!(idle_but_up.state, GatewayState::Up, "{idle_but_up:?}");
assert!(idle_but_up.evidence.contains("PID 4242"));
write_fake("echo 'something unparseable'");
let no_verdict = gateway_health(HarnessId::HERMES, true, None, None);
assert_eq!(no_verdict.state, GatewayState::Down);
assert!(no_verdict.evidence.contains("no verdict"));
crate::harness_command::TEST_PROGRAM_OVERRIDE.with(|slot| *slot.borrow_mut() = None);
let absent = gateway_health(HarnessId::HERMES, false, None, None);
assert_eq!(absent.state, GatewayState::Unknown);
let core = gateway_health(HarnessId::CODEX, true, None, Some("0.144.4"));
assert_eq!(core.state, GatewayState::Unknown);
assert!(core.evidence.contains("per session"));
}
#[test]
fn triggers_list_reads_both_stores_and_never_emits_secrets() {
let response = triggers_list(json!({"homes": jobs_fixture_homes()}));
let rows = response["result"]["triggers"]
.as_array()
.expect("triggers")
.clone();
let hermes: Vec<&Value> = rows.iter().filter(|r| r["harness"] == "hermes").collect();
assert!(
hermes.iter().any(|r| r["name"] == "deploys"
&& r["route"] == "/webhooks/deploys"
&& r["kind"] == "webhook"),
"{rows:#?}"
);
let openclaw: Vec<&Value> = rows.iter().filter(|r| r["harness"] == "openclaw").collect();
assert!(openclaw
.iter()
.any(|r| r["name"] == "wake" && r["kind"] == "builtin_wake"));
assert!(openclaw.iter().any(|r| r["name"] == "gmail"
&& r["kind"] == "hook_mapping"
&& r["target"]["action"] == "agent"));
let rendered = response.to_string();
for secret in [
"FAKE-WEBHOOK-HMAC-DO-NOT-EMIT",
"FAKE-HOOK-TOKEN-DO-NOT-EMIT",
] {
assert!(!rendered.contains(secret), "{rendered}");
}
let refused =
triggers_list(json!({"harness": "claude-code", "homes": jobs_fixture_homes()}));
assert_eq!(refused["error"]["code"], -32020, "{refused}");
}
fn triggers_list(params: Value) -> Value {
let mut service = HarnessSessionService::new();
service.handle(request(1, "harness.v1.triggers.list", params))
}
#[test]
fn routes_list_reads_both_gateway_configs_and_flags_the_defaults() {
let response = routes_list(json!({"homes": jobs_fixture_homes()}));
let rows = response["result"]["routes"]
.as_array()
.expect("routes")
.clone();
let hermes: Vec<&Value> = rows.iter().filter(|r| r["harness"] == "hermes").collect();
assert_eq!(hermes.len(), 2, "{rows:#?}");
assert_eq!(hermes[0]["target"], "coder");
assert_eq!(hermes[0]["match"]["platform"], "slack");
assert_eq!(hermes[0]["match"]["chat_id"], "C0FIXTURE");
assert_eq!(hermes[0]["specificity"], 4);
assert_eq!(hermes[1]["default"], true);
let openclaw: Vec<&Value> = rows.iter().filter(|r| r["harness"] == "openclaw").collect();
assert!(
openclaw.iter().any(|r| r["target"] == "design"
&& r["match"]["platform"] == "slack"
&& r["specificity"] == 1),
"{openclaw:#?}"
);
assert!(openclaw.iter().any(|r| r["default"] == true));
let refused = routes_list(json!({"harness": "codex", "homes": jobs_fixture_homes()}));
assert_eq!(refused["error"]["code"], -32020, "{refused}");
}
fn routes_list(params: Value) -> Value {
let mut service = HarnessSessionService::new();
service.handle(request(1, "harness.v1.routes.list", params))
}
#[test]
fn jobs_list_projects_every_fixture_store_onto_the_uniform_row() {
let response = jobs_list(json!({"homes": jobs_fixture_homes()}));
let result = &response["result"];
let ids: Vec<&str> = result["jobs"]
.as_array()
.unwrap()
.iter()
.map(|job| job["id"].as_str().unwrap())
.collect();
assert_eq!(
ids,
vec![
"release-watch",
"toolu_wake_recheck",
"digest-15m",
"nightly-audit",
"coder-standup",
"ops-once-boot",
"85ad7832-896f-42be-af31-3e1ed2fbdc4b",
"8bb7d938-ca46-4a6d-90eb-c92331155566",
"cron_standup",
"cron_reindex",
],
"{result}"
);
let health = job_row(result, "85ad7832-896f-42be-af31-3e1ed2fbdc4b");
assert_eq!(health["harness"], "openclaw");
assert_eq!(health["schedule"]["kind"], "interval");
assert_eq!(health["schedule"]["minutes"], 10.0);
assert_eq!(health["session_target"], "isolated");
assert_eq!(health["payload"]["kind"], "prompt");
assert_eq!(health["payload"]["text"], "nightly health check");
assert_eq!(health["deliver"]["mode"], "announce");
assert_eq!(health["deliver"]["target"], "last");
assert_eq!(health["next_run_at"], "2026-09-03T06:52:26Z");
let digest = job_row(result, "8bb7d938-ca46-4a6d-90eb-c92331155566");
assert_eq!(digest["schedule"]["kind"], "cron");
assert_eq!(digest["schedule"]["expr"], "0 9 * * 1");
assert_eq!(digest["session_target"], "main");
assert_eq!(digest["payload"]["kind"], "system_event");
let cron = job_row(result, "release-watch");
assert_eq!(cron["harness"], "claude-code");
assert_eq!(cron["scope"], "session");
assert_eq!(cron["session_id"], "7c1d2e3f-4a5b-6c7d-8e9f-0a1b2c3d4e5f");
assert_eq!(cron["schedule"]["kind"], "cron");
assert_eq!(cron["schedule"]["expr"], "*/10 * * * *");
assert_eq!(cron["schedule"]["display"], "*/10 * * * *");
assert_eq!(cron["payload"]["kind"], "prompt");
assert_eq!(cron["recurring"], true);
assert_eq!(cron["deliver"]["target"], "session");
let wakeup = job_row(result, "toolu_wake_recheck");
assert_eq!(wakeup["payload"]["kind"], "wakeup");
assert_eq!(wakeup["schedule"]["kind"], "once");
assert_eq!(wakeup["recurring"], false);
assert_eq!(wakeup["state"], "pending");
let interval = job_row(result, "digest-15m");
assert_eq!(interval["harness"], "hermes");
assert_eq!(interval["scope"], "install");
assert_eq!(interval["profile"], Value::Null);
assert_eq!(interval["schedule"]["kind"], "interval");
assert_eq!(interval["schedule"]["minutes"], 15.0);
assert_eq!(interval["schedule"]["display"], "every 15 min");
assert_eq!(interval["deliver"]["target"], "origin");
assert_eq!(interval["deliver"]["chat_id"], "-1002233445566");
assert_eq!(interval["next_run_at"], "2026-09-02T11:15:00Z");
assert_eq!(interval["last_status"], "ok");
let nightly = job_row(result, "nightly-audit");
assert_eq!(nightly["schedule"]["expr"], "0 3 * * *");
assert_eq!(nightly["deliver"]["target"], "local");
assert_eq!(nightly["enabled"], false);
assert_eq!(nightly["state"], "paused");
let profiled = job_row(result, "ops-once-boot");
assert_eq!(profiled["profile"], "ops");
assert_eq!(profiled["schedule"]["kind"], "once");
assert_eq!(profiled["schedule"]["run_at"], "2026-09-03T06:00:00Z");
assert_eq!(profiled["payload"]["kind"], "script");
assert_eq!(profiled["deliver"]["target"], "slack:C0429ABCD");
assert_eq!(profiled["deliver"]["chat_id"], "C0429ABCD");
assert_eq!(profiled["recurring"], false);
let standup_to_group = job_row(result, "coder-standup");
assert_eq!(standup_to_group["deliver"]["target"], "origin");
assert_eq!(standup_to_group["deliver"]["chat_id"], "-100777");
assert_eq!(standup_to_group["deliver"]["thread_id"], "55");
assert!(standup_to_group["deliver"]["mode"].is_null());
assert!(standup_to_group["deliver"]["account"].is_null());
let standup = job_row(result, "cron_standup");
assert_eq!(standup["harness"], "openclaw");
assert_eq!(standup["session_target"], "isolated");
assert_eq!(standup["deliver"]["mode"], "announce");
assert_eq!(standup["deliver"]["target"], "slack");
assert_eq!(standup["deliver"]["chat_id"], "C0429ABCD");
assert_eq!(standup["payload"]["kind"], "prompt");
assert_eq!(standup["profile"], "main");
let reindex = job_row(result, "cron_reindex");
assert_eq!(reindex["session_target"], "main");
assert_eq!(reindex["payload"]["kind"], "system_event");
assert_eq!(reindex["schedule"]["kind"], "interval");
assert_eq!(reindex["schedule"]["display"], "every 240 min");
assert_eq!(reindex["enabled"], false);
let states: Vec<(&str, &str)> = result["sources"]
.as_array()
.unwrap()
.iter()
.map(|source| {
(
source["harness"].as_str().unwrap(),
source["state"].as_str().unwrap(),
)
})
.collect();
assert_eq!(
states,
vec![
("claude-code", "scanned"),
("hermes", "read"),
("hermes", "absent_store"),
("hermes", "read"),
("openclaw", "read"),
("openclaw", "read"),
],
"{result}"
);
}
#[test]
fn jobs_list_filters_by_harness_session_and_profile() {
let by_harness = jobs_list(json!({"harness": "openclaw", "homes": jobs_fixture_homes()}));
let ids: Vec<&str> = by_harness["result"]["jobs"]
.as_array()
.unwrap()
.iter()
.map(|job| job["id"].as_str().unwrap())
.collect();
assert_eq!(
ids,
vec![
"85ad7832-896f-42be-af31-3e1ed2fbdc4b",
"8bb7d938-ca46-4a6d-90eb-c92331155566",
"cron_standup",
"cron_reindex",
]
);
let by_session = jobs_list(json!({
"session": "7c1d2e3f-4a5b-6c7d-8e9f-0a1b2c3d4e5f",
"homes": jobs_fixture_homes(),
}));
let jobs = by_session["result"]["jobs"].as_array().unwrap();
assert_eq!(jobs.len(), 2, "{by_session}");
assert!(jobs
.iter()
.all(|job| job["harness"] == "claude-code" && job["scope"] == "session"));
let by_profile = jobs_list(json!({
"harness": "hermes",
"profile": "ops",
"homes": jobs_fixture_homes(),
}));
let jobs = by_profile["result"]["jobs"].as_array().unwrap();
assert_eq!(jobs.len(), 1, "{by_profile}");
assert_eq!(jobs[0]["id"], "ops-once-boot");
}
#[test]
fn jobs_get_answers_with_the_row_and_the_verbatim_native_record() {
let mut service = HarnessSessionService::new();
let hermes = service.handle(request(
1,
"harness.v1.jobs.get",
json!({"harness": "hermes", "id": "digest-15m", "homes": jobs_fixture_homes()}),
));
assert_eq!(hermes["result"]["job"]["schedule"]["kind"], "interval");
assert_eq!(hermes["result"]["source"]["provider"], "nous");
assert_eq!(hermes["result"]["source"]["failure_deliver"], "local");
let claude = service.handle(request(
2,
"harness.v1.jobs.get",
json!({"harness": "claude-code", "id": "release-watch", "homes": jobs_fixture_homes()}),
));
assert_eq!(claude["result"]["job"]["payload"]["kind"], "prompt");
assert_eq!(
claude["result"]["source"]["tool_use_id"],
"toolu_cron_release_watch"
);
let missing = service.handle(request(
3,
"harness.v1.jobs.get",
json!({"harness": "hermes", "id": "no-such-job", "homes": jobs_fixture_homes()}),
));
assert!(missing["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("no scheduled job `no-such-job`")));
}
#[test]
fn jobs_refuse_a_harness_without_a_scheduled_job_concept() {
let mut service = HarnessSessionService::new();
for (id, method, params) in [
(
1,
"harness.v1.jobs.list",
json!({"harness": "codex", "homes": jobs_fixture_homes()}),
),
(
2,
"harness.v1.jobs.get",
json!({"harness": "codex", "id": "anything"}),
),
] {
let response = service.handle(request(id, method, params));
assert_eq!(response["error"]["code"], -32020, "{response}");
assert!(response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("has no scheduled jobs")));
assert!(response.get("result").is_none());
}
}
#[test]
fn jobs_list_reports_a_migrated_openclaw_store_as_absent_instead_of_failing() {
let scratch = std::env::temp_dir().join(format!(
"supercode-jobs-migrated-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(&scratch).unwrap();
let response = jobs_list(json!({
"harness": "openclaw",
"homes": {"openclaw": scratch.clone()},
}));
let result = &response["result"];
assert_eq!(result["jobs"].as_array().unwrap().len(), 0, "{result}");
assert_eq!(result["sources"][0]["state"], "absent_store");
assert_eq!(result["sources"][0]["harness"], "openclaw");
std::fs::remove_dir_all(&scratch).ok();
}
const OPENCLAW_HEALTH_JOB: &str = "85ad7832-896f-42be-af31-3e1ed2fbdc4b";
const OPENCLAW_DIGEST_JOB: &str = "8bb7d938-ca46-4a6d-90eb-c92331155566";
fn runs_list(params: Value) -> Value {
let mut service = HarnessSessionService::new();
service.handle(request(1, "harness.v1.runs.list", params))
}
fn run_row<'a>(result: &'a Value, id: &str) -> &'a Value {
result["runs"]
.as_array()
.expect("runs is an array")
.iter()
.find(|run| run["id"] == id)
.unwrap_or_else(|| panic!("no run `{id}` in {result}"))
}
#[test]
fn runs_list_projects_both_fixture_stores_onto_the_uniform_row() {
let response = runs_list(json!({"homes": jobs_fixture_homes()}));
let result = &response["result"];
let ids: Vec<&str> = result["runs"]
.as_array()
.expect("runs is an array")
.iter()
.map(|run| run["id"].as_str().unwrap())
.collect();
let digest_fire = format!("{OPENCLAW_DIGEST_JOB}#1");
assert_eq!(
ids,
vec![
"b2c3d4e5f60718293a4b5c6d7e8f9012",
"a1b2c3d4e5f60718293a4b5c6d7e8f90",
"c3d4e5f60718293a4b5c6d7e8f901234",
"f60718293a4b5c6d7e8f901234567890",
"e5f60718293a4b5c6d7e8f9012345678",
"d4e5f60718293a4b5c6d7e8f90123456",
"run_health_0002",
digest_fire.as_str(),
"run_health_0001",
],
"{result}"
);
let failed = run_row(result, "b2c3d4e5f60718293a4b5c6d7e8f9012");
assert_eq!(failed["harness"], "hermes");
assert_eq!(failed["job_id"], "job42");
assert_eq!(failed["status"], "failed");
assert_eq!(failed["error"], "provider returned 500 after 3 attempts");
assert_eq!(failed["claimed_at"], "2026-09-02T13:05:00.100442");
let abandoned = run_row(result, "d4e5f60718293a4b5c6d7e8f90123456");
assert_eq!(abandoned["status"], "unknown");
assert_eq!(abandoned["job_id"], "ops-once-boot");
let running = run_row(result, "c3d4e5f60718293a4b5c6d7e8f901234");
assert_eq!(running["status"], "running");
assert!(running["finished_at"].is_null(), "{running}");
assert!(running["session_id"].is_null(), "{running}");
let ok = run_row(result, "run_health_0001");
assert_eq!(ok["harness"], "openclaw");
assert_eq!(ok["job_id"], OPENCLAW_HEALTH_JOB);
assert_eq!(ok["status"], "ok");
assert_eq!(ok["started_at"], "2026-09-02T08:30:00.000Z");
assert_eq!(ok["finished_at"], "2026-09-02T08:30:30.000Z");
assert_eq!(ok["session_id"], "3dd577ae-a0a3-4b5b-8063-f402be4f5fd4");
assert!(ok["claimed_at"].is_null(), "{ok}");
assert_eq!(run_row(result, &digest_fire)["status"], "skipped");
for id in [
"b2c3d4e5f60718293a4b5c6d7e8f9012",
"d4e5f60718293a4b5c6d7e8f90123456",
] {
assert!(run_row(result, id)["delivery"].is_null(), "{id}");
}
let sources = result["sources"].as_array().unwrap();
let states: Vec<(&str, &str)> = sources
.iter()
.map(|source| {
(
source["harness"].as_str().unwrap(),
source["state"].as_str().unwrap(),
)
})
.collect();
assert_eq!(
states,
vec![
("hermes", "read"),
("hermes", "absent_store"),
("hermes", "read"),
("openclaw", "read"),
],
"{result}"
);
assert_eq!(sources[2]["profile"], "ops");
assert!(sources[3]["path"]
.as_str()
.is_some_and(|path| path.ends_with("state/openclaw.sqlite")));
}
#[test]
fn runs_list_joins_a_hermes_fire_to_the_session_it_opened() {
let response = runs_list(json!({
"harness": "hermes",
"job": "job42",
"homes": jobs_fixture_homes(),
}));
let result = &response["result"];
assert_eq!(result["runs"].as_array().unwrap().len(), 2, "{result}");
let ran = run_row(result, "a1b2c3d4e5f60718293a4b5c6d7e8f90");
assert_eq!(ran["session_id"], "cron_job42_20260902_120000");
let failed = run_row(result, "b2c3d4e5f60718293a4b5c6d7e8f9012");
assert!(failed["session_id"].is_null(), "{failed}");
}
#[test]
fn runs_list_reads_the_delivery_each_harness_recorded_for_a_fire() {
let response = runs_list(json!({"homes": jobs_fixture_homes()}));
let result = &response["result"];
let delivered = run_row(result, "e5f60718293a4b5c6d7e8f9012345678");
assert_eq!(delivered["status"], "completed");
assert_eq!(delivered["delivery"]["state"], "delivered");
assert_eq!(delivered["delivery"]["target"], "telegram:-100777:55");
assert_eq!(delivered["delivery"]["attempts"], 1);
assert!(delivered["delivery"]["last_error"].is_null(), "{delivered}");
assert_eq!(
delivered["delivery"]["delivered_at"],
"2026-09-02T09:00:30.400Z"
);
let undelivered = run_row(result, "f60718293a4b5c6d7e8f901234567890");
assert_eq!(undelivered["status"], "completed");
assert_eq!(undelivered["delivery"]["state"], "failed");
assert_eq!(undelivered["delivery"]["attempts"], 3);
assert_eq!(
undelivered["delivery"]["last_error"],
"telegram send failed: Bad Request: chat not found"
);
assert!(
undelivered["delivery"]["delivered_at"].is_null(),
"{undelivered}"
);
let announced = run_row(result, "run_health_0001");
assert_eq!(announced["delivery"]["state"], "delivered");
assert_eq!(announced["delivery"]["target"], "last");
assert!(announced["delivery"]["attempts"].is_null(), "{announced}");
assert!(
announced["delivery"]["delivered_at"].is_null(),
"{announced}"
);
let refused = run_row(result, "run_health_0002");
assert_eq!(refused["delivery"]["state"], "not-delivered");
assert_eq!(refused["delivery"]["last_error"], "channel_not_found");
let skipped = run_row(result, &format!("{OPENCLAW_DIGEST_JOB}#1"));
assert!(skipped["delivery"].is_null(), "{skipped}");
}
#[test]
fn runs_list_matches_a_hermes_obligation_by_the_session_key_first() {
let scratch = std::env::temp_dir().join(format!(
"supercode-runs-delivery-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(scratch.join("cron")).unwrap();
let fixture = jobs_fixture_root().join("hermes_home");
std::fs::copy(fixture.join("state.db"), scratch.join("state.db")).unwrap();
for name in ["cron/executions.db", "cron/jobs.json"] {
std::fs::copy(fixture.join(name), scratch.join(name)).unwrap();
}
{
let connection = rusqlite::Connection::open(scratch.join("state.db")).unwrap();
connection
.execute(
"UPDATE delivery_obligations SET platform = 'slack', chat_id = 'C0FALLBACK'",
[],
)
.unwrap();
connection
.execute(
"INSERT INTO sessions (id, source, session_key, started_at) VALUES \
('cron_coder-standup_20260902_090010', 'cron', \
'agent:coder:telegram:group:-100777:55', 1788339610.0)",
[],
)
.unwrap();
}
let response = runs_list(json!({
"harness": "hermes",
"job": "coder-standup",
"homes": {"hermes": scratch.join("state.db")},
}));
let result = &response["result"];
let matched = run_row(result, "e5f60718293a4b5c6d7e8f9012345678");
assert_eq!(
matched["session_id"], "cron_coder-standup_20260902_090010",
"{result}"
);
assert_eq!(matched["delivery"]["state"], "delivered", "{result}");
assert_eq!(
matched["delivery"]["target"], "slack:C0FALLBACK:55",
"{result}"
);
std::fs::remove_dir_all(&scratch).ok();
}
#[test]
fn runs_list_follows_a_compression_chain_to_the_readable_tip() {
let scratch = std::env::temp_dir().join(format!(
"supercode-runs-compressed-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(scratch.join("cron")).unwrap();
let fixture = jobs_fixture_root().join("hermes_home");
std::fs::copy(fixture.join("state.db"), scratch.join("state.db")).unwrap();
std::fs::copy(
fixture.join("cron/executions.db"),
scratch.join("cron/executions.db"),
)
.unwrap();
{
let connection = rusqlite::Connection::open(scratch.join("state.db")).unwrap();
connection
.execute(
"UPDATE sessions SET end_reason = 'compression' WHERE id = ?1",
["cron_job42_20260902_120000"],
)
.unwrap();
connection
.execute(
"INSERT INTO sessions (id, source, parent_session_id, started_at) \
VALUES ('job42-after-compaction', 'cron', \
'cron_job42_20260902_120000', 1788350000.0)",
[],
)
.unwrap();
}
let response = runs_list(json!({
"harness": "hermes",
"job": "job42",
"homes": {"hermes": scratch.join("state.db")},
}));
let result = &response["result"];
assert_eq!(
run_row(result, "a1b2c3d4e5f60718293a4b5c6d7e8f90")["session_id"],
"job42-after-compaction",
"{result}"
);
std::fs::remove_dir_all(&scratch).ok();
}
#[test]
fn runs_list_filters_by_job_and_caps_by_limit() {
let by_job = runs_list(json!({
"harness": "openclaw",
"job": OPENCLAW_HEALTH_JOB,
"homes": jobs_fixture_homes(),
}));
let ids: Vec<&str> = by_job["result"]["runs"]
.as_array()
.unwrap()
.iter()
.map(|run| run["id"].as_str().unwrap())
.collect();
assert_eq!(ids, vec!["run_health_0002", "run_health_0001"], "{by_job}");
let capped = runs_list(json!({
"harness": "openclaw",
"limit": 1,
"homes": jobs_fixture_homes(),
}));
let runs = capped["result"]["runs"].as_array().unwrap();
assert_eq!(runs.len(), 1, "{capped}");
assert_eq!(runs[0]["id"], "run_health_0002");
}
#[test]
fn runs_get_answers_with_the_row_and_the_verbatim_native_record() {
let mut service = HarnessSessionService::new();
let hermes = service.handle(request(
1,
"harness.v1.runs.get",
json!({
"harness": "hermes",
"id": "a1b2c3d4e5f60718293a4b5c6d7e8f90",
"homes": jobs_fixture_homes(),
}),
));
assert_eq!(hermes["result"]["run"]["status"], "completed");
assert_eq!(
hermes["result"]["run"]["session_id"],
"cron_job42_20260902_120000"
);
assert_eq!(hermes["result"]["source"]["source"], "scheduler");
assert_eq!(hermes["result"]["source"]["pid"], 4242);
assert_eq!(hermes["result"]["source"]["process_id"], "9f1c2d");
let openclaw = service.handle(request(
2,
"harness.v1.runs.get",
json!({
"harness": "openclaw",
"id": "run_health_0002",
"homes": jobs_fixture_homes(),
}),
));
assert_eq!(openclaw["result"]["run"]["status"], "error");
assert_eq!(
openclaw["result"]["source"]["delivery_status"],
"not-delivered"
);
assert_eq!(
openclaw["result"]["source"]["delivery_error"],
"channel_not_found"
);
assert_eq!(openclaw["result"]["source"]["delivered"], 0);
assert_eq!(
openclaw["result"]["run"]["delivery"]["state"],
"not-delivered"
);
assert_eq!(
openclaw["result"]["run"]["delivery"]["last_error"],
"channel_not_found"
);
let missing = service.handle(request(
3,
"harness.v1.runs.get",
json!({"harness": "hermes", "id": "no-such-run", "homes": jobs_fixture_homes()}),
));
assert!(missing["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("no run `no-such-run`")));
}
#[test]
fn runs_refuse_a_harness_that_keeps_no_run_store() {
let mut service = HarnessSessionService::new();
for (id, method, params) in [
(
1,
"harness.v1.runs.list",
json!({"harness": "claude-code", "homes": jobs_fixture_homes()}),
),
(
2,
"harness.v1.runs.get",
json!({"harness": "claude-code", "id": "anything"}),
),
(
3,
"harness.v1.runs.list",
json!({"harness": "codex", "homes": jobs_fixture_homes()}),
),
] {
let response = service.handle(request(id, method, params));
assert_eq!(response["error"]["code"], -32020, "{response}");
assert!(response["error"]["message"]
.as_str()
.is_some_and(|message| message.contains("keeps no run store")));
assert!(response.get("result").is_none());
}
}
#[test]
fn runs_list_reports_an_install_with_no_run_store_as_absent() {
let scratch = std::env::temp_dir().join(format!(
"supercode-runs-empty-{}-{}",
std::process::id(),
generated_session_id()
));
std::fs::create_dir_all(&scratch).unwrap();
let response = runs_list(json!({
"harness": "openclaw",
"homes": {"openclaw": scratch.clone()},
}));
let result = &response["result"];
assert_eq!(result["runs"].as_array().unwrap().len(), 0, "{result}");
assert_eq!(result["sources"][0]["state"], "absent_store");
assert!(result["sources"][0]["path"]
.as_str()
.is_some_and(|path| path.ends_with("state/openclaw.sqlite")));
std::fs::remove_dir_all(&scratch).ok();
}
}