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 let waiting = doors.all().iter().any(|live| {
3684 live.address.harness == session.locator.harness.as_str()
3685 && live.address.session_id == session.locator.session_id
3686 && live.status == "waiting"
3687 });
3688 if waiting && session.locator.harness.as_str() == "claude-code" {
3689 if let Some(request) = crate::mail_route::pending_request(session.locator.storage.path()) {
3690 value["pending_request"] = request;
3691 }
3692 }
3693 Ok(value)
3694}
3695
3696fn live_index_changes(
3697 changes: Vec<crate::session_index::SessionIndexChange>,
3698 homes: &HarnessHomes,
3699) -> std::result::Result<Vec<Value>, ServiceError> {
3700 use crate::session_index::SessionIndexChange;
3701 let doors = crate::mail_route::LiveSessions::read(homes);
3702 changes
3703 .into_iter()
3704 .map(|change| match change {
3705 SessionIndexChange::Added { descriptor } => Ok(json!({
3706 "kind": "added",
3707 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3708 })),
3709 SessionIndexChange::Updated { descriptor } => Ok(json!({
3710 "kind": "updated",
3711 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3712 })),
3713 SessionIndexChange::Removed { key } => Ok(json!({
3714 "kind": "removed",
3715 "key": key,
3716 })),
3717 })
3718 .collect()
3719}
3720
3721fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
3722 use crate::{SessionPresence, SessionTurnState};
3723 match (activity.presence, activity.turn) {
3724 (SessionPresence::Persisted, _) => None,
3725 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
3726 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
3727 (SessionPresence::Running, SessionTurnState::Unknown)
3731 if activity.evidence.native_state.is_none() =>
3732 {
3733 None
3734 }
3735 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
3736 }
3737}
3738
3739#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
3740#[serde(rename_all = "kebab-case")]
3741enum TransferFormat {
3742 ClaudeCode,
3743 Codex,
3744 #[serde(rename = "opencode", alias = "open-code")]
3745 OpenCode,
3746 Pi,
3747 Grok,
3748 Gemini,
3749 Goose,
3750 Hermes,
3754}
3755
3756impl TransferFormat {
3757 fn id(self) -> &'static str {
3758 match self {
3759 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
3760 Self::Codex => HarnessId::CODEX,
3761 Self::OpenCode => HarnessId::OPENCODE,
3762 Self::Pi => HarnessId::PI,
3763 Self::Grok => HarnessId::GROK,
3764 Self::Gemini => HarnessId::GEMINI,
3765 Self::Goose => HarnessId::GOOSE,
3766 Self::Hermes => HarnessId::HERMES,
3767 }
3768 }
3769}
3770
3771impl From<TransferFormat> for SessionFormat {
3772 fn from(value: TransferFormat) -> Self {
3773 match value {
3774 TransferFormat::ClaudeCode => Self::ClaudeCode,
3775 TransferFormat::Codex => Self::Codex,
3776 TransferFormat::OpenCode => Self::OpenCode,
3777 TransferFormat::Pi => Self::Pi,
3778 TransferFormat::Grok => Self::Grok,
3779 TransferFormat::Gemini => Self::Gemini,
3780 TransferFormat::Goose => Self::Goose,
3781 TransferFormat::Hermes => Self::Codex,
3783 }
3784 }
3785}
3786
3787#[derive(Deserialize)]
3788struct ImportSessionParams {
3789 source_harness: TransferFormat,
3790 content: String,
3791}
3792
3793#[derive(Deserialize)]
3794struct ExportSessionParams {
3795 locator: SessionLocator,
3796 target_harness: TransferFormat,
3797}
3798
3799#[derive(Deserialize)]
3800struct ReduceSessionParams {
3801 locator: SessionLocator,
3802 target_harness: TransferFormat,
3803 #[serde(default = "default_keep_last")]
3804 keep_last: usize,
3805}
3806
3807fn default_keep_last() -> usize {
3808 6
3809}
3810
3811#[derive(Deserialize)]
3812struct BranchSessionParams {
3813 locator: SessionLocator,
3814 #[serde(default)]
3815 target_harness: Option<TransferFormat>,
3816}
3817
3818#[derive(Deserialize)]
3819struct HandoffSessionParams {
3820 locator: SessionLocator,
3821 target_harness: TransferFormat,
3822 #[serde(default)]
3823 cwd: Option<PathBuf>,
3824}
3825
3826#[derive(Deserialize)]
3827struct MaterializeSessionParams {
3828 artifact: crate::native_materialize::MaterializeArtifact,
3829 cwd: PathBuf,
3830 #[serde(default)]
3832 homes: HarnessHomes,
3833}
3834
3835#[derive(Debug, Clone, Copy, Default, Deserialize)]
3838#[serde(rename_all = "snake_case")]
3839enum ResumePolicy {
3840 Default,
3841 #[default]
3842 Yolo,
3843}
3844
3845#[derive(Deserialize)]
3846struct ResumeInstructionsParams {
3847 locator: SessionLocator,
3848 #[serde(default)]
3849 cwd: Option<PathBuf>,
3850 #[serde(default)]
3851 policy: ResumePolicy,
3852}
3853
3854#[derive(Deserialize)]
3856struct WorkflowLoadParams {
3857 from: crate::workflow_doors::WorkflowHarness,
3858 home: PathBuf,
3859}
3860
3861#[derive(Deserialize)]
3864struct OrchestrationLoadParams {
3865 root: PathBuf,
3866 #[serde(default)]
3867 flavor: crate::orchestration_doors::HomeFlavor,
3868}
3869
3870#[derive(Deserialize)]
3873struct OrchestrationSaveParams {
3874 root: PathBuf,
3875 orchestration: crate::orchestration::Orchestration,
3876 #[serde(default)]
3877 vault: BTreeMap<String, String>,
3878}
3879
3880#[derive(Deserialize)]
3882struct OrchestrationCompileParams {
3883 from: crate::orchestration_doors::OrchestrationHarness,
3884 home: PathBuf,
3885}
3886
3887#[derive(Deserialize)]
3891struct OrchestrationDecompileParams {
3892 to: crate::orchestration_doors::OrchestrationHarness,
3893 orchestration: crate::orchestration::Orchestration,
3894 source: PathBuf,
3895 #[serde(default)]
3896 source_flavor: crate::orchestration_doors::SourceFlavor,
3897 dest: PathBuf,
3898 #[serde(default)]
3899 vault: BTreeMap<String, String>,
3900}
3901
3902#[derive(Deserialize)]
3905struct OrchestrationImportParams {
3906 from: crate::orchestration_doors::OrchestrationHarness,
3907 home: PathBuf,
3908 into: PathBuf,
3909}
3910
3911#[derive(Deserialize)]
3914struct OrchestrationExportParams {
3915 to: crate::orchestration_doors::OrchestrationHarness,
3916 root: PathBuf,
3917 dest: PathBuf,
3918}
3919
3920#[derive(Deserialize)]
3922struct JobsGetParams {
3923 harness: String,
3924 id: String,
3925 #[serde(default)]
3926 homes: crate::HarnessHomes,
3927}
3928
3929fn mutate_job(
3937 verb: crate::jobs_control::JobVerb,
3938 params: Value,
3939) -> std::result::Result<Value, ServiceError> {
3940 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
3941 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
3942 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
3943 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3944}
3945
3946fn mutate_skill(
3953 verb: crate::skills_control::SkillVerb,
3954 params: Value,
3955) -> std::result::Result<Value, ServiceError> {
3956 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
3957 if !crate::skills_control::supports_skill_control(&mutation.harness) {
3958 return Err(ServiceError::UnsupportedAction(format!(
3959 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
3960 mutation.harness,
3961 verb.as_str(),
3962 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
3963 )));
3964 }
3965 let outcome =
3966 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
3967 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3968}
3969
3970fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
3972 match error {
3973 crate::skills_control::SkillControlError::Unsupported(message) => {
3974 ServiceError::UnsupportedAction(message)
3975 }
3976 crate::skills_control::SkillControlError::Invalid(message) => {
3977 ServiceError::InvalidParams(message)
3978 }
3979 crate::skills_control::SkillControlError::Failed(message) => {
3980 ServiceError::Operation(message)
3981 }
3982 }
3983}
3984
3985fn mutate_profile(
3993 verb: crate::profiles_control::ProfileVerb,
3994 params: Value,
3995) -> std::result::Result<Value, ServiceError> {
3996 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
3997 let outcome =
3998 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
3999 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4000}
4001
4002fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4004 match error {
4005 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4006 ServiceError::UnsupportedAction(message)
4007 }
4008 crate::profiles_control::ProfileControlError::Invalid(message) => {
4009 ServiceError::InvalidParams(message)
4010 }
4011 crate::profiles_control::ProfileControlError::Failed(message) => {
4012 ServiceError::Operation(message)
4013 }
4014 }
4015}
4016
4017fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4021 match error {
4022 crate::jobs_control::JobControlError::Unsupported(message) => {
4023 ServiceError::UnsupportedAction(message)
4024 }
4025 crate::jobs_control::JobControlError::Invalid(message) => {
4026 ServiceError::InvalidParams(message)
4027 }
4028 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4029 }
4030}
4031
4032fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4037 match error {
4038 crate::SessionControlError::Unsupported(message) => {
4039 ServiceError::UnsupportedAction(message)
4040 }
4041 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4042 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4043 }
4044}
4045
4046fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4051 if crate::jobs::supports_jobs(harness) {
4052 return Ok(());
4053 }
4054 Err(ServiceError::UnsupportedAction(format!(
4055 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4056 crate::jobs::JOB_HARNESSES.join(", ")
4057 )))
4058}
4059
4060#[derive(Deserialize)]
4062struct RunsGetParams {
4063 harness: String,
4064 id: String,
4065 #[serde(default)]
4066 homes: crate::HarnessHomes,
4067}
4068
4069fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4074 if crate::runs::supports_runs(harness) {
4075 return Ok(());
4076 }
4077 Err(ServiceError::UnsupportedAction(format!(
4078 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4079 crate::runs::RUN_HARNESSES.join(", ")
4080 )))
4081}
4082
4083#[derive(Serialize)]
4084struct SessionArtifact {
4085 source_harness: HarnessId,
4086 target_harness: &'static str,
4087 session_id: Option<String>,
4088 content: String,
4089 suggested_filename: String,
4090 files: Vec<SessionArtifactFile>,
4091 fidelity: Fidelity,
4092 residue: Vec<String>,
4093}
4094
4095#[derive(Serialize)]
4096struct SessionArtifactFile {
4097 path: String,
4098 content: String,
4099 role: ArtifactFileRole,
4100}
4101
4102#[derive(Serialize)]
4103#[serde(rename_all = "snake_case")]
4104enum ArtifactFileRole {
4105 Primary,
4106 Subagent,
4107 Bundle,
4108 SourceRecovery,
4109}
4110
4111#[derive(Serialize)]
4112struct StructuredLaunch {
4113 cwd: PathBuf,
4114 program: String,
4115 arguments: Vec<String>,
4116 env: BTreeMap<String, String>,
4117}
4118
4119struct HandoffInstructions {
4120 launch: StructuredLaunch,
4121 materialize: Option<StructuredLaunch>,
4122 requires_materialization: bool,
4123 note: String,
4124}
4125
4126#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4127#[serde(rename_all = "snake_case")]
4128enum HarnessProbeLevel {
4129 #[default]
4130 Passive,
4131 Handshake,
4132}
4133
4134#[derive(Default, Deserialize)]
4135#[serde(default)]
4136struct HarnessInventoryParams {
4137 harness: Option<HarnessId>,
4138 harnesses: Vec<HarnessId>,
4139 workspace: Option<PathBuf>,
4140 probe: HarnessProbeLevel,
4141 include_sessions: bool,
4142 skip_versions: bool,
4144}
4145
4146#[derive(Deserialize)]
4147struct HarnessAuthenticationParams {
4148 harness: HarnessId,
4149}
4150
4151#[derive(Deserialize)]
4152struct BeginHarnessAuthenticationParams {
4153 harness: HarnessId,
4154 #[serde(default = "local_browser_authentication_environment")]
4155 environment: crate::HarnessAuthenticationEnvironment,
4156 #[serde(default)]
4157 method: Option<crate::HarnessAuthenticationMethodId>,
4158 #[serde(default)]
4159 cwd: Option<PathBuf>,
4160}
4161
4162fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4163 crate::HarnessAuthenticationEnvironment::LocalBrowser
4164}
4165
4166#[derive(Serialize)]
4167struct HarnessInventoryReport {
4168 probe: HarnessProbeLevel,
4169 workspace: Option<PathBuf>,
4170 harnesses: Vec<LocalHarness>,
4171}
4172
4173#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4174#[serde(rename_all = "snake_case")]
4175enum HarnessAuthState {
4176 Ready,
4177 Configured,
4178 Required,
4179 Unknown,
4180}
4181
4182#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4183#[serde(rename_all = "snake_case")]
4184enum HarnessRuntimeState {
4185 Ready,
4186 Degraded,
4187 Unavailable,
4188}
4189
4190#[derive(Serialize)]
4191struct HarnessSessionCounts {
4192 global: Option<usize>,
4193 workspace: Option<usize>,
4194}
4195
4196#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4207#[serde(rename_all = "snake_case")]
4208pub enum GatewayState {
4209 Up,
4210 Down,
4211 Unknown,
4212}
4213
4214#[derive(Debug, Clone, Serialize)]
4216pub struct GatewayHealth {
4217 pub state: GatewayState,
4218 #[serde(skip_serializing_if = "Option::is_none")]
4223 pub endpoint: Option<String>,
4224 #[serde(skip_serializing_if = "Option::is_none")]
4225 pub version: Option<String>,
4226 pub evidence: String,
4228 pub checked_at_ms: u64,
4229}
4230
4231fn openclaw_gateway_endpoint(home: &Path) -> String {
4235 let config_path = home.join(".openclaw/openclaw.json");
4236 let gateway = std::fs::read_to_string(&config_path)
4237 .ok()
4238 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4239 .and_then(|config| config.get("gateway").cloned());
4240 if let Some(url) = gateway
4241 .as_ref()
4242 .and_then(|gateway| gateway.get("url"))
4243 .and_then(serde_json::Value::as_str)
4244 {
4245 return url.to_string();
4246 }
4247 let port = gateway
4248 .as_ref()
4249 .and_then(|gateway| gateway.get("port"))
4250 .and_then(serde_json::Value::as_u64)
4251 .unwrap_or(18789);
4252 format!("ws://127.0.0.1:{port}")
4253}
4254
4255fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4262 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4263 let output = std::process::Command::new(&program)
4264 .args(["gateway", "status"])
4265 .stdin(std::process::Stdio::null())
4266 .output()
4267 .ok()?;
4268 let text = format!(
4269 "{}{}",
4270 String::from_utf8_lossy(&output.stdout),
4271 String::from_utf8_lossy(&output.stderr)
4272 );
4273 let verdict = text.lines().find_map(|line| {
4274 let l = line.trim();
4275 if l.contains("supervised by launchd (PID")
4276 || l.contains("supervised by systemd (PID")
4277 || l.contains("Gateway is running")
4278 || l.contains("process is running")
4279 {
4280 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4281 } else if l.contains("not running") || l.contains("not installed") {
4282 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4283 } else {
4284 None
4285 }
4286 });
4287 verdict
4288}
4289
4290fn gateway_health(
4291 id: &str,
4292 installed: bool,
4293 running: Option<&RunningInstance>,
4294 version: Option<&str>,
4295) -> GatewayHealth {
4296 let checked_at_ms = now_epoch_ms();
4297 let home = supercode_interchange::user_home()
4298 .map(std::path::PathBuf::into_os_string)
4299 .map(PathBuf::from);
4300 match id {
4301 HarnessId::HERMES | HarnessId::OPENCLAW => {
4302 let endpoint = (id == HarnessId::OPENCLAW)
4303 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4304 .flatten();
4305 let (state, evidence) = match running {
4306 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4307 None if !installed => (
4308 GatewayState::Unknown,
4309 format!("`{id}` is not installed; no gateway to probe"),
4310 ),
4311 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4312 Some((state, evidence)) => (state, evidence),
4315 None => (
4316 GatewayState::Down,
4317 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4318 ),
4319 },
4320 None => (
4321 GatewayState::Down,
4322 format!(
4323 "no TCP listener at {}",
4324 endpoint.as_deref().unwrap_or("the gateway endpoint")
4325 ),
4326 ),
4327 };
4328 GatewayHealth {
4329 state,
4330 endpoint,
4331 version: version.map(str::to_string),
4332 evidence,
4333 checked_at_ms,
4334 }
4335 }
4336 HarnessId::ORCHESTRATOR => {
4343 let root = crate::HarnessHomes::default().orchestrator;
4344 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4345 Some(lease) if lease.is_live() => (
4346 GatewayState::Up,
4347 format!(
4348 "`{}` names pid {} (started {}), which is live",
4349 crate::orchestrator::lock_path(&root).display(),
4350 lease.pid,
4351 lease.started_at
4352 ),
4353 ),
4354 Some(lease) => (
4355 GatewayState::Down,
4356 format!(
4357 "stale lease `{}`: pid {} is gone",
4358 crate::orchestrator::lock_path(&root).display(),
4359 lease.pid
4360 ),
4361 ),
4362 None => (
4363 GatewayState::Down,
4364 format!(
4365 "no lease at `{}`; `supercode orchestrator start` writes one",
4366 crate::orchestrator::lock_path(&root).display()
4367 ),
4368 ),
4369 };
4370 GatewayHealth {
4371 state,
4372 endpoint: None,
4373 version: version.map(str::to_string),
4374 evidence,
4375 checked_at_ms,
4376 }
4377 }
4378 _ => GatewayHealth {
4379 state: GatewayState::Unknown,
4380 endpoint: None,
4381 version: version.map(str::to_string),
4382 evidence: format!("`{id}` runs per session, not as a gateway"),
4383 checked_at_ms,
4384 },
4385 }
4386}
4387
4388#[derive(Debug, Clone, Serialize)]
4389struct RunningInstance {
4390 method: RunningInstanceMethod,
4392 evidence: String,
4394 checked_at_ms: u64,
4396}
4397
4398#[derive(Debug, Clone, Copy, Serialize)]
4399#[serde(rename_all = "snake_case")]
4400enum RunningInstanceMethod {
4401 GatewayConnect,
4404 StoreWalActivity,
4407}
4408
4409fn now_epoch_ms() -> u64 {
4410 std::time::SystemTime::now()
4411 .duration_since(std::time::UNIX_EPOCH)
4412 .map(|elapsed| elapsed.as_millis() as u64)
4413 .unwrap_or(0)
4414}
4415
4416fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4420 let config_path = home.join(".openclaw/openclaw.json");
4421 let text = std::fs::read_to_string(&config_path).ok();
4422 let gateway = text
4423 .as_deref()
4424 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4425 .and_then(|config| config.get("gateway").cloned());
4426 let address = gateway
4427 .as_ref()
4428 .and_then(|gateway| gateway.get("url"))
4429 .and_then(serde_json::Value::as_str)
4430 .and_then(|url| {
4431 url.split("://").nth(1).map(|rest| {
4432 rest.trim_end_matches('/')
4433 .split('/')
4434 .next()
4435 .unwrap_or(rest)
4436 .to_string()
4437 })
4438 })
4439 .unwrap_or_else(|| {
4440 let port = gateway
4441 .as_ref()
4442 .and_then(|gateway| gateway.get("port"))
4443 .and_then(serde_json::Value::as_u64)
4444 .unwrap_or(18789);
4445 format!("127.0.0.1:{port}")
4446 });
4447 let reachable = std::net::TcpStream::connect_timeout(
4448 &address.parse().ok()?,
4449 std::time::Duration::from_millis(400),
4450 )
4451 .is_ok();
4452 reachable.then(|| RunningInstance {
4453 method: RunningInstanceMethod::GatewayConnect,
4454 evidence: format!(
4455 "gateway endpoint {address} accepted a TCP connect (from {})",
4456 config_path.display()
4457 ),
4458 checked_at_ms: now_epoch_ms(),
4459 })
4460}
4461
4462fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4467 let wal = home.join(".hermes/state.db-wal");
4468 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4469 let age_ms = std::time::SystemTime::now()
4470 .duration_since(modified)
4471 .map(|age| age.as_millis() as u64)
4472 .unwrap_or(u64::MAX);
4473 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4474 method: RunningInstanceMethod::StoreWalActivity,
4475 evidence: format!(
4476 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4477 wal.display()
4478 ),
4479 checked_at_ms: now_epoch_ms(),
4480 })
4481}
4482
4483fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4485 let home = supercode_interchange::user_home()
4486 .map(std::path::PathBuf::into_os_string)
4487 .map(PathBuf::from)?;
4488 match id {
4489 HarnessId::OPENCLAW => probe_openclaw_running(&home),
4490 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4491 _ => None,
4492 }
4493}
4494
4495#[derive(Serialize)]
4496struct LocalHarness {
4497 id: HarnessId,
4498 display_name: String,
4499 supported: bool,
4500 installed: bool,
4501 executable: Option<String>,
4502 version: Option<String>,
4503 auth: HarnessAuthState,
4504 runtime: HarnessRuntimeState,
4505 protocol: String,
4506 capabilities: crate::RuntimeCapabilities,
4507 effective_capabilities: crate::RuntimeCapabilities,
4508 sessions: HarnessSessionCounts,
4509 #[serde(skip_serializing_if = "Option::is_none")]
4512 running: Option<RunningInstance>,
4513 gateway: GatewayHealth,
4515 reason: Option<String>,
4516 repair: Option<String>,
4517}
4518
4519#[derive(Clone, Deserialize)]
4520struct RuntimeBackendParams {
4521 harness: HarnessId,
4522 #[serde(default)]
4523 protocol: Option<String>,
4524 #[serde(default)]
4525 launch: Option<RuntimeLaunch>,
4526 #[serde(default)]
4527 base_url: Option<String>,
4528 #[serde(default)]
4529 policy: RuntimePolicy,
4530}
4531
4532#[derive(Debug, Clone, Copy, Default, Deserialize)]
4535#[serde(rename_all = "snake_case")]
4536enum RuntimePolicy {
4537 Default,
4538 #[default]
4539 Yolo,
4540}
4541
4542#[derive(Deserialize)]
4543struct RuntimeStartParams {
4544 #[serde(flatten)]
4545 backend: RuntimeBackendParams,
4546 cwd: PathBuf,
4547 #[serde(default)]
4550 mcp_servers: Vec<crate::McpServerLaunch>,
4551 #[serde(default)]
4553 approval_policy: Option<String>,
4554}
4555
4556#[derive(Deserialize)]
4557struct RuntimeAttachParams {
4558 #[serde(flatten)]
4559 backend: RuntimeBackendParams,
4560 runtime_id: String,
4561 #[serde(default)]
4562 cwd: Option<PathBuf>,
4563 #[serde(default)]
4566 mcp_servers: Vec<crate::McpServerLaunch>,
4567 #[serde(default)]
4569 approval_policy: Option<String>,
4570}
4571
4572#[derive(Deserialize)]
4573struct RuntimeConnectionParams {
4574 connection: String,
4575}
4576
4577#[derive(Deserialize)]
4578struct RuntimeInputParams {
4579 connection: String,
4580 text: String,
4581 #[serde(default)]
4582 image_urls: Vec<String>,
4583}
4584
4585const MAX_RUNTIME_IMAGES: usize = 4;
4586const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
4587const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
4588
4589fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
4590 if image_urls.len() > MAX_RUNTIME_IMAGES {
4591 return Err(ServiceError::InvalidParams(format!(
4592 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
4593 )));
4594 }
4595 let mut total = 0usize;
4596 for url in &image_urls {
4597 if !(url.starts_with("data:image/")
4598 || url.starts_with("https://")
4599 || url.starts_with("http://"))
4600 {
4601 return Err(ServiceError::InvalidParams(
4602 "runtime images must be image data URLs or HTTP(S) URLs".into(),
4603 ));
4604 }
4605 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
4606 return Err(ServiceError::InvalidParams(format!(
4607 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
4608 )));
4609 }
4610 total = total.saturating_add(url.len());
4611 }
4612 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
4613 return Err(ServiceError::InvalidParams(format!(
4614 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
4615 )));
4616 }
4617 Ok(image_urls)
4618}
4619
4620#[derive(Deserialize)]
4621struct RuntimeRespondParams {
4622 connection: String,
4623 request_id: Value,
4624 response: Value,
4625}
4626
4627fn default_reduction_store_root() -> PathBuf {
4628 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
4629 return PathBuf::from(root).join("sessions");
4630 }
4631 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
4632 return PathBuf::from(home).join(".supercode").join("sessions");
4633 }
4634 PathBuf::from(".supercode").join("sessions")
4635}
4636
4637fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
4638 let mut output = String::new();
4639 for message in messages {
4640 output.push_str(
4641 &serde_json::to_string(message)
4642 .map_err(|error| ServiceError::Operation(error.to_string()))?,
4643 );
4644 output.push('\n');
4645 }
4646 Ok(output)
4647}
4648
4649fn parse_messages_jsonl(
4650 content: &str,
4651) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
4652 content
4653 .lines()
4654 .enumerate()
4655 .filter(|(_, line)| !line.trim().is_empty())
4656 .map(|(index, line)| {
4657 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
4658 ServiceError::Operation(format!(
4659 "reduced transcript line {} is invalid: {error}",
4660 index + 1
4661 ))
4662 })
4663 })
4664 .collect()
4665}
4666
4667fn reduced_bootstrap_prompt(
4668 source: &SessionLocator,
4669 target: TransferFormat,
4670 view_jsonl: &str,
4671 sidecar_path: &Path,
4672 reduction_log_path: &Path,
4673) -> String {
4674 format!(
4675 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
4676 \n\
4677 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\
4678 \n\
4679 <supercode-reduced-session source-session=\"{source_id}\">\n\
4680 {view_jsonl}\
4681 </supercode-reduced-session>\n\
4682 \n\
4683 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
4684 source_harness = source.harness.as_str(),
4685 target_harness = target.id(),
4686 sidecar = sidecar_path.display(),
4687 log = reduction_log_path.display(),
4688 source_id = source.session_id,
4689 )
4690}
4691
4692fn session_artifact(
4693 locator: &SessionLocator,
4694 session: &Session,
4695 target: TransferFormat,
4696) -> std::result::Result<SessionArtifact, ServiceError> {
4697 session_artifact_with_id(locator, session, target, None)
4698}
4699
4700fn session_artifact_with_id(
4701 locator: &SessionLocator,
4702 session: &Session,
4703 target: TransferFormat,
4704 target_session_id: Option<&str>,
4705) -> std::result::Result<SessionArtifact, ServiceError> {
4706 let format: SessionFormat = target.into();
4707 let diagonal = format.source() == session.meta.source;
4708 crate::residue_store::store_segments(session);
4709 let has_appended_turns = session
4710 .imported_message_count
4711 .is_some_and(|imported| imported < session.messages.len());
4712 let mut restoration = None;
4713 let content = if let Some(id) = target_session_id {
4714 if diagonal && format != SessionFormat::OpenCode {
4715 session
4716 .to_jsonl_spliced(format, Some(id))
4717 .map_err(operation)?
4718 } else {
4719 let mut rewritten = session.clone();
4720 rewritten.meta.session_id = Some(id.to_string());
4721 rewritten.to_jsonl(format).map_err(operation)?
4722 }
4723 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
4724 session.raw_verbatim()
4725 } else if diagonal {
4726 session.to_jsonl_spliced(format, None).map_err(operation)?
4727 } else {
4728 match session
4731 .restore_residue(format, crate::residue_store::lookup)
4732 .map_err(operation)?
4733 {
4734 Some((content, report)) => {
4735 restoration = Some(report);
4736 content
4737 }
4738 None => session.to_jsonl(format).map_err(operation)?,
4739 }
4740 };
4741 let stem = sanitize_filename(
4742 target_session_id
4743 .or(session.meta.session_id.as_deref())
4744 .unwrap_or(&locator.session_id),
4745 );
4746 let suggested_filename = if diagonal && target == TransferFormat::Grok {
4747 "chat_history.jsonl".to_string()
4748 } else if target == TransferFormat::Goose {
4749 format!("{stem}.goose.json")
4750 } else {
4751 format!("{stem}.{}.jsonl", target.id())
4752 };
4753 let mut files = vec![SessionArtifactFile {
4754 path: suggested_filename.clone(),
4755 content: content.clone(),
4756 role: ArtifactFileRole::Primary,
4757 }];
4758 if target == TransferFormat::ClaudeCode {
4759 let bundle_stem = Path::new(&suggested_filename)
4760 .file_stem()
4761 .and_then(|stem| stem.to_str())
4762 .unwrap_or(&stem);
4763 let mut child_paths = BTreeSet::new();
4764 for (index, subagent) in session.subagents.iter().enumerate() {
4765 let agent_id = subagent
4766 .meta
4767 .agent_id
4768 .as_deref()
4769 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
4770 .map(sanitize_filename)
4771 .filter(|id| !id.is_empty())
4772 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4773 let child_has_appended_turns = subagent
4774 .imported_message_count
4775 .is_some_and(|imported| imported < subagent.messages.len());
4776 let child_content = if target_session_id.is_none()
4777 && subagent.meta.source == SessionSource::ClaudeCode
4778 && subagent.raw_is_verbatim
4779 && !child_has_appended_turns
4780 {
4781 subagent.raw_verbatim()
4782 } else if subagent.meta.source == SessionSource::ClaudeCode {
4783 subagent
4784 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
4785 .map_err(operation)?
4786 } else {
4787 let mut child = subagent.clone();
4788 if let Some(id) = target_session_id {
4789 child.meta.session_id = Some(id.to_string());
4790 }
4791 child
4792 .to_jsonl(SessionFormat::ClaudeCode)
4793 .map_err(operation)?
4794 };
4795 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
4796 if !child_paths.insert(path.clone()) {
4797 return Err(ServiceError::Operation(format!(
4798 "Claude subagent ids collide at artifact path `{path}`"
4799 )));
4800 }
4801 files.push(SessionArtifactFile {
4802 path,
4803 content: child_content,
4804 role: ArtifactFileRole::Subagent,
4805 });
4806 }
4807 }
4808 if diagonal && target == TransferFormat::Grok {
4809 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
4810 }
4811 if !diagonal || !session.raw_is_verbatim {
4812 files.push(SessionArtifactFile {
4813 path: "recovery/source.supercode.jsonl".into(),
4814 content: session.to_native_jsonl(),
4815 role: ArtifactFileRole::SourceRecovery,
4816 });
4817 for (index, subagent) in session.subagents.iter().enumerate() {
4818 let id = subagent
4819 .meta
4820 .agent_id
4821 .as_deref()
4822 .map(sanitize_filename)
4823 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4824 files.push(SessionArtifactFile {
4825 path: format!("recovery/subagents/{id}.supercode.jsonl"),
4826 content: subagent.to_native_jsonl(),
4827 role: ArtifactFileRole::SourceRecovery,
4828 });
4829 }
4830 }
4831 if !diagonal && session.meta.source == SessionSource::Grok {
4832 append_grok_bundle_files(
4833 locator,
4834 "recovery/grok/",
4835 ArtifactFileRole::SourceRecovery,
4836 &mut files,
4837 )?;
4838 }
4839 let (fidelity, residue) = if diagonal
4840 && target_session_id.is_none()
4841 && session.raw_is_verbatim
4842 && !has_appended_turns
4843 {
4844 (Fidelity::ByteLossless, Vec::new())
4845 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
4846 (
4847 Fidelity::ValueLossless,
4848 vec![if target_session_id.is_some() {
4849 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
4850 } else {
4851 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
4852 }],
4853 )
4854 } else {
4855 match restoration {
4856 Some(report) if report.rendered_messages == 0 => (
4857 Fidelity::ByteLossless,
4858 vec![format!(
4859 "restored verbatim from this conversation's {} source records in the residue store",
4860 target.id()
4861 )],
4862 ),
4863 Some(report) => (
4864 Fidelity::Semantic,
4865 vec![format!(
4866 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
4867 report.restored_messages,
4868 report.restored_messages + report.rendered_messages,
4869 report.rendered_messages,
4870 target.id()
4871 )],
4872 ),
4873 None => (
4874 Fidelity::Semantic,
4875 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
4876 ),
4877 }
4878 };
4879 Ok(SessionArtifact {
4880 source_harness: locator.harness.clone(),
4881 target_harness: target.id(),
4882 session_id: target_session_id
4883 .map(str::to_string)
4884 .or_else(|| session.meta.session_id.clone()),
4885 content,
4886 suggested_filename,
4887 files,
4888 fidelity,
4889 residue,
4890 })
4891}
4892
4893fn append_grok_bundle_files(
4894 locator: &SessionLocator,
4895 prefix: &str,
4896 role: ArtifactFileRole,
4897 files: &mut Vec<SessionArtifactFile>,
4898) -> std::result::Result<(), ServiceError> {
4899 let primary = locator.storage.path();
4900 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
4901 return Err(ServiceError::Operation(format!(
4902 "Grok bundle locator must name chat_history.jsonl, got {}",
4903 primary.display()
4904 )));
4905 }
4906 let parent = primary.parent().ok_or_else(|| {
4907 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
4908 })?;
4909 for name in ["summary.json", "updates.jsonl"] {
4910 let path = parent.join(name);
4911 let metadata = match std::fs::symlink_metadata(&path) {
4912 Ok(metadata) => metadata,
4913 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
4914 Err(error) => return Err(ServiceError::Operation(error.to_string())),
4915 };
4916 if metadata.file_type().is_symlink() || !metadata.is_file() {
4917 return Err(ServiceError::Operation(format!(
4918 "refusing non-regular Grok bundle member {}",
4919 path.display()
4920 )));
4921 }
4922 let content = std::fs::read_to_string(&path).map_err(|error| {
4923 ServiceError::Operation(format!(
4924 "Grok bundle member {} is not representable as UTF-8: {error}",
4925 path.display()
4926 ))
4927 })?;
4928 files.push(SessionArtifactFile {
4929 path: format!("{prefix}{name}"),
4930 content,
4931 role: match role {
4932 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
4933 _ => ArtifactFileRole::SourceRecovery,
4934 },
4935 });
4936 }
4937 Ok(())
4938}
4939
4940fn handoff_artifact(
4941 locator: &SessionLocator,
4942 session: &Session,
4943 target: TransferFormat,
4944) -> std::result::Result<SessionArtifact, ServiceError> {
4945 let target_session_id = target_session_id(target);
4946 session_artifact_with_id(locator, session, target, Some(&target_session_id))
4947}
4948
4949fn target_session_id(target: TransferFormat) -> String {
4950 let uuid = generated_session_id();
4951 match target {
4952 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
4953 TransferFormat::ClaudeCode
4954 | TransferFormat::Codex
4955 | TransferFormat::Pi
4956 | TransferFormat::Grok
4957 | TransferFormat::Gemini
4958 | TransferFormat::Goose
4959 | TransferFormat::Hermes => uuid,
4960 }
4961}
4962
4963fn sanitize_filename(value: &str) -> String {
4964 let value = value
4965 .chars()
4966 .map(|character| {
4967 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
4968 character
4969 } else {
4970 '-'
4971 }
4972 })
4973 .collect::<String>();
4974 let value = value.trim_matches('-');
4975 if value.is_empty() {
4976 "session".into()
4977 } else {
4978 value.chars().take(100).collect()
4979 }
4980}
4981
4982fn handoff_instructions(
4983 target: TransferFormat,
4984 session_id: &str,
4985 cwd: &Path,
4986) -> HandoffInstructions {
4987 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
4988 cwd: cwd.to_path_buf(),
4989 program: program.into(),
4990 arguments,
4991 env: BTreeMap::new(),
4992 };
4993 match target {
4994 TransferFormat::ClaudeCode => HandoffInstructions {
4995 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
4996 materialize: None,
4997 requires_materialization: true,
4998 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(),
4999 },
5000 TransferFormat::Hermes => HandoffInstructions {
5001 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5002 materialize: None,
5003 requires_materialization: true,
5004 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(),
5005 },
5006 TransferFormat::Codex => HandoffInstructions {
5007 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5008 materialize: None,
5009 requires_materialization: true,
5010 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5011 },
5012 TransferFormat::OpenCode => HandoffInstructions {
5013 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5014 materialize: Some(launch(
5015 "opencode",
5016 vec!["import".into(), "{artifact_path}".into()],
5017 )),
5018 requires_materialization: true,
5019 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5020 },
5021 TransferFormat::Pi => HandoffInstructions {
5022 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5023 materialize: None,
5024 requires_materialization: true,
5025 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5026 },
5027 TransferFormat::Grok => HandoffInstructions {
5028 launch: launch(
5029 "grok",
5030 vec!["--resume".into(), "{materialized_session_id}".into()],
5031 ),
5032 materialize: None,
5033 requires_materialization: true,
5034 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(),
5035 },
5036 TransferFormat::Gemini => HandoffInstructions {
5037 launch: launch(
5038 "gemini",
5039 vec!["--session-file".into(), "{artifact_path}".into()],
5040 ),
5041 materialize: None,
5042 requires_materialization: true,
5043 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(),
5044 },
5045 TransferFormat::Goose => HandoffInstructions {
5046 launch: launch(
5047 "goose",
5048 vec![
5049 "session".into(),
5050 "--resume".into(),
5051 "--session-id".into(),
5052 "{imported_session_id}".into(),
5053 ],
5054 ),
5055 materialize: Some(launch(
5056 "goose",
5057 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5058 )),
5059 requires_materialization: true,
5060 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(),
5061 },
5062 }
5063}
5064
5065fn resume_launch(
5066 harness: &str,
5067 session_id: &str,
5068 cwd: &Path,
5069 policy: ResumePolicy,
5070) -> std::result::Result<StructuredLaunch, ServiceError> {
5071 let mut arguments = Vec::new();
5072 let program = match harness {
5073 HarnessId::GROK => {
5074 if matches!(policy, ResumePolicy::Yolo) {
5075 if crate::support::self_sandbox_supported() {
5076 arguments.extend(["--sandbox".into(), "workspace".into()]);
5077 }
5078 arguments.push("--always-approve".into());
5079 }
5080 arguments.extend(["--resume".into(), session_id.into()]);
5081 "grok"
5082 }
5083 HarnessId::CODEX => {
5084 let cwd_key = serde_json::to_string(cwd.to_string_lossy().as_ref())
5085 .expect("a filesystem path always serializes as JSON text");
5086 arguments.extend([
5087 "-c".into(),
5088 "check_for_update_on_startup=false".into(),
5089 "-c".into(),
5090 format!("projects={{{cwd_key}={{trust_level=\"trusted\"}}}}"),
5093 ]);
5094 if matches!(policy, ResumePolicy::Yolo) {
5095 arguments.extend([
5096 "--dangerously-bypass-approvals-and-sandbox".into(),
5097 "--dangerously-bypass-hook-trust".into(),
5098 ]);
5099 }
5100 arguments.extend(["resume".into(), session_id.into()]);
5101 "codex"
5102 }
5103 HarnessId::CLAUDE_CODE => {
5104 if matches!(policy, ResumePolicy::Yolo) {
5105 arguments.push("--dangerously-skip-permissions".into());
5106 }
5107 arguments.extend(["--resume".into(), session_id.into()]);
5108 "claude"
5109 }
5110 HarnessId::GEMINI => {
5111 if matches!(policy, ResumePolicy::Yolo) {
5112 arguments.push("--yolo".into());
5113 }
5114 arguments.extend(["--resume".into(), session_id.into()]);
5115 "gemini"
5116 }
5117 HarnessId::GOOSE => {
5118 arguments.extend([
5119 "session".into(),
5120 "--resume".into(),
5121 "--session-id".into(),
5122 session_id.into(),
5123 ]);
5124 "goose"
5125 }
5126 HarnessId::PI => {
5127 if matches!(policy, ResumePolicy::Yolo) {
5128 arguments.push("--approve".into());
5129 }
5130 arguments.extend(["--session".into(), session_id.into()]);
5131 "pi"
5132 }
5133 HarnessId::OPENCODE => {
5134 arguments.extend(["--session".into(), session_id.into()]);
5135 "opencode"
5136 }
5137 HarnessId::SUPERCODE => {
5138 if matches!(policy, ResumePolicy::Yolo) {
5139 arguments.push("--dangerous".into());
5140 }
5141 arguments.extend(["resume".into(), session_id.into()]);
5142 "supercode"
5143 }
5144 other => {
5145 return Err(ServiceError::InvalidParams(format!(
5146 "no structured resume launch is registered for harness `{other}`"
5147 )))
5148 }
5149 };
5150 Ok(StructuredLaunch {
5151 cwd: cwd.to_path_buf(),
5152 env: if program == "grok" {
5153 crate::support::grok_home_env()
5154 } else {
5155 BTreeMap::new()
5156 },
5157 program: program.into(),
5158 arguments,
5159 })
5160}
5161
5162fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5168 let digest = blake3::hash(address.as_bytes()).to_hex();
5169 let path = std::env::temp_dir().join(format!(
5170 "supercode-openclaw-gateway-token-{}",
5171 &digest.as_str()[..16]
5172 ));
5173 #[cfg(unix)]
5174 {
5175 use std::io::Write;
5176 use std::os::unix::fs::OpenOptionsExt;
5177 let mut file = std::fs::OpenOptions::new()
5178 .write(true)
5179 .create(true)
5180 .truncate(true)
5181 .mode(0o600)
5182 .open(&path)?;
5183 file.write_all(secret.as_bytes())?;
5184 }
5185 #[cfg(not(unix))]
5186 std::fs::write(&path, secret)?;
5187 Ok(path)
5188}
5189
5190fn open_connect_descriptor(
5196 descriptor: &crate::HarnessSupportDescriptor,
5197 home: &Path,
5198) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5199 let Some(connect) = &descriptor.runtime.connect_launch else {
5200 return Err(ServiceError::InvalidParams(format!(
5201 "harness `{}` has no registered connect-mode launch",
5202 descriptor.id.as_str()
5203 )));
5204 };
5205 let resolved = connect
5206 .resolve(home)
5207 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5208 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5209 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5210 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5211 if let Some(token) = resolved.auth {
5212 backend = backend.with_bearer(token);
5213 }
5214 Ok(Box::new(backend))
5215 }
5216 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5217 let mut env = BTreeMap::new();
5228 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5229 if let Some(token) = resolved.auth {
5230 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5231 .map_err(|error| {
5232 ServiceError::UnsupportedAction(format!(
5233 "could not stage the gateway credential for the bridge: {error}"
5234 ))
5235 })?;
5236 arguments.push("--token-file".into());
5237 arguments.push(token_path.to_string_lossy().into_owned());
5238 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5239 }
5240 let program = descriptor
5245 .runtime
5246 .default_launch
5247 .as_ref()
5248 .map(|launch| launch.program.clone())
5249 .unwrap_or_else(|| "openclaw".into());
5250 let launch = RuntimeLaunch {
5251 program,
5252 arguments,
5253 env,
5254 };
5255 Ok(Box::new(
5256 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5257 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5258 ))
5259 }
5260 _ => Err(ServiceError::UnsupportedAction(format!(
5261 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5262 descriptor.id.as_str(),
5263 connect.protocol
5264 ))),
5265 }
5266}
5267
5268fn registry_connect_descriptor(
5271 params: &RuntimeBackendParams,
5272) -> Option<crate::HarnessSupportDescriptor> {
5273 if params.launch.is_some() || params.base_url.is_some() {
5274 return None;
5275 }
5276 harness_support_registry()
5277 .harnesses
5278 .into_iter()
5279 .find(|descriptor| descriptor.id == params.harness)
5280 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5281}
5282
5283fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5284 supercode_interchange::user_home()
5285 .map(std::path::PathBuf::into_os_string)
5286 .map(PathBuf::from)
5287 .ok_or_else(|| {
5288 ServiceError::UnsupportedAction(
5289 "connect-mode launches need HOME to locate the harness config".into(),
5290 )
5291 })
5292}
5293
5294pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5297 "harness.v1.runtimes.start",
5298 "harness.v1.runtimes.resume",
5299 "harness.v1.runtimes.attach",
5300 "harness.v1.runtimes.attach_existing",
5301];
5302
5303pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5308
5309pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5315
5316pub const DETACHED_METHODS: &[&str] = &[
5322 "harness.v1.harnesses.list",
5323 "harness.v1.harnesses.probe",
5324 "harness.v1.sessions.message",
5325 "harness.v1.sessions.new",
5326 "harness.v1.sessions.reset",
5327 "harness.v1.sessions.archive",
5328 "harness.v1.sessions.delete",
5329];
5330
5331pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5337
5338pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5346
5347async fn within_control_deadline<F: std::future::Future>(
5350 method: &str,
5351 call: F,
5352) -> std::result::Result<F::Output, ServiceError> {
5353 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5354 .await
5355 .map_err(|_| {
5356 ServiceError::Operation(format!(
5357 "`{method}` gave up after {}s: the runtime did not answer",
5358 RUNTIME_CONTROL_DEADLINE.as_secs()
5359 ))
5360 })
5361}
5362
5363pub struct RuntimeOpen {
5367 id: Value,
5368 method: String,
5369 params: Value,
5370}
5371
5372impl RuntimeOpen {
5373 pub async fn open(self) -> OpenedRuntime {
5377 let Self { id, method, params } = self;
5378 let outcome = open_runtime(&method, params).await;
5379 OpenedRuntime { id, outcome }
5380 }
5381}
5382
5383pub struct OpenedRuntime {
5386 id: Value,
5387 outcome: std::result::Result<OpenRuntime, ServiceError>,
5388}
5389
5390pub struct DetachedCall {
5395 id: Value,
5396 method: String,
5397 work: std::result::Result<Work, ServiceError>,
5398}
5399
5400impl DetachedCall {
5401 pub async fn run(self) -> DetachedAnswer {
5404 let Self { id, method, work } = self;
5405 match work {
5406 Ok(Work::Runtime(work)) => {
5411 let (result, returned) = work.run().await;
5412 DetachedAnswer {
5413 response: service_response(id, result),
5414 returned,
5415 }
5416 }
5417 Ok(Work::Free(work)) => {
5418 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5419 Ok(result) => result,
5420 Err(_) => Err(ServiceError::Operation(format!(
5421 "`{method}` gave up after {}s: the harness it waits on did not answer",
5422 DETACHED_CALL_DEADLINE.as_secs()
5423 ))),
5424 };
5425 DetachedAnswer {
5426 response: service_response(id, result),
5427 returned: None,
5428 }
5429 }
5430 Err(error) => DetachedAnswer {
5431 response: service_response(id, Err(error)),
5432 returned: None,
5433 },
5434 }
5435 }
5436}
5437
5438pub struct DetachedAnswer {
5442 response: Value,
5443 returned: Option<ReturnedRuntime>,
5444}
5445
5446impl DetachedAnswer {
5447 pub fn into_response(self) -> Value {
5450 self.response
5451 }
5452}
5453
5454pub struct ReturnedRuntime {
5457 connection: String,
5458 runtime: Box<dyn RuntimeConnection>,
5459}
5460
5461enum Work {
5464 Free(DetachedWork),
5465 Runtime(RuntimeWork),
5466}
5467
5468enum DetachedWork {
5471 Inventory(InventoryWork),
5475 Message(MessageSessionParams),
5477 SessionMutation {
5480 verb: crate::SessionVerb,
5481 mutation: crate::SessionMutation,
5482 },
5483}
5484
5485impl DetachedWork {
5486 async fn run(self) -> std::result::Result<Value, ServiceError> {
5487 match self {
5488 Self::Inventory(work) => run_inventory(work).await,
5489 Self::Message(params) => Ok(message_live_session(¶ms).await),
5490 Self::SessionMutation { verb, mutation } => {
5491 let outcome = run_session_mutation(verb, &mutation).await?;
5492 serde_json::to_value(outcome)
5493 .map_err(|error| ServiceError::Operation(error.to_string()))
5494 }
5495 }
5496 }
5497}
5498
5499enum RuntimeWork {
5501 Close {
5503 runtime: Box<dyn RuntimeConnection>,
5504 process_group: Option<u32>,
5505 },
5506 LiveCommand {
5509 connection: String,
5510 runtime: Box<dyn RuntimeConnection>,
5511 verb: crate::SessionVerb,
5512 mutation: crate::SessionMutation,
5513 command: &'static str,
5514 session: String,
5515 },
5516}
5517
5518type RuntimeWorkAnswer = (
5521 std::result::Result<Value, ServiceError>,
5522 Option<ReturnedRuntime>,
5523);
5524
5525impl RuntimeWork {
5526 async fn run(self) -> RuntimeWorkAnswer {
5527 match self {
5528 Self::Close {
5529 runtime,
5530 process_group,
5531 } => (close_runtime(runtime, process_group).await, None),
5532 Self::LiveCommand {
5533 connection,
5534 mut runtime,
5535 verb,
5536 mutation,
5537 command,
5538 session,
5539 } => {
5540 let result =
5541 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5542 (
5543 result,
5544 Some(ReturnedRuntime {
5545 connection,
5546 runtime,
5547 }),
5548 )
5549 }
5550 }
5551 }
5552}
5553
5554async fn close_runtime(
5557 mut runtime: Box<dyn RuntimeConnection>,
5558 process_group: Option<u32>,
5559) -> std::result::Result<Value, ServiceError> {
5560 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
5561 Ok(result) => {
5562 result.map_err(operation)?;
5563 Ok(json!({"closed": true}))
5564 }
5565 Err(deadline) => {
5566 let killed = kill_runtime_process_group(process_group);
5571 drop(runtime);
5572 Ok(json!({
5573 "closed": true,
5574 "killed": killed,
5575 "detail": error_message(deadline),
5576 }))
5577 }
5578 }
5579}
5580
5581fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
5584 mutation
5585 .session
5586 .clone()
5587 .filter(|value| !value.trim().is_empty())
5588 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
5589}
5590
5591async fn type_live_command(
5595 runtime: &mut dyn RuntimeConnection,
5596 verb: crate::SessionVerb,
5597 mutation: &crate::SessionMutation,
5598 command: &str,
5599 session: String,
5600) -> std::result::Result<Value, ServiceError> {
5601 within_control_deadline(
5602 &format!("sessions.{}", verb.as_str()),
5603 runtime.send_input(RuntimeInput {
5604 text: command.to_string(),
5605 image_urls: Vec::new(),
5606 }),
5607 )
5608 .await?
5609 .map_err(operation)?;
5610 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
5611 .map_err(session_control_error)?;
5612 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
5613}
5614
5615enum OpenRuntime {
5618 Hosted {
5621 runtime: Box<dyn RuntimeConnection>,
5622 capabilities: crate::RuntimeCapabilities,
5623 workspace: PathBuf,
5624 },
5625 Joined { runtime: Box<dyn RuntimeConnection> },
5628}
5629
5630async fn open_runtime(
5635 method: &str,
5636 params: Value,
5637) -> std::result::Result<OpenRuntime, ServiceError> {
5638 match tokio::time::timeout(
5639 RUNTIME_OPEN_DEADLINE,
5640 open_runtime_unbounded(method, params),
5641 )
5642 .await
5643 {
5644 Ok(result) => result,
5645 Err(_) => Err(ServiceError::Operation(format!(
5646 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
5647 RUNTIME_OPEN_DEADLINE.as_secs()
5648 ))),
5649 }
5650}
5651
5652async fn open_runtime_unbounded(
5653 method: &str,
5654 params: Value,
5655) -> std::result::Result<OpenRuntime, ServiceError> {
5656 match method {
5657 "harness.v1.runtimes.start" => {
5658 let params = decode::<RuntimeStartParams>(params)?;
5659 let backend = runtime_backend(¶ms.backend)?;
5660 let capabilities = backend.capabilities();
5661 let workspace = params.cwd.clone();
5662 let runtime = backend
5663 .start(RuntimeStartRequest {
5664 cwd: params.cwd,
5665 launch: runtime_launch(¶ms.backend),
5666 mcp_servers: params.mcp_servers,
5667 approval_policy: params.approval_policy,
5668 })
5669 .await
5670 .map_err(operation)?;
5671 Ok(OpenRuntime::Hosted {
5672 runtime,
5673 capabilities,
5674 workspace,
5675 })
5676 }
5677 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
5678 let params = decode::<RuntimeAttachParams>(params)?;
5679 let backend = runtime_backend(¶ms.backend)?;
5680 let capabilities = backend.capabilities();
5681 let workspace = params
5682 .cwd
5683 .clone()
5684 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
5685 let runtime = backend
5686 .attach(RuntimeAttachRequest {
5687 runtime_id: params.runtime_id,
5688 cwd: params.cwd,
5689 launch: runtime_launch(¶ms.backend),
5690 mcp_servers: params.mcp_servers,
5691 approval_policy: params.approval_policy,
5692 })
5693 .await
5694 .map_err(operation)?;
5695 Ok(OpenRuntime::Hosted {
5696 runtime,
5697 capabilities,
5698 workspace,
5699 })
5700 }
5701 "harness.v1.runtimes.attach_existing" => {
5702 let params = decode::<RuntimeAttachParams>(params)?;
5703 let backend: Box<dyn RuntimeBackend> = match params
5704 .backend
5705 .base_url
5706 .as_deref()
5707 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
5708 {
5709 Some(endpoint) => {
5710 #[cfg(not(feature = "adapter-api"))]
5711 {
5712 let _ = endpoint;
5713 return Err(ServiceError::UnsupportedAction(
5714 "live HTTP attachment adapter is not compiled".into(),
5715 ));
5716 }
5717 #[cfg(feature = "adapter-api")]
5718 {
5719 let workspace = params.cwd.clone().ok_or_else(|| {
5720 ServiceError::InvalidParams(
5721 "Volter Harness live attach requires the project cwd".into(),
5722 )
5723 })?;
5724 let source = LiveRuntimeSource {
5725 harness: params.backend.harness.as_str().to_string(),
5726 session_id: params.runtime_id.clone(),
5727 workspace,
5728 };
5729 let receipt = resolve_live_runtime(&endpoint, &source)
5730 .map_err(|error| ServiceError::Operation(error.to_string()))?;
5731 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
5732 }
5733 }
5734 None => runtime_backend(¶ms.backend)?,
5735 };
5736 let capabilities = backend.capabilities();
5737 if !capabilities.attach_existing_process {
5738 return Err(ServiceError::Operation(format!(
5739 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
5740 backend.harness().as_str()
5741 )));
5742 }
5743 let runtime = backend
5744 .attach_existing(RuntimeAttachRequest {
5745 runtime_id: params.runtime_id,
5746 cwd: params.cwd,
5747 launch: runtime_launch(¶ms.backend),
5748 mcp_servers: params.mcp_servers,
5749 approval_policy: params.approval_policy,
5750 })
5751 .await
5752 .map_err(operation)?;
5753 Ok(OpenRuntime::Joined { runtime })
5754 }
5755 _ => Err(ServiceError::MethodNotFound),
5756 }
5757}
5758
5759fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
5761 match result {
5762 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
5763 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
5764 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
5765 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
5766 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
5767 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
5768 }
5769}
5770
5771fn runtime_backend(
5772 params: &RuntimeBackendParams,
5773) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5774 if let Some(descriptor) = registry_connect_descriptor(params) {
5775 return open_connect_descriptor(&descriptor, &service_home()?);
5776 }
5777 if params.protocol.as_deref() == Some("acp") {
5778 let launch = params
5779 .launch
5780 .clone()
5781 .or_else(|| {
5782 harness_support_registry()
5783 .harnesses
5784 .into_iter()
5785 .find(|harness| harness.id == params.harness)
5786 .filter(|harness| {
5787 harness.runtime.implementation == ImplementationKind::GenericProtocol
5788 && harness.runtime.protocol.starts_with("acp")
5789 })
5790 .and_then(|harness| harness.runtime.default_launch)
5791 })
5792 .ok_or_else(|| {
5793 ServiceError::InvalidParams(
5794 "an ACP runtime requires `launch` unless the harness has a registered default"
5795 .into(),
5796 )
5797 })?;
5798 let resume_session = harness_support_registry()
5799 .harnesses
5800 .into_iter()
5801 .find(|harness| harness.id == params.harness)
5802 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
5803 return Ok(Box::new(
5804 AcpRuntimeBackend::new(params.harness.clone(), launch)
5805 .with_resume_support(resume_session),
5806 ));
5807 }
5808 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
5809 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
5810 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
5811 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
5812 HarnessId::OPENCODE => match ¶ms.base_url {
5813 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
5814 None => Box::new(OpenCodeRuntimeBackend::new()),
5815 },
5816 harness => {
5817 let descriptor = harness_support_registry()
5818 .harnesses
5819 .into_iter()
5820 .find(|descriptor| descriptor.id.as_str() == harness)
5821 .filter(|descriptor| {
5822 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
5823 && descriptor.runtime.protocol.starts_with("acp")
5824 });
5825 let Some(descriptor) = descriptor else {
5826 return Err(ServiceError::InvalidParams(format!(
5827 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
5828 )));
5829 };
5830 let resume = descriptor.runtime.capabilities.resume_session;
5831 Box::new(
5832 AcpRuntimeBackend::new(
5833 descriptor.id,
5834 descriptor
5835 .runtime
5836 .default_launch
5837 .expect("generic ACP registry entry includes its launch"),
5838 )
5839 .with_resume_support(resume),
5840 )
5841 }
5842 };
5843 Ok(backend)
5844}
5845
5846fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
5847 if let Some(launch) = ¶ms.launch {
5848 return Some(launch.clone());
5849 }
5850 if !matches!(params.policy, RuntimePolicy::Yolo) {
5851 return None;
5852 }
5853 let launch = match params.harness.as_str() {
5854 HarnessId::GROK => RuntimeLaunch {
5855 program: "grok".into(),
5856 arguments: {
5857 let mut arguments: Vec<String> = Vec::new();
5858 if crate::support::self_sandbox_supported() {
5859 arguments.extend(["--sandbox".into(), "workspace".into()]);
5860 }
5861 arguments.extend([
5862 "--always-approve".into(),
5863 "agent".into(),
5864 "--no-leader".into(),
5865 "stdio".into(),
5866 ]);
5867 arguments
5868 },
5869 env: crate::support::grok_env(),
5870 },
5871 HarnessId::CODEX => RuntimeLaunch {
5872 program: "codex".into(),
5873 arguments: vec![
5874 "--dangerously-bypass-approvals-and-sandbox".into(),
5875 "--dangerously-bypass-hook-trust".into(),
5876 "app-server".into(),
5877 ],
5878 env: BTreeMap::new(),
5879 },
5880 HarnessId::CLAUDE_CODE => RuntimeLaunch {
5881 program: "claude".into(),
5882 arguments: vec![
5883 "--dangerously-skip-permissions".into(),
5884 "--print".into(),
5885 "--input-format".into(),
5886 "stream-json".into(),
5887 "--output-format".into(),
5888 "stream-json".into(),
5889 "--verbose".into(),
5890 ],
5891 env: BTreeMap::new(),
5892 },
5893 HarnessId::PI => RuntimeLaunch {
5894 program: "pi".into(),
5895 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
5896 env: BTreeMap::new(),
5897 },
5898 HarnessId::OPENCODE => RuntimeLaunch {
5899 program: "opencode".into(),
5900 arguments: vec!["serve".into()],
5901 env: BTreeMap::new(),
5902 },
5903 HarnessId::GEMINI => RuntimeLaunch {
5904 program: "gemini".into(),
5905 arguments: vec!["--acp".into(), "--yolo".into()],
5906 env: BTreeMap::new(),
5907 },
5908 HarnessId::GOOSE => RuntimeLaunch {
5909 program: "goose".into(),
5910 arguments: vec!["acp".into()],
5911 env: BTreeMap::new(),
5912 },
5913 HarnessId::SUPERCODE => RuntimeLaunch {
5914 program: "supercode".into(),
5915 arguments: vec!["acp".into(), "--dangerous".into()],
5916 env: BTreeMap::new(),
5917 },
5918 _ => return None,
5919 };
5920 Some(launch)
5921}
5922
5923struct IsolatedProbeHome {
5929 launch: RuntimeLaunch,
5930 root: PathBuf,
5931}
5932
5933impl IsolatedProbeHome {
5934 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
5935 let root = std::env::temp_dir().join(format!(
5936 "supercode-harness-probe-{harness}-{}",
5937 generated_session_id()
5938 ));
5939 std::fs::create_dir_all(&root)?;
5940 set_private_dir_permissions(&root)?;
5941
5942 if let Some(source_home) = supercode_interchange::user_home()
5943 .map(std::path::PathBuf::into_os_string)
5944 .map(PathBuf::from)
5945 {
5946 for relative in probe_auth_files(harness) {
5947 copy_probe_file(&source_home, &root, relative)?;
5948 }
5949 }
5950 if harness == HarnessId::SUPERCODE {
5955 let config_home = crate::agent::global_instructions_dir();
5956 for file in ["config.toml", "credentials.toml"] {
5957 copy_probe_path(
5958 &config_home.join(file),
5959 &root.join(".config/supercode").join(file),
5960 )?;
5961 }
5962 }
5963 configure_isolated_probe_auth(harness, &root)?;
5964
5965 let root_text = root.to_string_lossy().into_owned();
5966 for (key, value) in [
5967 ("HOME", root_text.clone()),
5968 (
5969 "XDG_CACHE_HOME",
5970 root.join(".cache").to_string_lossy().into_owned(),
5971 ),
5972 (
5973 "XDG_CONFIG_HOME",
5974 root.join(".config").to_string_lossy().into_owned(),
5975 ),
5976 (
5977 "XDG_DATA_HOME",
5978 root.join(".local/share").to_string_lossy().into_owned(),
5979 ),
5980 ] {
5981 launch.env.insert(key.into(), value);
5982 }
5983 let scoped = match harness {
5984 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
5985 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
5986 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
5987 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
5988 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
5989 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
5990 _ => None,
5991 };
5992 if let Some((key, value)) = scoped {
5993 launch
5994 .env
5995 .insert(key.into(), value.to_string_lossy().into_owned());
5996 }
5997 Ok(Self { launch, root })
5998 }
5999
6000 fn cleanup(&self) -> std::io::Result<()> {
6001 match std::fs::remove_dir_all(&self.root) {
6002 Ok(()) => Ok(()),
6003 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6004 Err(error) => Err(error),
6005 }
6006 }
6007}
6008
6009impl Drop for IsolatedProbeHome {
6010 fn drop(&mut self) {
6011 let _ = self.cleanup();
6012 }
6013}
6014
6015fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6016 match harness {
6017 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6018 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6022 HarnessId::CODEX => &[".codex/auth.json"],
6023 HarnessId::GEMINI => &[
6024 ".gemini/google_accounts.json",
6025 ".gemini/oauth_creds.json",
6026 ".gemini/settings.json",
6027 ],
6028 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6029 HarnessId::OPENCODE => &[
6030 ".config/opencode/auth.json",
6031 ".local/share/opencode/auth.json",
6032 ],
6033 HarnessId::PI => &[".pi/agent/auth.json"],
6034 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6039 _ => &[],
6040 }
6041}
6042
6043fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6044 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6045}
6046
6047fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6048 if !source.is_file() {
6049 return Ok(());
6050 }
6051 if let Some(parent) = destination.parent() {
6052 std::fs::create_dir_all(parent)?;
6053 set_private_dir_permissions(parent)?;
6054 }
6055 std::fs::copy(source, destination)?;
6056 set_private_file_permissions(destination)
6057}
6058
6059fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6060 if harness != HarnessId::GEMINI {
6061 return Ok(());
6062 }
6063 let oauth = probe_home.join(".gemini/oauth_creds.json");
6064 if !oauth.is_file() {
6065 return Ok(());
6066 }
6067 let settings_path = probe_home.join(".gemini/settings.json");
6068 let mut settings = std::fs::read_to_string(&settings_path)
6069 .ok()
6070 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6071 .unwrap_or_else(|| json!({}));
6072 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6073 std::fs::write(
6074 &settings_path,
6075 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6076 )?;
6077 set_private_file_permissions(&settings_path)
6078}
6079
6080#[cfg(unix)]
6081fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6082 use std::os::unix::fs::PermissionsExt;
6083 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6084}
6085
6086#[cfg(not(unix))]
6087fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6088 Ok(())
6089}
6090
6091#[cfg(unix)]
6092fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6093 use std::os::unix::fs::PermissionsExt;
6094 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6095}
6096
6097#[cfg(not(unix))]
6098fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6099 Ok(())
6100}
6101
6102fn find_executable(program: &str) -> Option<PathBuf> {
6103 let candidate = PathBuf::from(program);
6104 if candidate.components().count() > 1 {
6105 return candidate.is_file().then_some(candidate);
6106 }
6107 let path = std::env::var_os("PATH")?;
6108 for directory in std::env::split_paths(&path) {
6109 let candidate = directory.join(program);
6110 if candidate.is_file() {
6111 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6112 }
6113 #[cfg(windows)]
6114 {
6115 for extension in ["exe", "cmd", "bat"] {
6116 let candidate = directory.join(format!("{program}.{extension}"));
6117 if candidate.is_file() {
6118 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6119 }
6120 }
6121 }
6122 }
6123 None
6124}
6125
6126async fn executable_version(executable: &Path) -> Option<String> {
6127 let mut command = tokio::process::Command::new(executable);
6128 command
6129 .arg("--version")
6130 .stdin(std::process::Stdio::null())
6131 .stdout(std::process::Stdio::piped())
6132 .stderr(std::process::Stdio::piped())
6133 .kill_on_drop(true);
6134 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6135 .await
6136 .ok()?
6137 .ok()?;
6138 let stdout = String::from_utf8_lossy(&output.stdout);
6139 let stderr = String::from_utf8_lossy(&output.stderr);
6140 stdout
6141 .lines()
6142 .chain(stderr.lines())
6143 .map(str::trim)
6144 .find(|line| !line.is_empty())
6145 .map(|line| truncate_text(line, 200))
6146}
6147
6148pub(crate) fn auth_evidence(harness: &str) -> bool {
6149 let env_names: &[&str] = match harness {
6150 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6151 HarnessId::CODEX => &["OPENAI_API_KEY"],
6152 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6153 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6154 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6155 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6156 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6157 _ => &[],
6158 };
6159 if env_names
6160 .iter()
6161 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6162 {
6163 return true;
6164 }
6165 let Some(home) = supercode_interchange::user_home()
6166 .map(std::path::PathBuf::into_os_string)
6167 .map(PathBuf::from)
6168 else {
6169 return false;
6170 };
6171 let files: Vec<PathBuf> = match harness {
6172 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6173 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6174 HarnessId::OPENCODE => vec![
6175 home.join(".local/share/opencode/auth.json"),
6176 home.join(".config/opencode/auth.json"),
6177 ],
6178 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6179 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6180 HarnessId::GEMINI => vec![
6181 home.join(".gemini/oauth_creds.json"),
6182 home.join(".gemini/google_accounts.json"),
6183 ],
6184 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6185 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6186 _ => Vec::new(),
6187 };
6188 if files.into_iter().any(|path| {
6189 std::fs::metadata(path)
6190 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6191 .unwrap_or(false)
6192 }) {
6193 return true;
6194 }
6195 if harness == HarnessId::CLAUDE_CODE {
6202 return std::fs::read_to_string(home.join(".claude.json"))
6203 .map(|text| text.contains("\"oauthAccount\""))
6204 .unwrap_or(false);
6205 }
6206 false
6207}
6208
6209fn looks_like_auth_error(message: &str) -> bool {
6210 let message = message.to_ascii_lowercase();
6211 [
6212 "auth",
6213 "login",
6214 "sign in",
6215 "sign-in",
6216 "credential",
6217 "unauthorized",
6218 "forbidden",
6219 "token",
6220 ]
6221 .iter()
6222 .any(|needle| message.contains(needle))
6223}
6224
6225fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6226 crate::RuntimeCapabilities {
6227 start_session: false,
6228 resume_session: false,
6229 attach_existing_process: false,
6230 send_input: false,
6231 stream_events: false,
6232 interrupt: false,
6233 steer: false,
6234 respond_to_requests: false,
6235 }
6236}
6237
6238fn truncate_text(text: &str, max_chars: usize) -> String {
6239 let mut chars = text.chars();
6240 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6241 if chars.next().is_some() {
6242 format!("{truncated}…")
6243 } else {
6244 truncated
6245 }
6246}
6247
6248fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6255 match &handle.endpoint {
6256 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6257 crate::RuntimeEndpoint::Http { .. } => None,
6258 }
6259}
6260
6261fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6267 match process_group {
6268 #[cfg(unix)]
6269 Some(pid) => {
6270 crate::lsp::kill_process_group(pid);
6271 true
6272 }
6273 #[cfg(not(unix))]
6274 Some(_) => false,
6275 None => false,
6276 }
6277}
6278
6279fn error_message(error: ServiceError) -> String {
6280 match error {
6281 ServiceError::InvalidParams(message)
6282 | ServiceError::Operation(message)
6283 | ServiceError::UnsupportedAction(message) => message,
6284 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6285 ServiceError::Sdk(error) => error.to_string(),
6286 }
6287}
6288
6289#[derive(Debug)]
6290enum ServiceError {
6291 InvalidParams(String),
6292 MethodNotFound,
6293 UnsupportedAction(String),
6294 Operation(String),
6295 Sdk(SdkError),
6296}
6297
6298fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6299 match error {
6300 ServiceError::InvalidParams(message) => {
6301 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6302 }
6303 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6304 SdkError::unsupported(operation)
6305 }
6306 ServiceError::Operation(message) => {
6307 let code = if message.contains("already in progress") {
6308 SdkErrorCode::Busy
6309 } else if message.contains("not supported by this runtime") {
6310 SdkErrorCode::UnsupportedAction
6311 } else if message.contains("unknown runtime connection") {
6312 SdkErrorCode::NotFound
6313 } else {
6314 SdkErrorCode::Execution
6315 };
6316 SdkError::new(code, operation, message)
6317 }
6318 ServiceError::Sdk(error) => error,
6319 }
6320}
6321
6322fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6323 let error_code = error.code();
6324 let code = match error_code {
6325 SdkErrorCode::Unauthenticated => -32030,
6326 SdkErrorCode::Unauthorized => -32031,
6327 SdkErrorCode::ControllerRequired => -32032,
6328 SdkErrorCode::LeaseExpired => -32033,
6329 SdkErrorCode::InvalidArgument => -32602,
6330 SdkErrorCode::NotFound => -32004,
6331 SdkErrorCode::Busy => -32000,
6332 SdkErrorCode::UnsupportedAction => -32020,
6333 SdkErrorCode::Execution => -32002,
6334 SdkErrorCode::Transport => -32003,
6335 };
6336 json!({
6337 "jsonrpc": "2.0",
6338 "id": id,
6339 "error": {
6340 "code": code,
6341 "name": error_code,
6342 "operation": error.operation(),
6343 "message": error.to_string(),
6344 },
6345 })
6346}
6347
6348fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6349 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6350}
6351
6352fn operation(error: impl Into<crate::Error>) -> ServiceError {
6353 let error = error.into();
6354 match error {
6355 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6356 error => ServiceError::Operation(error.to_string()),
6357 }
6358}
6359
6360#[derive(Debug, Clone, Deserialize, Default)]
6364#[serde(default)]
6365struct MemoryRequest {
6366 harness: Option<String>,
6368 query: Option<String>,
6370 profile: Option<String>,
6372 session: Option<String>,
6374 full: bool,
6376 regex: bool,
6378 cwd: Option<std::path::PathBuf>,
6380 homes: crate::HarnessHomes,
6382}
6383
6384fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6387 let request = decode::<MemoryRequest>(params)?;
6388 let harness = request
6389 .harness
6390 .clone()
6391 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6392 let to_service = |error: crate::memory::MemoryError| match error {
6393 crate::memory::MemoryError::UnsupportedHarness { .. }
6394 | crate::memory::MemoryError::SessionNotScoped { .. } => {
6395 ServiceError::UnsupportedAction(error.to_string())
6396 }
6397 other => ServiceError::InvalidParams(other.to_string()),
6398 };
6399 match method {
6400 "harness.v1.memory.show" => {
6401 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6402 harness,
6403 profile: request.profile,
6404 session: request.session,
6405 full: request.full,
6406 cwd: request.cwd,
6407 homes: request.homes,
6408 })
6409 .map_err(to_service)?;
6410 Ok(json!({
6411 "schema": crate::memory::MEMORY_SCHEMA,
6412 "documents": documents,
6413 }))
6414 }
6415 "harness.v1.memory.search" => {
6416 let query = request
6417 .query
6418 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6419 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6420 harness,
6421 query,
6422 profile: request.profile,
6423 regex: request.regex,
6424 cwd: request.cwd,
6425 homes: request.homes,
6426 })
6427 .map_err(to_service)?;
6428 Ok(json!({
6429 "schema": crate::memory::MEMORY_SCHEMA,
6430 "matches": matches,
6431 }))
6432 }
6433 _ => Err(ServiceError::MethodNotFound),
6434 }
6435}
6436
6437#[derive(Debug, Clone, Deserialize)]
6441#[serde(default)]
6442struct ProfilesQuery {
6443 harness: Option<String>,
6445 name: Option<String>,
6447 homes: crate::HarnessHomes,
6449}
6450
6451impl Default for ProfilesQuery {
6452 fn default() -> Self {
6453 Self {
6454 harness: None,
6455 name: None,
6456 homes: crate::HarnessHomes::default(),
6457 }
6458 }
6459}
6460
6461fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6464 let query = decode::<ProfilesQuery>(params)?;
6465 let to_service = |error: crate::profiles::ProfileError| match error {
6466 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6467 ServiceError::UnsupportedAction(error.to_string())
6468 }
6469 crate::profiles::ProfileError::NotFound { .. } => {
6470 ServiceError::InvalidParams(error.to_string())
6471 }
6472 };
6473 match method {
6474 "harness.v1.profiles.list" => {
6475 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6476 .map_err(to_service)?;
6477 Ok(json!({
6478 "schema": crate::profiles::PROFILES_SCHEMA,
6479 "profiles": profiles,
6480 }))
6481 }
6482 "harness.v1.profiles.get" => {
6483 let harness = query
6484 .harness
6485 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6486 let name = query
6487 .name
6488 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6489 let profile =
6490 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6491 Ok(json!({
6492 "schema": crate::profiles::PROFILES_SCHEMA,
6493 "profile": profile,
6494 }))
6495 }
6496 _ => Err(ServiceError::MethodNotFound),
6497 }
6498}
6499
6500#[derive(Debug, Clone, Deserialize)]
6504#[serde(default)]
6505struct ChannelsQuery {
6506 harness: Option<String>,
6508 name: Option<String>,
6510 homes: crate::HarnessHomes,
6512}
6513
6514impl Default for ChannelsQuery {
6515 fn default() -> Self {
6516 Self {
6517 harness: None,
6518 name: None,
6519 homes: crate::HarnessHomes::default(),
6520 }
6521 }
6522}
6523
6524#[derive(Debug, Clone, Deserialize)]
6528#[serde(default)]
6529struct RoutesQuery {
6530 harness: Option<String>,
6531 profile: Option<String>,
6533 homes: crate::HarnessHomes,
6534}
6535
6536impl Default for RoutesQuery {
6537 fn default() -> Self {
6538 Self {
6539 harness: None,
6540 profile: None,
6541 homes: crate::HarnessHomes::default(),
6542 }
6543 }
6544}
6545
6546#[derive(Debug, Clone, Deserialize)]
6547#[serde(default)]
6548struct TriggersQuery {
6549 harness: Option<String>,
6550 homes: crate::HarnessHomes,
6551}
6552
6553impl Default for TriggersQuery {
6554 fn default() -> Self {
6555 Self {
6556 harness: None,
6557 homes: crate::HarnessHomes::default(),
6558 }
6559 }
6560}
6561
6562fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
6563 let query = decode::<TriggersQuery>(params)?;
6564 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
6565 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6566 Ok(json!({
6567 "schema": crate::triggers::TRIGGERS_SCHEMA,
6568 "triggers": triggers,
6569 }))
6570}
6571
6572fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
6573 let query = decode::<RoutesQuery>(params)?;
6574 let routes = crate::routes::list_routes(
6575 &query.homes,
6576 query.harness.as_deref(),
6577 query.profile.as_deref(),
6578 )
6579 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6580 Ok(json!({
6581 "schema": crate::routes::ROUTES_SCHEMA,
6582 "routes": routes,
6583 }))
6584}
6585
6586fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6587 let query = decode::<ChannelsQuery>(params)?;
6588 let to_service = |error: crate::channels::ChannelError| match error {
6589 crate::channels::ChannelError::UnsupportedHarness { .. } => {
6590 ServiceError::UnsupportedAction(error.to_string())
6591 }
6592 crate::channels::ChannelError::NotFound { .. } => {
6593 ServiceError::InvalidParams(error.to_string())
6594 }
6595 };
6596 match method {
6597 "harness.v1.channels.list" => {
6598 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
6599 .map_err(to_service)?;
6600 Ok(json!({
6601 "schema": crate::channels::CHANNELS_SCHEMA,
6602 "channels": channels,
6603 }))
6604 }
6605 "harness.v1.channels.status" => {
6606 let harness = query
6607 .harness
6608 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6609 let name = query
6610 .name
6611 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6612 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
6613 .map_err(to_service)?;
6614 Ok(json!({
6615 "schema": crate::channels::CHANNELS_SCHEMA,
6616 "channel": channel,
6617 }))
6618 }
6619 _ => Err(ServiceError::MethodNotFound),
6620 }
6621}
6622
6623fn rpc_error(id: Value, code: i64, message: &str) -> Value {
6624 json!({
6625 "jsonrpc": "2.0",
6626 "id": id,
6627 "error": {"code": code, "message": message},
6628 })
6629}