1use std::collections::{BTreeMap, BTreeSet};
8use std::path::{Path, PathBuf};
9use std::sync::Arc;
10use std::time::Duration;
11
12use serde::{Deserialize, Serialize};
13use serde_json::{json, Value};
14use tokio::sync::Notify;
15
16use crate::reduce;
17use crate::runtime::generated_session_id;
18#[cfg(feature = "adapter-api")]
19use crate::runtime::{HostedHarnessConnection, HostedHarnessRuntime};
20use crate::sdk::{
21 discover_session_page, load_session, load_session_with_fidelity, SdkCapabilities, SdkError,
22 SdkErrorCode, SdkEvent, SdkOperation, SdkRequest, SdkRuntimeEvent, SdkService,
23};
24use crate::Fidelity;
25#[cfg(feature = "adapter-api")]
26use crate::SupercodeHttpRuntimeBackend;
27use crate::{
28 discover_live_runtime, harness_support_registry, AcpRuntimeBackend, ClaudeCodeRuntimeBackend,
29 CodexRuntimeBackend, DiscoveryQuery, HarnessCatalog, HarnessHomes, HarnessId,
30 ImplementationKind, LiveRuntimeEndpoint, LiveRuntimeSource, OpenCodeRuntimeBackend,
31 PiRuntimeBackend, Role, RuntimeAttachRequest, RuntimeBackend, RuntimeConnection, RuntimeInput,
32 RuntimeLaunch, RuntimeStartRequest, Session, SessionDescriptor, SessionFollower, SessionFormat,
33 SessionLocator, SessionSource,
34};
35#[cfg(feature = "adapter-api")]
36use crate::{register_live_runtime, resolve_live_runtime, LiveRuntimeRegistration};
37use supercode_interchange::watch::{bound_session_view, message_json, normalized_session_json};
38
39pub const HARNESS_SERVICE_METHODS: &[&str] = &[
42 "harness.v1.support.report",
43 "harness.v1.harnesses.list",
44 "harness.v1.harnesses.probe",
45 "harness.v1.harnesses.settings",
46 "harness.v1.harnesses.configure",
47 "harness.v1.harnesses.auth.methods",
48 "harness.v1.harnesses.auth.begin",
49 "harness.v1.harnesses.auth.verify",
50 "harness.v1.sessions.discover",
51 "harness.v1.sessions.load",
52 "harness.v1.sessions.follow",
53 "harness.v1.sessions.unfollow",
54 "harness.v1.sessions.activity.subscribe",
55 "harness.v1.sessions.activity.unsubscribe",
56 "harness.v1.sessions.activity_under",
57 "harness.v1.sessions.index.subscribe",
58 "harness.v1.sessions.index.resize",
59 "harness.v1.sessions.index.unsubscribe",
60 "harness.v1.sessions.message",
61 "harness.v1.sessions.inbox",
62 "harness.v1.sessions.import",
63 "harness.v1.sessions.export",
64 "harness.v1.sessions.translate",
65 "harness.v1.sessions.reduce",
66 "harness.v1.sessions.branch",
67 "harness.v1.sessions.handoff",
68 "harness.v1.sessions.materialize",
69 "harness.v1.sessions.resume_instructions",
70 "harness.v1.skills.list",
71 "harness.v1.skills.install",
72 "harness.v1.skills.remove",
73 "harness.v1.memory.show",
74 "harness.v1.memory.search",
75 "harness.v1.jobs.list",
76 "harness.v1.jobs.get",
77 "harness.v1.jobs.create",
78 "harness.v1.jobs.update",
79 "harness.v1.jobs.pause",
80 "harness.v1.jobs.resume",
81 "harness.v1.jobs.run",
82 "harness.v1.jobs.delete",
83 "harness.v1.jobs.notepad",
84 "harness.v1.jobs.notepad_set",
85 "harness.v1.jobs.notepad_delete",
86 "harness.v1.sessions.new",
87 "harness.v1.sessions.reset",
88 "harness.v1.sessions.archive",
89 "harness.v1.sessions.delete",
90 "harness.v1.runs.list",
91 "harness.v1.runs.get",
92 "harness.v1.approvals.list",
93 "harness.v1.approvals.resolve",
94 "harness.v1.runtimes.capabilities",
95 "harness.v1.runtimes.start",
96 "harness.v1.runtimes.resume",
97 "harness.v1.runtimes.attach_existing",
98 "harness.v1.runtimes.attach",
99 "harness.v1.runtimes.send_input",
100 "harness.v1.runtimes.interrupt",
101 "harness.v1.runtimes.steer",
102 "harness.v1.runtimes.respond",
103 "harness.v1.runtimes.terminal_instructions",
104 "harness.v1.runtimes.acquire_control",
105 "harness.v1.runtimes.heartbeat",
106 "harness.v1.runtimes.detach",
107 "harness.v1.runtimes.close",
108 "harness.v1.profiles.list",
109 "harness.v1.profiles.get",
110 "harness.v1.profiles.create",
111 "harness.v1.profiles.delete",
112 "harness.v1.channels.list",
113 "harness.v1.routes.list",
114 "harness.v1.triggers.list",
115 "harness.v1.channels.status",
116 "harness.v1.orchestration.load",
117 "harness.v1.orchestration.save",
118 "harness.v1.orchestration.compile",
119 "harness.v1.orchestration.decompile",
120 "harness.v1.orchestration.import",
121 "harness.v1.orchestration.export",
122 "harness.v1.workflow.load",
123];
124
125pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
127pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
129pub const SESSION_ACTIVITY_EVENT_METHOD: &str = "harness.v1.sessions.activity_event";
131pub const SESSION_INDEX_EVENT_METHOD: &str = "harness.v1.sessions.index_event";
133pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
135
136pub struct HarnessSessionService {
139 catalog: HarnessCatalog,
140 followers: BTreeMap<String, SessionFollower>,
141 followed_sources: BTreeMap<String, FollowedSource>,
142 activity_subscriptions: BTreeMap<String, ActivitySubscription>,
143 index_subscriptions: BTreeMap<String, crate::session_index::SessionIndexSubscription>,
144 index_notifier: Arc<Notify>,
145 #[cfg(feature = "adapter-api")]
146 activity_monitor: crate::session_activity::SessionActivityMonitor,
147 next_subscription: u64,
148 runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
149 runtimes_in_flight: BTreeSet<String>,
154 terminal_launches: BTreeMap<String, StructuredLaunch>,
155 runtime_sequences: BTreeMap<String, u64>,
156 next_runtime: u64,
157 reduction_store_root: Option<PathBuf>,
158 approvals: crate::approvals::ApprovalRegistry,
162 subagent_approvals: Option<Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>>,
165}
166
167impl Default for HarnessSessionService {
168 fn default() -> Self {
169 Self::new()
170 }
171}
172
173impl HarnessSessionService {
174 pub fn new() -> Self {
176 Self {
177 catalog: HarnessCatalog::new(),
178 followers: BTreeMap::new(),
179 followed_sources: BTreeMap::new(),
180 activity_subscriptions: BTreeMap::new(),
181 index_subscriptions: BTreeMap::new(),
182 index_notifier: Arc::new(Notify::new()),
183 #[cfg(feature = "adapter-api")]
184 activity_monitor: Default::default(),
185 next_subscription: 1,
186 runtimes: BTreeMap::new(),
187 runtimes_in_flight: BTreeSet::new(),
188 terminal_launches: BTreeMap::new(),
189 runtime_sequences: BTreeMap::new(),
190 next_runtime: 1,
191 reduction_store_root: None,
192 approvals: crate::approvals::ApprovalRegistry::new(),
193 subagent_approvals: None,
194 }
195 }
196
197 pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
202 self.reduction_store_root = Some(root.into());
203 self
204 }
205
206 pub fn observe_subagent_approvals(
214 &mut self,
215 queue: Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>,
216 ) {
217 self.subagent_approvals = Some(queue);
218 }
219
220 pub fn approvals(&self, query: &crate::approvals::ApprovalsQuery) -> Vec<crate::ApprovalRow> {
227 let now = crate::approvals::now_ms();
228 let mut rows = self.approvals.rows(now);
229 if let Some(queue) = self.subagent_approvals.as_ref() {
230 let queued = queue
231 .lock()
232 .unwrap_or_else(std::sync::PoisonError::into_inner)
233 .clone();
234 rows.extend(crate::approvals::subagent_rows(&queued, now));
235 }
236 rows.retain(|row| query.matches(row));
237 rows.sort_by(|left, right| {
238 left.requested_at_ms
239 .cmp(&right.requested_at_ms)
240 .then_with(|| left.id.cmp(&right.id))
241 });
242 rows
243 }
244
245 async fn approvals_resolve(
255 &mut self,
256 params: Value,
257 ) -> std::result::Result<Value, ServiceError> {
258 let params = decode::<crate::approvals::ApprovalsResolveParams>(params)?;
259 if params.id.trim().is_empty() {
260 return Err(ServiceError::InvalidParams(
261 "approvals resolve requires the `id` of a listed approval row".into(),
262 ));
263 }
264 let choice = match (params.decision, params.option_id.as_deref()) {
265 (Some(_), Some(_)) => {
266 return Err(ServiceError::InvalidParams(
267 "approvals resolve takes either `decision` or `option_id`, not both".into(),
268 ))
269 }
270 (Some(decision), None) => crate::approvals::ApprovalChoice::Decision(decision),
271 (None, Some(option)) => crate::approvals::ApprovalChoice::Option(option.to_string()),
272 (None, None) => {
273 return Err(ServiceError::InvalidParams(format!(
274 "approvals resolve requires `decision` ({}) or an explicit `option_id`",
275 crate::approvals::ApprovalDecision::ALL
276 .map(|decision| decision.as_str())
277 .join(" | "),
278 )))
279 }
280 };
281 let resolution = self
282 .approvals
283 .resolution(¶ms.id, &choice)
284 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?;
285 self.runtime_call(
289 "harness.v1.runtimes.respond",
290 json!({
291 "connection": resolution.connection,
292 "request_id": resolution.request_id,
293 "response": resolution.response,
294 }),
295 )
296 .await?;
297 Ok(json!({
298 "id": params.id,
299 "decision": params.decision.map(|decision| decision.as_str()),
300 "option_id": resolution.option_id,
301 "resolved": true,
302 }))
303 }
304
305 #[cfg(feature = "adapter-api")]
308 pub fn session_index_notifier(&self) -> Arc<Notify> {
309 Arc::clone(&self.index_notifier)
310 }
311
312 #[cfg(feature = "adapter-api")]
314 pub fn handle(&mut self, request: Value) -> Value {
315 let id = request.get("id").cloned().unwrap_or(Value::Null);
316 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
317 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
318 }
319 let Some(method) = request.get("method").and_then(Value::as_str) else {
320 return rpc_error(id, -32600, "request is missing `method`");
321 };
322 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
323 match self.call(method, params) {
324 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
325 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
326 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
327 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
328 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
329 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
330 }
331 }
332
333 #[cfg(feature = "adapter-api")]
336 pub async fn handle_async(&mut self, request: Value) -> Value {
337 let method = request
338 .get("method")
339 .and_then(Value::as_str)
340 .unwrap_or_default();
341 if matches!(
342 method,
343 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
344 ) {
345 let id = request.get("id").cloned().unwrap_or(Value::Null);
346 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
347 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
348 }
349 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
350 return match self.inventory_call(method, params).await {
351 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
352 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
353 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
354 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
355 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
356 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
357 };
358 }
359 if matches!(
360 method,
361 "harness.v1.harnesses.auth.methods"
362 | "harness.v1.harnesses.auth.begin"
363 | "harness.v1.harnesses.auth.verify"
364 ) {
365 let id = request.get("id").cloned().unwrap_or(Value::Null);
366 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
367 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
368 }
369 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
370 return match self.harness_authentication_call(method, params).await {
371 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
372 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
373 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
374 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
375 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
376 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
377 };
378 }
379 if matches!(
385 method,
386 "harness.v1.sessions.new"
387 | "harness.v1.sessions.reset"
388 | "harness.v1.sessions.archive"
389 | "harness.v1.sessions.delete"
390 ) {
391 let id = request.get("id").cloned().unwrap_or(Value::Null);
392 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
393 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
394 }
395 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
396 let verb = match method {
397 "harness.v1.sessions.new" => crate::SessionVerb::New,
398 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
399 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
400 _ => crate::SessionVerb::Delete,
401 };
402 return match self.mutate_session(verb, params).await {
403 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
404 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
405 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
406 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
407 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
408 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
409 };
410 }
411 if method == "harness.v1.sessions.message" {
412 let id = request.get("id").cloned().unwrap_or(Value::Null);
413 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
414 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
415 }
416 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
417 return match self.message_call(params).await {
418 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
419 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
420 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
421 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
422 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
423 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
424 };
425 }
426 if matches!(
427 method,
428 "harness.v1.harnesses.settings" | "harness.v1.harnesses.configure"
429 ) {
430 let id = request.get("id").cloned().unwrap_or(Value::Null);
431 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
432 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
433 }
434 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
435 return match self.harness_settings_call(method, params) {
436 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
437 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
438 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
439 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
440 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
441 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
442 };
443 }
444 if method == "harness.v1.sessions.activity_under" {
445 let id = request.get("id").cloned().unwrap_or(Value::Null);
446 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
447 return match activity_under_call(params).await {
448 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
449 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
450 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
451 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
452 Err(_) => rpc_error(id, -32000, "sessions.activity_under failed"),
453 };
454 }
455 if method == "harness.v1.sessions.activity.subscribe" {
456 let id = request.get("id").cloned().unwrap_or(Value::Null);
457 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
458 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
459 }
460 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
461 return match self.subscribe_session_activity(params).await {
462 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
463 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
464 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
465 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
466 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
467 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
468 };
469 }
470 if let Some(operation) = SdkOperation::from_method(method) {
471 let id = request.get("id").cloned().unwrap_or(Value::Null);
472 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
473 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
474 }
475 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
476 return match self.execute(SdkRequest { operation, params }).await {
477 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
478 Err(error) => sdk_rpc_error(id, &error),
479 };
480 }
481 if !method.starts_with("harness.v1.runtimes.") {
482 return self.handle(request);
483 }
484 let id = request.get("id").cloned().unwrap_or(Value::Null);
485 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
486 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
487 }
488 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
489 match self.runtime_call(method, params).await {
490 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
491 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
492 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
493 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
494 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
495 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
496 }
497 }
498
499 #[cfg(feature = "adapter-api")]
502 pub fn poll(&mut self) -> Vec<Value> {
503 let mut notifications = Vec::new();
504 for (subscription, follower) in &mut self.followers {
505 match follower.poll() {
506 Ok(Some(event)) => notifications.push(json!({
507 "jsonrpc": "2.0",
508 "method": SESSION_EVENT_METHOD,
509 "params": {
510 "subscription": subscription,
511 "event": event.to_json(),
512 }
513 })),
514 Ok(None) => {}
515 Err(error) => notifications.push(json!({
516 "jsonrpc": "2.0",
517 "method": SESSION_EVENT_METHOD,
518 "params": {
519 "subscription": subscription,
520 "event": {
521 "type": "watch_error",
522 "recoverable": true,
523 "message": error.to_string(),
524 },
525 }
526 })),
527 }
528 }
529 notifications
530 }
531
532 #[cfg(feature = "adapter-api")]
543 pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
544 let registry = crate::LocalRuntimeRegistry::new();
545 let authorization = crate::RuntimeAuthorization::observer();
546 let mut notifications = Vec::new();
547 for (subscription, source) in &mut self.followed_sources {
548 let state = match registry
549 .source_state(&source.harness, &source.session_id, &authorization)
550 .await
551 {
552 Ok(Some(state)) => state,
553 Ok(None) => crate::RuntimeRegistryState::Persisted,
554 Err(_) => continue,
556 };
557 if source.reported.as_deref() == Some(state.as_str()) {
558 continue;
559 }
560 source.reported = Some(state.as_str().to_string());
561 notifications.push(json!({
562 "jsonrpc": "2.0",
563 "method": SESSION_EVENT_METHOD,
564 "params": {
565 "subscription": subscription,
566 "event": {"type": "runtime_state", "state": state.as_str()},
567 },
568 }));
569 }
570 notifications
571 }
572
573 #[cfg(feature = "adapter-api")]
577 pub async fn poll_session_activities(&mut self) -> Vec<Value> {
578 let subscriptions = self
579 .activity_subscriptions
580 .iter()
581 .map(|(id, subscription)| {
582 (
583 id.clone(),
584 subscription.locators.clone(),
585 subscription.homes.clone(),
586 )
587 })
588 .collect::<Vec<_>>();
589 let mut notifications = Vec::new();
590 for (subscription_id, locators, homes) in subscriptions {
591 let Ok(activities) = self.activity_monitor.resolve(&locators, &homes).await else {
592 continue;
595 };
596 let Some(subscription) = self.activity_subscriptions.get_mut(&subscription_id) else {
597 continue;
598 };
599 let mut changed = Vec::new();
600 for activity in activities {
601 let key = activity.key();
602 if subscription
603 .reported
604 .get(&key)
605 .is_some_and(|previous| previous.same_state(&activity))
606 {
607 continue;
608 }
609 subscription.reported.insert(key, activity.clone());
610 changed.push(activity);
611 }
612 if !changed.is_empty() {
613 notifications.push(json!({
614 "jsonrpc": "2.0",
615 "method": SESSION_ACTIVITY_EVENT_METHOD,
616 "params": {
617 "subscription": subscription_id,
618 "activities": changed,
619 },
620 }));
621 }
622 }
623 notifications
624 }
625
626 #[cfg(feature = "adapter-api")]
630 pub fn poll_session_indexes(&mut self) -> Vec<Value> {
631 let mut notifications = Vec::new();
632 for (subscription, index) in &mut self.index_subscriptions {
633 let homes = index.homes().clone();
634 match index.poll() {
635 Ok(Some(delta)) => match live_index_changes(delta.changes, &homes) {
636 Ok(changes) => notifications.push(json!({
637 "jsonrpc": "2.0",
638 "method": SESSION_INDEX_EVENT_METHOD,
639 "params": {
640 "subscription": subscription,
641 "revision": delta.revision,
642 "changes": changes,
643 },
644 })),
645 Err(error) => notifications.push(json!({
646 "jsonrpc": "2.0",
647 "method": SESSION_INDEX_EVENT_METHOD,
648 "params": {
649 "subscription": subscription,
650 "error": {"recoverable": true, "message": error_message(error)},
651 },
652 })),
653 },
654 Ok(None) => {}
655 Err(error) => notifications.push(json!({
656 "jsonrpc": "2.0",
657 "method": SESSION_INDEX_EVENT_METHOD,
658 "params": {
659 "subscription": subscription,
660 "error": {"recoverable": true, "message": error},
661 },
662 })),
663 }
664 }
665 notifications
666 }
667
668 #[cfg(feature = "adapter-api")]
669 async fn subscribe_session_activity(
670 &mut self,
671 params: Value,
672 ) -> std::result::Result<Value, ServiceError> {
673 let params = decode::<ActivitySubscribeParams>(params)?;
674 if params.locators.is_empty() {
675 return Err(ServiceError::InvalidParams(
676 "sessions.activity.subscribe requires at least one locator".into(),
677 ));
678 }
679 if params.locators.len() > 2_048 {
680 return Err(ServiceError::InvalidParams(
681 "sessions.activity.subscribe accepts at most 2048 locators".into(),
682 ));
683 }
684 let initial = self
685 .activity_monitor
686 .resolve(¶ms.locators, ¶ms.homes)
687 .await
688 .map_err(ServiceError::Sdk)?;
689 let subscription = format!("activity-sub-{}", self.next_subscription);
690 self.next_subscription += 1;
691 let reported = initial
692 .iter()
693 .cloned()
694 .map(|activity| (activity.key(), activity))
695 .collect();
696 self.activity_subscriptions.insert(
697 subscription.clone(),
698 ActivitySubscription {
699 locators: params.locators,
700 homes: params.homes,
701 reported,
702 },
703 );
704 Ok(json!({"subscription": subscription, "initial": initial}))
705 }
706
707 #[cfg(feature = "adapter-api")]
709 pub async fn poll_runtimes(&mut self) -> Vec<Value> {
710 self.poll_sdk_events()
711 .await
712 .into_iter()
713 .map(|(connection, runtime_event)| {
714 json!({
715 "jsonrpc": "2.0",
716 "method": RUNTIME_EVENT_METHOD,
717 "params": {
718 "connection": connection,
719 "session_id": runtime_event.session_id,
720 "sequence": runtime_event.event.sequence,
721 "event": {
722 "kind": runtime_event.event.kind,
723 "payload": runtime_event.event.payload,
724 },
725 },
726 })
727 })
728 .collect()
729 }
730
731 async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
732 let mut events = Vec::new();
733 let mut closed = Vec::new();
734 let now_ms = crate::approvals::now_ms();
735 for (connection, runtime) in &mut self.runtimes {
736 let session_id = runtime.handle().runtime_id.clone();
737 let harness = runtime.handle().harness.clone();
738 for _ in 0..256 {
743 match tokio::time::timeout(Duration::ZERO, runtime.next_event()).await {
744 Ok(Ok(Some(event))) => {
745 let terminal = event.kind == "transport_closed";
746 self.approvals
750 .observe(connection, &harness, &session_id, &event, now_ms);
751 let next_sequence = self
752 .runtime_sequences
753 .entry(session_id.clone())
754 .or_insert(0);
755 let sequence = event.sequence.unwrap_or_else(|| {
756 *next_sequence = next_sequence.saturating_add(1);
757 *next_sequence
758 });
759 *next_sequence = (*next_sequence).max(sequence);
760 events.push((
761 connection.clone(),
762 SdkRuntimeEvent {
763 session_id: session_id.clone(),
764 event: SdkEvent {
765 sequence,
766 kind: event.kind,
767 payload: event.payload,
768 },
769 },
770 ));
771 if terminal {
772 closed.push(connection.clone());
773 break;
774 }
775 }
776 Ok(Ok(None)) => {
777 let sequence = self
778 .runtime_sequences
779 .entry(session_id.clone())
780 .or_insert(0);
781 *sequence = sequence.saturating_add(1);
782 events.push((
783 connection.clone(),
784 SdkRuntimeEvent {
785 session_id,
786 event: SdkEvent {
787 sequence: *sequence,
788 kind: "transport_closed".into(),
789 payload: json!({"message": "Harness runtime transport closed."}),
790 },
791 },
792 ));
793 closed.push(connection.clone());
794 break;
795 }
796 Err(_) => break,
797 Ok(Err(error)) => {
798 let sequence = self
799 .runtime_sequences
800 .entry(session_id.clone())
801 .or_insert(0);
802 *sequence = sequence.saturating_add(1);
803 events.push((
804 connection.clone(),
805 SdkRuntimeEvent {
806 session_id,
807 event: SdkEvent {
808 sequence: *sequence,
809 kind: "transport_error".into(),
810 payload: json!({"message": error.to_string(), "terminal": true}),
811 },
812 },
813 ));
814 closed.push(connection.clone());
815 break;
816 }
817 }
818 }
819 }
820 for connection in closed {
821 if let Some(runtime) = self.runtimes.remove(&connection) {
822 self.runtime_sequences.remove(&runtime.handle().runtime_id);
823 }
824 self.terminal_launches.remove(&connection);
825 self.approvals.forget(&connection);
828 }
829 events
830 }
831
832 fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
833 match method {
834 "harness.v1.capabilities" => Ok(json!({
835 "version": HARNESS_SERVICE_VERSION,
836 "sdk": self.capabilities(),
837 "methods": HARNESS_SERVICE_METHODS,
838 "notifications": [
839 SESSION_EVENT_METHOD,
840 SESSION_ACTIVITY_EVENT_METHOD,
841 SESSION_INDEX_EVENT_METHOD,
842 RUNTIME_EVENT_METHOD
843 ],
844 "harnesses": harness_support_registry()
845 .harnesses
846 .into_iter()
847 .map(|harness| harness.id)
848 .collect::<Vec<_>>(),
849 })),
850 "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
851 .map_err(|error| ServiceError::Operation(error.to_string())),
852 "harness.v1.profiles.list" | "harness.v1.profiles.get" => profiles_call(method, params),
853 "harness.v1.profiles.create" => {
859 mutate_profile(crate::profiles_control::ProfileVerb::Create, params)
860 }
861 "harness.v1.profiles.delete" => {
862 mutate_profile(crate::profiles_control::ProfileVerb::Delete, params)
863 }
864 "harness.v1.channels.list" | "harness.v1.channels.status" => {
865 channels_call(method, params)
866 }
867 "harness.v1.routes.list" => routes_call(params),
870 "harness.v1.triggers.list" => triggers_call(params),
872 "harness.v1.workflow.load" => {
882 let params = decode::<WorkflowLoadParams>(params)?;
883 let read =
884 crate::workflow_doors::load(params.from, ¶ms.home).map_err(operation)?;
885 serde_json::to_value(read)
886 .map_err(|error| ServiceError::Operation(error.to_string()))
887 }
888 "harness.v1.orchestration.load" => {
889 let params = decode::<OrchestrationLoadParams>(params)?;
890 let read = crate::orchestration_doors::load(¶ms.root, params.flavor)
891 .map_err(operation)?;
892 serde_json::to_value(read)
893 .map_err(|error| ServiceError::Operation(error.to_string()))
894 }
895 "harness.v1.orchestration.save" => {
896 let params = decode::<OrchestrationSaveParams>(params)?;
897 let saved = crate::orchestration_doors::save(
898 ¶ms.root,
899 params.orchestration,
900 params.vault,
901 )
902 .map_err(operation)?;
903 serde_json::to_value(saved)
904 .map_err(|error| ServiceError::Operation(error.to_string()))
905 }
906 "harness.v1.orchestration.compile" => {
907 let params = decode::<OrchestrationCompileParams>(params)?;
908 let read = crate::orchestration_doors::compile(params.from, ¶ms.home)
909 .map_err(operation)?;
910 serde_json::to_value(read)
911 .map_err(|error| ServiceError::Operation(error.to_string()))
912 }
913 "harness.v1.orchestration.decompile" => {
914 let params = decode::<OrchestrationDecompileParams>(params)?;
915 let report = crate::orchestration_doors::decompile(
916 params.to,
917 params.orchestration,
918 ¶ms.source,
919 params.source_flavor,
920 ¶ms.dest,
921 params.vault,
922 )
923 .map_err(operation)?;
924 serde_json::to_value(report)
925 .map_err(|error| ServiceError::Operation(error.to_string()))
926 }
927 "harness.v1.orchestration.import" => {
932 let params = decode::<OrchestrationImportParams>(params)?;
933 let imported =
934 crate::orchestration_doors::import(params.from, ¶ms.home, ¶ms.into)
935 .map_err(operation)?;
936 serde_json::to_value(imported)
937 .map_err(|error| ServiceError::Operation(error.to_string()))
938 }
939 "harness.v1.orchestration.export" => {
940 let params = decode::<OrchestrationExportParams>(params)?;
941 let report =
942 crate::orchestration_doors::export(params.to, ¶ms.root, ¶ms.dest)
943 .map_err(operation)?;
944 serde_json::to_value(report)
945 .map_err(|error| ServiceError::Operation(error.to_string()))
946 }
947 "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
953 "harness.v1.skills.list" => {
958 let query = decode::<crate::skills::SkillsQuery>(params)?;
959 if let Some(harness) = query.harness.as_deref() {
960 if !crate::skills::SKILL_HARNESSES.contains(&harness) {
961 return Err(ServiceError::UnsupportedAction(format!(
962 "`{harness}` has no skills root Volter Harness reads"
963 )));
964 }
965 }
966 serde_json::to_value(crate::skills::list_skills(&query))
967 .map_err(|error| ServiceError::Operation(error.to_string()))
968 }
969 "harness.v1.skills.install" => {
977 mutate_skill(crate::skills_control::SkillVerb::Install, params)
978 }
979 "harness.v1.skills.remove" => {
980 mutate_skill(crate::skills_control::SkillVerb::Remove, params)
981 }
982 "harness.v1.approvals.list" => {
990 let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
991 if let Some(harness) = query.harness.as_deref() {
992 if !crate::approvals::lists_approvals(harness) {
993 return Err(ServiceError::UnsupportedAction(format!(
994 "`{harness}` has no runtime door that carries an approval request"
995 )));
996 }
997 }
998 serde_json::to_value(self.approvals(&query))
999 .map_err(|error| ServiceError::Operation(error.to_string()))
1000 }
1001 "harness.v1.sessions.discover" => {
1002 let query = decode::<DiscoveryQuery>(params)?;
1003 let mut page = discover_session_page(&query).map_err(operation)?;
1004 let doors = crate::mail_route::LiveSessions::read(&query.homes);
1009 if let Some(text) = query
1013 .query
1014 .as_deref()
1015 .map(str::trim)
1016 .filter(|text| !text.is_empty())
1017 {
1018 let needle = text.to_lowercase();
1019 for live in doors.all() {
1020 let name = live
1021 .name
1022 .split('@')
1023 .next()
1024 .unwrap_or(&live.name)
1025 .to_lowercase();
1026 if !name.contains(&needle)
1027 || page.sessions.iter().any(|session| {
1028 session.locator.session_id == live.address.session_id
1029 })
1030 {
1031 continue;
1032 }
1033 if query
1034 .limit
1035 .is_some_and(|limit| page.sessions.len() >= limit)
1036 {
1037 break;
1038 }
1039 let mut by_id = query.clone();
1041 by_id.query = Some(live.address.session_id.clone());
1042 by_id.cursor = None;
1043 if let Some(cwd) = &live.cwd {
1044 by_id.workspace = Some(cwd.clone());
1045 by_id.workspace_subtree = false;
1046 by_id.workspace_family = None;
1047 }
1048 if let Ok(found) = discover_session_page(&by_id) {
1049 page.sessions
1050 .extend(found.sessions.into_iter().filter(|session| {
1051 session.locator.session_id == live.address.session_id
1052 }));
1053 }
1054 }
1055 }
1056 let activities = crate::session_activity::resolve_stock_session_activities(
1057 &page
1058 .sessions
1059 .iter()
1060 .map(|session| session.locator.clone())
1061 .collect::<Vec<_>>(),
1062 &query.homes,
1063 )
1064 .into_iter()
1065 .map(|activity| (activity.key(), activity))
1066 .collect::<BTreeMap<_, _>>();
1067 let sessions = page
1068 .sessions
1069 .into_iter()
1070 .map(|session| {
1071 let mut value = live_descriptor_value(&session, &doors)?;
1072 let activity_key = (
1073 session.locator.harness.as_str().to_string(),
1074 session.locator.session_id.clone(),
1075 );
1076 if let Some(activity) = activities.get(&activity_key) {
1077 value["activity"] = serde_json::to_value(activity)
1078 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1079 if let Some(status) = legacy_live_status(activity) {
1080 value["live_status"] = json!(status);
1081 }
1082 }
1083 Ok(value)
1084 })
1085 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1086 let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1087 if query.search_previews {
1090 result["receipt"] = serde_json::to_value(page.receipt)
1091 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1092 }
1093 Ok(result)
1094 }
1095 "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1096 "harness.v1.sessions.load" => {
1097 let params = decode::<LoadSessionParams>(params)?;
1098 if let Some(options) = ¶ms.options {
1099 options.validate()?;
1100 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
1101 return Ok(result);
1102 }
1103 return load_session(¶ms.read.locator)
1104 .map(|session| projected_session_result(&session, options))
1105 .map_err(operation);
1106 }
1107 let mut session = if params.read.display_history() {
1108 self.catalog
1109 .load_display_view(
1110 ¶ms.read.locator,
1111 params.read.read_fidelity(),
1112 params.read.tail_messages().unwrap_or(500),
1113 )
1114 .map_err(crate::Error::from)
1115 } else if params.read.include_subagents() {
1116 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
1117 } else {
1118 self.catalog
1119 .load_parent_with_fidelity(
1120 ¶ms.read.locator,
1121 params.read.read_fidelity(),
1122 )
1123 .map_err(crate::Error::from)
1124 }
1125 .map_err(operation)?;
1126 params.read.bound_session(&mut session);
1127 Ok(json!({"session": normalized_session_json(&session)}))
1128 }
1129 "harness.v1.sessions.follow" => {
1130 let params = decode::<LocatorParams>(params)?;
1131 let mut follower = self
1132 .catalog
1133 .follow_read_view(
1134 ¶ms.locator,
1135 params.read_fidelity(),
1136 params.include_subagents(),
1137 params.tail_messages(),
1138 params.max_message_chars(),
1139 params.display_history(),
1140 )
1141 .map_err(operation)?;
1142 let initial = follower
1143 .poll()
1144 .map_err(operation)?
1145 .map(|event| event.to_json());
1146 let subscription = format!("sub-{}", self.next_subscription);
1147 self.next_subscription += 1;
1148 self.followers.insert(subscription.clone(), follower);
1149 self.followed_sources.insert(
1150 subscription.clone(),
1151 FollowedSource {
1152 harness: params.locator.harness.as_str().to_string(),
1153 session_id: params.locator.session_id.clone(),
1154 reported: None,
1155 },
1156 );
1157 Ok(json!({"subscription": subscription, "initial": initial}))
1158 }
1159 "harness.v1.sessions.unfollow" => {
1160 let params = decode::<UnfollowParams>(params)?;
1161 self.followed_sources.remove(¶ms.subscription);
1162 Ok(json!({
1163 "removed": self.followers.remove(¶ms.subscription).is_some()
1164 }))
1165 }
1166 "harness.v1.sessions.activity.unsubscribe" => {
1167 let params = decode::<UnfollowParams>(params)?;
1168 Ok(json!({
1169 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1170 }))
1171 }
1172 "harness.v1.sessions.index.subscribe" => {
1173 let query = decode::<DiscoveryQuery>(params)?;
1174 crate::session_index::validate_query(&query)
1175 .map_err(ServiceError::InvalidParams)?;
1176 let homes = query.homes.clone();
1177 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1178 query,
1179 Arc::clone(&self.index_notifier),
1180 )
1181 .map_err(ServiceError::Operation)?;
1182 let doors = crate::mail_route::LiveSessions::read(&homes);
1183 let initial = initial
1184 .iter()
1185 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1186 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1187 let subscription = format!("index-sub-{}", self.next_subscription);
1188 self.next_subscription += 1;
1189 self.index_subscriptions.insert(subscription.clone(), index);
1190 Ok(json!({
1191 "subscription": subscription,
1192 "revision": 1,
1193 "initial": initial,
1194 }))
1195 }
1196 "harness.v1.sessions.index.resize" => {
1197 let params = decode::<IndexResizeParams>(params)?;
1198 crate::session_index::validate_limit(params.limit)
1199 .map_err(ServiceError::InvalidParams)?;
1200 let index = self
1201 .index_subscriptions
1202 .get_mut(¶ms.subscription)
1203 .ok_or_else(|| {
1204 ServiceError::InvalidParams("unknown session index subscription".into())
1205 })?;
1206 let prepared = index
1207 .prepare_resize(params.limit)
1208 .map_err(ServiceError::Operation)?;
1209 let doors = crate::mail_route::LiveSessions::read(index.homes());
1210 let initial = prepared
1211 .page
1212 .sessions
1213 .iter()
1214 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1215 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1216 let response = json!({
1217 "subscription": params.subscription,
1218 "revision": prepared.revision,
1219 "initial": initial,
1220 "receipt": prepared.page.receipt,
1221 });
1222 index.commit_resize(prepared);
1223 Ok(response)
1224 }
1225 "harness.v1.sessions.index.unsubscribe" => {
1226 let params = decode::<UnfollowParams>(params)?;
1227 Ok(json!({
1228 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1229 }))
1230 }
1231 "harness.v1.sessions.import" => {
1232 let params = decode::<ImportSessionParams>(params)?;
1233 let session = Session::load_str(¶ms.content, params.source_harness.into())
1234 .map_err(operation)?;
1235 Ok(json!({"session": normalized_session_json(&session)}))
1236 }
1237 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1238 let params = decode::<ExportSessionParams>(params)?;
1239 let session = load_session(¶ms.locator).map_err(operation)?;
1240 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1241 if method == "harness.v1.sessions.export"
1242 && params.target_harness == TransferFormat::Hermes
1243 {
1244 let imported = crate::hermes_import::import_into_hermes(&session, None)
1246 .map_err(operation)?;
1247 return Ok(json!({"artifact": artifact, "imported": imported}));
1248 }
1249 Ok(json!({"artifact": artifact}))
1250 }
1251 "harness.v1.sessions.reduce" => {
1252 let params = decode::<ReduceSessionParams>(params)?;
1253 self.reduce_session(params)
1254 }
1255 "harness.v1.sessions.branch" => {
1256 let params = decode::<BranchSessionParams>(params)?;
1257 let session = load_session(¶ms.locator).map_err(operation)?;
1258 let storage = params.locator.storage.path().display().to_string();
1259 let bootstrap_prompt = format!(
1260 "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.",
1261 params.locator.harness.as_str(), params.locator.session_id, storage
1262 );
1263 let artifact = params
1264 .target_harness
1265 .map(|target| session_artifact(¶ms.locator, &session, target))
1266 .transpose()?;
1267 Ok(json!({
1268 "parent": params.locator,
1269 "session": normalized_session_json(&session),
1270 "bootstrap_prompt": bootstrap_prompt,
1271 "artifact": artifact,
1272 }))
1273 }
1274 "harness.v1.sessions.handoff" => {
1275 let params = decode::<HandoffSessionParams>(params)?;
1276 let session = load_session(¶ms.locator).map_err(operation)?;
1277 let cwd = params
1278 .cwd
1279 .or_else(|| session.meta.cwd.clone())
1280 .unwrap_or_else(|| PathBuf::from("."));
1281 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1282 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1283 ServiceError::Operation(
1284 "handoff artifact omitted target session identity".into(),
1285 )
1286 })?;
1287 let instructions =
1288 handoff_instructions(params.target_harness, target_session_id, &cwd);
1289 Ok(json!({
1290 "artifact": artifact,
1291 "launch": instructions.launch,
1292 "materialize": instructions.materialize,
1293 "requires_materialization": instructions.requires_materialization,
1294 "note": instructions.note,
1295 }))
1296 }
1297 "harness.v1.sessions.materialize" => {
1298 let params = decode::<MaterializeSessionParams>(params)?;
1299 for file in ¶ms.artifact.files {
1303 if file.role == "source_recovery"
1304 && file.path == "recovery/source.supercode.jsonl"
1305 {
1306 if let Ok(source) = Session::from_native_str(&file.content) {
1307 crate::residue_store::store_segments(&source);
1308 }
1309 }
1310 }
1311 let locator = crate::native_materialize::materialize_native_artifact(
1312 params.artifact,
1313 ¶ms.cwd,
1314 ¶ms.homes,
1315 )
1316 .map_err(ServiceError::Operation)?;
1317 Ok(json!({"locator": locator}))
1318 }
1319 "harness.v1.jobs.list" => {
1323 let query = decode::<crate::jobs::JobsQuery>(params)?;
1324 if let Some(harness) = query.harness.as_deref() {
1325 refuse_harness_without_jobs(harness, "jobs.list")?;
1326 }
1327 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1328 serde_json::to_value(listing)
1329 .map_err(|error| ServiceError::Operation(error.to_string()))
1330 }
1331 "harness.v1.jobs.get" => {
1332 let params = decode::<JobsGetParams>(params)?;
1333 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1334 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1335 .map_err(operation)?
1336 {
1337 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1338 None => Err(ServiceError::Operation(format!(
1339 "`{}` has no scheduled job `{}`",
1340 params.harness, params.id
1341 ))),
1342 }
1343 }
1344 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1350 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1351 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1352 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1353 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1354 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1355 "harness.v1.jobs.notepad"
1356 | "harness.v1.jobs.notepad_set"
1357 | "harness.v1.jobs.notepad_delete" => {
1358 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1359 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1360 let answer = match method {
1361 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1362 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1363 _ => crate::jobs_notepad::read(&request),
1364 }
1365 .map_err(job_control_error)?;
1366 serde_json::to_value(answer)
1367 .map_err(|error| ServiceError::Operation(error.to_string()))
1368 }
1369 "harness.v1.runs.list" => {
1373 let query = decode::<crate::runs::RunsQuery>(params)?;
1374 if let Some(harness) = query.harness.as_deref() {
1375 refuse_harness_without_runs(harness, "runs.list")?;
1376 }
1377 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1378 serde_json::to_value(listing)
1379 .map_err(|error| ServiceError::Operation(error.to_string()))
1380 }
1381 "harness.v1.runs.get" => {
1382 let params = decode::<RunsGetParams>(params)?;
1383 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1384 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1385 .map_err(operation)?
1386 {
1387 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1388 None => Err(ServiceError::Operation(format!(
1389 "`{}` has no run `{}`",
1390 params.harness, params.id
1391 ))),
1392 }
1393 }
1394 "harness.v1.sessions.resume_instructions" => {
1395 let params = decode::<ResumeInstructionsParams>(params)?;
1396 let session = load_session(¶ms.locator).map_err(operation)?;
1397 let cwd = params
1398 .cwd
1399 .or(session.meta.cwd)
1400 .unwrap_or_else(|| PathBuf::from("."));
1401 let launch = resume_launch(
1402 params.locator.harness.as_str(),
1403 ¶ms.locator.session_id,
1404 &cwd,
1405 params.policy,
1406 )?;
1407 Ok(json!({"launch": launch}))
1408 }
1409 _ => Err(ServiceError::MethodNotFound),
1410 }
1411 }
1412
1413 fn reduce_session(
1414 &self,
1415 params: ReduceSessionParams,
1416 ) -> std::result::Result<Value, ServiceError> {
1417 let session = load_session(¶ms.locator).map_err(operation)?;
1418 if session.messages.is_empty() {
1419 return Err(ServiceError::InvalidParams(
1420 "cannot reduce an empty session".into(),
1421 ));
1422 }
1423 let keep_last = params.keep_last.clamp(1, 128);
1424 let policy = reduce::ReductionPolicy {
1425 clear_turns_older_than: Some(keep_last),
1426 ..Default::default()
1427 };
1428 let (view, log) =
1429 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1430 if log.reductions.is_empty() {
1431 return Err(ServiceError::UnsupportedAction(format!(
1432 "session `{}` is already too small for a meaningful reversible reduction",
1433 params.locator.session_id
1434 )));
1435 }
1436 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1437 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1438 if reduced_tokens >= source_tokens {
1439 return Err(ServiceError::UnsupportedAction(format!(
1440 "session `{}` has no token-reducing reversible projection",
1441 params.locator.session_id
1442 )));
1443 }
1444
1445 let store_root = self
1446 .reduction_store_root
1447 .clone()
1448 .unwrap_or_else(default_reduction_store_root);
1449 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1450 let rescue_id = format!("rescue-{}", generated_session_id());
1451 let imported = session
1452 .imported_message_count
1453 .unwrap_or(session.messages.len())
1454 .min(session.messages.len());
1455 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1456 let view_jsonl = messages_jsonl(&view)?;
1457 let title = format!(
1458 "Reduced {} continuation from {}",
1459 params.target_harness.id(),
1460 params.locator.session_id
1461 );
1462
1463 store
1468 .save_sidecar(&rescue_id, &sidecar_jsonl)
1469 .map_err(operation)?;
1470 store
1471 .save_reduction_log(&rescue_id, &log)
1472 .map_err(operation)?;
1473 store
1474 .save(&rescue_id, &title, &view_jsonl)
1475 .map_err(operation)?;
1476
1477 let source_bytes = serde_json::to_vec(&session.messages)
1478 .map_err(|error| ServiceError::Operation(error.to_string()))?
1479 .len() as u64;
1480 let reduced_bytes = serde_json::to_vec(&view)
1481 .map_err(|error| ServiceError::Operation(error.to_string()))?
1482 .len() as u64;
1483 store
1484 .set_reduction_stats(
1485 &rescue_id,
1486 &title,
1487 source_bytes,
1488 reduced_bytes,
1489 log.reductions.len() as u32,
1490 )
1491 .map_err(operation)?;
1492
1493 let reloaded_sidecar = store
1497 .load_sidecar(&rescue_id)
1498 .map_err(operation)?
1499 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1500 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1501 let reloaded_log = store
1502 .load_reduction_log(&rescue_id)
1503 .map_err(operation)?
1504 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1505 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1506 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1507 let (restamped_view, restamped_log) =
1514 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1515 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1516 return Err(ServiceError::Operation(
1517 "persisted reduction view does not match its durable log and sidecar".into(),
1518 ));
1519 }
1520 if restamped_log != reloaded_log {
1521 return Err(ServiceError::Operation(
1522 "reapplying the durable reduction log changed its identity".into(),
1523 ));
1524 }
1525 let inverted =
1526 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1527 if inverted != session.messages {
1528 return Err(ServiceError::Operation(
1529 "reduction inversion did not restore the source messages byte-exactly".into(),
1530 ));
1531 }
1532
1533 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1534 let sidecar_path = store.sidecar_path(&rescue_id);
1535 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1536 let bootstrap_prompt = reduced_bootstrap_prompt(
1537 ¶ms.locator,
1538 params.target_harness,
1539 &view_jsonl,
1540 &sidecar_path,
1541 &reduction_log_path,
1542 );
1543 let mut reduced_session = session.clone();
1544 reduced_session.meta.session_id = Some(rescue_id.clone());
1545 reduced_session.messages = view;
1546
1547 Ok(json!({
1548 "session": normalized_session_json(&reduced_session),
1549 "bootstrap_prompt": bootstrap_prompt,
1550 "receipt": {
1551 "id": rescue_id,
1552 "sidecar_id": rescue_id,
1553 "source_harness": params.locator.harness,
1554 "target_harness": params.target_harness.id(),
1555 "source_tokens": source_tokens,
1556 "reduced_tokens": reduced_tokens,
1557 "ratio": ratio,
1558 "source_bytes": source_bytes,
1559 "reduced_bytes": reduced_bytes,
1560 "reductions": reloaded_log.reductions.len(),
1561 "sidecar_path": sidecar_path,
1562 "reduction_log_path": reduction_log_path,
1563 "verified": true,
1564 "reversible": true,
1565 }
1566 }))
1567 }
1568
1569 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1587 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1588 return None;
1589 }
1590 let method = request.get("method").and_then(Value::as_str)?;
1591 if !RUNTIME_OPEN_METHODS.contains(&method) {
1592 return None;
1593 }
1594 Some(RuntimeOpen {
1595 id: request.get("id").cloned().unwrap_or(Value::Null),
1596 method: method.to_string(),
1597 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1598 })
1599 }
1600
1601 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1623 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1624 return None;
1625 }
1626 let method = request.get("method").and_then(Value::as_str)?;
1627 if !DETACHED_METHODS.contains(&method) {
1628 return None;
1629 }
1630 let id = request.get("id").cloned().unwrap_or(Value::Null);
1631 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1632 let work = match method {
1633 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1634 .inventory_work(method, params)
1635 .map(DetachedWork::Inventory),
1636 "harness.v1.sessions.message" => {
1637 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1638 }
1639 _ => {
1640 let verb = match method {
1641 "harness.v1.sessions.new" => crate::SessionVerb::New,
1642 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1643 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1644 _ => crate::SessionVerb::Delete,
1645 };
1646 match decode::<crate::SessionMutation>(params) {
1647 Ok(mutation) => {
1648 match crate::sessions_control::door(&mutation.harness, verb) {
1649 Ok(crate::SessionDoor::Live(_)) => return None,
1652 Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1653 Err(error) => Err(session_control_error(error)),
1654 }
1655 }
1656 Err(error) => Err(error),
1657 }
1658 }
1659 };
1660 Some(DetachedCall {
1661 id,
1662 method: method.to_string(),
1663 work: work.map(Work::Free),
1664 })
1665 }
1666
1667 pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1681 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1682 return None;
1683 }
1684 let method = request.get("method").and_then(Value::as_str)?;
1685 let id = request.get("id").cloned().unwrap_or(Value::Null);
1686 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1687 let work = match method {
1688 "harness.v1.runtimes.close" => decode::<RuntimeConnectionParams>(params)
1689 .and_then(|params| self.surrender_runtime(¶ms.connection))
1690 .map(|(runtime, process_group)| {
1691 Work::Runtime(RuntimeWork::Close {
1692 runtime,
1693 process_group,
1694 })
1695 }),
1696 "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1697 let verb = if method == "harness.v1.sessions.new" {
1698 crate::SessionVerb::New
1699 } else {
1700 crate::SessionVerb::Reset
1701 };
1702 let mutation = decode::<crate::SessionMutation>(params).ok()?;
1703 let Ok(crate::SessionDoor::Live(command)) =
1707 crate::sessions_control::door(&mutation.harness, verb)
1708 else {
1709 return None;
1710 };
1711 let connection = mutation
1712 .connection
1713 .clone()
1714 .filter(|value| !value.trim().is_empty())?;
1715 self.lend_runtime(&connection).map(|runtime| {
1716 let session = live_session_name(runtime.as_ref(), &mutation);
1717 Work::Runtime(RuntimeWork::LiveCommand {
1718 connection,
1719 runtime,
1720 verb,
1721 mutation,
1722 command,
1723 session,
1724 })
1725 })
1726 }
1727 _ => return None,
1728 };
1729 Some(DetachedCall {
1730 id,
1731 method: method.to_string(),
1732 work,
1733 })
1734 }
1735
1736 pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1741 let DetachedAnswer { response, returned } = answer;
1742 if let Some(ReturnedRuntime {
1743 connection,
1744 runtime,
1745 }) = returned
1746 {
1747 self.runtimes_in_flight.remove(&connection);
1748 self.runtimes.insert(connection, runtime);
1749 }
1750 response
1751 }
1752
1753 pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1757 let OpenedRuntime { id, outcome } = opened;
1758 let result = match outcome {
1759 Ok(open) => self.register_open_runtime(open).await,
1760 Err(error) => Err(error),
1761 };
1762 service_response(id, result)
1763 }
1764
1765 async fn register_open_runtime(
1767 &mut self,
1768 open: OpenRuntime,
1769 ) -> std::result::Result<Value, ServiceError> {
1770 match open {
1771 OpenRuntime::Hosted {
1772 runtime,
1773 capabilities,
1774 workspace,
1775 } => {
1776 self.insert_hosted_runtime(runtime, capabilities, workspace)
1777 .await
1778 }
1779 OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1780 }
1781 }
1782
1783 async fn runtime_call(
1784 &mut self,
1785 method: &str,
1786 params: Value,
1787 ) -> std::result::Result<Value, ServiceError> {
1788 match method {
1789 "harness.v1.runtimes.capabilities" => {
1790 let params = decode::<RuntimeBackendParams>(params)?;
1791 let backend = runtime_backend(¶ms)?;
1792 Ok(json!({
1793 "harness": backend.harness(),
1794 "capabilities": backend.capabilities(),
1795 }))
1796 }
1797 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1798 self.register_open_runtime(open_runtime(method, params).await?)
1799 .await
1800 }
1801 "harness.v1.runtimes.send_input" => {
1802 let params = decode::<RuntimeInputParams>(params)?;
1803 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1804 let runtime = self.runtime_mut(¶ms.connection)?;
1805 let turn_id = within_control_deadline(
1806 method,
1807 runtime.send_input(RuntimeInput {
1808 text: params.text,
1809 image_urls,
1810 }),
1811 )
1812 .await?
1813 .map_err(operation)?;
1814 Ok(json!({"turn_id": turn_id}))
1815 }
1816 "harness.v1.runtimes.interrupt" => {
1817 let params = decode::<RuntimeConnectionParams>(params)?;
1818 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1819 .await?
1820 .map_err(operation)?;
1821 Ok(json!({}))
1822 }
1823 "harness.v1.runtimes.steer" => {
1824 let params = decode::<RuntimeInputParams>(params)?;
1825 if !params.image_urls.is_empty() {
1826 return Err(ServiceError::InvalidParams(
1827 "runtime steering accepts text only".into(),
1828 ));
1829 }
1830 let text = params.text.trim();
1831 if text.is_empty() || text.chars().count() > 50_000 {
1832 return Err(ServiceError::InvalidParams(
1833 "runtime steering requires 1 to 50,000 text characters".into(),
1834 ));
1835 }
1836 within_control_deadline(
1837 method,
1838 self.runtime_mut(¶ms.connection)?
1839 .steer(text.to_string()),
1840 )
1841 .await?
1842 .map_err(operation)?;
1843 Ok(json!({}))
1844 }
1845 "harness.v1.runtimes.respond" => {
1846 let params = decode::<RuntimeRespondParams>(params)?;
1847 let request_id = params.request_id.clone();
1848 within_control_deadline(
1849 method,
1850 self.runtime_mut(¶ms.connection)?
1851 .respond(params.request_id, params.response),
1852 )
1853 .await?
1854 .map_err(operation)?;
1855 self.approvals.answered(¶ms.connection, &request_id);
1857 Ok(json!({}))
1858 }
1859 "harness.v1.runtimes.acquire_control" => {
1860 let params = decode::<RuntimeConnectionParams>(params)?;
1861 let snapshot = within_control_deadline(
1862 method,
1863 self.runtime_mut(¶ms.connection)?.acquire_control(),
1864 )
1865 .await?
1866 .map_err(operation)?;
1867 serde_json::to_value(snapshot)
1868 .map_err(|error| ServiceError::Operation(error.to_string()))
1869 }
1870 "harness.v1.runtimes.heartbeat" => {
1871 let params = decode::<RuntimeConnectionParams>(params)?;
1872 let snapshot = within_control_deadline(
1873 method,
1874 self.runtime_mut(¶ms.connection)?.heartbeat(),
1875 )
1876 .await?
1877 .map_err(operation)?;
1878 serde_json::to_value(snapshot)
1879 .map_err(|error| ServiceError::Operation(error.to_string()))
1880 }
1881 "harness.v1.runtimes.detach" => {
1882 let params = decode::<RuntimeConnectionParams>(params)?;
1883 let snapshot =
1884 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1885 .await?
1886 .map_err(operation)?;
1887 serde_json::to_value(snapshot)
1888 .map_err(|error| ServiceError::Operation(error.to_string()))
1889 }
1890 "harness.v1.runtimes.terminal_instructions" => {
1891 let params = decode::<RuntimeConnectionParams>(params)?;
1892 let launch = self
1893 .terminal_launches
1894 .get(¶ms.connection)
1895 .ok_or_else(|| {
1896 ServiceError::Operation(
1897 "this runtime is not hosted for terminal attachment".into(),
1898 )
1899 })?;
1900 Ok(json!({"launch":launch}))
1901 }
1902 "harness.v1.runtimes.close" => {
1903 let params = decode::<RuntimeConnectionParams>(params)?;
1904 let (runtime, process_group) = self.surrender_runtime(¶ms.connection)?;
1905 close_runtime(runtime, process_group).await
1906 }
1907 _ => Err(ServiceError::MethodNotFound),
1908 }
1909 }
1910
1911 #[cfg(feature = "adapter-api")]
1913 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1914 let params = decode::<MessageSessionParams>(params)?;
1915 Ok(message_live_session(¶ms).await)
1916 }
1917
1918 #[cfg(feature = "adapter-api")]
1919 fn harness_settings_call(
1920 &self,
1921 method: &str,
1922 params: Value,
1923 ) -> std::result::Result<Value, ServiceError> {
1924 let homes = crate::HarnessHomes::default();
1925 match method {
1926 "harness.v1.harnesses.settings" => {
1927 let params = decode::<HarnessSettingsParams>(params)?;
1928 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1929 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1930 serde_json::to_value(report)
1931 .map_err(|error| ServiceError::Operation(error.to_string()))
1932 }
1933 "harness.v1.harnesses.configure" => {
1934 let params = decode::<ConfigureHarnessParams>(params)?;
1935 let report = crate::configure_harness_interop_settings(
1936 &homes,
1937 ¶ms.harness,
1938 ¶ms.changes,
1939 params.expected_revision.as_deref(),
1940 )
1941 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1942 serde_json::to_value(report)
1943 .map_err(|error| ServiceError::Operation(error.to_string()))
1944 }
1945 _ => Err(ServiceError::MethodNotFound),
1946 }
1947 }
1948
1949 fn insert_runtime(
1950 &mut self,
1951 runtime: Box<dyn RuntimeConnection>,
1952 ) -> std::result::Result<Value, ServiceError> {
1953 let connection = format!("runtime-{}", self.next_runtime);
1954 self.next_runtime += 1;
1955 let handle = runtime.handle().clone();
1956 self.runtime_sequences
1957 .entry(handle.runtime_id.clone())
1958 .or_insert(0);
1959 self.runtimes.insert(connection.clone(), runtime);
1960 Ok(json!({"connection": connection, "handle": handle}))
1961 }
1962
1963 #[cfg(feature = "adapter-api")]
1964 async fn insert_hosted_runtime(
1965 &mut self,
1966 runtime: Box<dyn RuntimeConnection>,
1967 capabilities: crate::RuntimeCapabilities,
1968 workspace: PathBuf,
1969 ) -> std::result::Result<Value, ServiceError> {
1970 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities);
1971 let token: std::sync::Arc<str> = crate::server::generate_token().into();
1972 let server = crate::server::run_frontend_http(
1973 host.clone(),
1974 host.frontend_sender(),
1975 "127.0.0.1:0",
1976 token.clone(),
1977 connection.handle().runtime_id.clone(),
1978 )
1979 .await
1980 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1981 let source = LiveRuntimeSource {
1982 harness: connection.handle().harness.as_str().to_string(),
1983 session_id: connection.handle().runtime_id.clone(),
1984 workspace: workspace.clone(),
1985 };
1986 let registration = register_live_runtime(
1987 connection.handle().runtime_id.clone(),
1988 source.clone(),
1989 format!("http://{}", server.address()),
1990 token.to_string(),
1991 )
1992 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1993 let endpoint = registration.endpoint().to_string();
1994 let launch = StructuredLaunch {
1995 cwd: workspace,
1996 program: std::env::current_exe()
2000 .ok()
2001 .map(|path| path.to_string_lossy().into_owned())
2002 .unwrap_or_else(|| "supercode".into()),
2003 arguments: vec![
2004 "open".into(),
2005 endpoint,
2006 "--harness".into(),
2007 source.harness,
2008 "--session".into(),
2009 source.session_id,
2010 ],
2011 env: BTreeMap::new(),
2012 };
2013 let lease = HostedRuntimeLease {
2014 connection,
2015 _host: host,
2016 _registration: registration,
2017 _server: server,
2018 };
2019 let opened = self.insert_runtime(Box::new(lease))?;
2020 let connection_id = opened["connection"]
2021 .as_str()
2022 .expect("insert_runtime returns a connection id")
2023 .to_string();
2024 self.terminal_launches.insert(connection_id, launch);
2025 Ok(opened)
2026 }
2027
2028 #[cfg(not(feature = "adapter-api"))]
2029 async fn insert_hosted_runtime(
2030 &mut self,
2031 runtime: Box<dyn RuntimeConnection>,
2032 _capabilities: crate::RuntimeCapabilities,
2033 _workspace: PathBuf,
2034 ) -> std::result::Result<Value, ServiceError> {
2035 self.insert_runtime(runtime)
2036 }
2037
2038 fn runtime_mut(
2039 &mut self,
2040 connection: &str,
2041 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2042 if self.runtimes_in_flight.contains(connection) {
2043 return Err(self.lent_out(connection));
2044 }
2045 self.runtimes.get_mut(connection).ok_or_else(|| {
2046 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2047 })
2048 }
2049
2050 fn lent_out(&self, connection: &str) -> ServiceError {
2054 ServiceError::Operation(format!(
2055 "runtime connection `{connection}`: a harness turn is already in progress"
2056 ))
2057 }
2058
2059 fn lend_runtime(
2062 &mut self,
2063 connection: &str,
2064 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2065 if self.runtimes_in_flight.contains(connection) {
2066 return Err(self.lent_out(connection));
2067 }
2068 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2069 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2070 })?;
2071 self.runtimes_in_flight.insert(connection.to_string());
2072 Ok(runtime)
2073 }
2074
2075 fn surrender_runtime(
2087 &mut self,
2088 connection: &str,
2089 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2090 if self.runtimes_in_flight.contains(connection) {
2091 return Err(self.lent_out(connection));
2092 }
2093 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2094 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2095 })?;
2096 let process_group = runtime_process_group(runtime.handle());
2097 let runtime_id = runtime.handle().runtime_id.clone();
2098 self.terminal_launches.remove(connection);
2099 self.runtime_sequences.remove(&runtime_id);
2100 self.approvals.forget(connection);
2101 Ok((runtime, process_group))
2102 }
2103
2104 pub fn kill_all_runtime_groups(&self) -> usize {
2115 self.runtimes
2116 .values()
2117 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2118 .count()
2119 }
2120
2121 async fn mutate_session(
2132 &mut self,
2133 verb: crate::SessionVerb,
2134 params: Value,
2135 ) -> std::result::Result<Value, ServiceError> {
2136 let mutation = decode::<crate::SessionMutation>(params)?;
2137 let door = crate::sessions_control::door(&mutation.harness, verb)
2138 .map_err(session_control_error)?;
2139 let outcome = match door {
2140 #[cfg(not(feature = "adapter-api"))]
2144 crate::SessionDoor::Live(command) => {
2145 return Err(ServiceError::Operation(format!(
2146 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2147 session, which needs this build's `adapter-api` feature",
2148 mutation.harness,
2149 verb.as_str()
2150 )));
2151 }
2152 #[cfg(feature = "adapter-api")]
2153 crate::SessionDoor::Live(command) => {
2154 let connection = mutation
2155 .connection
2156 .clone()
2157 .filter(|value| !value.trim().is_empty())
2158 .ok_or_else(|| {
2159 ServiceError::InvalidParams(format!(
2160 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2161 driven session: pass the `connection` of an open runtime \
2162 (`harness.v1.runtimes.start`)",
2163 mutation.harness,
2164 verb.as_str()
2165 ))
2166 })?;
2167 let runtime = self.runtime_mut(&connection)?;
2168 let session = live_session_name(runtime.as_ref(), &mutation);
2169 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2175 .await;
2176 }
2177 _ => run_session_mutation(verb, &mutation).await?,
2178 };
2179 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2180 }
2181
2182 async fn inventory_call(
2186 &self,
2187 method: &str,
2188 params: Value,
2189 ) -> std::result::Result<Value, ServiceError> {
2190 run_inventory(self.inventory_work(method, params)?).await
2191 }
2192
2193 fn inventory_work(
2199 &self,
2200 method: &str,
2201 params: Value,
2202 ) -> std::result::Result<InventoryWork, ServiceError> {
2203 let mut params = decode::<HarnessInventoryParams>(params)?;
2204 if method == "harness.v1.harnesses.probe" {
2205 let harness = params.harness.take().ok_or_else(|| {
2206 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2207 })?;
2208 params.harnesses = vec![harness];
2209 }
2210 let selected = params
2211 .harnesses
2212 .iter()
2213 .map(HarnessId::as_str)
2214 .collect::<std::collections::BTreeSet<_>>();
2215 let supported = harness_support_registry()
2216 .harnesses
2217 .into_iter()
2218 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2219 .collect::<Vec<_>>();
2220 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2221 let known = supported
2222 .iter()
2223 .map(|harness| harness.id.as_str())
2224 .collect::<std::collections::BTreeSet<_>>();
2225 let missing = params
2226 .harnesses
2227 .iter()
2228 .filter(|id| !known.contains(id.as_str()))
2229 .map(HarnessId::as_str)
2230 .collect::<Vec<_>>();
2231 return Err(ServiceError::InvalidParams(format!(
2232 "unknown harness(es): {}",
2233 missing.join(", ")
2234 )));
2235 }
2236 let global_counts = params
2237 .include_sessions
2238 .then(|| self.session_counts(None, ¶ms.harnesses));
2239 let workspace_counts = params
2240 .include_sessions
2241 .then(|| {
2242 params
2243 .workspace
2244 .as_deref()
2245 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2246 })
2247 .flatten();
2248 Ok(InventoryWork {
2249 params,
2250 supported,
2251 global_counts,
2252 workspace_counts,
2253 })
2254 }
2255
2256 #[cfg(feature = "adapter-api")]
2257 async fn harness_authentication_call(
2258 &self,
2259 method: &str,
2260 params: Value,
2261 ) -> std::result::Result<Value, ServiceError> {
2262 match method {
2263 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2264 let params = decode::<HarnessAuthenticationParams>(params)?;
2265 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2266 .map_err(|error| ServiceError::Operation(error.to_string()))
2267 }
2268 "harness.v1.harnesses.auth.begin" => {
2269 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2270 let cwd = params
2271 .cwd
2272 .or_else(|| std::env::current_dir().ok())
2273 .unwrap_or_else(|| PathBuf::from("."));
2274 let plan = crate::harness_authentication_plan(
2275 ¶ms.harness,
2276 params.environment,
2277 params.method,
2278 &cwd,
2279 )
2280 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2281 serde_json::to_value(plan)
2282 .map_err(|error| ServiceError::Operation(error.to_string()))
2283 }
2284 _ => Err(ServiceError::MethodNotFound),
2285 }
2286 }
2287
2288 fn session_counts(
2289 &self,
2290 workspace: Option<&Path>,
2291 harnesses: &[HarnessId],
2292 ) -> BTreeMap<String, usize> {
2293 let mut counts = BTreeMap::new();
2294 for session in self
2295 .catalog
2296 .discover(&DiscoveryQuery {
2297 workspace: workspace.map(Path::to_path_buf),
2298 harnesses: harnesses.to_vec(),
2299 ..DiscoveryQuery::default()
2300 })
2301 .unwrap_or_default()
2302 {
2303 *counts
2304 .entry(session.locator.harness.as_str().to_string())
2305 .or_insert(0) += 1;
2306 }
2307 counts
2308 }
2309}
2310
2311#[async_trait::async_trait]
2312impl SdkService for HarnessSessionService {
2313 fn capabilities(&self) -> SdkCapabilities {
2314 SdkCapabilities::default()
2315 }
2316
2317 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2318 if request.operation == SdkOperation::Events {
2319 let events = self
2320 .poll_sdk_events()
2321 .await
2322 .into_iter()
2323 .map(|(_, event)| event)
2324 .collect::<Vec<_>>();
2325 return serde_json::to_value(events).map_err(|error| {
2326 SdkError::new(
2327 SdkErrorCode::Execution,
2328 request.operation,
2329 error.to_string(),
2330 )
2331 });
2332 }
2333 if self.runtimes.is_empty()
2334 && matches!(
2335 request.operation,
2336 SdkOperation::Input
2337 | SdkOperation::Interrupt
2338 | SdkOperation::Steer
2339 | SdkOperation::Respond
2340 | SdkOperation::Close
2341 )
2342 {
2343 return Err(SdkError::unsupported(request.operation));
2344 }
2345 let method = request
2346 .operation
2347 .method()
2348 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2349 let result = match request.operation {
2350 SdkOperation::Discover
2351 | SdkOperation::Load
2352 | SdkOperation::Export
2353 | SdkOperation::ProfilesList
2354 | SdkOperation::ProfilesGet
2355 | SdkOperation::ProfilesCreate
2356 | SdkOperation::ProfilesDelete
2357 | SdkOperation::SkillsList
2358 | SdkOperation::SkillsInstall
2359 | SdkOperation::SkillsRemove
2360 | SdkOperation::ChannelsList
2361 | SdkOperation::RoutesList
2362 | SdkOperation::TriggersList
2363 | SdkOperation::ChannelsStatus
2364 | SdkOperation::MemoryShow
2365 | SdkOperation::MemorySearch
2366 | SdkOperation::JobsList
2367 | SdkOperation::JobsGet
2368 | SdkOperation::JobsCreate
2369 | SdkOperation::JobsUpdate
2370 | SdkOperation::JobsPause
2371 | SdkOperation::JobsResume
2372 | SdkOperation::JobsRun
2373 | SdkOperation::JobsDelete
2374 | SdkOperation::JobsNotepad
2375 | SdkOperation::JobsNotepadSet
2376 | SdkOperation::JobsNotepadDelete
2377 | SdkOperation::RunsList
2378 | SdkOperation::RunsGet
2379 | SdkOperation::ApprovalsList
2380 | SdkOperation::OrchestrationLoad
2381 | SdkOperation::OrchestrationSave
2382 | SdkOperation::OrchestrationCompile
2383 | SdkOperation::OrchestrationDecompile
2384 | SdkOperation::OrchestrationImport
2385 | SdkOperation::OrchestrationExport
2386 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2387 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2390 SdkOperation::Start
2391 | SdkOperation::Resume
2392 | SdkOperation::Input
2393 | SdkOperation::Interrupt
2394 | SdkOperation::Steer
2395 | SdkOperation::Respond
2396 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2397 SdkOperation::SessionsNew => {
2402 self.mutate_session(crate::SessionVerb::New, request.params)
2403 .await
2404 }
2405 SdkOperation::SessionsReset => {
2406 self.mutate_session(crate::SessionVerb::Reset, request.params)
2407 .await
2408 }
2409 SdkOperation::SessionsArchive => {
2410 self.mutate_session(crate::SessionVerb::Archive, request.params)
2411 .await
2412 }
2413 SdkOperation::SessionsDelete => {
2414 self.mutate_session(crate::SessionVerb::Delete, request.params)
2415 .await
2416 }
2417 SdkOperation::Events => unreachable!("handled before method dispatch"),
2418 };
2419 result.map_err(|error| sdk_error(request.operation, error))
2420 }
2421
2422 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2423 Ok(self
2424 .poll_sdk_events()
2425 .await
2426 .into_iter()
2427 .map(|(_, event)| event)
2428 .collect())
2429 }
2430}
2431
2432#[cfg(feature = "adapter-api")]
2433struct HostedRuntimeLease {
2434 connection: HostedHarnessConnection,
2435 _host: std::sync::Arc<HostedHarnessRuntime>,
2436 _registration: LiveRuntimeRegistration,
2437 _server: crate::server::FrontendHttpServer,
2438}
2439
2440#[async_trait::async_trait]
2441#[cfg(feature = "adapter-api")]
2442impl RuntimeConnection for HostedRuntimeLease {
2443 fn handle(&self) -> &crate::RuntimeHandle {
2444 self.connection.handle()
2445 }
2446
2447 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2448 self.connection.send_input(input).await
2449 }
2450
2451 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2452 self.connection.next_event().await
2453 }
2454
2455 async fn interrupt(&mut self) -> crate::Result<()> {
2456 self.connection.interrupt().await
2457 }
2458
2459 async fn steer(&mut self, text: String) -> crate::Result<()> {
2462 self.connection.steer(text).await
2463 }
2464
2465 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2466 self.connection.respond(request_id, response).await
2467 }
2468
2469 async fn close(&mut self) -> crate::Result<()> {
2470 self.connection.close().await
2471 }
2472}
2473
2474struct InventoryWork {
2477 params: HarnessInventoryParams,
2478 supported: Vec<crate::HarnessSupportDescriptor>,
2479 global_counts: Option<BTreeMap<String, usize>>,
2480 workspace_counts: Option<BTreeMap<String, usize>>,
2481}
2482
2483async fn run_session_mutation(
2490 verb: crate::SessionVerb,
2491 mutation: &crate::SessionMutation,
2492) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2493 let door =
2500 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2501 if let crate::SessionDoor::Http = door {
2502 return crate::sessions_control::mutate(verb, mutation)
2503 .await
2504 .map_err(session_control_error);
2505 }
2506 let mutation = mutation.clone();
2507 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2508 .await
2509 .map_err(|error| {
2510 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2511 })?
2512 .map_err(session_control_error)
2513}
2514
2515async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2518 let InventoryWork {
2519 params,
2520 supported,
2521 global_counts,
2522 workspace_counts,
2523 } = work;
2524 let probes = supported.into_iter().map(|descriptor| {
2525 let global = global_counts
2526 .as_ref()
2527 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2528 let workspace = workspace_counts
2529 .as_ref()
2530 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2531 probe_harness(descriptor, ¶ms, global, workspace)
2532 });
2533 let harnesses = futures::future::join_all(probes).await;
2534 serde_json::to_value(HarnessInventoryReport {
2535 probe: params.probe,
2536 workspace: params.workspace,
2537 harnesses,
2538 })
2539 .map_err(|error| ServiceError::Operation(error.to_string()))
2540}
2541
2542async fn probe_harness(
2543 descriptor: crate::HarnessSupportDescriptor,
2544 params: &HarnessInventoryParams,
2545 global: Option<usize>,
2546 workspace: Option<usize>,
2547) -> LocalHarness {
2548 let launch = descriptor.runtime.default_launch.as_ref();
2549 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2554 .then(crate::orchestrator::daemon_entry)
2555 .and_then(Result::ok);
2556 let executable = match &orchestrator_entry {
2557 Some(entry) => Some(entry.clone()),
2558 None => launch.and_then(|launch| find_executable(&launch.program)),
2559 };
2560 let installed = executable.is_some();
2561 let version = if params.skip_versions || orchestrator_entry.is_some() {
2562 None
2565 } else {
2566 match executable.as_deref() {
2567 Some(path) => executable_version(path).await,
2568 None => None,
2569 }
2570 };
2571 let configured = auth_evidence(descriptor.id.as_str());
2572 let mut auth = if configured {
2573 HarnessAuthState::Configured
2574 } else if matches!(
2575 descriptor.id.as_str(),
2576 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2577 ) {
2578 HarnessAuthState::Required
2583 } else {
2584 HarnessAuthState::Unknown
2585 };
2586 let mut runtime = if installed {
2587 HarnessRuntimeState::Degraded
2588 } else {
2589 HarnessRuntimeState::Unavailable
2590 };
2591 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2592 let mut reason = (!installed).then(|| {
2593 if is_orchestrator {
2594 format!(
2595 "{} is supported but its daemon entry `{}` was not found",
2596 descriptor.display_name,
2597 crate::orchestrator::DAEMON_ENTRY
2598 )
2599 } else {
2600 format!(
2601 "{} is supported but `{}` was not found on PATH",
2602 descriptor.display_name,
2603 launch
2604 .map(|launch| launch.program.as_str())
2605 .unwrap_or("executable")
2606 )
2607 }
2608 });
2609 let mut repair = (!installed).then(|| {
2610 if is_orchestrator {
2611 format!(
2612 "Install the `supercode-orchestrator` package so `{}` resolves.",
2613 crate::orchestrator::DAEMON_ENTRY
2614 )
2615 } else {
2616 format!(
2617 "Install {} and ensure `{}` is on PATH.",
2618 descriptor.display_name,
2619 launch
2620 .map(|launch| launch.program.as_str())
2621 .unwrap_or("its executable")
2622 )
2623 }
2624 });
2625
2626 if installed && params.probe == HarnessProbeLevel::Handshake {
2627 let backend_params = RuntimeBackendParams {
2628 harness: descriptor.id.clone(),
2629 protocol: None,
2630 launch: None,
2631 base_url: None,
2632 policy: RuntimePolicy::Default,
2633 };
2634 match runtime_backend(&backend_params) {
2635 Ok(backend) => {
2636 let cwd = params
2637 .workspace
2638 .clone()
2639 .or_else(|| std::env::current_dir().ok())
2640 .unwrap_or_else(|| PathBuf::from("."));
2641 let isolated = descriptor
2642 .runtime
2643 .default_launch
2644 .clone()
2645 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2646 let Some(isolated) = isolated else {
2647 reason = Some(
2648 "No-prompt runtime handshake could not create its isolated harness home."
2649 .into(),
2650 );
2651 repair = Some(
2652 "Check temporary-directory permissions, then run the handshake probe again."
2653 .into(),
2654 );
2655 let running = probe_running_instance(descriptor.id.as_str());
2656 return LocalHarness {
2657 gateway: gateway_health(
2658 descriptor.id.as_str(),
2659 installed,
2660 running.as_ref(),
2661 version.as_deref(),
2662 ),
2663 id: descriptor.id,
2664 display_name: descriptor.display_name,
2665 supported: true,
2666 installed,
2667 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2668 version,
2669 auth,
2670 runtime,
2671 protocol: descriptor.runtime.protocol,
2672 capabilities: descriptor.runtime.capabilities.clone(),
2673 effective_capabilities: descriptor.runtime.capabilities,
2674 sessions: HarnessSessionCounts { global, workspace },
2675 running,
2676 reason,
2677 repair,
2678 };
2679 };
2680 match tokio::time::timeout(
2681 Duration::from_secs(30),
2682 backend.start(RuntimeStartRequest {
2683 cwd,
2684 launch: Some(isolated.launch.clone()),
2685 mcp_servers: Vec::new(),
2686 approval_policy: None,
2687 }),
2688 )
2689 .await
2690 {
2691 Ok(Ok(mut connection)) => {
2692 match stabilize_handshake(connection.as_mut()).await {
2693 Ok(()) => {
2694 auth = HarnessAuthState::Ready;
2695 runtime = HarnessRuntimeState::Ready;
2696 reason = Some(
2697 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2698 .into(),
2699 );
2700 repair = None;
2701 }
2702 Err(message) => {
2703 auth = if looks_like_auth_error(&message) {
2704 HarnessAuthState::Required
2705 } else if configured {
2706 HarnessAuthState::Configured
2707 } else {
2708 HarnessAuthState::Unknown
2709 };
2710 reason = Some(format!(
2711 "No-prompt runtime handshake became unhealthy during startup: {message}"
2712 ));
2713 repair = Some(if auth == HarnessAuthState::Required {
2714 format!(
2715 "Run `{}` interactively once and complete sign-in, then probe again.",
2716 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2717 )
2718 } else {
2719 "Run the harness directly to inspect its startup failure, then probe again."
2720 .into()
2721 });
2722 }
2723 }
2724 let _ =
2725 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2726 }
2727 Ok(Err(error)) => {
2728 let message = truncate_text(&error.to_string(), 500);
2729 auth = if looks_like_auth_error(&message) {
2730 HarnessAuthState::Required
2731 } else if configured {
2732 HarnessAuthState::Configured
2733 } else {
2734 HarnessAuthState::Unknown
2735 };
2736 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2737 repair = Some(if auth == HarnessAuthState::Required {
2738 format!(
2739 "Run `{}` interactively once and complete sign-in, then probe again.",
2740 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2741 )
2742 } else {
2743 "Check the harness installation and run the handshake probe again."
2744 .into()
2745 });
2746 }
2747 Err(_) => {
2748 reason =
2749 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2750 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2751 }
2752 }
2753 let _ = isolated.cleanup();
2762 tokio::time::sleep(Duration::from_millis(250)).await;
2763 if let Err(error) = isolated.cleanup() {
2764 auth = if configured {
2765 HarnessAuthState::Configured
2766 } else {
2767 HarnessAuthState::Unknown
2768 };
2769 runtime = HarnessRuntimeState::Degraded;
2770 reason = Some(format!(
2771 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2772 ));
2773 repair = Some(
2774 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2775 .into(),
2776 );
2777 }
2778 }
2779 Err(error) => {
2780 reason = Some(error_message(error));
2781 }
2782 }
2783 } else if installed && configured {
2784 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2785 } else if installed && auth == HarnessAuthState::Required {
2786 reason = Some("Executable found, but no native authentication evidence is present.".into());
2787 repair = Some(format!(
2788 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2789 descriptor.id.as_str()
2790 ));
2791 } else if installed {
2792 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2793 repair = Some(format!(
2794 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2795 launch
2796 .map(|launch| launch.program.as_str())
2797 .unwrap_or("the harness")
2798 ));
2799 }
2800
2801 let effective_capabilities = if installed {
2802 descriptor.runtime.capabilities.clone()
2803 } else {
2804 unavailable_capabilities()
2805 };
2806 let running = probe_running_instance(descriptor.id.as_str());
2807 LocalHarness {
2808 gateway: gateway_health(
2809 descriptor.id.as_str(),
2810 installed,
2811 running.as_ref(),
2812 version.as_deref(),
2813 ),
2814 id: descriptor.id,
2815 display_name: descriptor.display_name,
2816 supported: true,
2817 installed,
2818 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2819 version,
2820 auth,
2821 runtime,
2822 protocol: descriptor.runtime.protocol,
2823 capabilities: descriptor.runtime.capabilities,
2824 effective_capabilities,
2825 sessions: HarnessSessionCounts { global, workspace },
2826 running,
2827 reason,
2828 repair,
2829 }
2830}
2831
2832async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2833 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2834 loop {
2835 let now = tokio::time::Instant::now();
2836 if now >= deadline {
2837 return Ok(());
2838 }
2839 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2840 Err(_) => return Ok(()),
2841 Ok(Ok(Some(event))) => {
2842 if let Some(message) = handshake_event_failure(&event) {
2843 return Err(truncate_text(&message, 500));
2844 }
2845 }
2846 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2847 Ok(Err(error)) => return Err(error.to_string()),
2848 }
2849 }
2850}
2851
2852fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2853 let detail = event
2854 .payload
2855 .get("message")
2856 .or_else(|| event.payload.get("line"))
2857 .and_then(Value::as_str)
2858 .unwrap_or(event.kind.as_str());
2859 match event.kind.as_str() {
2860 "transport_closed" => Some("runtime transport closed during startup".into()),
2861 "transport_error" => Some(format!("runtime transport error: {detail}")),
2862 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2863 _ => None,
2868 }
2869}
2870
2871fn indexed_claude_window(
2872 locator: &SessionLocator,
2873 options: &SessionLoadOptions,
2874) -> std::result::Result<Option<Value>, ServiceError> {
2875 use supercode_interchange::session::ClaudeReadIndex;
2876 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2879 || options.include_subagents != Some(false)
2880 {
2881 return Ok(None);
2882 }
2883 let crate::StorageLocator::File { path } = &locator.storage else {
2884 return Ok(None);
2885 };
2886 if !ClaudeReadIndex::supports(path)
2887 .map_err(|error| ServiceError::Operation(error.to_string()))?
2888 {
2889 return Ok(None);
2890 }
2891 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2892 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2893 let total = index.len();
2894 let (offset, end) = projected_message_window(total, options);
2895 let session = index
2896 .read_messages(offset..end)
2897 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2898 let summary = index
2899 .read_summary()
2900 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2901 let selected_options = SessionLoadOptions {
2902 message_offset: None,
2903 message_limit: None,
2904 message_tail: None,
2905 ..options.clone()
2906 };
2907 let mut selected = projected_session_json(&session, &selected_options);
2908 selected["raw_record_count"] = json!(index.raw_record_count());
2909 Ok(Some(json!({
2910 "session": selected,
2911 "summary": projected_session_summary(&summary, options),
2912 "window": {
2913 "has_more": offset > 0 || end < total, "has_newer": end < total,
2914 "has_older": offset > 0, "newer_items": index.item_count(end..total),
2915 "offset": offset, "older_items": index.item_count(0..offset),
2916 "returned": end - offset, "total_messages": total,
2917 }
2918 })))
2919}
2920
2921fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
2922 let total_messages = session.messages.len();
2923 let (offset, end) = projected_message_window(total_messages, options);
2924 json!({
2925 "session": projected_session_json(session, options),
2926 "summary": projected_session_summary(session, options),
2927 "window": {
2928 "has_more": offset > 0 || end < total_messages,
2929 "has_newer": end < total_messages,
2930 "has_older": offset > 0,
2931 "newer_items": normalized_item_count(&session.messages[end..]),
2932 "offset": offset,
2933 "older_items": normalized_item_count(&session.messages[..offset]),
2934 "returned": end.saturating_sub(offset),
2935 "total_messages": total_messages,
2936 }
2937 })
2938}
2939
2940fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
2941 messages
2942 .iter()
2943 .map(|message| {
2944 let conversation = usize::from(
2945 matches!(message.role, Role::Assistant | Role::User)
2946 && message_has_content(message),
2947 );
2948 let tool_result =
2949 usize::from(message.role == Role::Tool && message_has_content(message));
2950 conversation + tool_result + message.tool_calls().len()
2951 })
2952 .sum()
2953}
2954
2955fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
2956 let mut conversational = session.messages.iter().filter(|message| {
2957 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
2958 });
2959 let first_message = conversational.clone().next();
2960 let last_message = conversational.next_back();
2961 let mut assistant = session
2962 .messages
2963 .iter()
2964 .filter(|message| message.role == Role::Assistant && message_has_content(message));
2965 let first_assistant_message = assistant.clone().next();
2966 let last_assistant_message = assistant.next_back();
2967 let end_of_turn = session
2968 .messages
2969 .iter()
2970 .rev()
2971 .find(|message| message.role != Role::System)
2972 .is_some_and(|message| {
2973 message.role == Role::Assistant
2974 && message_has_content(message)
2975 && message.tool_calls().is_empty()
2976 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
2978 });
2979 let project = |message: Option<&crate::ChatMessage>| {
2980 message.map(|message| project_inline_media(message_json(message), options))
2981 };
2982 json!({
2983 "end_of_turn": end_of_turn,
2984 "first_assistant_message": project(first_assistant_message),
2985 "first_message": project(first_message),
2986 "last_assistant_message": project(last_assistant_message),
2987 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
2988 "last_message": project(last_message),
2989 })
2990}
2991
2992fn message_has_content(message: &crate::ChatMessage) -> bool {
2993 message
2994 .content
2995 .as_deref()
2996 .is_some_and(|content| !content.trim().is_empty())
2997 || message
2998 .content_parts
2999 .as_ref()
3000 .is_some_and(|parts| !parts.is_empty())
3001}
3002
3003fn message_text(message: &crate::ChatMessage) -> String {
3004 if let Some(content) = &message.content {
3005 return content.clone();
3006 }
3007 message
3008 .content_parts
3009 .as_ref()
3010 .into_iter()
3011 .flatten()
3012 .filter_map(|part| part.get("text").and_then(Value::as_str))
3013 .collect::<Vec<_>>()
3014 .join("\n")
3015}
3016
3017fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3018 let (offset, end) = projected_message_window(session.messages.len(), options);
3019 let messages = session.messages[offset..end]
3020 .iter()
3021 .map(|message| project_inline_media(message_json(message), options))
3022 .collect::<Vec<_>>();
3023 let subagents = if options.include_subagents.unwrap_or(true) {
3024 let subagent_options = SessionLoadOptions {
3029 message_limit: None,
3030 message_offset: None,
3031 message_tail: None,
3032 ..options.clone()
3033 };
3034 session
3035 .subagents
3036 .iter()
3037 .map(|subagent| projected_session_json(subagent, &subagent_options))
3038 .collect::<Vec<_>>()
3039 } else {
3040 Vec::new()
3041 };
3042 json!({
3043 "source": match session.meta.source {
3044 SessionSource::ClaudeCode => "claude_code",
3045 SessionSource::Codex => "codex",
3046 SessionSource::Gemini => "gemini",
3047 SessionSource::Goose => "goose",
3048 SessionSource::Grok => "grok",
3049 SessionSource::Native => "native",
3050 SessionSource::OpenClaw => "openclaw",
3051 SessionSource::Hermes => "hermes",
3052 SessionSource::OpenCode => "opencode",
3053 SessionSource::Pi => "pi",
3054 },
3055 "session_id": session.meta.session_id,
3056 "ended_at": session.meta.ended_at,
3057 "end_reason": session.meta.end_reason,
3058 "model": session.meta.model,
3059 "cwd": session.meta.cwd,
3060 "system_prompt": session.meta.system_prompt,
3061 "agent_id": session.meta.agent_id,
3062 "parent_tool_use_id": session.meta.parent_tool_use_id,
3063 "lineage": session.meta.lineage,
3064 "messages": messages,
3065 "subagents": subagents,
3066 "raw_record_count": session.raw.len(),
3067 "parse_error_lines": session.parse_error_lines,
3068 })
3069}
3070
3071fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3072 if let Some(tail) = options.message_tail {
3073 return (total.saturating_sub(tail), total);
3074 }
3075 let offset = options.message_offset.unwrap_or(0).min(total);
3076 let end = options
3077 .message_limit
3078 .map(|limit| offset.saturating_add(limit).min(total))
3079 .unwrap_or(total);
3080 (offset, end)
3081}
3082
3083fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3084 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3085 return message;
3086 };
3087 for part in parts {
3088 let Some(url) = part
3089 .get("image_url")
3090 .and_then(|image| image.get("url"))
3091 .and_then(Value::as_str)
3092 else {
3093 continue;
3094 };
3095 let Some(rest) = url.strip_prefix("data:") else {
3096 continue;
3097 };
3098 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3099 continue;
3100 };
3101 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3102 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3103 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3104 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3105 || options
3106 .max_inline_media_bytes
3107 .is_some_and(|limit| decoded_bytes > limit);
3108 if should_elide {
3109 *part = json!({
3110 "type": "media_reference",
3111 "media_type": media_type,
3112 "encoding": "base64",
3113 "encoded_bytes": encoded.len(),
3114 "decoded_bytes": decoded_bytes,
3115 "omitted": true,
3116 });
3117 }
3118 }
3119 message
3120}
3121
3122#[derive(Deserialize)]
3123struct LocatorParams {
3124 locator: SessionLocator,
3125 #[serde(default)]
3138 fidelity: Option<Fidelity>,
3139 #[serde(default)]
3142 view: Option<SessionReadView>,
3143}
3144
3145#[derive(Deserialize)]
3146struct SessionReadView {
3147 #[serde(default)]
3150 tail_messages: Option<usize>,
3151 #[serde(default)]
3154 include_subagents: bool,
3155 #[serde(default)]
3157 display_history: bool,
3158 #[serde(default)]
3161 max_message_chars: Option<usize>,
3162}
3163
3164impl LocatorParams {
3165 fn read_fidelity(&self) -> Fidelity {
3166 self.fidelity.unwrap_or(Fidelity::Semantic)
3167 }
3168
3169 fn include_subagents(&self) -> bool {
3170 self.view
3171 .as_ref()
3172 .map(|view| view.include_subagents)
3173 .unwrap_or(true)
3174 }
3175
3176 fn tail_messages(&self) -> Option<usize> {
3177 self.view
3178 .as_ref()
3179 .and_then(|view| view.tail_messages)
3180 .map(|limit| limit.clamp(1, 5_000))
3181 }
3182
3183 fn display_history(&self) -> bool {
3184 self.view.as_ref().is_some_and(|view| view.display_history)
3185 }
3186
3187 fn max_message_chars(&self) -> Option<usize> {
3188 self.view
3189 .as_ref()
3190 .and_then(|view| view.max_message_chars)
3191 .map(|limit| limit.clamp(256, 64_000))
3192 }
3193
3194 fn bound_session(&self, session: &mut Session) {
3195 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3196 }
3197}
3198
3199#[derive(Debug, Clone, Copy, Default, Deserialize)]
3200#[serde(rename_all = "snake_case")]
3201enum InlineMediaMode {
3202 #[default]
3203 Full,
3204 Metadata,
3205}
3206
3207#[derive(Debug, Clone, Default, Deserialize)]
3208#[serde(default)]
3209struct SessionLoadOptions {
3210 include_subagents: Option<bool>,
3211 inline_media: InlineMediaMode,
3212 max_inline_media_bytes: Option<usize>,
3213 message_limit: Option<usize>,
3214 message_offset: Option<usize>,
3215 message_tail: Option<usize>,
3216}
3217
3218impl SessionLoadOptions {
3219 fn validate(&self) -> std::result::Result<(), ServiceError> {
3220 if self.message_tail.is_some()
3221 && (self.message_limit.is_some() || self.message_offset.is_some())
3222 {
3223 return Err(ServiceError::InvalidParams(
3224 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3225 .into(),
3226 ));
3227 }
3228 Ok(())
3229 }
3230}
3231
3232#[derive(Deserialize)]
3233struct LoadSessionParams {
3234 #[serde(flatten)]
3235 read: LocatorParams,
3236 #[serde(default)]
3237 options: Option<SessionLoadOptions>,
3238}
3239
3240#[derive(Deserialize)]
3241struct UnfollowParams {
3242 subscription: String,
3243}
3244
3245#[derive(Debug, Deserialize)]
3246#[serde(deny_unknown_fields)]
3247struct IndexResizeParams {
3248 subscription: String,
3249 limit: usize,
3250}
3251
3252#[derive(Deserialize)]
3253struct ActivitySubscribeParams {
3254 locators: Vec<SessionLocator>,
3255 #[serde(default)]
3256 homes: crate::HarnessHomes,
3257}
3258
3259#[derive(Deserialize)]
3260struct MessageSessionParams {
3261 locator: SessionLocator,
3262 text: String,
3263 #[serde(default)]
3267 from_name: Option<String>,
3268 #[serde(default)]
3270 in_reply_to: Option<String>,
3271 #[serde(default)]
3273 notify_when_idle: bool,
3274 #[serde(default)]
3279 as_user: bool,
3280 #[serde(default)]
3283 homes: crate::HarnessHomes,
3284}
3285
3286#[derive(Deserialize)]
3287#[serde(deny_unknown_fields)]
3288struct InboxParams {
3289 #[serde(default)]
3291 from_name: Option<String>,
3292 #[serde(default)]
3294 address: Option<String>,
3295 #[serde(default)]
3297 all: bool,
3298}
3299
3300#[derive(Deserialize)]
3301#[serde(deny_unknown_fields)]
3302struct HarnessSettingsParams {
3303 harness: String,
3304}
3305
3306#[derive(Deserialize)]
3307#[serde(deny_unknown_fields)]
3308struct ConfigureHarnessParams {
3309 harness: String,
3310 #[serde(default)]
3311 changes: Vec<crate::HarnessSettingChange>,
3312 #[serde(default)]
3313 expected_revision: Option<String>,
3314}
3315
3316fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3317 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3318 Ok(report) => (
3319 serde_json::to_value(report).unwrap_or(Value::Null),
3320 Value::Null,
3321 ),
3322 Err(error) => (
3323 Value::Null,
3324 Value::String(format!(
3325 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3326 )),
3327 ),
3328 }
3329}
3330
3331async fn message_live_session(params: &MessageSessionParams) -> Value {
3345 use crate::mail_route::{Delivered, NoDoor, Refused};
3346 let (inbound_controls, inbound_controls_error) =
3347 claude_inbound_controls_or_error(¶ms.homes);
3348 let refused = |reason: &str, message: String| {
3349 json!({
3350 "delivered_to_bus": false,
3351 "refusal": {"reason": reason, "message": message},
3352 "inbound_controls": inbound_controls,
3353 "inbound_controls_error": inbound_controls_error,
3354 })
3355 };
3356 if params.text.trim().is_empty() {
3357 return refused(
3358 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3359 "refusing to deliver an empty message".into(),
3360 );
3361 }
3362 let sender = match operator_address(params.from_name.as_deref()) {
3363 Ok(sender) => sender,
3364 Err(message) => return refused("invalid_sender", message),
3365 };
3366 let receiver = match crate::mailbox::MailAddress::new(
3367 crate::mailbox::local_machine_name(),
3368 params.locator.harness.as_str(),
3369 ¶ms.locator.session_id,
3370 ) {
3371 Ok(receiver) => receiver,
3372 Err(error) => return refused("delivery_failed", error.to_string()),
3373 };
3374 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3376 && crate::runtime_mail::controlled_runtime(
3377 HarnessId::CLAUDE_CODE,
3378 ¶ms.locator.session_id,
3379 )
3380 .is_none()
3381 {
3382 if let Err(refusal) =
3383 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3384 {
3385 return refused(refusal.reason.as_str(), refusal.message);
3386 }
3387 }
3388 if params.as_user {
3389 return message_as_user(
3390 params,
3391 sender,
3392 receiver,
3393 inbound_controls,
3394 inbound_controls_error,
3395 )
3396 .await;
3397 }
3398 let door = match crate::mail_route::door_for(¶ms.homes, &receiver) {
3399 Ok(door) => door,
3400 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3401 return refused(
3402 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3403 format!(
3404 "no running `{}` session `{}` is reachable; its transcript is persisted only",
3405 params.locator.harness.as_str(),
3406 params.locator.session_id
3407 ),
3408 )
3409 }
3410 };
3411 let mut envelope = match crate::mailbox::Envelope::new(
3412 sender.clone(),
3413 format!("{}@{}", sender.session_id, sender.machine),
3414 crate::mailbox::MailKind::Peer,
3415 crate::mailbox::ReplyVia::Command,
3416 params.text.clone(),
3417 ) {
3418 Ok(envelope) => envelope,
3419 Err(error) => return refused("delivery_failed", error.to_string()),
3420 };
3421 envelope.in_reply_to = params.in_reply_to.clone();
3422 let how = match crate::mail_route::deliver(
3423 &envelope,
3424 &receiver,
3425 &door,
3426 true,
3427 params.notify_when_idle,
3428 )
3429 .await
3430 {
3431 Err(detail) => {
3432 return refused(
3433 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3434 detail,
3435 )
3436 }
3437 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3438 Ok(Err(Refused::TooLong(bytes))) => {
3439 return refused(
3440 "too_long",
3441 format!(
3442 "the message is {bytes} bytes; the limit is {}",
3443 crate::mail_route::MAX_RELAYED_BYTES
3444 ),
3445 )
3446 }
3447 Ok(Ok(delivered)) => match delivered {
3448 Delivered::Steered => "steered",
3449 Delivered::Started => "started",
3450 Delivered::Native { busy: true } => "next_tool_call",
3451 Delivered::Native { busy: false } => "started",
3452 Delivered::Hooked => "hook",
3453 Delivered::Queued => "queued",
3454 Delivered::Stored => "stored",
3455 Delivered::Operator => "filed",
3456 },
3457 };
3458 json!({
3459 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3460 "message_id": envelope.id,
3461 "reply_to": sender.to_string(),
3462 "target": {
3463 "harness": params.locator.harness.as_str(),
3464 "session_id": params.locator.session_id,
3465 "name": match &door {
3466 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
3467 _ => None,
3468 },
3469 },
3470 "delivery": {"door": door.name(), "how": how},
3471 "inbound_controls": inbound_controls,
3472 "inbound_controls_error": inbound_controls_error,
3473 })
3474}
3475
3476async fn message_as_user(
3480 params: &MessageSessionParams,
3481 sender: crate::mailbox::MailAddress,
3482 receiver: crate::mailbox::MailAddress,
3483 inbound_controls: Value,
3484 inbound_controls_error: Value,
3485) -> Value {
3486 let envelope = match crate::mailbox::Envelope::new(
3487 sender.clone(),
3488 format!("{}@{}", sender.session_id, sender.machine),
3489 crate::mailbox::MailKind::User,
3490 crate::mailbox::ReplyVia::None,
3491 params.text.clone(),
3492 ) {
3493 Ok(envelope) => envelope,
3494 Err(error) => {
3495 return json!({
3496 "delivered_to_bus": false,
3497 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
3498 "inbound_controls": inbound_controls,
3499 "inbound_controls_error": inbound_controls_error,
3500 })
3501 }
3502 };
3503 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
3504 Ok(turn) => json!({
3505 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
3506 "message_id": envelope.id,
3507 "target": {
3508 "harness": params.locator.harness.as_str(),
3509 "session_id": params.locator.session_id,
3510 },
3511 "delivery": {
3512 "door": match turn {
3513 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
3514 _ => "pane",
3515 },
3516 "how": turn.as_str(),
3517 },
3518 "inbound_controls": inbound_controls,
3519 "inbound_controls_error": inbound_controls_error,
3520 }),
3521 Err(message) => json!({
3522 "delivered_to_bus": false,
3523 "refusal": {"reason": "no_user_door", "message": message},
3524 "inbound_controls": inbound_controls,
3525 "inbound_controls_error": inbound_controls_error,
3526 }),
3527 }
3528}
3529
3530#[derive(Deserialize)]
3531#[serde(deny_unknown_fields)]
3532struct ActivityUnderParams {
3533 pids: Vec<u32>,
3535 #[serde(default)]
3536 homes: crate::HarnessHomes,
3537}
3538
3539async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
3543 let params = decode::<ActivityUnderParams>(params)?;
3544 if params.pids.len() > 1024 {
3545 return Err(ServiceError::InvalidParams(
3546 "sessions.activity_under accepts at most 1024 pids".into(),
3547 ));
3548 }
3549 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
3550 .await
3551 .map_err(ServiceError::Sdk)?;
3552 Ok(json!({
3553 "activities": found
3554 .into_iter()
3555 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
3556 .collect::<Vec<_>>(),
3557 }))
3558}
3559
3560fn operator_address(
3563 name: Option<&str>,
3564) -> std::result::Result<crate::mailbox::MailAddress, String> {
3565 let name: String = name
3566 .unwrap_or("supercode")
3567 .trim()
3568 .chars()
3569 .map(|character| {
3570 if character.is_whitespace() || character == '@' {
3571 '-'
3572 } else {
3573 character
3574 }
3575 })
3576 .collect();
3577 if name.is_empty() {
3578 return Err("from_name must not be empty".into());
3579 }
3580 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
3581 .map_err(|error| error.to_string())
3582}
3583
3584fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
3588 let address = match (¶ms.address, ¶ms.from_name) {
3589 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
3590 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
3591 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
3592 (Some(_), Some(_)) => {
3593 return Err(ServiceError::InvalidParams(
3594 "sessions.inbox takes from_name or address, not both".into(),
3595 ))
3596 }
3597 };
3598 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
3599 let mailbox =
3600 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
3601 let claimed = mailbox.claim_unread().map_err(operation)?;
3602 let mut messages: Vec<Value> = Vec::new();
3603 if params.all {
3604 for stored in mailbox.list().map_err(operation)? {
3605 if stored.state == crate::mailbox::MailState::Read {
3606 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3607 }
3608 }
3609 }
3610 for stored in &claimed {
3611 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3612 }
3613 for stored in &claimed {
3614 mailbox.acknowledge(stored).map_err(operation)?;
3615 }
3616 Ok(json!({"address": address.to_string(), "messages": messages}))
3617}
3618
3619#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3624struct FollowedSource {
3625 harness: String,
3626 session_id: String,
3627 reported: Option<String>,
3628}
3629
3630#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3631struct ActivitySubscription {
3632 locators: Vec<SessionLocator>,
3633 homes: crate::HarnessHomes,
3634 reported: BTreeMap<(String, String), crate::SessionActivity>,
3635}
3636
3637fn live_descriptor_value(
3644 session: &SessionDescriptor,
3645 doors: &crate::mail_route::LiveSessions,
3646) -> std::result::Result<Value, ServiceError> {
3647 let mut value = serde_json::to_value(session)
3648 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3649 if value.get("title").is_none_or(Value::is_null) {
3651 if let Some(live) = doors.all().iter().find(|live| {
3652 live.address.harness == session.locator.harness.as_str()
3653 && live.address.session_id == session.locator.session_id
3654 }) {
3655 let name = live.name.split('@').next().unwrap_or(&live.name);
3656 if !name.is_empty() {
3657 value["title"] = json!(name);
3658 }
3659 }
3660 }
3661 if let Some(workspace) = &session.cwd {
3662 let source = LiveRuntimeSource {
3663 harness: session.locator.harness.as_str().to_string(),
3664 session_id: session.locator.session_id.clone(),
3665 workspace: workspace.clone(),
3666 };
3667 if let Some(endpoint) = discover_live_runtime(&source)
3668 .map_err(|error| ServiceError::Operation(error.to_string()))?
3669 {
3670 value["live_endpoint"] = json!(endpoint.as_str());
3671 }
3672 }
3673 if let Some(door) = doors.door(
3677 session.locator.harness.as_str(),
3678 &session.locator.session_id,
3679 ) {
3680 value["delivery"] = json!(door);
3681 }
3682 Ok(value)
3683}
3684
3685fn live_index_changes(
3686 changes: Vec<crate::session_index::SessionIndexChange>,
3687 homes: &HarnessHomes,
3688) -> std::result::Result<Vec<Value>, ServiceError> {
3689 use crate::session_index::SessionIndexChange;
3690 let doors = crate::mail_route::LiveSessions::read(homes);
3691 changes
3692 .into_iter()
3693 .map(|change| match change {
3694 SessionIndexChange::Added { descriptor } => Ok(json!({
3695 "kind": "added",
3696 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3697 })),
3698 SessionIndexChange::Updated { descriptor } => Ok(json!({
3699 "kind": "updated",
3700 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3701 })),
3702 SessionIndexChange::Removed { key } => Ok(json!({
3703 "kind": "removed",
3704 "key": key,
3705 })),
3706 })
3707 .collect()
3708}
3709
3710fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
3711 use crate::{SessionPresence, SessionTurnState};
3712 match (activity.presence, activity.turn) {
3713 (SessionPresence::Persisted, _) => None,
3714 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
3715 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
3716 (SessionPresence::Running, SessionTurnState::Unknown)
3720 if activity.evidence.native_state.is_none() =>
3721 {
3722 None
3723 }
3724 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
3725 }
3726}
3727
3728#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
3729#[serde(rename_all = "kebab-case")]
3730enum TransferFormat {
3731 ClaudeCode,
3732 Codex,
3733 #[serde(rename = "opencode", alias = "open-code")]
3734 OpenCode,
3735 Pi,
3736 Grok,
3737 Gemini,
3738 Goose,
3739 Hermes,
3743}
3744
3745impl TransferFormat {
3746 fn id(self) -> &'static str {
3747 match self {
3748 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
3749 Self::Codex => HarnessId::CODEX,
3750 Self::OpenCode => HarnessId::OPENCODE,
3751 Self::Pi => HarnessId::PI,
3752 Self::Grok => HarnessId::GROK,
3753 Self::Gemini => HarnessId::GEMINI,
3754 Self::Goose => HarnessId::GOOSE,
3755 Self::Hermes => HarnessId::HERMES,
3756 }
3757 }
3758}
3759
3760impl From<TransferFormat> for SessionFormat {
3761 fn from(value: TransferFormat) -> Self {
3762 match value {
3763 TransferFormat::ClaudeCode => Self::ClaudeCode,
3764 TransferFormat::Codex => Self::Codex,
3765 TransferFormat::OpenCode => Self::OpenCode,
3766 TransferFormat::Pi => Self::Pi,
3767 TransferFormat::Grok => Self::Grok,
3768 TransferFormat::Gemini => Self::Gemini,
3769 TransferFormat::Goose => Self::Goose,
3770 TransferFormat::Hermes => Self::Codex,
3772 }
3773 }
3774}
3775
3776#[derive(Deserialize)]
3777struct ImportSessionParams {
3778 source_harness: TransferFormat,
3779 content: String,
3780}
3781
3782#[derive(Deserialize)]
3783struct ExportSessionParams {
3784 locator: SessionLocator,
3785 target_harness: TransferFormat,
3786}
3787
3788#[derive(Deserialize)]
3789struct ReduceSessionParams {
3790 locator: SessionLocator,
3791 target_harness: TransferFormat,
3792 #[serde(default = "default_keep_last")]
3793 keep_last: usize,
3794}
3795
3796fn default_keep_last() -> usize {
3797 6
3798}
3799
3800#[derive(Deserialize)]
3801struct BranchSessionParams {
3802 locator: SessionLocator,
3803 #[serde(default)]
3804 target_harness: Option<TransferFormat>,
3805}
3806
3807#[derive(Deserialize)]
3808struct HandoffSessionParams {
3809 locator: SessionLocator,
3810 target_harness: TransferFormat,
3811 #[serde(default)]
3812 cwd: Option<PathBuf>,
3813}
3814
3815#[derive(Deserialize)]
3816struct MaterializeSessionParams {
3817 artifact: crate::native_materialize::MaterializeArtifact,
3818 cwd: PathBuf,
3819 #[serde(default)]
3821 homes: HarnessHomes,
3822}
3823
3824#[derive(Debug, Clone, Copy, Default, Deserialize)]
3827#[serde(rename_all = "snake_case")]
3828enum ResumePolicy {
3829 Default,
3830 #[default]
3831 Yolo,
3832}
3833
3834#[derive(Deserialize)]
3835struct ResumeInstructionsParams {
3836 locator: SessionLocator,
3837 #[serde(default)]
3838 cwd: Option<PathBuf>,
3839 #[serde(default)]
3840 policy: ResumePolicy,
3841}
3842
3843#[derive(Deserialize)]
3845struct WorkflowLoadParams {
3846 from: crate::workflow_doors::WorkflowHarness,
3847 home: PathBuf,
3848}
3849
3850#[derive(Deserialize)]
3853struct OrchestrationLoadParams {
3854 root: PathBuf,
3855 #[serde(default)]
3856 flavor: crate::orchestration_doors::HomeFlavor,
3857}
3858
3859#[derive(Deserialize)]
3862struct OrchestrationSaveParams {
3863 root: PathBuf,
3864 orchestration: crate::orchestration::Orchestration,
3865 #[serde(default)]
3866 vault: BTreeMap<String, String>,
3867}
3868
3869#[derive(Deserialize)]
3871struct OrchestrationCompileParams {
3872 from: crate::orchestration_doors::OrchestrationHarness,
3873 home: PathBuf,
3874}
3875
3876#[derive(Deserialize)]
3880struct OrchestrationDecompileParams {
3881 to: crate::orchestration_doors::OrchestrationHarness,
3882 orchestration: crate::orchestration::Orchestration,
3883 source: PathBuf,
3884 #[serde(default)]
3885 source_flavor: crate::orchestration_doors::SourceFlavor,
3886 dest: PathBuf,
3887 #[serde(default)]
3888 vault: BTreeMap<String, String>,
3889}
3890
3891#[derive(Deserialize)]
3894struct OrchestrationImportParams {
3895 from: crate::orchestration_doors::OrchestrationHarness,
3896 home: PathBuf,
3897 into: PathBuf,
3898}
3899
3900#[derive(Deserialize)]
3903struct OrchestrationExportParams {
3904 to: crate::orchestration_doors::OrchestrationHarness,
3905 root: PathBuf,
3906 dest: PathBuf,
3907}
3908
3909#[derive(Deserialize)]
3911struct JobsGetParams {
3912 harness: String,
3913 id: String,
3914 #[serde(default)]
3915 homes: crate::HarnessHomes,
3916}
3917
3918fn mutate_job(
3926 verb: crate::jobs_control::JobVerb,
3927 params: Value,
3928) -> std::result::Result<Value, ServiceError> {
3929 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
3930 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
3931 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
3932 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3933}
3934
3935fn mutate_skill(
3942 verb: crate::skills_control::SkillVerb,
3943 params: Value,
3944) -> std::result::Result<Value, ServiceError> {
3945 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
3946 if !crate::skills_control::supports_skill_control(&mutation.harness) {
3947 return Err(ServiceError::UnsupportedAction(format!(
3948 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
3949 mutation.harness,
3950 verb.as_str(),
3951 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
3952 )));
3953 }
3954 let outcome =
3955 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
3956 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3957}
3958
3959fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
3961 match error {
3962 crate::skills_control::SkillControlError::Unsupported(message) => {
3963 ServiceError::UnsupportedAction(message)
3964 }
3965 crate::skills_control::SkillControlError::Invalid(message) => {
3966 ServiceError::InvalidParams(message)
3967 }
3968 crate::skills_control::SkillControlError::Failed(message) => {
3969 ServiceError::Operation(message)
3970 }
3971 }
3972}
3973
3974fn mutate_profile(
3982 verb: crate::profiles_control::ProfileVerb,
3983 params: Value,
3984) -> std::result::Result<Value, ServiceError> {
3985 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
3986 let outcome =
3987 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
3988 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3989}
3990
3991fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
3993 match error {
3994 crate::profiles_control::ProfileControlError::Unsupported(message) => {
3995 ServiceError::UnsupportedAction(message)
3996 }
3997 crate::profiles_control::ProfileControlError::Invalid(message) => {
3998 ServiceError::InvalidParams(message)
3999 }
4000 crate::profiles_control::ProfileControlError::Failed(message) => {
4001 ServiceError::Operation(message)
4002 }
4003 }
4004}
4005
4006fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4010 match error {
4011 crate::jobs_control::JobControlError::Unsupported(message) => {
4012 ServiceError::UnsupportedAction(message)
4013 }
4014 crate::jobs_control::JobControlError::Invalid(message) => {
4015 ServiceError::InvalidParams(message)
4016 }
4017 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4018 }
4019}
4020
4021fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4026 match error {
4027 crate::SessionControlError::Unsupported(message) => {
4028 ServiceError::UnsupportedAction(message)
4029 }
4030 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4031 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4032 }
4033}
4034
4035fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4040 if crate::jobs::supports_jobs(harness) {
4041 return Ok(());
4042 }
4043 Err(ServiceError::UnsupportedAction(format!(
4044 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4045 crate::jobs::JOB_HARNESSES.join(", ")
4046 )))
4047}
4048
4049#[derive(Deserialize)]
4051struct RunsGetParams {
4052 harness: String,
4053 id: String,
4054 #[serde(default)]
4055 homes: crate::HarnessHomes,
4056}
4057
4058fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4063 if crate::runs::supports_runs(harness) {
4064 return Ok(());
4065 }
4066 Err(ServiceError::UnsupportedAction(format!(
4067 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4068 crate::runs::RUN_HARNESSES.join(", ")
4069 )))
4070}
4071
4072#[derive(Serialize)]
4073struct SessionArtifact {
4074 source_harness: HarnessId,
4075 target_harness: &'static str,
4076 session_id: Option<String>,
4077 content: String,
4078 suggested_filename: String,
4079 files: Vec<SessionArtifactFile>,
4080 fidelity: Fidelity,
4081 residue: Vec<String>,
4082}
4083
4084#[derive(Serialize)]
4085struct SessionArtifactFile {
4086 path: String,
4087 content: String,
4088 role: ArtifactFileRole,
4089}
4090
4091#[derive(Serialize)]
4092#[serde(rename_all = "snake_case")]
4093enum ArtifactFileRole {
4094 Primary,
4095 Subagent,
4096 Bundle,
4097 SourceRecovery,
4098}
4099
4100#[derive(Serialize)]
4101struct StructuredLaunch {
4102 cwd: PathBuf,
4103 program: String,
4104 arguments: Vec<String>,
4105 env: BTreeMap<String, String>,
4106}
4107
4108struct HandoffInstructions {
4109 launch: StructuredLaunch,
4110 materialize: Option<StructuredLaunch>,
4111 requires_materialization: bool,
4112 note: String,
4113}
4114
4115#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4116#[serde(rename_all = "snake_case")]
4117enum HarnessProbeLevel {
4118 #[default]
4119 Passive,
4120 Handshake,
4121}
4122
4123#[derive(Default, Deserialize)]
4124#[serde(default)]
4125struct HarnessInventoryParams {
4126 harness: Option<HarnessId>,
4127 harnesses: Vec<HarnessId>,
4128 workspace: Option<PathBuf>,
4129 probe: HarnessProbeLevel,
4130 include_sessions: bool,
4131 skip_versions: bool,
4133}
4134
4135#[derive(Deserialize)]
4136struct HarnessAuthenticationParams {
4137 harness: HarnessId,
4138}
4139
4140#[derive(Deserialize)]
4141struct BeginHarnessAuthenticationParams {
4142 harness: HarnessId,
4143 #[serde(default = "local_browser_authentication_environment")]
4144 environment: crate::HarnessAuthenticationEnvironment,
4145 #[serde(default)]
4146 method: Option<crate::HarnessAuthenticationMethodId>,
4147 #[serde(default)]
4148 cwd: Option<PathBuf>,
4149}
4150
4151fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4152 crate::HarnessAuthenticationEnvironment::LocalBrowser
4153}
4154
4155#[derive(Serialize)]
4156struct HarnessInventoryReport {
4157 probe: HarnessProbeLevel,
4158 workspace: Option<PathBuf>,
4159 harnesses: Vec<LocalHarness>,
4160}
4161
4162#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4163#[serde(rename_all = "snake_case")]
4164enum HarnessAuthState {
4165 Ready,
4166 Configured,
4167 Required,
4168 Unknown,
4169}
4170
4171#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4172#[serde(rename_all = "snake_case")]
4173enum HarnessRuntimeState {
4174 Ready,
4175 Degraded,
4176 Unavailable,
4177}
4178
4179#[derive(Serialize)]
4180struct HarnessSessionCounts {
4181 global: Option<usize>,
4182 workspace: Option<usize>,
4183}
4184
4185#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4196#[serde(rename_all = "snake_case")]
4197pub enum GatewayState {
4198 Up,
4199 Down,
4200 Unknown,
4201}
4202
4203#[derive(Debug, Clone, Serialize)]
4205pub struct GatewayHealth {
4206 pub state: GatewayState,
4207 #[serde(skip_serializing_if = "Option::is_none")]
4212 pub endpoint: Option<String>,
4213 #[serde(skip_serializing_if = "Option::is_none")]
4214 pub version: Option<String>,
4215 pub evidence: String,
4217 pub checked_at_ms: u64,
4218}
4219
4220fn openclaw_gateway_endpoint(home: &Path) -> String {
4224 let config_path = home.join(".openclaw/openclaw.json");
4225 let gateway = std::fs::read_to_string(&config_path)
4226 .ok()
4227 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4228 .and_then(|config| config.get("gateway").cloned());
4229 if let Some(url) = gateway
4230 .as_ref()
4231 .and_then(|gateway| gateway.get("url"))
4232 .and_then(serde_json::Value::as_str)
4233 {
4234 return url.to_string();
4235 }
4236 let port = gateway
4237 .as_ref()
4238 .and_then(|gateway| gateway.get("port"))
4239 .and_then(serde_json::Value::as_u64)
4240 .unwrap_or(18789);
4241 format!("ws://127.0.0.1:{port}")
4242}
4243
4244fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4251 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4252 let output = std::process::Command::new(&program)
4253 .args(["gateway", "status"])
4254 .stdin(std::process::Stdio::null())
4255 .output()
4256 .ok()?;
4257 let text = format!(
4258 "{}{}",
4259 String::from_utf8_lossy(&output.stdout),
4260 String::from_utf8_lossy(&output.stderr)
4261 );
4262 let verdict = text.lines().find_map(|line| {
4263 let l = line.trim();
4264 if l.contains("supervised by launchd (PID")
4265 || l.contains("supervised by systemd (PID")
4266 || l.contains("Gateway is running")
4267 || l.contains("process is running")
4268 {
4269 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4270 } else if l.contains("not running") || l.contains("not installed") {
4271 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4272 } else {
4273 None
4274 }
4275 });
4276 verdict
4277}
4278
4279fn gateway_health(
4280 id: &str,
4281 installed: bool,
4282 running: Option<&RunningInstance>,
4283 version: Option<&str>,
4284) -> GatewayHealth {
4285 let checked_at_ms = now_epoch_ms();
4286 let home = supercode_interchange::user_home()
4287 .map(std::path::PathBuf::into_os_string)
4288 .map(PathBuf::from);
4289 match id {
4290 HarnessId::HERMES | HarnessId::OPENCLAW => {
4291 let endpoint = (id == HarnessId::OPENCLAW)
4292 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4293 .flatten();
4294 let (state, evidence) = match running {
4295 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4296 None if !installed => (
4297 GatewayState::Unknown,
4298 format!("`{id}` is not installed; no gateway to probe"),
4299 ),
4300 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4301 Some((state, evidence)) => (state, evidence),
4304 None => (
4305 GatewayState::Down,
4306 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4307 ),
4308 },
4309 None => (
4310 GatewayState::Down,
4311 format!(
4312 "no TCP listener at {}",
4313 endpoint.as_deref().unwrap_or("the gateway endpoint")
4314 ),
4315 ),
4316 };
4317 GatewayHealth {
4318 state,
4319 endpoint,
4320 version: version.map(str::to_string),
4321 evidence,
4322 checked_at_ms,
4323 }
4324 }
4325 HarnessId::ORCHESTRATOR => {
4332 let root = crate::HarnessHomes::default().orchestrator;
4333 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4334 Some(lease) if lease.is_live() => (
4335 GatewayState::Up,
4336 format!(
4337 "`{}` names pid {} (started {}), which is live",
4338 crate::orchestrator::lock_path(&root).display(),
4339 lease.pid,
4340 lease.started_at
4341 ),
4342 ),
4343 Some(lease) => (
4344 GatewayState::Down,
4345 format!(
4346 "stale lease `{}`: pid {} is gone",
4347 crate::orchestrator::lock_path(&root).display(),
4348 lease.pid
4349 ),
4350 ),
4351 None => (
4352 GatewayState::Down,
4353 format!(
4354 "no lease at `{}`; `supercode orchestrator start` writes one",
4355 crate::orchestrator::lock_path(&root).display()
4356 ),
4357 ),
4358 };
4359 GatewayHealth {
4360 state,
4361 endpoint: None,
4362 version: version.map(str::to_string),
4363 evidence,
4364 checked_at_ms,
4365 }
4366 }
4367 _ => GatewayHealth {
4368 state: GatewayState::Unknown,
4369 endpoint: None,
4370 version: version.map(str::to_string),
4371 evidence: format!("`{id}` runs per session, not as a gateway"),
4372 checked_at_ms,
4373 },
4374 }
4375}
4376
4377#[derive(Debug, Clone, Serialize)]
4378struct RunningInstance {
4379 method: RunningInstanceMethod,
4381 evidence: String,
4383 checked_at_ms: u64,
4385}
4386
4387#[derive(Debug, Clone, Copy, Serialize)]
4388#[serde(rename_all = "snake_case")]
4389enum RunningInstanceMethod {
4390 GatewayConnect,
4393 StoreWalActivity,
4396}
4397
4398fn now_epoch_ms() -> u64 {
4399 std::time::SystemTime::now()
4400 .duration_since(std::time::UNIX_EPOCH)
4401 .map(|elapsed| elapsed.as_millis() as u64)
4402 .unwrap_or(0)
4403}
4404
4405fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4409 let config_path = home.join(".openclaw/openclaw.json");
4410 let text = std::fs::read_to_string(&config_path).ok();
4411 let gateway = text
4412 .as_deref()
4413 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4414 .and_then(|config| config.get("gateway").cloned());
4415 let address = gateway
4416 .as_ref()
4417 .and_then(|gateway| gateway.get("url"))
4418 .and_then(serde_json::Value::as_str)
4419 .and_then(|url| {
4420 url.split("://").nth(1).map(|rest| {
4421 rest.trim_end_matches('/')
4422 .split('/')
4423 .next()
4424 .unwrap_or(rest)
4425 .to_string()
4426 })
4427 })
4428 .unwrap_or_else(|| {
4429 let port = gateway
4430 .as_ref()
4431 .and_then(|gateway| gateway.get("port"))
4432 .and_then(serde_json::Value::as_u64)
4433 .unwrap_or(18789);
4434 format!("127.0.0.1:{port}")
4435 });
4436 let reachable = std::net::TcpStream::connect_timeout(
4437 &address.parse().ok()?,
4438 std::time::Duration::from_millis(400),
4439 )
4440 .is_ok();
4441 reachable.then(|| RunningInstance {
4442 method: RunningInstanceMethod::GatewayConnect,
4443 evidence: format!(
4444 "gateway endpoint {address} accepted a TCP connect (from {})",
4445 config_path.display()
4446 ),
4447 checked_at_ms: now_epoch_ms(),
4448 })
4449}
4450
4451fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4456 let wal = home.join(".hermes/state.db-wal");
4457 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4458 let age_ms = std::time::SystemTime::now()
4459 .duration_since(modified)
4460 .map(|age| age.as_millis() as u64)
4461 .unwrap_or(u64::MAX);
4462 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4463 method: RunningInstanceMethod::StoreWalActivity,
4464 evidence: format!(
4465 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4466 wal.display()
4467 ),
4468 checked_at_ms: now_epoch_ms(),
4469 })
4470}
4471
4472fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4474 let home = supercode_interchange::user_home()
4475 .map(std::path::PathBuf::into_os_string)
4476 .map(PathBuf::from)?;
4477 match id {
4478 HarnessId::OPENCLAW => probe_openclaw_running(&home),
4479 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4480 _ => None,
4481 }
4482}
4483
4484#[derive(Serialize)]
4485struct LocalHarness {
4486 id: HarnessId,
4487 display_name: String,
4488 supported: bool,
4489 installed: bool,
4490 executable: Option<String>,
4491 version: Option<String>,
4492 auth: HarnessAuthState,
4493 runtime: HarnessRuntimeState,
4494 protocol: String,
4495 capabilities: crate::RuntimeCapabilities,
4496 effective_capabilities: crate::RuntimeCapabilities,
4497 sessions: HarnessSessionCounts,
4498 #[serde(skip_serializing_if = "Option::is_none")]
4501 running: Option<RunningInstance>,
4502 gateway: GatewayHealth,
4504 reason: Option<String>,
4505 repair: Option<String>,
4506}
4507
4508#[derive(Clone, Deserialize)]
4509struct RuntimeBackendParams {
4510 harness: HarnessId,
4511 #[serde(default)]
4512 protocol: Option<String>,
4513 #[serde(default)]
4514 launch: Option<RuntimeLaunch>,
4515 #[serde(default)]
4516 base_url: Option<String>,
4517 #[serde(default)]
4518 policy: RuntimePolicy,
4519}
4520
4521#[derive(Debug, Clone, Copy, Default, Deserialize)]
4524#[serde(rename_all = "snake_case")]
4525enum RuntimePolicy {
4526 Default,
4527 #[default]
4528 Yolo,
4529}
4530
4531#[derive(Deserialize)]
4532struct RuntimeStartParams {
4533 #[serde(flatten)]
4534 backend: RuntimeBackendParams,
4535 cwd: PathBuf,
4536 #[serde(default)]
4539 mcp_servers: Vec<crate::McpServerLaunch>,
4540 #[serde(default)]
4542 approval_policy: Option<String>,
4543}
4544
4545#[derive(Deserialize)]
4546struct RuntimeAttachParams {
4547 #[serde(flatten)]
4548 backend: RuntimeBackendParams,
4549 runtime_id: String,
4550 #[serde(default)]
4551 cwd: Option<PathBuf>,
4552 #[serde(default)]
4555 mcp_servers: Vec<crate::McpServerLaunch>,
4556 #[serde(default)]
4558 approval_policy: Option<String>,
4559}
4560
4561#[derive(Deserialize)]
4562struct RuntimeConnectionParams {
4563 connection: String,
4564}
4565
4566#[derive(Deserialize)]
4567struct RuntimeInputParams {
4568 connection: String,
4569 text: String,
4570 #[serde(default)]
4571 image_urls: Vec<String>,
4572}
4573
4574const MAX_RUNTIME_IMAGES: usize = 4;
4575const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
4576const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
4577
4578fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
4579 if image_urls.len() > MAX_RUNTIME_IMAGES {
4580 return Err(ServiceError::InvalidParams(format!(
4581 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
4582 )));
4583 }
4584 let mut total = 0usize;
4585 for url in &image_urls {
4586 if !(url.starts_with("data:image/")
4587 || url.starts_with("https://")
4588 || url.starts_with("http://"))
4589 {
4590 return Err(ServiceError::InvalidParams(
4591 "runtime images must be image data URLs or HTTP(S) URLs".into(),
4592 ));
4593 }
4594 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
4595 return Err(ServiceError::InvalidParams(format!(
4596 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
4597 )));
4598 }
4599 total = total.saturating_add(url.len());
4600 }
4601 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
4602 return Err(ServiceError::InvalidParams(format!(
4603 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
4604 )));
4605 }
4606 Ok(image_urls)
4607}
4608
4609#[derive(Deserialize)]
4610struct RuntimeRespondParams {
4611 connection: String,
4612 request_id: Value,
4613 response: Value,
4614}
4615
4616fn default_reduction_store_root() -> PathBuf {
4617 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
4618 return PathBuf::from(root).join("sessions");
4619 }
4620 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
4621 return PathBuf::from(home).join(".supercode").join("sessions");
4622 }
4623 PathBuf::from(".supercode").join("sessions")
4624}
4625
4626fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
4627 let mut output = String::new();
4628 for message in messages {
4629 output.push_str(
4630 &serde_json::to_string(message)
4631 .map_err(|error| ServiceError::Operation(error.to_string()))?,
4632 );
4633 output.push('\n');
4634 }
4635 Ok(output)
4636}
4637
4638fn parse_messages_jsonl(
4639 content: &str,
4640) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
4641 content
4642 .lines()
4643 .enumerate()
4644 .filter(|(_, line)| !line.trim().is_empty())
4645 .map(|(index, line)| {
4646 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
4647 ServiceError::Operation(format!(
4648 "reduced transcript line {} is invalid: {error}",
4649 index + 1
4650 ))
4651 })
4652 })
4653 .collect()
4654}
4655
4656fn reduced_bootstrap_prompt(
4657 source: &SessionLocator,
4658 target: TransferFormat,
4659 view_jsonl: &str,
4660 sidecar_path: &Path,
4661 reduction_log_path: &Path,
4662) -> String {
4663 format!(
4664 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
4665 \n\
4666 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 Volter Harness 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\
4667 \n\
4668 <supercode-reduced-session source-session=\"{source_id}\">\n\
4669 {view_jsonl}\
4670 </supercode-reduced-session>\n\
4671 \n\
4672 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
4673 source_harness = source.harness.as_str(),
4674 target_harness = target.id(),
4675 sidecar = sidecar_path.display(),
4676 log = reduction_log_path.display(),
4677 source_id = source.session_id,
4678 )
4679}
4680
4681fn session_artifact(
4682 locator: &SessionLocator,
4683 session: &Session,
4684 target: TransferFormat,
4685) -> std::result::Result<SessionArtifact, ServiceError> {
4686 session_artifact_with_id(locator, session, target, None)
4687}
4688
4689fn session_artifact_with_id(
4690 locator: &SessionLocator,
4691 session: &Session,
4692 target: TransferFormat,
4693 target_session_id: Option<&str>,
4694) -> std::result::Result<SessionArtifact, ServiceError> {
4695 let format: SessionFormat = target.into();
4696 let diagonal = format.source() == session.meta.source;
4697 crate::residue_store::store_segments(session);
4698 let has_appended_turns = session
4699 .imported_message_count
4700 .is_some_and(|imported| imported < session.messages.len());
4701 let mut restoration = None;
4702 let content = if let Some(id) = target_session_id {
4703 if diagonal && format != SessionFormat::OpenCode {
4704 session
4705 .to_jsonl_spliced(format, Some(id))
4706 .map_err(operation)?
4707 } else {
4708 let mut rewritten = session.clone();
4709 rewritten.meta.session_id = Some(id.to_string());
4710 rewritten.to_jsonl(format).map_err(operation)?
4711 }
4712 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
4713 session.raw_verbatim()
4714 } else if diagonal {
4715 session.to_jsonl_spliced(format, None).map_err(operation)?
4716 } else {
4717 match session
4720 .restore_residue(format, crate::residue_store::lookup)
4721 .map_err(operation)?
4722 {
4723 Some((content, report)) => {
4724 restoration = Some(report);
4725 content
4726 }
4727 None => session.to_jsonl(format).map_err(operation)?,
4728 }
4729 };
4730 let stem = sanitize_filename(
4731 target_session_id
4732 .or(session.meta.session_id.as_deref())
4733 .unwrap_or(&locator.session_id),
4734 );
4735 let suggested_filename = if diagonal && target == TransferFormat::Grok {
4736 "chat_history.jsonl".to_string()
4737 } else if target == TransferFormat::Goose {
4738 format!("{stem}.goose.json")
4739 } else {
4740 format!("{stem}.{}.jsonl", target.id())
4741 };
4742 let mut files = vec![SessionArtifactFile {
4743 path: suggested_filename.clone(),
4744 content: content.clone(),
4745 role: ArtifactFileRole::Primary,
4746 }];
4747 if target == TransferFormat::ClaudeCode {
4748 let bundle_stem = Path::new(&suggested_filename)
4749 .file_stem()
4750 .and_then(|stem| stem.to_str())
4751 .unwrap_or(&stem);
4752 let mut child_paths = BTreeSet::new();
4753 for (index, subagent) in session.subagents.iter().enumerate() {
4754 let agent_id = subagent
4755 .meta
4756 .agent_id
4757 .as_deref()
4758 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
4759 .map(sanitize_filename)
4760 .filter(|id| !id.is_empty())
4761 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4762 let child_has_appended_turns = subagent
4763 .imported_message_count
4764 .is_some_and(|imported| imported < subagent.messages.len());
4765 let child_content = if target_session_id.is_none()
4766 && subagent.meta.source == SessionSource::ClaudeCode
4767 && subagent.raw_is_verbatim
4768 && !child_has_appended_turns
4769 {
4770 subagent.raw_verbatim()
4771 } else if subagent.meta.source == SessionSource::ClaudeCode {
4772 subagent
4773 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
4774 .map_err(operation)?
4775 } else {
4776 let mut child = subagent.clone();
4777 if let Some(id) = target_session_id {
4778 child.meta.session_id = Some(id.to_string());
4779 }
4780 child
4781 .to_jsonl(SessionFormat::ClaudeCode)
4782 .map_err(operation)?
4783 };
4784 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
4785 if !child_paths.insert(path.clone()) {
4786 return Err(ServiceError::Operation(format!(
4787 "Claude subagent ids collide at artifact path `{path}`"
4788 )));
4789 }
4790 files.push(SessionArtifactFile {
4791 path,
4792 content: child_content,
4793 role: ArtifactFileRole::Subagent,
4794 });
4795 }
4796 }
4797 if diagonal && target == TransferFormat::Grok {
4798 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
4799 }
4800 if !diagonal || !session.raw_is_verbatim {
4801 files.push(SessionArtifactFile {
4802 path: "recovery/source.supercode.jsonl".into(),
4803 content: session.to_native_jsonl(),
4804 role: ArtifactFileRole::SourceRecovery,
4805 });
4806 for (index, subagent) in session.subagents.iter().enumerate() {
4807 let id = subagent
4808 .meta
4809 .agent_id
4810 .as_deref()
4811 .map(sanitize_filename)
4812 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4813 files.push(SessionArtifactFile {
4814 path: format!("recovery/subagents/{id}.supercode.jsonl"),
4815 content: subagent.to_native_jsonl(),
4816 role: ArtifactFileRole::SourceRecovery,
4817 });
4818 }
4819 }
4820 if !diagonal && session.meta.source == SessionSource::Grok {
4821 append_grok_bundle_files(
4822 locator,
4823 "recovery/grok/",
4824 ArtifactFileRole::SourceRecovery,
4825 &mut files,
4826 )?;
4827 }
4828 let (fidelity, residue) = if diagonal
4829 && target_session_id.is_none()
4830 && session.raw_is_verbatim
4831 && !has_appended_turns
4832 {
4833 (Fidelity::ByteLossless, Vec::new())
4834 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
4835 (
4836 Fidelity::ValueLossless,
4837 vec![if target_session_id.is_some() {
4838 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
4839 } else {
4840 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
4841 }],
4842 )
4843 } else {
4844 match restoration {
4845 Some(report) if report.rendered_messages == 0 => (
4846 Fidelity::ByteLossless,
4847 vec![format!(
4848 "restored verbatim from this conversation's {} source records in the residue store",
4849 target.id()
4850 )],
4851 ),
4852 Some(report) => (
4853 Fidelity::Semantic,
4854 vec![format!(
4855 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
4856 report.restored_messages,
4857 report.restored_messages + report.rendered_messages,
4858 report.rendered_messages,
4859 target.id()
4860 )],
4861 ),
4862 None => (
4863 Fidelity::Semantic,
4864 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
4865 ),
4866 }
4867 };
4868 Ok(SessionArtifact {
4869 source_harness: locator.harness.clone(),
4870 target_harness: target.id(),
4871 session_id: target_session_id
4872 .map(str::to_string)
4873 .or_else(|| session.meta.session_id.clone()),
4874 content,
4875 suggested_filename,
4876 files,
4877 fidelity,
4878 residue,
4879 })
4880}
4881
4882fn append_grok_bundle_files(
4883 locator: &SessionLocator,
4884 prefix: &str,
4885 role: ArtifactFileRole,
4886 files: &mut Vec<SessionArtifactFile>,
4887) -> std::result::Result<(), ServiceError> {
4888 let primary = locator.storage.path();
4889 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
4890 return Err(ServiceError::Operation(format!(
4891 "Grok bundle locator must name chat_history.jsonl, got {}",
4892 primary.display()
4893 )));
4894 }
4895 let parent = primary.parent().ok_or_else(|| {
4896 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
4897 })?;
4898 for name in ["summary.json", "updates.jsonl"] {
4899 let path = parent.join(name);
4900 let metadata = match std::fs::symlink_metadata(&path) {
4901 Ok(metadata) => metadata,
4902 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
4903 Err(error) => return Err(ServiceError::Operation(error.to_string())),
4904 };
4905 if metadata.file_type().is_symlink() || !metadata.is_file() {
4906 return Err(ServiceError::Operation(format!(
4907 "refusing non-regular Grok bundle member {}",
4908 path.display()
4909 )));
4910 }
4911 let content = std::fs::read_to_string(&path).map_err(|error| {
4912 ServiceError::Operation(format!(
4913 "Grok bundle member {} is not representable as UTF-8: {error}",
4914 path.display()
4915 ))
4916 })?;
4917 files.push(SessionArtifactFile {
4918 path: format!("{prefix}{name}"),
4919 content,
4920 role: match role {
4921 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
4922 _ => ArtifactFileRole::SourceRecovery,
4923 },
4924 });
4925 }
4926 Ok(())
4927}
4928
4929fn handoff_artifact(
4930 locator: &SessionLocator,
4931 session: &Session,
4932 target: TransferFormat,
4933) -> std::result::Result<SessionArtifact, ServiceError> {
4934 let target_session_id = target_session_id(target);
4935 session_artifact_with_id(locator, session, target, Some(&target_session_id))
4936}
4937
4938fn target_session_id(target: TransferFormat) -> String {
4939 let uuid = generated_session_id();
4940 match target {
4941 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
4942 TransferFormat::ClaudeCode
4943 | TransferFormat::Codex
4944 | TransferFormat::Pi
4945 | TransferFormat::Grok
4946 | TransferFormat::Gemini
4947 | TransferFormat::Goose
4948 | TransferFormat::Hermes => uuid,
4949 }
4950}
4951
4952fn sanitize_filename(value: &str) -> String {
4953 let value = value
4954 .chars()
4955 .map(|character| {
4956 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
4957 character
4958 } else {
4959 '-'
4960 }
4961 })
4962 .collect::<String>();
4963 let value = value.trim_matches('-');
4964 if value.is_empty() {
4965 "session".into()
4966 } else {
4967 value.chars().take(100).collect()
4968 }
4969}
4970
4971fn handoff_instructions(
4972 target: TransferFormat,
4973 session_id: &str,
4974 cwd: &Path,
4975) -> HandoffInstructions {
4976 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
4977 cwd: cwd.to_path_buf(),
4978 program: program.into(),
4979 arguments,
4980 env: BTreeMap::new(),
4981 };
4982 match target {
4983 TransferFormat::ClaudeCode => HandoffInstructions {
4984 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
4985 materialize: None,
4986 requires_materialization: true,
4987 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(),
4988 },
4989 TransferFormat::Hermes => HandoffInstructions {
4990 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
4991 materialize: None,
4992 requires_materialization: true,
4993 note: "Hand the artifact (a Codex rollout) to `hermes sessions import --from codex <file>` — `sessions.export --to hermes` does exactly that — and resume the id Hermes prints: Hermes mints its own id and writes its own store.".into(),
4994 },
4995 TransferFormat::Codex => HandoffInstructions {
4996 launch: launch("codex", vec!["resume".into(), session_id.into()]),
4997 materialize: None,
4998 requires_materialization: true,
4999 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5000 },
5001 TransferFormat::OpenCode => HandoffInstructions {
5002 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5003 materialize: Some(launch(
5004 "opencode",
5005 vec!["import".into(), "{artifact_path}".into()],
5006 )),
5007 requires_materialization: true,
5008 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5009 },
5010 TransferFormat::Pi => HandoffInstructions {
5011 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5012 materialize: None,
5013 requires_materialization: true,
5014 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5015 },
5016 TransferFormat::Grok => HandoffInstructions {
5017 launch: launch(
5018 "grok",
5019 vec!["--resume".into(), "{materialized_session_id}".into()],
5020 ),
5021 materialize: None,
5022 requires_materialization: true,
5023 note: "Grok has no import command. Materialize the artifact through `harness.v1.sessions.materialize` (target `grok`, `value_lossless`, the destination cwd): it writes Grok's store entry (`chat_history.jsonl` and the `summary.json` `--resume` requires) under a fresh id; replace {materialized_session_id} with the id it returns.".into(),
5024 },
5025 TransferFormat::Gemini => HandoffInstructions {
5026 launch: launch(
5027 "gemini",
5028 vec!["--session-file".into(), "{artifact_path}".into()],
5029 ),
5030 materialize: None,
5031 requires_materialization: true,
5032 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(),
5033 },
5034 TransferFormat::Goose => HandoffInstructions {
5035 launch: launch(
5036 "goose",
5037 vec![
5038 "session".into(),
5039 "--resume".into(),
5040 "--session-id".into(),
5041 "{imported_session_id}".into(),
5042 ],
5043 ),
5044 materialize: Some(launch(
5045 "goose",
5046 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5047 )),
5048 requires_materialization: true,
5049 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(),
5050 },
5051 }
5052}
5053
5054fn resume_launch(
5055 harness: &str,
5056 session_id: &str,
5057 cwd: &Path,
5058 policy: ResumePolicy,
5059) -> std::result::Result<StructuredLaunch, ServiceError> {
5060 let mut arguments = Vec::new();
5061 let program = match harness {
5062 HarnessId::GROK => {
5063 if matches!(policy, ResumePolicy::Yolo) {
5064 if crate::support::self_sandbox_supported() {
5065 arguments.extend(["--sandbox".into(), "workspace".into()]);
5066 }
5067 arguments.push("--always-approve".into());
5068 }
5069 arguments.extend(["--resume".into(), session_id.into()]);
5070 "grok"
5071 }
5072 HarnessId::CODEX => {
5073 let cwd_key = serde_json::to_string(cwd.to_string_lossy().as_ref())
5074 .expect("a filesystem path always serializes as JSON text");
5075 arguments.extend([
5076 "-c".into(),
5077 "check_for_update_on_startup=false".into(),
5078 "-c".into(),
5079 format!("projects={{{cwd_key}={{trust_level=\"trusted\"}}}}"),
5082 ]);
5083 if matches!(policy, ResumePolicy::Yolo) {
5084 arguments.extend([
5085 "--dangerously-bypass-approvals-and-sandbox".into(),
5086 "--dangerously-bypass-hook-trust".into(),
5087 ]);
5088 }
5089 arguments.extend(["resume".into(), session_id.into()]);
5090 "codex"
5091 }
5092 HarnessId::CLAUDE_CODE => {
5093 if matches!(policy, ResumePolicy::Yolo) {
5094 arguments.push("--dangerously-skip-permissions".into());
5095 }
5096 arguments.extend(["--resume".into(), session_id.into()]);
5097 "claude"
5098 }
5099 HarnessId::GEMINI => {
5100 if matches!(policy, ResumePolicy::Yolo) {
5101 arguments.push("--yolo".into());
5102 }
5103 arguments.extend(["--resume".into(), session_id.into()]);
5104 "gemini"
5105 }
5106 HarnessId::GOOSE => {
5107 arguments.extend([
5108 "session".into(),
5109 "--resume".into(),
5110 "--session-id".into(),
5111 session_id.into(),
5112 ]);
5113 "goose"
5114 }
5115 HarnessId::PI => {
5116 if matches!(policy, ResumePolicy::Yolo) {
5117 arguments.push("--approve".into());
5118 }
5119 arguments.extend(["--session".into(), session_id.into()]);
5120 "pi"
5121 }
5122 HarnessId::OPENCODE => {
5123 arguments.extend(["--session".into(), session_id.into()]);
5124 "opencode"
5125 }
5126 HarnessId::SUPERCODE => {
5127 if matches!(policy, ResumePolicy::Yolo) {
5128 arguments.push("--dangerous".into());
5129 }
5130 arguments.extend(["resume".into(), session_id.into()]);
5131 "supercode"
5132 }
5133 other => {
5134 return Err(ServiceError::InvalidParams(format!(
5135 "no structured resume launch is registered for harness `{other}`"
5136 )))
5137 }
5138 };
5139 Ok(StructuredLaunch {
5140 cwd: cwd.to_path_buf(),
5141 env: if program == "grok" {
5142 crate::support::grok_home_env()
5143 } else {
5144 BTreeMap::new()
5145 },
5146 program: program.into(),
5147 arguments,
5148 })
5149}
5150
5151fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5157 let digest = blake3::hash(address.as_bytes()).to_hex();
5158 let path = std::env::temp_dir().join(format!(
5159 "supercode-openclaw-gateway-token-{}",
5160 &digest.as_str()[..16]
5161 ));
5162 #[cfg(unix)]
5163 {
5164 use std::io::Write;
5165 use std::os::unix::fs::OpenOptionsExt;
5166 let mut file = std::fs::OpenOptions::new()
5167 .write(true)
5168 .create(true)
5169 .truncate(true)
5170 .mode(0o600)
5171 .open(&path)?;
5172 file.write_all(secret.as_bytes())?;
5173 }
5174 #[cfg(not(unix))]
5175 std::fs::write(&path, secret)?;
5176 Ok(path)
5177}
5178
5179fn open_connect_descriptor(
5185 descriptor: &crate::HarnessSupportDescriptor,
5186 home: &Path,
5187) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5188 let Some(connect) = &descriptor.runtime.connect_launch else {
5189 return Err(ServiceError::InvalidParams(format!(
5190 "harness `{}` has no registered connect-mode launch",
5191 descriptor.id.as_str()
5192 )));
5193 };
5194 let resolved = connect
5195 .resolve(home)
5196 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5197 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5198 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5199 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5200 if let Some(token) = resolved.auth {
5201 backend = backend.with_bearer(token);
5202 }
5203 Ok(Box::new(backend))
5204 }
5205 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5206 let mut env = BTreeMap::new();
5217 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5218 if let Some(token) = resolved.auth {
5219 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5220 .map_err(|error| {
5221 ServiceError::UnsupportedAction(format!(
5222 "could not stage the gateway credential for the bridge: {error}"
5223 ))
5224 })?;
5225 arguments.push("--token-file".into());
5226 arguments.push(token_path.to_string_lossy().into_owned());
5227 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5228 }
5229 let program = descriptor
5234 .runtime
5235 .default_launch
5236 .as_ref()
5237 .map(|launch| launch.program.clone())
5238 .unwrap_or_else(|| "openclaw".into());
5239 let launch = RuntimeLaunch {
5240 program,
5241 arguments,
5242 env,
5243 };
5244 Ok(Box::new(
5245 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5246 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5247 ))
5248 }
5249 _ => Err(ServiceError::UnsupportedAction(format!(
5250 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5251 descriptor.id.as_str(),
5252 connect.protocol
5253 ))),
5254 }
5255}
5256
5257fn registry_connect_descriptor(
5260 params: &RuntimeBackendParams,
5261) -> Option<crate::HarnessSupportDescriptor> {
5262 if params.launch.is_some() || params.base_url.is_some() {
5263 return None;
5264 }
5265 harness_support_registry()
5266 .harnesses
5267 .into_iter()
5268 .find(|descriptor| descriptor.id == params.harness)
5269 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5270}
5271
5272fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5273 supercode_interchange::user_home()
5274 .map(std::path::PathBuf::into_os_string)
5275 .map(PathBuf::from)
5276 .ok_or_else(|| {
5277 ServiceError::UnsupportedAction(
5278 "connect-mode launches need HOME to locate the harness config".into(),
5279 )
5280 })
5281}
5282
5283pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5286 "harness.v1.runtimes.start",
5287 "harness.v1.runtimes.resume",
5288 "harness.v1.runtimes.attach",
5289 "harness.v1.runtimes.attach_existing",
5290];
5291
5292pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5297
5298pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5304
5305pub const DETACHED_METHODS: &[&str] = &[
5311 "harness.v1.harnesses.list",
5312 "harness.v1.harnesses.probe",
5313 "harness.v1.sessions.message",
5314 "harness.v1.sessions.new",
5315 "harness.v1.sessions.reset",
5316 "harness.v1.sessions.archive",
5317 "harness.v1.sessions.delete",
5318];
5319
5320pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5326
5327pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5335
5336async fn within_control_deadline<F: std::future::Future>(
5339 method: &str,
5340 call: F,
5341) -> std::result::Result<F::Output, ServiceError> {
5342 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5343 .await
5344 .map_err(|_| {
5345 ServiceError::Operation(format!(
5346 "`{method}` gave up after {}s: the runtime did not answer",
5347 RUNTIME_CONTROL_DEADLINE.as_secs()
5348 ))
5349 })
5350}
5351
5352pub struct RuntimeOpen {
5356 id: Value,
5357 method: String,
5358 params: Value,
5359}
5360
5361impl RuntimeOpen {
5362 pub async fn open(self) -> OpenedRuntime {
5366 let Self { id, method, params } = self;
5367 let outcome = open_runtime(&method, params).await;
5368 OpenedRuntime { id, outcome }
5369 }
5370}
5371
5372pub struct OpenedRuntime {
5375 id: Value,
5376 outcome: std::result::Result<OpenRuntime, ServiceError>,
5377}
5378
5379pub struct DetachedCall {
5384 id: Value,
5385 method: String,
5386 work: std::result::Result<Work, ServiceError>,
5387}
5388
5389impl DetachedCall {
5390 pub async fn run(self) -> DetachedAnswer {
5393 let Self { id, method, work } = self;
5394 match work {
5395 Ok(Work::Runtime(work)) => {
5400 let (result, returned) = work.run().await;
5401 DetachedAnswer {
5402 response: service_response(id, result),
5403 returned,
5404 }
5405 }
5406 Ok(Work::Free(work)) => {
5407 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5408 Ok(result) => result,
5409 Err(_) => Err(ServiceError::Operation(format!(
5410 "`{method}` gave up after {}s: the harness it waits on did not answer",
5411 DETACHED_CALL_DEADLINE.as_secs()
5412 ))),
5413 };
5414 DetachedAnswer {
5415 response: service_response(id, result),
5416 returned: None,
5417 }
5418 }
5419 Err(error) => DetachedAnswer {
5420 response: service_response(id, Err(error)),
5421 returned: None,
5422 },
5423 }
5424 }
5425}
5426
5427pub struct DetachedAnswer {
5431 response: Value,
5432 returned: Option<ReturnedRuntime>,
5433}
5434
5435impl DetachedAnswer {
5436 pub fn into_response(self) -> Value {
5439 self.response
5440 }
5441}
5442
5443pub struct ReturnedRuntime {
5446 connection: String,
5447 runtime: Box<dyn RuntimeConnection>,
5448}
5449
5450enum Work {
5453 Free(DetachedWork),
5454 Runtime(RuntimeWork),
5455}
5456
5457enum DetachedWork {
5460 Inventory(InventoryWork),
5464 Message(MessageSessionParams),
5466 SessionMutation {
5469 verb: crate::SessionVerb,
5470 mutation: crate::SessionMutation,
5471 },
5472}
5473
5474impl DetachedWork {
5475 async fn run(self) -> std::result::Result<Value, ServiceError> {
5476 match self {
5477 Self::Inventory(work) => run_inventory(work).await,
5478 Self::Message(params) => Ok(message_live_session(¶ms).await),
5479 Self::SessionMutation { verb, mutation } => {
5480 let outcome = run_session_mutation(verb, &mutation).await?;
5481 serde_json::to_value(outcome)
5482 .map_err(|error| ServiceError::Operation(error.to_string()))
5483 }
5484 }
5485 }
5486}
5487
5488enum RuntimeWork {
5490 Close {
5492 runtime: Box<dyn RuntimeConnection>,
5493 process_group: Option<u32>,
5494 },
5495 LiveCommand {
5498 connection: String,
5499 runtime: Box<dyn RuntimeConnection>,
5500 verb: crate::SessionVerb,
5501 mutation: crate::SessionMutation,
5502 command: &'static str,
5503 session: String,
5504 },
5505}
5506
5507type RuntimeWorkAnswer = (
5510 std::result::Result<Value, ServiceError>,
5511 Option<ReturnedRuntime>,
5512);
5513
5514impl RuntimeWork {
5515 async fn run(self) -> RuntimeWorkAnswer {
5516 match self {
5517 Self::Close {
5518 runtime,
5519 process_group,
5520 } => (close_runtime(runtime, process_group).await, None),
5521 Self::LiveCommand {
5522 connection,
5523 mut runtime,
5524 verb,
5525 mutation,
5526 command,
5527 session,
5528 } => {
5529 let result =
5530 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5531 (
5532 result,
5533 Some(ReturnedRuntime {
5534 connection,
5535 runtime,
5536 }),
5537 )
5538 }
5539 }
5540 }
5541}
5542
5543async fn close_runtime(
5546 mut runtime: Box<dyn RuntimeConnection>,
5547 process_group: Option<u32>,
5548) -> std::result::Result<Value, ServiceError> {
5549 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
5550 Ok(result) => {
5551 result.map_err(operation)?;
5552 Ok(json!({"closed": true}))
5553 }
5554 Err(deadline) => {
5555 let killed = kill_runtime_process_group(process_group);
5560 drop(runtime);
5561 Ok(json!({
5562 "closed": true,
5563 "killed": killed,
5564 "detail": error_message(deadline),
5565 }))
5566 }
5567 }
5568}
5569
5570fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
5573 mutation
5574 .session
5575 .clone()
5576 .filter(|value| !value.trim().is_empty())
5577 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
5578}
5579
5580async fn type_live_command(
5584 runtime: &mut dyn RuntimeConnection,
5585 verb: crate::SessionVerb,
5586 mutation: &crate::SessionMutation,
5587 command: &str,
5588 session: String,
5589) -> std::result::Result<Value, ServiceError> {
5590 within_control_deadline(
5591 &format!("sessions.{}", verb.as_str()),
5592 runtime.send_input(RuntimeInput {
5593 text: command.to_string(),
5594 image_urls: Vec::new(),
5595 }),
5596 )
5597 .await?
5598 .map_err(operation)?;
5599 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
5600 .map_err(session_control_error)?;
5601 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
5602}
5603
5604enum OpenRuntime {
5607 Hosted {
5610 runtime: Box<dyn RuntimeConnection>,
5611 capabilities: crate::RuntimeCapabilities,
5612 workspace: PathBuf,
5613 },
5614 Joined { runtime: Box<dyn RuntimeConnection> },
5617}
5618
5619async fn open_runtime(
5624 method: &str,
5625 params: Value,
5626) -> std::result::Result<OpenRuntime, ServiceError> {
5627 match tokio::time::timeout(
5628 RUNTIME_OPEN_DEADLINE,
5629 open_runtime_unbounded(method, params),
5630 )
5631 .await
5632 {
5633 Ok(result) => result,
5634 Err(_) => Err(ServiceError::Operation(format!(
5635 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
5636 RUNTIME_OPEN_DEADLINE.as_secs()
5637 ))),
5638 }
5639}
5640
5641async fn open_runtime_unbounded(
5642 method: &str,
5643 params: Value,
5644) -> std::result::Result<OpenRuntime, ServiceError> {
5645 match method {
5646 "harness.v1.runtimes.start" => {
5647 let params = decode::<RuntimeStartParams>(params)?;
5648 let backend = runtime_backend(¶ms.backend)?;
5649 let capabilities = backend.capabilities();
5650 let workspace = params.cwd.clone();
5651 let runtime = backend
5652 .start(RuntimeStartRequest {
5653 cwd: params.cwd,
5654 launch: runtime_launch(¶ms.backend),
5655 mcp_servers: params.mcp_servers,
5656 approval_policy: params.approval_policy,
5657 })
5658 .await
5659 .map_err(operation)?;
5660 Ok(OpenRuntime::Hosted {
5661 runtime,
5662 capabilities,
5663 workspace,
5664 })
5665 }
5666 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
5667 let params = decode::<RuntimeAttachParams>(params)?;
5668 let backend = runtime_backend(¶ms.backend)?;
5669 let capabilities = backend.capabilities();
5670 let workspace = params
5671 .cwd
5672 .clone()
5673 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
5674 let runtime = backend
5675 .attach(RuntimeAttachRequest {
5676 runtime_id: params.runtime_id,
5677 cwd: params.cwd,
5678 launch: runtime_launch(¶ms.backend),
5679 mcp_servers: params.mcp_servers,
5680 approval_policy: params.approval_policy,
5681 })
5682 .await
5683 .map_err(operation)?;
5684 Ok(OpenRuntime::Hosted {
5685 runtime,
5686 capabilities,
5687 workspace,
5688 })
5689 }
5690 "harness.v1.runtimes.attach_existing" => {
5691 let params = decode::<RuntimeAttachParams>(params)?;
5692 let backend: Box<dyn RuntimeBackend> = match params
5693 .backend
5694 .base_url
5695 .as_deref()
5696 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
5697 {
5698 Some(endpoint) => {
5699 #[cfg(not(feature = "adapter-api"))]
5700 {
5701 let _ = endpoint;
5702 return Err(ServiceError::UnsupportedAction(
5703 "live HTTP attachment adapter is not compiled".into(),
5704 ));
5705 }
5706 #[cfg(feature = "adapter-api")]
5707 {
5708 let workspace = params.cwd.clone().ok_or_else(|| {
5709 ServiceError::InvalidParams(
5710 "Volter Harness live attach requires the project cwd".into(),
5711 )
5712 })?;
5713 let source = LiveRuntimeSource {
5714 harness: params.backend.harness.as_str().to_string(),
5715 session_id: params.runtime_id.clone(),
5716 workspace,
5717 };
5718 let receipt = resolve_live_runtime(&endpoint, &source)
5719 .map_err(|error| ServiceError::Operation(error.to_string()))?;
5720 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
5721 }
5722 }
5723 None => runtime_backend(¶ms.backend)?,
5724 };
5725 let capabilities = backend.capabilities();
5726 if !capabilities.attach_existing_process {
5727 return Err(ServiceError::Operation(format!(
5728 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
5729 backend.harness().as_str()
5730 )));
5731 }
5732 let runtime = backend
5733 .attach_existing(RuntimeAttachRequest {
5734 runtime_id: params.runtime_id,
5735 cwd: params.cwd,
5736 launch: runtime_launch(¶ms.backend),
5737 mcp_servers: params.mcp_servers,
5738 approval_policy: params.approval_policy,
5739 })
5740 .await
5741 .map_err(operation)?;
5742 Ok(OpenRuntime::Joined { runtime })
5743 }
5744 _ => Err(ServiceError::MethodNotFound),
5745 }
5746}
5747
5748fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
5750 match result {
5751 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
5752 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
5753 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
5754 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
5755 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
5756 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
5757 }
5758}
5759
5760fn runtime_backend(
5761 params: &RuntimeBackendParams,
5762) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5763 if let Some(descriptor) = registry_connect_descriptor(params) {
5764 return open_connect_descriptor(&descriptor, &service_home()?);
5765 }
5766 if params.protocol.as_deref() == Some("acp") {
5767 let launch = params
5768 .launch
5769 .clone()
5770 .or_else(|| {
5771 harness_support_registry()
5772 .harnesses
5773 .into_iter()
5774 .find(|harness| harness.id == params.harness)
5775 .filter(|harness| {
5776 harness.runtime.implementation == ImplementationKind::GenericProtocol
5777 && harness.runtime.protocol.starts_with("acp")
5778 })
5779 .and_then(|harness| harness.runtime.default_launch)
5780 })
5781 .ok_or_else(|| {
5782 ServiceError::InvalidParams(
5783 "an ACP runtime requires `launch` unless the harness has a registered default"
5784 .into(),
5785 )
5786 })?;
5787 let resume_session = harness_support_registry()
5788 .harnesses
5789 .into_iter()
5790 .find(|harness| harness.id == params.harness)
5791 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
5792 return Ok(Box::new(
5793 AcpRuntimeBackend::new(params.harness.clone(), launch)
5794 .with_resume_support(resume_session),
5795 ));
5796 }
5797 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
5798 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
5799 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
5800 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
5801 HarnessId::OPENCODE => match ¶ms.base_url {
5802 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
5803 None => Box::new(OpenCodeRuntimeBackend::new()),
5804 },
5805 harness => {
5806 let descriptor = harness_support_registry()
5807 .harnesses
5808 .into_iter()
5809 .find(|descriptor| descriptor.id.as_str() == harness)
5810 .filter(|descriptor| {
5811 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
5812 && descriptor.runtime.protocol.starts_with("acp")
5813 });
5814 let Some(descriptor) = descriptor else {
5815 return Err(ServiceError::InvalidParams(format!(
5816 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
5817 )));
5818 };
5819 let resume = descriptor.runtime.capabilities.resume_session;
5820 Box::new(
5821 AcpRuntimeBackend::new(
5822 descriptor.id,
5823 descriptor
5824 .runtime
5825 .default_launch
5826 .expect("generic ACP registry entry includes its launch"),
5827 )
5828 .with_resume_support(resume),
5829 )
5830 }
5831 };
5832 Ok(backend)
5833}
5834
5835fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
5836 if let Some(launch) = ¶ms.launch {
5837 return Some(launch.clone());
5838 }
5839 if !matches!(params.policy, RuntimePolicy::Yolo) {
5840 return None;
5841 }
5842 let launch = match params.harness.as_str() {
5843 HarnessId::GROK => RuntimeLaunch {
5844 program: "grok".into(),
5845 arguments: {
5846 let mut arguments: Vec<String> = Vec::new();
5847 if crate::support::self_sandbox_supported() {
5848 arguments.extend(["--sandbox".into(), "workspace".into()]);
5849 }
5850 arguments.extend([
5851 "--always-approve".into(),
5852 "agent".into(),
5853 "--no-leader".into(),
5854 "stdio".into(),
5855 ]);
5856 arguments
5857 },
5858 env: crate::support::grok_env(),
5859 },
5860 HarnessId::CODEX => RuntimeLaunch {
5861 program: "codex".into(),
5862 arguments: vec![
5863 "--dangerously-bypass-approvals-and-sandbox".into(),
5864 "--dangerously-bypass-hook-trust".into(),
5865 "app-server".into(),
5866 ],
5867 env: BTreeMap::new(),
5868 },
5869 HarnessId::CLAUDE_CODE => RuntimeLaunch {
5870 program: "claude".into(),
5871 arguments: vec![
5872 "--dangerously-skip-permissions".into(),
5873 "--print".into(),
5874 "--input-format".into(),
5875 "stream-json".into(),
5876 "--output-format".into(),
5877 "stream-json".into(),
5878 "--verbose".into(),
5879 ],
5880 env: BTreeMap::new(),
5881 },
5882 HarnessId::PI => RuntimeLaunch {
5883 program: "pi".into(),
5884 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
5885 env: BTreeMap::new(),
5886 },
5887 HarnessId::OPENCODE => RuntimeLaunch {
5888 program: "opencode".into(),
5889 arguments: vec!["serve".into()],
5890 env: BTreeMap::new(),
5891 },
5892 HarnessId::GEMINI => RuntimeLaunch {
5893 program: "gemini".into(),
5894 arguments: vec!["--acp".into(), "--yolo".into()],
5895 env: BTreeMap::new(),
5896 },
5897 HarnessId::GOOSE => RuntimeLaunch {
5898 program: "goose".into(),
5899 arguments: vec!["acp".into()],
5900 env: BTreeMap::new(),
5901 },
5902 HarnessId::SUPERCODE => RuntimeLaunch {
5903 program: "supercode".into(),
5904 arguments: vec!["acp".into(), "--dangerous".into()],
5905 env: BTreeMap::new(),
5906 },
5907 _ => return None,
5908 };
5909 Some(launch)
5910}
5911
5912struct IsolatedProbeHome {
5918 launch: RuntimeLaunch,
5919 root: PathBuf,
5920}
5921
5922impl IsolatedProbeHome {
5923 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
5924 let root = std::env::temp_dir().join(format!(
5925 "supercode-harness-probe-{harness}-{}",
5926 generated_session_id()
5927 ));
5928 std::fs::create_dir_all(&root)?;
5929 set_private_dir_permissions(&root)?;
5930
5931 if let Some(source_home) = supercode_interchange::user_home()
5932 .map(std::path::PathBuf::into_os_string)
5933 .map(PathBuf::from)
5934 {
5935 for relative in probe_auth_files(harness) {
5936 copy_probe_file(&source_home, &root, relative)?;
5937 }
5938 }
5939 if harness == HarnessId::SUPERCODE {
5944 let config_home = crate::agent::global_instructions_dir();
5945 for file in ["config.toml", "credentials.toml"] {
5946 copy_probe_path(
5947 &config_home.join(file),
5948 &root.join(".config/supercode").join(file),
5949 )?;
5950 }
5951 }
5952 configure_isolated_probe_auth(harness, &root)?;
5953
5954 let root_text = root.to_string_lossy().into_owned();
5955 for (key, value) in [
5956 ("HOME", root_text.clone()),
5957 (
5958 "XDG_CACHE_HOME",
5959 root.join(".cache").to_string_lossy().into_owned(),
5960 ),
5961 (
5962 "XDG_CONFIG_HOME",
5963 root.join(".config").to_string_lossy().into_owned(),
5964 ),
5965 (
5966 "XDG_DATA_HOME",
5967 root.join(".local/share").to_string_lossy().into_owned(),
5968 ),
5969 ] {
5970 launch.env.insert(key.into(), value);
5971 }
5972 let scoped = match harness {
5973 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
5974 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
5975 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
5976 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
5977 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
5978 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
5979 _ => None,
5980 };
5981 if let Some((key, value)) = scoped {
5982 launch
5983 .env
5984 .insert(key.into(), value.to_string_lossy().into_owned());
5985 }
5986 Ok(Self { launch, root })
5987 }
5988
5989 fn cleanup(&self) -> std::io::Result<()> {
5990 match std::fs::remove_dir_all(&self.root) {
5991 Ok(()) => Ok(()),
5992 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
5993 Err(error) => Err(error),
5994 }
5995 }
5996}
5997
5998impl Drop for IsolatedProbeHome {
5999 fn drop(&mut self) {
6000 let _ = self.cleanup();
6001 }
6002}
6003
6004fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6005 match harness {
6006 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6007 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6011 HarnessId::CODEX => &[".codex/auth.json"],
6012 HarnessId::GEMINI => &[
6013 ".gemini/google_accounts.json",
6014 ".gemini/oauth_creds.json",
6015 ".gemini/settings.json",
6016 ],
6017 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6018 HarnessId::OPENCODE => &[
6019 ".config/opencode/auth.json",
6020 ".local/share/opencode/auth.json",
6021 ],
6022 HarnessId::PI => &[".pi/agent/auth.json"],
6023 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6028 _ => &[],
6029 }
6030}
6031
6032fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6033 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6034}
6035
6036fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6037 if !source.is_file() {
6038 return Ok(());
6039 }
6040 if let Some(parent) = destination.parent() {
6041 std::fs::create_dir_all(parent)?;
6042 set_private_dir_permissions(parent)?;
6043 }
6044 std::fs::copy(source, destination)?;
6045 set_private_file_permissions(destination)
6046}
6047
6048fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6049 if harness != HarnessId::GEMINI {
6050 return Ok(());
6051 }
6052 let oauth = probe_home.join(".gemini/oauth_creds.json");
6053 if !oauth.is_file() {
6054 return Ok(());
6055 }
6056 let settings_path = probe_home.join(".gemini/settings.json");
6057 let mut settings = std::fs::read_to_string(&settings_path)
6058 .ok()
6059 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6060 .unwrap_or_else(|| json!({}));
6061 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6062 std::fs::write(
6063 &settings_path,
6064 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6065 )?;
6066 set_private_file_permissions(&settings_path)
6067}
6068
6069#[cfg(unix)]
6070fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6071 use std::os::unix::fs::PermissionsExt;
6072 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6073}
6074
6075#[cfg(not(unix))]
6076fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6077 Ok(())
6078}
6079
6080#[cfg(unix)]
6081fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6082 use std::os::unix::fs::PermissionsExt;
6083 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6084}
6085
6086#[cfg(not(unix))]
6087fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6088 Ok(())
6089}
6090
6091fn find_executable(program: &str) -> Option<PathBuf> {
6092 let candidate = PathBuf::from(program);
6093 if candidate.components().count() > 1 {
6094 return candidate.is_file().then_some(candidate);
6095 }
6096 let path = std::env::var_os("PATH")?;
6097 for directory in std::env::split_paths(&path) {
6098 let candidate = directory.join(program);
6099 if candidate.is_file() {
6100 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6101 }
6102 #[cfg(windows)]
6103 {
6104 for extension in ["exe", "cmd", "bat"] {
6105 let candidate = directory.join(format!("{program}.{extension}"));
6106 if candidate.is_file() {
6107 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6108 }
6109 }
6110 }
6111 }
6112 None
6113}
6114
6115async fn executable_version(executable: &Path) -> Option<String> {
6116 let mut command = tokio::process::Command::new(executable);
6117 command
6118 .arg("--version")
6119 .stdin(std::process::Stdio::null())
6120 .stdout(std::process::Stdio::piped())
6121 .stderr(std::process::Stdio::piped())
6122 .kill_on_drop(true);
6123 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6124 .await
6125 .ok()?
6126 .ok()?;
6127 let stdout = String::from_utf8_lossy(&output.stdout);
6128 let stderr = String::from_utf8_lossy(&output.stderr);
6129 stdout
6130 .lines()
6131 .chain(stderr.lines())
6132 .map(str::trim)
6133 .find(|line| !line.is_empty())
6134 .map(|line| truncate_text(line, 200))
6135}
6136
6137pub(crate) fn auth_evidence(harness: &str) -> bool {
6138 let env_names: &[&str] = match harness {
6139 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6140 HarnessId::CODEX => &["OPENAI_API_KEY"],
6141 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6142 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6143 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6144 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6145 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6146 _ => &[],
6147 };
6148 if env_names
6149 .iter()
6150 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6151 {
6152 return true;
6153 }
6154 let Some(home) = supercode_interchange::user_home()
6155 .map(std::path::PathBuf::into_os_string)
6156 .map(PathBuf::from)
6157 else {
6158 return false;
6159 };
6160 let files: Vec<PathBuf> = match harness {
6161 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6162 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6163 HarnessId::OPENCODE => vec![
6164 home.join(".local/share/opencode/auth.json"),
6165 home.join(".config/opencode/auth.json"),
6166 ],
6167 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6168 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6169 HarnessId::GEMINI => vec![
6170 home.join(".gemini/oauth_creds.json"),
6171 home.join(".gemini/google_accounts.json"),
6172 ],
6173 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6174 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6175 _ => Vec::new(),
6176 };
6177 if files.into_iter().any(|path| {
6178 std::fs::metadata(path)
6179 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6180 .unwrap_or(false)
6181 }) {
6182 return true;
6183 }
6184 if harness == HarnessId::CLAUDE_CODE {
6191 return std::fs::read_to_string(home.join(".claude.json"))
6192 .map(|text| text.contains("\"oauthAccount\""))
6193 .unwrap_or(false);
6194 }
6195 false
6196}
6197
6198fn looks_like_auth_error(message: &str) -> bool {
6199 let message = message.to_ascii_lowercase();
6200 [
6201 "auth",
6202 "login",
6203 "sign in",
6204 "sign-in",
6205 "credential",
6206 "unauthorized",
6207 "forbidden",
6208 "token",
6209 ]
6210 .iter()
6211 .any(|needle| message.contains(needle))
6212}
6213
6214fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6215 crate::RuntimeCapabilities {
6216 start_session: false,
6217 resume_session: false,
6218 attach_existing_process: false,
6219 send_input: false,
6220 stream_events: false,
6221 interrupt: false,
6222 steer: false,
6223 respond_to_requests: false,
6224 }
6225}
6226
6227fn truncate_text(text: &str, max_chars: usize) -> String {
6228 let mut chars = text.chars();
6229 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6230 if chars.next().is_some() {
6231 format!("{truncated}…")
6232 } else {
6233 truncated
6234 }
6235}
6236
6237fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6244 match &handle.endpoint {
6245 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6246 crate::RuntimeEndpoint::Http { .. } => None,
6247 }
6248}
6249
6250fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6256 match process_group {
6257 #[cfg(unix)]
6258 Some(pid) => {
6259 crate::lsp::kill_process_group(pid);
6260 true
6261 }
6262 #[cfg(not(unix))]
6263 Some(_) => false,
6264 None => false,
6265 }
6266}
6267
6268fn error_message(error: ServiceError) -> String {
6269 match error {
6270 ServiceError::InvalidParams(message)
6271 | ServiceError::Operation(message)
6272 | ServiceError::UnsupportedAction(message) => message,
6273 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6274 ServiceError::Sdk(error) => error.to_string(),
6275 }
6276}
6277
6278#[derive(Debug)]
6279enum ServiceError {
6280 InvalidParams(String),
6281 MethodNotFound,
6282 UnsupportedAction(String),
6283 Operation(String),
6284 Sdk(SdkError),
6285}
6286
6287fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6288 match error {
6289 ServiceError::InvalidParams(message) => {
6290 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6291 }
6292 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6293 SdkError::unsupported(operation)
6294 }
6295 ServiceError::Operation(message) => {
6296 let code = if message.contains("already in progress") {
6297 SdkErrorCode::Busy
6298 } else if message.contains("not supported by this runtime") {
6299 SdkErrorCode::UnsupportedAction
6300 } else if message.contains("unknown runtime connection") {
6301 SdkErrorCode::NotFound
6302 } else {
6303 SdkErrorCode::Execution
6304 };
6305 SdkError::new(code, operation, message)
6306 }
6307 ServiceError::Sdk(error) => error,
6308 }
6309}
6310
6311fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6312 let error_code = error.code();
6313 let code = match error_code {
6314 SdkErrorCode::Unauthenticated => -32030,
6315 SdkErrorCode::Unauthorized => -32031,
6316 SdkErrorCode::ControllerRequired => -32032,
6317 SdkErrorCode::LeaseExpired => -32033,
6318 SdkErrorCode::InvalidArgument => -32602,
6319 SdkErrorCode::NotFound => -32004,
6320 SdkErrorCode::Busy => -32000,
6321 SdkErrorCode::UnsupportedAction => -32020,
6322 SdkErrorCode::Execution => -32002,
6323 SdkErrorCode::Transport => -32003,
6324 };
6325 json!({
6326 "jsonrpc": "2.0",
6327 "id": id,
6328 "error": {
6329 "code": code,
6330 "name": error_code,
6331 "operation": error.operation(),
6332 "message": error.to_string(),
6333 },
6334 })
6335}
6336
6337fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6338 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6339}
6340
6341fn operation(error: impl Into<crate::Error>) -> ServiceError {
6342 let error = error.into();
6343 match error {
6344 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6345 error => ServiceError::Operation(error.to_string()),
6346 }
6347}
6348
6349#[derive(Debug, Clone, Deserialize, Default)]
6353#[serde(default)]
6354struct MemoryRequest {
6355 harness: Option<String>,
6357 query: Option<String>,
6359 profile: Option<String>,
6361 session: Option<String>,
6363 full: bool,
6365 regex: bool,
6367 cwd: Option<std::path::PathBuf>,
6369 homes: crate::HarnessHomes,
6371}
6372
6373fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6376 let request = decode::<MemoryRequest>(params)?;
6377 let harness = request
6378 .harness
6379 .clone()
6380 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6381 let to_service = |error: crate::memory::MemoryError| match error {
6382 crate::memory::MemoryError::UnsupportedHarness { .. }
6383 | crate::memory::MemoryError::SessionNotScoped { .. } => {
6384 ServiceError::UnsupportedAction(error.to_string())
6385 }
6386 other => ServiceError::InvalidParams(other.to_string()),
6387 };
6388 match method {
6389 "harness.v1.memory.show" => {
6390 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6391 harness,
6392 profile: request.profile,
6393 session: request.session,
6394 full: request.full,
6395 cwd: request.cwd,
6396 homes: request.homes,
6397 })
6398 .map_err(to_service)?;
6399 Ok(json!({
6400 "schema": crate::memory::MEMORY_SCHEMA,
6401 "documents": documents,
6402 }))
6403 }
6404 "harness.v1.memory.search" => {
6405 let query = request
6406 .query
6407 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6408 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6409 harness,
6410 query,
6411 profile: request.profile,
6412 regex: request.regex,
6413 cwd: request.cwd,
6414 homes: request.homes,
6415 })
6416 .map_err(to_service)?;
6417 Ok(json!({
6418 "schema": crate::memory::MEMORY_SCHEMA,
6419 "matches": matches,
6420 }))
6421 }
6422 _ => Err(ServiceError::MethodNotFound),
6423 }
6424}
6425
6426#[derive(Debug, Clone, Deserialize)]
6430#[serde(default)]
6431struct ProfilesQuery {
6432 harness: Option<String>,
6434 name: Option<String>,
6436 homes: crate::HarnessHomes,
6438}
6439
6440impl Default for ProfilesQuery {
6441 fn default() -> Self {
6442 Self {
6443 harness: None,
6444 name: None,
6445 homes: crate::HarnessHomes::default(),
6446 }
6447 }
6448}
6449
6450fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6453 let query = decode::<ProfilesQuery>(params)?;
6454 let to_service = |error: crate::profiles::ProfileError| match error {
6455 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6456 ServiceError::UnsupportedAction(error.to_string())
6457 }
6458 crate::profiles::ProfileError::NotFound { .. } => {
6459 ServiceError::InvalidParams(error.to_string())
6460 }
6461 };
6462 match method {
6463 "harness.v1.profiles.list" => {
6464 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6465 .map_err(to_service)?;
6466 Ok(json!({
6467 "schema": crate::profiles::PROFILES_SCHEMA,
6468 "profiles": profiles,
6469 }))
6470 }
6471 "harness.v1.profiles.get" => {
6472 let harness = query
6473 .harness
6474 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6475 let name = query
6476 .name
6477 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6478 let profile =
6479 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6480 Ok(json!({
6481 "schema": crate::profiles::PROFILES_SCHEMA,
6482 "profile": profile,
6483 }))
6484 }
6485 _ => Err(ServiceError::MethodNotFound),
6486 }
6487}
6488
6489#[derive(Debug, Clone, Deserialize)]
6493#[serde(default)]
6494struct ChannelsQuery {
6495 harness: Option<String>,
6497 name: Option<String>,
6499 homes: crate::HarnessHomes,
6501}
6502
6503impl Default for ChannelsQuery {
6504 fn default() -> Self {
6505 Self {
6506 harness: None,
6507 name: None,
6508 homes: crate::HarnessHomes::default(),
6509 }
6510 }
6511}
6512
6513#[derive(Debug, Clone, Deserialize)]
6517#[serde(default)]
6518struct RoutesQuery {
6519 harness: Option<String>,
6520 profile: Option<String>,
6522 homes: crate::HarnessHomes,
6523}
6524
6525impl Default for RoutesQuery {
6526 fn default() -> Self {
6527 Self {
6528 harness: None,
6529 profile: None,
6530 homes: crate::HarnessHomes::default(),
6531 }
6532 }
6533}
6534
6535#[derive(Debug, Clone, Deserialize)]
6536#[serde(default)]
6537struct TriggersQuery {
6538 harness: Option<String>,
6539 homes: crate::HarnessHomes,
6540}
6541
6542impl Default for TriggersQuery {
6543 fn default() -> Self {
6544 Self {
6545 harness: None,
6546 homes: crate::HarnessHomes::default(),
6547 }
6548 }
6549}
6550
6551fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
6552 let query = decode::<TriggersQuery>(params)?;
6553 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
6554 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6555 Ok(json!({
6556 "schema": crate::triggers::TRIGGERS_SCHEMA,
6557 "triggers": triggers,
6558 }))
6559}
6560
6561fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
6562 let query = decode::<RoutesQuery>(params)?;
6563 let routes = crate::routes::list_routes(
6564 &query.homes,
6565 query.harness.as_deref(),
6566 query.profile.as_deref(),
6567 )
6568 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6569 Ok(json!({
6570 "schema": crate::routes::ROUTES_SCHEMA,
6571 "routes": routes,
6572 }))
6573}
6574
6575fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6576 let query = decode::<ChannelsQuery>(params)?;
6577 let to_service = |error: crate::channels::ChannelError| match error {
6578 crate::channels::ChannelError::UnsupportedHarness { .. } => {
6579 ServiceError::UnsupportedAction(error.to_string())
6580 }
6581 crate::channels::ChannelError::NotFound { .. } => {
6582 ServiceError::InvalidParams(error.to_string())
6583 }
6584 };
6585 match method {
6586 "harness.v1.channels.list" => {
6587 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
6588 .map_err(to_service)?;
6589 Ok(json!({
6590 "schema": crate::channels::CHANNELS_SCHEMA,
6591 "channels": channels,
6592 }))
6593 }
6594 "harness.v1.channels.status" => {
6595 let harness = query
6596 .harness
6597 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6598 let name = query
6599 .name
6600 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6601 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
6602 .map_err(to_service)?;
6603 Ok(json!({
6604 "schema": crate::channels::CHANNELS_SCHEMA,
6605 "channel": channel,
6606 }))
6607 }
6608 _ => Err(ServiceError::MethodNotFound),
6609 }
6610}
6611
6612fn rpc_error(id: Value, code: i64, message: &str) -> Value {
6613 json!({
6614 "jsonrpc": "2.0",
6615 "id": id,
6616 "error": {"code": code, "message": message},
6617 })
6618}