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 let mut activity = activity.clone();
1080 let waiting = doors.all().iter().any(|live| {
1081 live.address.harness == activity_key.0
1082 && live.address.session_id == activity_key.1
1083 && live.status == "waiting"
1084 });
1085 if waiting && activity.turn != crate::SessionTurnState::NeedsInput {
1086 activity.turn = crate::SessionTurnState::NeedsInput;
1087 activity.evidence.source = "pane_prompt".into();
1088 }
1089 let activity = &activity;
1090 value["activity"] = serde_json::to_value(activity)
1091 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1092 if let Some(status) = legacy_live_status(activity) {
1093 value["live_status"] = json!(status);
1094 }
1095 }
1096 Ok(value)
1097 })
1098 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1099 let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1100 if query.search_previews {
1103 result["receipt"] = serde_json::to_value(page.receipt)
1104 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1105 }
1106 Ok(result)
1107 }
1108 "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1109 "harness.v1.sessions.load" => {
1110 let params = decode::<LoadSessionParams>(params)?;
1111 if let Some(options) = ¶ms.options {
1112 options.validate()?;
1113 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
1114 return Ok(result);
1115 }
1116 return load_session(¶ms.read.locator)
1117 .map(|session| projected_session_result(&session, options))
1118 .map_err(operation);
1119 }
1120 let mut session = if params.read.display_history() {
1121 self.catalog
1122 .load_display_view(
1123 ¶ms.read.locator,
1124 params.read.read_fidelity(),
1125 params.read.tail_messages().unwrap_or(500),
1126 )
1127 .map_err(crate::Error::from)
1128 } else if params.read.include_subagents() {
1129 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
1130 } else {
1131 self.catalog
1132 .load_parent_with_fidelity(
1133 ¶ms.read.locator,
1134 params.read.read_fidelity(),
1135 )
1136 .map_err(crate::Error::from)
1137 }
1138 .map_err(operation)?;
1139 params.read.bound_session(&mut session);
1140 Ok(json!({"session": normalized_session_json(&session)}))
1141 }
1142 "harness.v1.sessions.follow" => {
1143 let params = decode::<LocatorParams>(params)?;
1144 let mut follower = self
1145 .catalog
1146 .follow_read_view(
1147 ¶ms.locator,
1148 params.read_fidelity(),
1149 params.include_subagents(),
1150 params.tail_messages(),
1151 params.max_message_chars(),
1152 params.display_history(),
1153 )
1154 .map_err(operation)?;
1155 let initial = follower
1156 .poll()
1157 .map_err(operation)?
1158 .map(|event| event.to_json());
1159 let subscription = format!("sub-{}", self.next_subscription);
1160 self.next_subscription += 1;
1161 self.followers.insert(subscription.clone(), follower);
1162 self.followed_sources.insert(
1163 subscription.clone(),
1164 FollowedSource {
1165 harness: params.locator.harness.as_str().to_string(),
1166 session_id: params.locator.session_id.clone(),
1167 reported: None,
1168 },
1169 );
1170 Ok(json!({"subscription": subscription, "initial": initial}))
1171 }
1172 "harness.v1.sessions.unfollow" => {
1173 let params = decode::<UnfollowParams>(params)?;
1174 self.followed_sources.remove(¶ms.subscription);
1175 Ok(json!({
1176 "removed": self.followers.remove(¶ms.subscription).is_some()
1177 }))
1178 }
1179 "harness.v1.sessions.activity.unsubscribe" => {
1180 let params = decode::<UnfollowParams>(params)?;
1181 Ok(json!({
1182 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1183 }))
1184 }
1185 "harness.v1.sessions.index.subscribe" => {
1186 let query = decode::<DiscoveryQuery>(params)?;
1187 crate::session_index::validate_query(&query)
1188 .map_err(ServiceError::InvalidParams)?;
1189 let homes = query.homes.clone();
1190 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1191 query,
1192 Arc::clone(&self.index_notifier),
1193 )
1194 .map_err(ServiceError::Operation)?;
1195 let doors = crate::mail_route::LiveSessions::read(&homes);
1196 let initial = initial
1197 .iter()
1198 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1199 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1200 let subscription = format!("index-sub-{}", self.next_subscription);
1201 self.next_subscription += 1;
1202 self.index_subscriptions.insert(subscription.clone(), index);
1203 Ok(json!({
1204 "subscription": subscription,
1205 "revision": 1,
1206 "initial": initial,
1207 }))
1208 }
1209 "harness.v1.sessions.index.resize" => {
1210 let params = decode::<IndexResizeParams>(params)?;
1211 crate::session_index::validate_limit(params.limit)
1212 .map_err(ServiceError::InvalidParams)?;
1213 let index = self
1214 .index_subscriptions
1215 .get_mut(¶ms.subscription)
1216 .ok_or_else(|| {
1217 ServiceError::InvalidParams("unknown session index subscription".into())
1218 })?;
1219 let prepared = index
1220 .prepare_resize(params.limit)
1221 .map_err(ServiceError::Operation)?;
1222 let doors = crate::mail_route::LiveSessions::read(index.homes());
1223 let initial = prepared
1224 .page
1225 .sessions
1226 .iter()
1227 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1228 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1229 let response = json!({
1230 "subscription": params.subscription,
1231 "revision": prepared.revision,
1232 "initial": initial,
1233 "receipt": prepared.page.receipt,
1234 });
1235 index.commit_resize(prepared);
1236 Ok(response)
1237 }
1238 "harness.v1.sessions.index.unsubscribe" => {
1239 let params = decode::<UnfollowParams>(params)?;
1240 Ok(json!({
1241 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1242 }))
1243 }
1244 "harness.v1.sessions.import" => {
1245 let params = decode::<ImportSessionParams>(params)?;
1246 let session = Session::load_str(¶ms.content, params.source_harness.into())
1247 .map_err(operation)?;
1248 Ok(json!({"session": normalized_session_json(&session)}))
1249 }
1250 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1251 let params = decode::<ExportSessionParams>(params)?;
1252 let session = load_session(¶ms.locator).map_err(operation)?;
1253 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1254 if method == "harness.v1.sessions.export"
1255 && params.target_harness == TransferFormat::Hermes
1256 {
1257 let imported = crate::hermes_import::import_into_hermes(&session, None)
1259 .map_err(operation)?;
1260 return Ok(json!({"artifact": artifact, "imported": imported}));
1261 }
1262 Ok(json!({"artifact": artifact}))
1263 }
1264 "harness.v1.sessions.reduce" => {
1265 let params = decode::<ReduceSessionParams>(params)?;
1266 self.reduce_session(params)
1267 }
1268 "harness.v1.sessions.branch" => {
1269 let params = decode::<BranchSessionParams>(params)?;
1270 let session = load_session(¶ms.locator).map_err(operation)?;
1271 let storage = params.locator.storage.path().display().to_string();
1272 let bootstrap_prompt = format!(
1273 "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.",
1274 params.locator.harness.as_str(), params.locator.session_id, storage
1275 );
1276 let artifact = params
1277 .target_harness
1278 .map(|target| session_artifact(¶ms.locator, &session, target))
1279 .transpose()?;
1280 Ok(json!({
1281 "parent": params.locator,
1282 "session": normalized_session_json(&session),
1283 "bootstrap_prompt": bootstrap_prompt,
1284 "artifact": artifact,
1285 }))
1286 }
1287 "harness.v1.sessions.handoff" => {
1288 let params = decode::<HandoffSessionParams>(params)?;
1289 let session = load_session(¶ms.locator).map_err(operation)?;
1290 let cwd = params
1291 .cwd
1292 .or_else(|| session.meta.cwd.clone())
1293 .unwrap_or_else(|| PathBuf::from("."));
1294 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1295 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1296 ServiceError::Operation(
1297 "handoff artifact omitted target session identity".into(),
1298 )
1299 })?;
1300 let instructions =
1301 handoff_instructions(params.target_harness, target_session_id, &cwd);
1302 Ok(json!({
1303 "artifact": artifact,
1304 "launch": instructions.launch,
1305 "materialize": instructions.materialize,
1306 "requires_materialization": instructions.requires_materialization,
1307 "note": instructions.note,
1308 }))
1309 }
1310 "harness.v1.sessions.materialize" => {
1311 let params = decode::<MaterializeSessionParams>(params)?;
1312 for file in ¶ms.artifact.files {
1316 if file.role == "source_recovery"
1317 && file.path == "recovery/source.supercode.jsonl"
1318 {
1319 if let Ok(source) = Session::from_native_str(&file.content) {
1320 crate::residue_store::store_segments(&source);
1321 }
1322 }
1323 }
1324 let locator = crate::native_materialize::materialize_native_artifact(
1325 params.artifact,
1326 ¶ms.cwd,
1327 ¶ms.homes,
1328 )
1329 .map_err(ServiceError::Operation)?;
1330 Ok(json!({"locator": locator}))
1331 }
1332 "harness.v1.jobs.list" => {
1336 let query = decode::<crate::jobs::JobsQuery>(params)?;
1337 if let Some(harness) = query.harness.as_deref() {
1338 refuse_harness_without_jobs(harness, "jobs.list")?;
1339 }
1340 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1341 serde_json::to_value(listing)
1342 .map_err(|error| ServiceError::Operation(error.to_string()))
1343 }
1344 "harness.v1.jobs.get" => {
1345 let params = decode::<JobsGetParams>(params)?;
1346 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1347 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1348 .map_err(operation)?
1349 {
1350 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1351 None => Err(ServiceError::Operation(format!(
1352 "`{}` has no scheduled job `{}`",
1353 params.harness, params.id
1354 ))),
1355 }
1356 }
1357 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1363 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1364 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1365 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1366 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1367 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1368 "harness.v1.jobs.notepad"
1369 | "harness.v1.jobs.notepad_set"
1370 | "harness.v1.jobs.notepad_delete" => {
1371 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1372 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1373 let answer = match method {
1374 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1375 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1376 _ => crate::jobs_notepad::read(&request),
1377 }
1378 .map_err(job_control_error)?;
1379 serde_json::to_value(answer)
1380 .map_err(|error| ServiceError::Operation(error.to_string()))
1381 }
1382 "harness.v1.runs.list" => {
1386 let query = decode::<crate::runs::RunsQuery>(params)?;
1387 if let Some(harness) = query.harness.as_deref() {
1388 refuse_harness_without_runs(harness, "runs.list")?;
1389 }
1390 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1391 serde_json::to_value(listing)
1392 .map_err(|error| ServiceError::Operation(error.to_string()))
1393 }
1394 "harness.v1.runs.get" => {
1395 let params = decode::<RunsGetParams>(params)?;
1396 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1397 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1398 .map_err(operation)?
1399 {
1400 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1401 None => Err(ServiceError::Operation(format!(
1402 "`{}` has no run `{}`",
1403 params.harness, params.id
1404 ))),
1405 }
1406 }
1407 "harness.v1.sessions.resume_instructions" => {
1408 let params = decode::<ResumeInstructionsParams>(params)?;
1409 let session = load_session(¶ms.locator).map_err(operation)?;
1410 let cwd = params
1411 .cwd
1412 .or(session.meta.cwd)
1413 .unwrap_or_else(|| PathBuf::from("."));
1414 let launch = resume_launch(
1415 params.locator.harness.as_str(),
1416 ¶ms.locator.session_id,
1417 &cwd,
1418 params.policy,
1419 )?;
1420 Ok(json!({"launch": launch}))
1421 }
1422 _ => Err(ServiceError::MethodNotFound),
1423 }
1424 }
1425
1426 fn reduce_session(
1427 &self,
1428 params: ReduceSessionParams,
1429 ) -> std::result::Result<Value, ServiceError> {
1430 let session = load_session(¶ms.locator).map_err(operation)?;
1431 if session.messages.is_empty() {
1432 return Err(ServiceError::InvalidParams(
1433 "cannot reduce an empty session".into(),
1434 ));
1435 }
1436 let keep_last = params.keep_last.clamp(1, 128);
1437 let policy = reduce::ReductionPolicy {
1438 clear_turns_older_than: Some(keep_last),
1439 ..Default::default()
1440 };
1441 let (view, log) =
1442 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1443 if log.reductions.is_empty() {
1444 return Err(ServiceError::UnsupportedAction(format!(
1445 "session `{}` is already too small for a meaningful reversible reduction",
1446 params.locator.session_id
1447 )));
1448 }
1449 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1450 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1451 if reduced_tokens >= source_tokens {
1452 return Err(ServiceError::UnsupportedAction(format!(
1453 "session `{}` has no token-reducing reversible projection",
1454 params.locator.session_id
1455 )));
1456 }
1457
1458 let store_root = self
1459 .reduction_store_root
1460 .clone()
1461 .unwrap_or_else(default_reduction_store_root);
1462 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1463 let rescue_id = format!("rescue-{}", generated_session_id());
1464 let imported = session
1465 .imported_message_count
1466 .unwrap_or(session.messages.len())
1467 .min(session.messages.len());
1468 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1469 let view_jsonl = messages_jsonl(&view)?;
1470 let title = format!(
1471 "Reduced {} continuation from {}",
1472 params.target_harness.id(),
1473 params.locator.session_id
1474 );
1475
1476 store
1481 .save_sidecar(&rescue_id, &sidecar_jsonl)
1482 .map_err(operation)?;
1483 store
1484 .save_reduction_log(&rescue_id, &log)
1485 .map_err(operation)?;
1486 store
1487 .save(&rescue_id, &title, &view_jsonl)
1488 .map_err(operation)?;
1489
1490 let source_bytes = serde_json::to_vec(&session.messages)
1491 .map_err(|error| ServiceError::Operation(error.to_string()))?
1492 .len() as u64;
1493 let reduced_bytes = serde_json::to_vec(&view)
1494 .map_err(|error| ServiceError::Operation(error.to_string()))?
1495 .len() as u64;
1496 store
1497 .set_reduction_stats(
1498 &rescue_id,
1499 &title,
1500 source_bytes,
1501 reduced_bytes,
1502 log.reductions.len() as u32,
1503 )
1504 .map_err(operation)?;
1505
1506 let reloaded_sidecar = store
1510 .load_sidecar(&rescue_id)
1511 .map_err(operation)?
1512 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1513 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1514 let reloaded_log = store
1515 .load_reduction_log(&rescue_id)
1516 .map_err(operation)?
1517 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1518 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1519 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1520 let (restamped_view, restamped_log) =
1527 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1528 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1529 return Err(ServiceError::Operation(
1530 "persisted reduction view does not match its durable log and sidecar".into(),
1531 ));
1532 }
1533 if restamped_log != reloaded_log {
1534 return Err(ServiceError::Operation(
1535 "reapplying the durable reduction log changed its identity".into(),
1536 ));
1537 }
1538 let inverted =
1539 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1540 if inverted != session.messages {
1541 return Err(ServiceError::Operation(
1542 "reduction inversion did not restore the source messages byte-exactly".into(),
1543 ));
1544 }
1545
1546 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1547 let sidecar_path = store.sidecar_path(&rescue_id);
1548 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1549 let bootstrap_prompt = reduced_bootstrap_prompt(
1550 ¶ms.locator,
1551 params.target_harness,
1552 &view_jsonl,
1553 &sidecar_path,
1554 &reduction_log_path,
1555 );
1556 let mut reduced_session = session.clone();
1557 reduced_session.meta.session_id = Some(rescue_id.clone());
1558 reduced_session.messages = view;
1559
1560 Ok(json!({
1561 "session": normalized_session_json(&reduced_session),
1562 "bootstrap_prompt": bootstrap_prompt,
1563 "receipt": {
1564 "id": rescue_id,
1565 "sidecar_id": rescue_id,
1566 "source_harness": params.locator.harness,
1567 "target_harness": params.target_harness.id(),
1568 "source_tokens": source_tokens,
1569 "reduced_tokens": reduced_tokens,
1570 "ratio": ratio,
1571 "source_bytes": source_bytes,
1572 "reduced_bytes": reduced_bytes,
1573 "reductions": reloaded_log.reductions.len(),
1574 "sidecar_path": sidecar_path,
1575 "reduction_log_path": reduction_log_path,
1576 "verified": true,
1577 "reversible": true,
1578 }
1579 }))
1580 }
1581
1582 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1600 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1601 return None;
1602 }
1603 let method = request.get("method").and_then(Value::as_str)?;
1604 if !RUNTIME_OPEN_METHODS.contains(&method) {
1605 return None;
1606 }
1607 Some(RuntimeOpen {
1608 id: request.get("id").cloned().unwrap_or(Value::Null),
1609 method: method.to_string(),
1610 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1611 })
1612 }
1613
1614 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1636 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1637 return None;
1638 }
1639 let method = request.get("method").and_then(Value::as_str)?;
1640 if !DETACHED_METHODS.contains(&method) {
1641 return None;
1642 }
1643 let id = request.get("id").cloned().unwrap_or(Value::Null);
1644 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1645 let work = match method {
1646 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1647 .inventory_work(method, params)
1648 .map(DetachedWork::Inventory),
1649 "harness.v1.sessions.message" => {
1650 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1651 }
1652 _ => {
1653 let verb = match method {
1654 "harness.v1.sessions.new" => crate::SessionVerb::New,
1655 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1656 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1657 _ => crate::SessionVerb::Delete,
1658 };
1659 match decode::<crate::SessionMutation>(params) {
1660 Ok(mutation) => {
1661 match crate::sessions_control::door(&mutation.harness, verb) {
1662 Ok(crate::SessionDoor::Live(_)) => return None,
1665 Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1666 Err(error) => Err(session_control_error(error)),
1667 }
1668 }
1669 Err(error) => Err(error),
1670 }
1671 }
1672 };
1673 Some(DetachedCall {
1674 id,
1675 method: method.to_string(),
1676 work: work.map(Work::Free),
1677 })
1678 }
1679
1680 pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1694 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1695 return None;
1696 }
1697 let method = request.get("method").and_then(Value::as_str)?;
1698 let id = request.get("id").cloned().unwrap_or(Value::Null);
1699 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1700 let work = match method {
1701 "harness.v1.runtimes.close" => decode::<RuntimeConnectionParams>(params)
1702 .and_then(|params| self.surrender_runtime(¶ms.connection))
1703 .map(|(runtime, process_group)| {
1704 Work::Runtime(RuntimeWork::Close {
1705 runtime,
1706 process_group,
1707 })
1708 }),
1709 "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1710 let verb = if method == "harness.v1.sessions.new" {
1711 crate::SessionVerb::New
1712 } else {
1713 crate::SessionVerb::Reset
1714 };
1715 let mutation = decode::<crate::SessionMutation>(params).ok()?;
1716 let Ok(crate::SessionDoor::Live(command)) =
1720 crate::sessions_control::door(&mutation.harness, verb)
1721 else {
1722 return None;
1723 };
1724 let connection = mutation
1725 .connection
1726 .clone()
1727 .filter(|value| !value.trim().is_empty())?;
1728 self.lend_runtime(&connection).map(|runtime| {
1729 let session = live_session_name(runtime.as_ref(), &mutation);
1730 Work::Runtime(RuntimeWork::LiveCommand {
1731 connection,
1732 runtime,
1733 verb,
1734 mutation,
1735 command,
1736 session,
1737 })
1738 })
1739 }
1740 _ => return None,
1741 };
1742 Some(DetachedCall {
1743 id,
1744 method: method.to_string(),
1745 work,
1746 })
1747 }
1748
1749 pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1754 let DetachedAnswer { response, returned } = answer;
1755 if let Some(ReturnedRuntime {
1756 connection,
1757 runtime,
1758 }) = returned
1759 {
1760 self.runtimes_in_flight.remove(&connection);
1761 self.runtimes.insert(connection, runtime);
1762 }
1763 response
1764 }
1765
1766 pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1770 let OpenedRuntime { id, outcome } = opened;
1771 let result = match outcome {
1772 Ok(open) => self.register_open_runtime(open).await,
1773 Err(error) => Err(error),
1774 };
1775 service_response(id, result)
1776 }
1777
1778 async fn register_open_runtime(
1780 &mut self,
1781 open: OpenRuntime,
1782 ) -> std::result::Result<Value, ServiceError> {
1783 match open {
1784 OpenRuntime::Hosted {
1785 runtime,
1786 capabilities,
1787 workspace,
1788 fresh,
1789 } => {
1790 self.insert_hosted_runtime(runtime, capabilities, workspace, fresh)
1791 .await
1792 }
1793 OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1794 }
1795 }
1796
1797 async fn runtime_call(
1798 &mut self,
1799 method: &str,
1800 params: Value,
1801 ) -> std::result::Result<Value, ServiceError> {
1802 match method {
1803 "harness.v1.runtimes.capabilities" => {
1804 let params = decode::<RuntimeBackendParams>(params)?;
1805 let backend = runtime_backend(¶ms)?;
1806 Ok(json!({
1807 "harness": backend.harness(),
1808 "capabilities": backend.capabilities(),
1809 }))
1810 }
1811 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1812 self.register_open_runtime(open_runtime(method, params).await?)
1813 .await
1814 }
1815 "harness.v1.runtimes.send_input" => {
1816 let params = decode::<RuntimeInputParams>(params)?;
1817 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1818 let runtime = self.runtime_mut(¶ms.connection)?;
1819 let turn_id = within_control_deadline(
1820 method,
1821 runtime.send_input(RuntimeInput {
1822 text: params.text,
1823 image_urls,
1824 }),
1825 )
1826 .await?
1827 .map_err(operation)?;
1828 Ok(json!({"turn_id": turn_id}))
1829 }
1830 "harness.v1.runtimes.interrupt" => {
1831 let params = decode::<RuntimeConnectionParams>(params)?;
1832 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1833 .await?
1834 .map_err(operation)?;
1835 Ok(json!({}))
1836 }
1837 "harness.v1.runtimes.steer" => {
1838 let params = decode::<RuntimeInputParams>(params)?;
1839 if !params.image_urls.is_empty() {
1840 return Err(ServiceError::InvalidParams(
1841 "runtime steering accepts text only".into(),
1842 ));
1843 }
1844 let text = params.text.trim();
1845 if text.is_empty() || text.chars().count() > 50_000 {
1846 return Err(ServiceError::InvalidParams(
1847 "runtime steering requires 1 to 50,000 text characters".into(),
1848 ));
1849 }
1850 within_control_deadline(
1851 method,
1852 self.runtime_mut(¶ms.connection)?
1853 .steer(text.to_string()),
1854 )
1855 .await?
1856 .map_err(operation)?;
1857 Ok(json!({}))
1858 }
1859 "harness.v1.runtimes.respond" => {
1860 let params = decode::<RuntimeRespondParams>(params)?;
1861 let request_id = params.request_id.clone();
1862 within_control_deadline(
1863 method,
1864 self.runtime_mut(¶ms.connection)?
1865 .respond(params.request_id, params.response),
1866 )
1867 .await?
1868 .map_err(operation)?;
1869 self.approvals.answered(¶ms.connection, &request_id);
1871 Ok(json!({}))
1872 }
1873 "harness.v1.runtimes.acquire_control" => {
1874 let params = decode::<RuntimeConnectionParams>(params)?;
1875 let snapshot = within_control_deadline(
1876 method,
1877 self.runtime_mut(¶ms.connection)?.acquire_control(),
1878 )
1879 .await?
1880 .map_err(operation)?;
1881 serde_json::to_value(snapshot)
1882 .map_err(|error| ServiceError::Operation(error.to_string()))
1883 }
1884 "harness.v1.runtimes.heartbeat" => {
1885 let params = decode::<RuntimeConnectionParams>(params)?;
1886 let snapshot = within_control_deadline(
1887 method,
1888 self.runtime_mut(¶ms.connection)?.heartbeat(),
1889 )
1890 .await?
1891 .map_err(operation)?;
1892 serde_json::to_value(snapshot)
1893 .map_err(|error| ServiceError::Operation(error.to_string()))
1894 }
1895 "harness.v1.runtimes.detach" => {
1896 let params = decode::<RuntimeConnectionParams>(params)?;
1897 let snapshot =
1898 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1899 .await?
1900 .map_err(operation)?;
1901 serde_json::to_value(snapshot)
1902 .map_err(|error| ServiceError::Operation(error.to_string()))
1903 }
1904 "harness.v1.runtimes.terminal_instructions" => {
1905 let params = decode::<RuntimeConnectionParams>(params)?;
1906 let launch = self
1907 .terminal_launches
1908 .get(¶ms.connection)
1909 .ok_or_else(|| {
1910 ServiceError::Operation(
1911 "this runtime is not hosted for terminal attachment".into(),
1912 )
1913 })?;
1914 Ok(json!({"launch":launch}))
1915 }
1916 "harness.v1.runtimes.close" => {
1917 let params = decode::<RuntimeConnectionParams>(params)?;
1918 let (runtime, process_group) = self.surrender_runtime(¶ms.connection)?;
1919 close_runtime(runtime, process_group).await
1920 }
1921 _ => Err(ServiceError::MethodNotFound),
1922 }
1923 }
1924
1925 #[cfg(feature = "adapter-api")]
1927 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1928 let params = decode::<MessageSessionParams>(params)?;
1929 Ok(message_live_session(¶ms).await)
1930 }
1931
1932 #[cfg(feature = "adapter-api")]
1933 fn harness_settings_call(
1934 &self,
1935 method: &str,
1936 params: Value,
1937 ) -> std::result::Result<Value, ServiceError> {
1938 let homes = crate::HarnessHomes::default();
1939 match method {
1940 "harness.v1.harnesses.settings" => {
1941 let params = decode::<HarnessSettingsParams>(params)?;
1942 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1943 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1944 serde_json::to_value(report)
1945 .map_err(|error| ServiceError::Operation(error.to_string()))
1946 }
1947 "harness.v1.harnesses.configure" => {
1948 let params = decode::<ConfigureHarnessParams>(params)?;
1949 let report = crate::configure_harness_interop_settings(
1950 &homes,
1951 ¶ms.harness,
1952 ¶ms.changes,
1953 params.expected_revision.as_deref(),
1954 )
1955 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1956 serde_json::to_value(report)
1957 .map_err(|error| ServiceError::Operation(error.to_string()))
1958 }
1959 _ => Err(ServiceError::MethodNotFound),
1960 }
1961 }
1962
1963 fn insert_runtime(
1964 &mut self,
1965 runtime: Box<dyn RuntimeConnection>,
1966 ) -> std::result::Result<Value, ServiceError> {
1967 let connection = format!("runtime-{}", self.next_runtime);
1968 self.next_runtime += 1;
1969 let handle = runtime.handle().clone();
1970 self.runtime_sequences
1971 .entry(handle.runtime_id.clone())
1972 .or_insert(0);
1973 self.runtimes.insert(connection.clone(), runtime);
1974 Ok(json!({"connection": connection, "handle": handle}))
1975 }
1976
1977 #[cfg(feature = "adapter-api")]
1978 async fn insert_hosted_runtime(
1979 &mut self,
1980 runtime: Box<dyn RuntimeConnection>,
1981 capabilities: crate::RuntimeCapabilities,
1982 workspace: PathBuf,
1983 fresh: bool,
1984 ) -> std::result::Result<Value, ServiceError> {
1985 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
1986 let token: std::sync::Arc<str> = crate::server::generate_token().into();
1987 let server = crate::server::run_frontend_http(
1988 host.clone(),
1989 host.frontend_sender(),
1990 "127.0.0.1:0",
1991 token.clone(),
1992 connection.handle().runtime_id.clone(),
1993 )
1994 .await
1995 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1996 let source = LiveRuntimeSource {
1997 harness: connection.handle().harness.as_str().to_string(),
1998 session_id: connection.handle().runtime_id.clone(),
1999 workspace: workspace.clone(),
2000 };
2001 let registration = register_live_runtime(
2002 connection.handle().runtime_id.clone(),
2003 source.clone(),
2004 format!("http://{}", server.address()),
2005 token.to_string(),
2006 )
2007 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2008 let endpoint = registration.endpoint().to_string();
2009 let launch = StructuredLaunch {
2010 cwd: workspace,
2011 program: std::env::current_exe()
2015 .ok()
2016 .map(|path| path.to_string_lossy().into_owned())
2017 .unwrap_or_else(|| "supercode".into()),
2018 arguments: vec![
2019 "open".into(),
2020 endpoint,
2021 "--harness".into(),
2022 source.harness,
2023 "--session".into(),
2024 source.session_id,
2025 ],
2026 env: BTreeMap::new(),
2027 };
2028 let lease = HostedRuntimeLease {
2029 connection,
2030 _host: host,
2031 _registration: registration,
2032 _server: server,
2033 };
2034 let opened = self.insert_runtime(Box::new(lease))?;
2035 let connection_id = opened["connection"]
2036 .as_str()
2037 .expect("insert_runtime returns a connection id")
2038 .to_string();
2039 self.terminal_launches.insert(connection_id, launch);
2040 Ok(opened)
2041 }
2042
2043 #[cfg(not(feature = "adapter-api"))]
2044 async fn insert_hosted_runtime(
2045 &mut self,
2046 runtime: Box<dyn RuntimeConnection>,
2047 _capabilities: crate::RuntimeCapabilities,
2048 _workspace: PathBuf,
2049 _fresh: bool,
2050 ) -> std::result::Result<Value, ServiceError> {
2051 self.insert_runtime(runtime)
2052 }
2053
2054 fn runtime_mut(
2055 &mut self,
2056 connection: &str,
2057 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2058 if self.runtimes_in_flight.contains(connection) {
2059 return Err(self.lent_out(connection));
2060 }
2061 self.runtimes.get_mut(connection).ok_or_else(|| {
2062 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2063 })
2064 }
2065
2066 fn lent_out(&self, connection: &str) -> ServiceError {
2070 ServiceError::Operation(format!(
2071 "runtime connection `{connection}`: a harness turn is already in progress"
2072 ))
2073 }
2074
2075 fn lend_runtime(
2078 &mut self,
2079 connection: &str,
2080 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2081 if self.runtimes_in_flight.contains(connection) {
2082 return Err(self.lent_out(connection));
2083 }
2084 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2085 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2086 })?;
2087 self.runtimes_in_flight.insert(connection.to_string());
2088 Ok(runtime)
2089 }
2090
2091 fn surrender_runtime(
2103 &mut self,
2104 connection: &str,
2105 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2106 if self.runtimes_in_flight.contains(connection) {
2107 return Err(self.lent_out(connection));
2108 }
2109 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2110 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2111 })?;
2112 let process_group = runtime_process_group(runtime.handle());
2113 let runtime_id = runtime.handle().runtime_id.clone();
2114 self.terminal_launches.remove(connection);
2115 self.runtime_sequences.remove(&runtime_id);
2116 self.approvals.forget(connection);
2117 Ok((runtime, process_group))
2118 }
2119
2120 pub fn kill_all_runtime_groups(&self) -> usize {
2131 self.runtimes
2132 .values()
2133 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2134 .count()
2135 }
2136
2137 async fn mutate_session(
2148 &mut self,
2149 verb: crate::SessionVerb,
2150 params: Value,
2151 ) -> std::result::Result<Value, ServiceError> {
2152 let mutation = decode::<crate::SessionMutation>(params)?;
2153 let door = crate::sessions_control::door(&mutation.harness, verb)
2154 .map_err(session_control_error)?;
2155 let outcome = match door {
2156 #[cfg(not(feature = "adapter-api"))]
2160 crate::SessionDoor::Live(command) => {
2161 return Err(ServiceError::Operation(format!(
2162 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2163 session, which needs this build's `adapter-api` feature",
2164 mutation.harness,
2165 verb.as_str()
2166 )));
2167 }
2168 #[cfg(feature = "adapter-api")]
2169 crate::SessionDoor::Live(command) => {
2170 let connection = mutation
2171 .connection
2172 .clone()
2173 .filter(|value| !value.trim().is_empty())
2174 .ok_or_else(|| {
2175 ServiceError::InvalidParams(format!(
2176 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2177 driven session: pass the `connection` of an open runtime \
2178 (`harness.v1.runtimes.start`)",
2179 mutation.harness,
2180 verb.as_str()
2181 ))
2182 })?;
2183 let runtime = self.runtime_mut(&connection)?;
2184 let session = live_session_name(runtime.as_ref(), &mutation);
2185 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2191 .await;
2192 }
2193 _ => run_session_mutation(verb, &mutation).await?,
2194 };
2195 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2196 }
2197
2198 async fn inventory_call(
2202 &self,
2203 method: &str,
2204 params: Value,
2205 ) -> std::result::Result<Value, ServiceError> {
2206 run_inventory(self.inventory_work(method, params)?).await
2207 }
2208
2209 fn inventory_work(
2215 &self,
2216 method: &str,
2217 params: Value,
2218 ) -> std::result::Result<InventoryWork, ServiceError> {
2219 let mut params = decode::<HarnessInventoryParams>(params)?;
2220 if method == "harness.v1.harnesses.probe" {
2221 let harness = params.harness.take().ok_or_else(|| {
2222 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2223 })?;
2224 params.harnesses = vec![harness];
2225 }
2226 let selected = params
2227 .harnesses
2228 .iter()
2229 .map(HarnessId::as_str)
2230 .collect::<std::collections::BTreeSet<_>>();
2231 let supported = harness_support_registry()
2232 .harnesses
2233 .into_iter()
2234 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2235 .collect::<Vec<_>>();
2236 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2237 let known = supported
2238 .iter()
2239 .map(|harness| harness.id.as_str())
2240 .collect::<std::collections::BTreeSet<_>>();
2241 let missing = params
2242 .harnesses
2243 .iter()
2244 .filter(|id| !known.contains(id.as_str()))
2245 .map(HarnessId::as_str)
2246 .collect::<Vec<_>>();
2247 return Err(ServiceError::InvalidParams(format!(
2248 "unknown harness(es): {}",
2249 missing.join(", ")
2250 )));
2251 }
2252 let global_counts = params
2253 .include_sessions
2254 .then(|| self.session_counts(None, ¶ms.harnesses));
2255 let workspace_counts = params
2256 .include_sessions
2257 .then(|| {
2258 params
2259 .workspace
2260 .as_deref()
2261 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2262 })
2263 .flatten();
2264 Ok(InventoryWork {
2265 params,
2266 supported,
2267 global_counts,
2268 workspace_counts,
2269 })
2270 }
2271
2272 #[cfg(feature = "adapter-api")]
2273 async fn harness_authentication_call(
2274 &self,
2275 method: &str,
2276 params: Value,
2277 ) -> std::result::Result<Value, ServiceError> {
2278 match method {
2279 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2280 let params = decode::<HarnessAuthenticationParams>(params)?;
2281 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2282 .map_err(|error| ServiceError::Operation(error.to_string()))
2283 }
2284 "harness.v1.harnesses.auth.begin" => {
2285 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2286 let cwd = params
2287 .cwd
2288 .or_else(|| std::env::current_dir().ok())
2289 .unwrap_or_else(|| PathBuf::from("."));
2290 let plan = crate::harness_authentication_plan(
2291 ¶ms.harness,
2292 params.environment,
2293 params.method,
2294 &cwd,
2295 )
2296 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2297 serde_json::to_value(plan)
2298 .map_err(|error| ServiceError::Operation(error.to_string()))
2299 }
2300 _ => Err(ServiceError::MethodNotFound),
2301 }
2302 }
2303
2304 fn session_counts(
2305 &self,
2306 workspace: Option<&Path>,
2307 harnesses: &[HarnessId],
2308 ) -> BTreeMap<String, usize> {
2309 let mut counts = BTreeMap::new();
2310 for session in self
2311 .catalog
2312 .discover(&DiscoveryQuery {
2313 workspace: workspace.map(Path::to_path_buf),
2314 harnesses: harnesses.to_vec(),
2315 ..DiscoveryQuery::default()
2316 })
2317 .unwrap_or_default()
2318 {
2319 *counts
2320 .entry(session.locator.harness.as_str().to_string())
2321 .or_insert(0) += 1;
2322 }
2323 counts
2324 }
2325}
2326
2327#[async_trait::async_trait]
2328impl SdkService for HarnessSessionService {
2329 fn capabilities(&self) -> SdkCapabilities {
2330 SdkCapabilities::default()
2331 }
2332
2333 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2334 if request.operation == SdkOperation::Events {
2335 let events = self
2336 .poll_sdk_events()
2337 .await
2338 .into_iter()
2339 .map(|(_, event)| event)
2340 .collect::<Vec<_>>();
2341 return serde_json::to_value(events).map_err(|error| {
2342 SdkError::new(
2343 SdkErrorCode::Execution,
2344 request.operation,
2345 error.to_string(),
2346 )
2347 });
2348 }
2349 if self.runtimes.is_empty()
2350 && matches!(
2351 request.operation,
2352 SdkOperation::Input
2353 | SdkOperation::Interrupt
2354 | SdkOperation::Steer
2355 | SdkOperation::Respond
2356 | SdkOperation::Close
2357 )
2358 {
2359 return Err(SdkError::unsupported(request.operation));
2360 }
2361 let method = request
2362 .operation
2363 .method()
2364 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2365 let result = match request.operation {
2366 SdkOperation::Discover
2367 | SdkOperation::Load
2368 | SdkOperation::Export
2369 | SdkOperation::ProfilesList
2370 | SdkOperation::ProfilesGet
2371 | SdkOperation::ProfilesCreate
2372 | SdkOperation::ProfilesDelete
2373 | SdkOperation::SkillsList
2374 | SdkOperation::SkillsInstall
2375 | SdkOperation::SkillsRemove
2376 | SdkOperation::ChannelsList
2377 | SdkOperation::RoutesList
2378 | SdkOperation::TriggersList
2379 | SdkOperation::ChannelsStatus
2380 | SdkOperation::MemoryShow
2381 | SdkOperation::MemorySearch
2382 | SdkOperation::JobsList
2383 | SdkOperation::JobsGet
2384 | SdkOperation::JobsCreate
2385 | SdkOperation::JobsUpdate
2386 | SdkOperation::JobsPause
2387 | SdkOperation::JobsResume
2388 | SdkOperation::JobsRun
2389 | SdkOperation::JobsDelete
2390 | SdkOperation::JobsNotepad
2391 | SdkOperation::JobsNotepadSet
2392 | SdkOperation::JobsNotepadDelete
2393 | SdkOperation::RunsList
2394 | SdkOperation::RunsGet
2395 | SdkOperation::ApprovalsList
2396 | SdkOperation::OrchestrationLoad
2397 | SdkOperation::OrchestrationSave
2398 | SdkOperation::OrchestrationCompile
2399 | SdkOperation::OrchestrationDecompile
2400 | SdkOperation::OrchestrationImport
2401 | SdkOperation::OrchestrationExport
2402 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2403 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2406 SdkOperation::Start
2407 | SdkOperation::Resume
2408 | SdkOperation::Input
2409 | SdkOperation::Interrupt
2410 | SdkOperation::Steer
2411 | SdkOperation::Respond
2412 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2413 SdkOperation::SessionsNew => {
2418 self.mutate_session(crate::SessionVerb::New, request.params)
2419 .await
2420 }
2421 SdkOperation::SessionsReset => {
2422 self.mutate_session(crate::SessionVerb::Reset, request.params)
2423 .await
2424 }
2425 SdkOperation::SessionsArchive => {
2426 self.mutate_session(crate::SessionVerb::Archive, request.params)
2427 .await
2428 }
2429 SdkOperation::SessionsDelete => {
2430 self.mutate_session(crate::SessionVerb::Delete, request.params)
2431 .await
2432 }
2433 SdkOperation::Events => unreachable!("handled before method dispatch"),
2434 };
2435 result.map_err(|error| sdk_error(request.operation, error))
2436 }
2437
2438 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2439 Ok(self
2440 .poll_sdk_events()
2441 .await
2442 .into_iter()
2443 .map(|(_, event)| event)
2444 .collect())
2445 }
2446}
2447
2448#[cfg(feature = "adapter-api")]
2449struct HostedRuntimeLease {
2450 connection: HostedHarnessConnection,
2451 _host: std::sync::Arc<HostedHarnessRuntime>,
2452 _registration: LiveRuntimeRegistration,
2453 _server: crate::server::FrontendHttpServer,
2454}
2455
2456#[async_trait::async_trait]
2457#[cfg(feature = "adapter-api")]
2458impl RuntimeConnection for HostedRuntimeLease {
2459 fn handle(&self) -> &crate::RuntimeHandle {
2460 self.connection.handle()
2461 }
2462
2463 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2464 self.connection.send_input(input).await
2465 }
2466
2467 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2468 self.connection.next_event().await
2469 }
2470
2471 async fn interrupt(&mut self) -> crate::Result<()> {
2472 self.connection.interrupt().await
2473 }
2474
2475 async fn steer(&mut self, text: String) -> crate::Result<()> {
2478 self.connection.steer(text).await
2479 }
2480
2481 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2482 self.connection.respond(request_id, response).await
2483 }
2484
2485 async fn close(&mut self) -> crate::Result<()> {
2486 self.connection.close().await
2487 }
2488}
2489
2490struct InventoryWork {
2493 params: HarnessInventoryParams,
2494 supported: Vec<crate::HarnessSupportDescriptor>,
2495 global_counts: Option<BTreeMap<String, usize>>,
2496 workspace_counts: Option<BTreeMap<String, usize>>,
2497}
2498
2499async fn run_session_mutation(
2506 verb: crate::SessionVerb,
2507 mutation: &crate::SessionMutation,
2508) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2509 let door =
2516 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2517 if let crate::SessionDoor::Http = door {
2518 return crate::sessions_control::mutate(verb, mutation)
2519 .await
2520 .map_err(session_control_error);
2521 }
2522 let mutation = mutation.clone();
2523 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2524 .await
2525 .map_err(|error| {
2526 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2527 })?
2528 .map_err(session_control_error)
2529}
2530
2531async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2534 let InventoryWork {
2535 params,
2536 supported,
2537 global_counts,
2538 workspace_counts,
2539 } = work;
2540 let probes = supported.into_iter().map(|descriptor| {
2541 let global = global_counts
2542 .as_ref()
2543 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2544 let workspace = workspace_counts
2545 .as_ref()
2546 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2547 probe_harness(descriptor, ¶ms, global, workspace)
2548 });
2549 let harnesses = futures::future::join_all(probes).await;
2550 serde_json::to_value(HarnessInventoryReport {
2551 probe: params.probe,
2552 workspace: params.workspace,
2553 harnesses,
2554 })
2555 .map_err(|error| ServiceError::Operation(error.to_string()))
2556}
2557
2558async fn probe_harness(
2559 descriptor: crate::HarnessSupportDescriptor,
2560 params: &HarnessInventoryParams,
2561 global: Option<usize>,
2562 workspace: Option<usize>,
2563) -> LocalHarness {
2564 let launch = descriptor.runtime.default_launch.as_ref();
2565 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2570 .then(crate::orchestrator::daemon_entry)
2571 .and_then(Result::ok);
2572 let executable = match &orchestrator_entry {
2573 Some(entry) => Some(entry.clone()),
2574 None => launch.and_then(|launch| find_executable(&launch.program)),
2575 };
2576 let installed = executable.is_some();
2577 let version = if params.skip_versions || orchestrator_entry.is_some() {
2578 None
2581 } else {
2582 match executable.as_deref() {
2583 Some(path) => executable_version(path).await,
2584 None => None,
2585 }
2586 };
2587 let configured = auth_evidence(descriptor.id.as_str());
2588 let mut auth = if configured {
2589 HarnessAuthState::Configured
2590 } else if matches!(
2591 descriptor.id.as_str(),
2592 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2593 ) {
2594 HarnessAuthState::Required
2599 } else {
2600 HarnessAuthState::Unknown
2601 };
2602 let mut runtime = if installed {
2603 HarnessRuntimeState::Degraded
2604 } else {
2605 HarnessRuntimeState::Unavailable
2606 };
2607 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2608 let mut reason = (!installed).then(|| {
2609 if is_orchestrator {
2610 format!(
2611 "{} is supported but its daemon entry `{}` was not found",
2612 descriptor.display_name,
2613 crate::orchestrator::DAEMON_ENTRY
2614 )
2615 } else {
2616 format!(
2617 "{} is supported but `{}` was not found on PATH",
2618 descriptor.display_name,
2619 launch
2620 .map(|launch| launch.program.as_str())
2621 .unwrap_or("executable")
2622 )
2623 }
2624 });
2625 let mut repair = (!installed).then(|| {
2626 if is_orchestrator {
2627 format!(
2628 "Install the `supercode-orchestrator` package so `{}` resolves.",
2629 crate::orchestrator::DAEMON_ENTRY
2630 )
2631 } else {
2632 format!(
2633 "Install {} and ensure `{}` is on PATH.",
2634 descriptor.display_name,
2635 launch
2636 .map(|launch| launch.program.as_str())
2637 .unwrap_or("its executable")
2638 )
2639 }
2640 });
2641
2642 if installed && params.probe == HarnessProbeLevel::Handshake {
2643 let backend_params = RuntimeBackendParams {
2644 harness: descriptor.id.clone(),
2645 protocol: None,
2646 launch: None,
2647 base_url: None,
2648 policy: RuntimePolicy::Default,
2649 };
2650 match runtime_backend(&backend_params) {
2651 Ok(backend) => {
2652 let cwd = params
2653 .workspace
2654 .clone()
2655 .or_else(|| std::env::current_dir().ok())
2656 .unwrap_or_else(|| PathBuf::from("."));
2657 let isolated = descriptor
2658 .runtime
2659 .default_launch
2660 .clone()
2661 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2662 let Some(isolated) = isolated else {
2663 reason = Some(
2664 "No-prompt runtime handshake could not create its isolated harness home."
2665 .into(),
2666 );
2667 repair = Some(
2668 "Check temporary-directory permissions, then run the handshake probe again."
2669 .into(),
2670 );
2671 let running = probe_running_instance(descriptor.id.as_str());
2672 return LocalHarness {
2673 gateway: gateway_health(
2674 descriptor.id.as_str(),
2675 installed,
2676 running.as_ref(),
2677 version.as_deref(),
2678 ),
2679 id: descriptor.id,
2680 display_name: descriptor.display_name,
2681 supported: true,
2682 installed,
2683 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2684 version,
2685 auth,
2686 runtime,
2687 protocol: descriptor.runtime.protocol,
2688 capabilities: descriptor.runtime.capabilities.clone(),
2689 effective_capabilities: descriptor.runtime.capabilities,
2690 sessions: HarnessSessionCounts { global, workspace },
2691 running,
2692 reason,
2693 repair,
2694 };
2695 };
2696 match tokio::time::timeout(
2697 Duration::from_secs(30),
2698 backend.start(RuntimeStartRequest {
2699 cwd,
2700 launch: Some(isolated.launch.clone()),
2701 mcp_servers: Vec::new(),
2702 approval_policy: None,
2703 }),
2704 )
2705 .await
2706 {
2707 Ok(Ok(mut connection)) => {
2708 match stabilize_handshake(connection.as_mut()).await {
2709 Ok(()) => {
2710 auth = HarnessAuthState::Ready;
2711 runtime = HarnessRuntimeState::Ready;
2712 reason = Some(
2713 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2714 .into(),
2715 );
2716 repair = None;
2717 }
2718 Err(message) => {
2719 auth = if looks_like_auth_error(&message) {
2720 HarnessAuthState::Required
2721 } else if configured {
2722 HarnessAuthState::Configured
2723 } else {
2724 HarnessAuthState::Unknown
2725 };
2726 reason = Some(format!(
2727 "No-prompt runtime handshake became unhealthy during startup: {message}"
2728 ));
2729 repair = Some(if auth == HarnessAuthState::Required {
2730 format!(
2731 "Run `{}` interactively once and complete sign-in, then probe again.",
2732 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2733 )
2734 } else {
2735 "Run the harness directly to inspect its startup failure, then probe again."
2736 .into()
2737 });
2738 }
2739 }
2740 let _ =
2741 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2742 }
2743 Ok(Err(error)) => {
2744 let message = truncate_text(&error.to_string(), 500);
2745 auth = if looks_like_auth_error(&message) {
2746 HarnessAuthState::Required
2747 } else if configured {
2748 HarnessAuthState::Configured
2749 } else {
2750 HarnessAuthState::Unknown
2751 };
2752 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2753 repair = Some(if auth == HarnessAuthState::Required {
2754 format!(
2755 "Run `{}` interactively once and complete sign-in, then probe again.",
2756 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2757 )
2758 } else {
2759 "Check the harness installation and run the handshake probe again."
2760 .into()
2761 });
2762 }
2763 Err(_) => {
2764 reason =
2765 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2766 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2767 }
2768 }
2769 let _ = isolated.cleanup();
2778 tokio::time::sleep(Duration::from_millis(250)).await;
2779 if let Err(error) = isolated.cleanup() {
2780 auth = if configured {
2781 HarnessAuthState::Configured
2782 } else {
2783 HarnessAuthState::Unknown
2784 };
2785 runtime = HarnessRuntimeState::Degraded;
2786 reason = Some(format!(
2787 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2788 ));
2789 repair = Some(
2790 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2791 .into(),
2792 );
2793 }
2794 }
2795 Err(error) => {
2796 reason = Some(error_message(error));
2797 }
2798 }
2799 } else if installed && configured {
2800 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2801 } else if installed && auth == HarnessAuthState::Required {
2802 reason = Some("Executable found, but no native authentication evidence is present.".into());
2803 repair = Some(format!(
2804 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2805 descriptor.id.as_str()
2806 ));
2807 } else if installed {
2808 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2809 repair = Some(format!(
2810 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2811 launch
2812 .map(|launch| launch.program.as_str())
2813 .unwrap_or("the harness")
2814 ));
2815 }
2816
2817 let effective_capabilities = if installed {
2818 descriptor.runtime.capabilities.clone()
2819 } else {
2820 unavailable_capabilities()
2821 };
2822 let running = probe_running_instance(descriptor.id.as_str());
2823 LocalHarness {
2824 gateway: gateway_health(
2825 descriptor.id.as_str(),
2826 installed,
2827 running.as_ref(),
2828 version.as_deref(),
2829 ),
2830 id: descriptor.id,
2831 display_name: descriptor.display_name,
2832 supported: true,
2833 installed,
2834 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2835 version,
2836 auth,
2837 runtime,
2838 protocol: descriptor.runtime.protocol,
2839 capabilities: descriptor.runtime.capabilities,
2840 effective_capabilities,
2841 sessions: HarnessSessionCounts { global, workspace },
2842 running,
2843 reason,
2844 repair,
2845 }
2846}
2847
2848async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2849 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2850 loop {
2851 let now = tokio::time::Instant::now();
2852 if now >= deadline {
2853 return Ok(());
2854 }
2855 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2856 Err(_) => return Ok(()),
2857 Ok(Ok(Some(event))) => {
2858 if let Some(message) = handshake_event_failure(&event) {
2859 return Err(truncate_text(&message, 500));
2860 }
2861 }
2862 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2863 Ok(Err(error)) => return Err(error.to_string()),
2864 }
2865 }
2866}
2867
2868fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2869 let detail = event
2870 .payload
2871 .get("message")
2872 .or_else(|| event.payload.get("line"))
2873 .and_then(Value::as_str)
2874 .unwrap_or(event.kind.as_str());
2875 match event.kind.as_str() {
2876 "transport_closed" => Some("runtime transport closed during startup".into()),
2877 "transport_error" => Some(format!("runtime transport error: {detail}")),
2878 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2879 _ => None,
2884 }
2885}
2886
2887fn indexed_claude_window(
2888 locator: &SessionLocator,
2889 options: &SessionLoadOptions,
2890) -> std::result::Result<Option<Value>, ServiceError> {
2891 use supercode_interchange::session::ClaudeReadIndex;
2892 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2895 || options.include_subagents != Some(false)
2896 {
2897 return Ok(None);
2898 }
2899 let crate::StorageLocator::File { path } = &locator.storage else {
2900 return Ok(None);
2901 };
2902 if !ClaudeReadIndex::supports(path)
2903 .map_err(|error| ServiceError::Operation(error.to_string()))?
2904 {
2905 return Ok(None);
2906 }
2907 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2908 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2909 let total = index.len();
2910 let (offset, end) = projected_message_window(total, options);
2911 let session = index
2912 .read_messages(offset..end)
2913 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2914 let summary = index
2915 .read_summary()
2916 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2917 let selected_options = SessionLoadOptions {
2918 message_offset: None,
2919 message_limit: None,
2920 message_tail: None,
2921 ..options.clone()
2922 };
2923 let mut selected = projected_session_json(&session, &selected_options);
2924 selected["raw_record_count"] = json!(index.raw_record_count());
2925 Ok(Some(json!({
2926 "session": selected,
2927 "summary": projected_session_summary(&summary, options),
2928 "window": {
2929 "has_more": offset > 0 || end < total, "has_newer": end < total,
2930 "has_older": offset > 0, "newer_items": index.item_count(end..total),
2931 "offset": offset, "older_items": index.item_count(0..offset),
2932 "returned": end - offset, "total_messages": total,
2933 }
2934 })))
2935}
2936
2937fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
2938 let total_messages = session.messages.len();
2939 let (offset, end) = projected_message_window(total_messages, options);
2940 json!({
2941 "session": projected_session_json(session, options),
2942 "summary": projected_session_summary(session, options),
2943 "window": {
2944 "has_more": offset > 0 || end < total_messages,
2945 "has_newer": end < total_messages,
2946 "has_older": offset > 0,
2947 "newer_items": normalized_item_count(&session.messages[end..]),
2948 "offset": offset,
2949 "older_items": normalized_item_count(&session.messages[..offset]),
2950 "returned": end.saturating_sub(offset),
2951 "total_messages": total_messages,
2952 }
2953 })
2954}
2955
2956fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
2957 messages
2958 .iter()
2959 .map(|message| {
2960 let conversation = usize::from(
2961 matches!(message.role, Role::Assistant | Role::User)
2962 && message_has_content(message),
2963 );
2964 let tool_result =
2965 usize::from(message.role == Role::Tool && message_has_content(message));
2966 conversation + tool_result + message.tool_calls().len()
2967 })
2968 .sum()
2969}
2970
2971fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
2972 let mut conversational = session.messages.iter().filter(|message| {
2973 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
2974 });
2975 let first_message = conversational.clone().next();
2976 let last_message = conversational.next_back();
2977 let mut assistant = session
2978 .messages
2979 .iter()
2980 .filter(|message| message.role == Role::Assistant && message_has_content(message));
2981 let first_assistant_message = assistant.clone().next();
2982 let last_assistant_message = assistant.next_back();
2983 let end_of_turn = session
2984 .messages
2985 .iter()
2986 .rev()
2987 .find(|message| message.role != Role::System)
2988 .is_some_and(|message| {
2989 message.role == Role::Assistant
2990 && message_has_content(message)
2991 && message.tool_calls().is_empty()
2992 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
2994 });
2995 let project = |message: Option<&crate::ChatMessage>| {
2996 message.map(|message| project_inline_media(message_json(message), options))
2997 };
2998 json!({
2999 "end_of_turn": end_of_turn,
3000 "first_assistant_message": project(first_assistant_message),
3001 "first_message": project(first_message),
3002 "last_assistant_message": project(last_assistant_message),
3003 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3004 "last_message": project(last_message),
3005 })
3006}
3007
3008fn message_has_content(message: &crate::ChatMessage) -> bool {
3009 message
3010 .content
3011 .as_deref()
3012 .is_some_and(|content| !content.trim().is_empty())
3013 || message
3014 .content_parts
3015 .as_ref()
3016 .is_some_and(|parts| !parts.is_empty())
3017}
3018
3019fn message_text(message: &crate::ChatMessage) -> String {
3020 if let Some(content) = &message.content {
3021 return content.clone();
3022 }
3023 message
3024 .content_parts
3025 .as_ref()
3026 .into_iter()
3027 .flatten()
3028 .filter_map(|part| part.get("text").and_then(Value::as_str))
3029 .collect::<Vec<_>>()
3030 .join("\n")
3031}
3032
3033fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3034 let (offset, end) = projected_message_window(session.messages.len(), options);
3035 let messages = session.messages[offset..end]
3036 .iter()
3037 .map(|message| project_inline_media(message_json(message), options))
3038 .collect::<Vec<_>>();
3039 let subagents = if options.include_subagents.unwrap_or(true) {
3040 let subagent_options = SessionLoadOptions {
3045 message_limit: None,
3046 message_offset: None,
3047 message_tail: None,
3048 ..options.clone()
3049 };
3050 session
3051 .subagents
3052 .iter()
3053 .map(|subagent| projected_session_json(subagent, &subagent_options))
3054 .collect::<Vec<_>>()
3055 } else {
3056 Vec::new()
3057 };
3058 json!({
3059 "source": match session.meta.source {
3060 SessionSource::ClaudeCode => "claude_code",
3061 SessionSource::Codex => "codex",
3062 SessionSource::Gemini => "gemini",
3063 SessionSource::Goose => "goose",
3064 SessionSource::Grok => "grok",
3065 SessionSource::Native => "native",
3066 SessionSource::OpenClaw => "openclaw",
3067 SessionSource::Hermes => "hermes",
3068 SessionSource::OpenCode => "opencode",
3069 SessionSource::Pi => "pi",
3070 },
3071 "session_id": session.meta.session_id,
3072 "ended_at": session.meta.ended_at,
3073 "end_reason": session.meta.end_reason,
3074 "model": session.meta.model,
3075 "cwd": session.meta.cwd,
3076 "system_prompt": session.meta.system_prompt,
3077 "agent_id": session.meta.agent_id,
3078 "parent_tool_use_id": session.meta.parent_tool_use_id,
3079 "lineage": session.meta.lineage,
3080 "messages": messages,
3081 "subagents": subagents,
3082 "raw_record_count": session.raw.len(),
3083 "parse_error_lines": session.parse_error_lines,
3084 })
3085}
3086
3087fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3088 if let Some(tail) = options.message_tail {
3089 return (total.saturating_sub(tail), total);
3090 }
3091 let offset = options.message_offset.unwrap_or(0).min(total);
3092 let end = options
3093 .message_limit
3094 .map(|limit| offset.saturating_add(limit).min(total))
3095 .unwrap_or(total);
3096 (offset, end)
3097}
3098
3099fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3100 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3101 return message;
3102 };
3103 for part in parts {
3104 let Some(url) = part
3105 .get("image_url")
3106 .and_then(|image| image.get("url"))
3107 .and_then(Value::as_str)
3108 else {
3109 continue;
3110 };
3111 let Some(rest) = url.strip_prefix("data:") else {
3112 continue;
3113 };
3114 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3115 continue;
3116 };
3117 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3118 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3119 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3120 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3121 || options
3122 .max_inline_media_bytes
3123 .is_some_and(|limit| decoded_bytes > limit);
3124 if should_elide {
3125 *part = json!({
3126 "type": "media_reference",
3127 "media_type": media_type,
3128 "encoding": "base64",
3129 "encoded_bytes": encoded.len(),
3130 "decoded_bytes": decoded_bytes,
3131 "omitted": true,
3132 });
3133 }
3134 }
3135 message
3136}
3137
3138#[derive(Deserialize)]
3139struct LocatorParams {
3140 locator: SessionLocator,
3141 #[serde(default)]
3154 fidelity: Option<Fidelity>,
3155 #[serde(default)]
3158 view: Option<SessionReadView>,
3159}
3160
3161#[derive(Deserialize)]
3162struct SessionReadView {
3163 #[serde(default)]
3166 tail_messages: Option<usize>,
3167 #[serde(default)]
3170 include_subagents: bool,
3171 #[serde(default)]
3173 display_history: bool,
3174 #[serde(default)]
3177 max_message_chars: Option<usize>,
3178}
3179
3180impl LocatorParams {
3181 fn read_fidelity(&self) -> Fidelity {
3182 self.fidelity.unwrap_or(Fidelity::Semantic)
3183 }
3184
3185 fn include_subagents(&self) -> bool {
3186 self.view
3187 .as_ref()
3188 .map(|view| view.include_subagents)
3189 .unwrap_or(true)
3190 }
3191
3192 fn tail_messages(&self) -> Option<usize> {
3193 self.view
3194 .as_ref()
3195 .and_then(|view| view.tail_messages)
3196 .map(|limit| limit.clamp(1, 5_000))
3197 }
3198
3199 fn display_history(&self) -> bool {
3200 self.view.as_ref().is_some_and(|view| view.display_history)
3201 }
3202
3203 fn max_message_chars(&self) -> Option<usize> {
3204 self.view
3205 .as_ref()
3206 .and_then(|view| view.max_message_chars)
3207 .map(|limit| limit.clamp(256, 64_000))
3208 }
3209
3210 fn bound_session(&self, session: &mut Session) {
3211 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3212 }
3213}
3214
3215#[derive(Debug, Clone, Copy, Default, Deserialize)]
3216#[serde(rename_all = "snake_case")]
3217enum InlineMediaMode {
3218 #[default]
3219 Full,
3220 Metadata,
3221}
3222
3223#[derive(Debug, Clone, Default, Deserialize)]
3224#[serde(default)]
3225struct SessionLoadOptions {
3226 include_subagents: Option<bool>,
3227 inline_media: InlineMediaMode,
3228 max_inline_media_bytes: Option<usize>,
3229 message_limit: Option<usize>,
3230 message_offset: Option<usize>,
3231 message_tail: Option<usize>,
3232}
3233
3234impl SessionLoadOptions {
3235 fn validate(&self) -> std::result::Result<(), ServiceError> {
3236 if self.message_tail.is_some()
3237 && (self.message_limit.is_some() || self.message_offset.is_some())
3238 {
3239 return Err(ServiceError::InvalidParams(
3240 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3241 .into(),
3242 ));
3243 }
3244 Ok(())
3245 }
3246}
3247
3248#[derive(Deserialize)]
3249struct LoadSessionParams {
3250 #[serde(flatten)]
3251 read: LocatorParams,
3252 #[serde(default)]
3253 options: Option<SessionLoadOptions>,
3254}
3255
3256#[derive(Deserialize)]
3257struct UnfollowParams {
3258 subscription: String,
3259}
3260
3261#[derive(Debug, Deserialize)]
3262#[serde(deny_unknown_fields)]
3263struct IndexResizeParams {
3264 subscription: String,
3265 limit: usize,
3266}
3267
3268#[derive(Deserialize)]
3269struct ActivitySubscribeParams {
3270 locators: Vec<SessionLocator>,
3271 #[serde(default)]
3272 homes: crate::HarnessHomes,
3273}
3274
3275#[derive(Deserialize)]
3276struct MessageSessionParams {
3277 locator: SessionLocator,
3278 text: String,
3279 #[serde(default)]
3283 from_name: Option<String>,
3284 #[serde(default)]
3286 in_reply_to: Option<String>,
3287 #[serde(default)]
3289 notify_when_idle: bool,
3290 #[serde(default)]
3295 as_user: bool,
3296 #[serde(default)]
3299 homes: crate::HarnessHomes,
3300}
3301
3302#[derive(Deserialize)]
3303#[serde(deny_unknown_fields)]
3304struct InboxParams {
3305 #[serde(default)]
3307 from_name: Option<String>,
3308 #[serde(default)]
3310 address: Option<String>,
3311 #[serde(default)]
3313 all: bool,
3314}
3315
3316#[derive(Deserialize)]
3317#[serde(deny_unknown_fields)]
3318struct HarnessSettingsParams {
3319 harness: String,
3320}
3321
3322#[derive(Deserialize)]
3323#[serde(deny_unknown_fields)]
3324struct ConfigureHarnessParams {
3325 harness: String,
3326 #[serde(default)]
3327 changes: Vec<crate::HarnessSettingChange>,
3328 #[serde(default)]
3329 expected_revision: Option<String>,
3330}
3331
3332fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3333 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3334 Ok(report) => (
3335 serde_json::to_value(report).unwrap_or(Value::Null),
3336 Value::Null,
3337 ),
3338 Err(error) => (
3339 Value::Null,
3340 Value::String(format!(
3341 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3342 )),
3343 ),
3344 }
3345}
3346
3347async fn message_live_session(params: &MessageSessionParams) -> Value {
3361 use crate::mail_route::{Delivered, NoDoor, Refused};
3362 let (inbound_controls, inbound_controls_error) =
3363 claude_inbound_controls_or_error(¶ms.homes);
3364 let refused = |reason: &str, message: String| {
3365 json!({
3366 "delivered_to_bus": false,
3367 "refusal": {"reason": reason, "message": message},
3368 "inbound_controls": inbound_controls,
3369 "inbound_controls_error": inbound_controls_error,
3370 })
3371 };
3372 if params.text.trim().is_empty() {
3373 return refused(
3374 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3375 "refusing to deliver an empty message".into(),
3376 );
3377 }
3378 let sender = match operator_address(params.from_name.as_deref()) {
3379 Ok(sender) => sender,
3380 Err(message) => return refused("invalid_sender", message),
3381 };
3382 let receiver = match crate::mailbox::MailAddress::new(
3383 crate::mailbox::local_machine_name(),
3384 params.locator.harness.as_str(),
3385 ¶ms.locator.session_id,
3386 ) {
3387 Ok(receiver) => receiver,
3388 Err(error) => return refused("delivery_failed", error.to_string()),
3389 };
3390 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3392 && crate::runtime_mail::controlled_runtime(
3393 HarnessId::CLAUDE_CODE,
3394 ¶ms.locator.session_id,
3395 )
3396 .is_none()
3397 {
3398 if let Err(refusal) =
3399 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3400 {
3401 return refused(refusal.reason.as_str(), refusal.message);
3402 }
3403 }
3404 if params.as_user {
3405 return message_as_user(
3406 params,
3407 sender,
3408 receiver,
3409 inbound_controls,
3410 inbound_controls_error,
3411 )
3412 .await;
3413 }
3414 let door = match crate::mail_route::door_for(¶ms.homes, &receiver) {
3415 Ok(door) => door,
3416 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3417 return refused(
3418 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3419 format!(
3420 "no running `{}` session `{}` is reachable; its transcript is persisted only",
3421 params.locator.harness.as_str(),
3422 params.locator.session_id
3423 ),
3424 )
3425 }
3426 };
3427 let mut envelope = match crate::mailbox::Envelope::new(
3428 sender.clone(),
3429 format!("{}@{}", sender.session_id, sender.machine),
3430 crate::mailbox::MailKind::Peer,
3431 crate::mailbox::ReplyVia::Command,
3432 params.text.clone(),
3433 ) {
3434 Ok(envelope) => envelope,
3435 Err(error) => return refused("delivery_failed", error.to_string()),
3436 };
3437 envelope.in_reply_to = params.in_reply_to.clone();
3438 let how = match crate::mail_route::deliver(
3439 &envelope,
3440 &receiver,
3441 &door,
3442 true,
3443 params.notify_when_idle,
3444 )
3445 .await
3446 {
3447 Err(detail) => {
3448 return refused(
3449 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3450 detail,
3451 )
3452 }
3453 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3454 Ok(Err(Refused::TooLong(bytes))) => {
3455 return refused(
3456 "too_long",
3457 format!(
3458 "the message is {bytes} bytes; the limit is {}",
3459 crate::mail_route::MAX_RELAYED_BYTES
3460 ),
3461 )
3462 }
3463 Ok(Ok(delivered)) => match delivered {
3464 Delivered::Steered => "steered",
3465 Delivered::Started => "started",
3466 Delivered::Native { busy: true } => "next_tool_call",
3467 Delivered::Native { busy: false } => "started",
3468 Delivered::Hooked => "hook",
3469 Delivered::HookWoken => "woken",
3470 Delivered::Queued => "queued",
3471 Delivered::Stored => "stored",
3472 Delivered::Operator => "filed",
3473 },
3474 };
3475 json!({
3476 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3477 "message_id": envelope.id,
3478 "reply_to": sender.to_string(),
3479 "target": {
3480 "harness": params.locator.harness.as_str(),
3481 "session_id": params.locator.session_id,
3482 "name": match &door {
3483 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
3484 _ => None,
3485 },
3486 },
3487 "delivery": {"door": door.name(), "how": how},
3488 "inbound_controls": inbound_controls,
3489 "inbound_controls_error": inbound_controls_error,
3490 })
3491}
3492
3493async fn message_as_user(
3497 params: &MessageSessionParams,
3498 sender: crate::mailbox::MailAddress,
3499 receiver: crate::mailbox::MailAddress,
3500 inbound_controls: Value,
3501 inbound_controls_error: Value,
3502) -> Value {
3503 let envelope = match crate::mailbox::Envelope::new(
3504 sender.clone(),
3505 format!("{}@{}", sender.session_id, sender.machine),
3506 crate::mailbox::MailKind::User,
3507 crate::mailbox::ReplyVia::None,
3508 params.text.clone(),
3509 ) {
3510 Ok(envelope) => envelope,
3511 Err(error) => {
3512 return json!({
3513 "delivered_to_bus": false,
3514 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
3515 "inbound_controls": inbound_controls,
3516 "inbound_controls_error": inbound_controls_error,
3517 })
3518 }
3519 };
3520 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
3521 Ok(turn) => json!({
3522 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
3523 "message_id": envelope.id,
3524 "target": {
3525 "harness": params.locator.harness.as_str(),
3526 "session_id": params.locator.session_id,
3527 },
3528 "delivery": {
3529 "door": match turn {
3530 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
3531 _ => "pane",
3532 },
3533 "how": turn.as_str(),
3534 },
3535 "inbound_controls": inbound_controls,
3536 "inbound_controls_error": inbound_controls_error,
3537 }),
3538 Err(message) => json!({
3539 "delivered_to_bus": false,
3540 "refusal": {"reason": "no_user_door", "message": message},
3541 "inbound_controls": inbound_controls,
3542 "inbound_controls_error": inbound_controls_error,
3543 }),
3544 }
3545}
3546
3547#[derive(Deserialize)]
3548#[serde(deny_unknown_fields)]
3549struct ActivityUnderParams {
3550 pids: Vec<u32>,
3552 #[serde(default)]
3553 homes: crate::HarnessHomes,
3554}
3555
3556async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
3560 let params = decode::<ActivityUnderParams>(params)?;
3561 if params.pids.len() > 1024 {
3562 return Err(ServiceError::InvalidParams(
3563 "sessions.activity_under accepts at most 1024 pids".into(),
3564 ));
3565 }
3566 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
3567 .await
3568 .map_err(ServiceError::Sdk)?;
3569 Ok(json!({
3570 "activities": found
3571 .into_iter()
3572 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
3573 .collect::<Vec<_>>(),
3574 }))
3575}
3576
3577fn operator_address(
3580 name: Option<&str>,
3581) -> std::result::Result<crate::mailbox::MailAddress, String> {
3582 let name: String = name
3583 .unwrap_or("supercode")
3584 .trim()
3585 .chars()
3586 .map(|character| {
3587 if character.is_whitespace() || character == '@' {
3588 '-'
3589 } else {
3590 character
3591 }
3592 })
3593 .collect();
3594 if name.is_empty() {
3595 return Err("from_name must not be empty".into());
3596 }
3597 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
3598 .map_err(|error| error.to_string())
3599}
3600
3601fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
3605 let address = match (¶ms.address, ¶ms.from_name) {
3606 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
3607 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
3608 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
3609 (Some(_), Some(_)) => {
3610 return Err(ServiceError::InvalidParams(
3611 "sessions.inbox takes from_name or address, not both".into(),
3612 ))
3613 }
3614 };
3615 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
3616 let mailbox =
3617 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
3618 let claimed = mailbox.claim_unread().map_err(operation)?;
3619 let mut messages: Vec<Value> = Vec::new();
3620 if params.all {
3621 for stored in mailbox.list().map_err(operation)? {
3622 if stored.state == crate::mailbox::MailState::Read {
3623 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3624 }
3625 }
3626 }
3627 for stored in &claimed {
3628 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3629 }
3630 for stored in &claimed {
3631 mailbox.acknowledge(stored).map_err(operation)?;
3632 }
3633 Ok(json!({"address": address.to_string(), "messages": messages}))
3634}
3635
3636#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3641struct FollowedSource {
3642 harness: String,
3643 session_id: String,
3644 reported: Option<String>,
3645}
3646
3647#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3648struct ActivitySubscription {
3649 locators: Vec<SessionLocator>,
3650 homes: crate::HarnessHomes,
3651 reported: BTreeMap<(String, String), crate::SessionActivity>,
3652}
3653
3654fn live_descriptor_value(
3661 session: &SessionDescriptor,
3662 doors: &crate::mail_route::LiveSessions,
3663) -> std::result::Result<Value, ServiceError> {
3664 let mut value = serde_json::to_value(session)
3665 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3666 if value.get("title").is_none_or(Value::is_null) {
3668 if let Some(live) = doors.all().iter().find(|live| {
3669 live.address.harness == session.locator.harness.as_str()
3670 && live.address.session_id == session.locator.session_id
3671 }) {
3672 let name = live.name.split('@').next().unwrap_or(&live.name);
3673 if !name.is_empty() {
3674 value["title"] = json!(name);
3675 }
3676 }
3677 }
3678 if let Some(workspace) = &session.cwd {
3679 let source = LiveRuntimeSource {
3680 harness: session.locator.harness.as_str().to_string(),
3681 session_id: session.locator.session_id.clone(),
3682 workspace: workspace.clone(),
3683 };
3684 if let Some(endpoint) = discover_live_runtime(&source)
3685 .map_err(|error| ServiceError::Operation(error.to_string()))?
3686 {
3687 value["live_endpoint"] = json!(endpoint.as_str());
3688 }
3689 }
3690 if let Some(door) = doors.door(
3694 session.locator.harness.as_str(),
3695 &session.locator.session_id,
3696 ) {
3697 value["delivery"] = json!(door);
3698 }
3699 let waiting = doors.all().iter().any(|live| {
3701 live.address.harness == session.locator.harness.as_str()
3702 && live.address.session_id == session.locator.session_id
3703 && live.status == "waiting"
3704 });
3705 if waiting && session.locator.harness.as_str() == "claude-code" {
3706 if let Some(request) = crate::mail_route::pending_request(session.locator.storage.path()) {
3707 value["pending_request"] = request;
3708 }
3709 }
3710 Ok(value)
3711}
3712
3713fn live_index_changes(
3714 changes: Vec<crate::session_index::SessionIndexChange>,
3715 homes: &HarnessHomes,
3716) -> std::result::Result<Vec<Value>, ServiceError> {
3717 use crate::session_index::SessionIndexChange;
3718 let doors = crate::mail_route::LiveSessions::read(homes);
3719 changes
3720 .into_iter()
3721 .map(|change| match change {
3722 SessionIndexChange::Added { descriptor } => Ok(json!({
3723 "kind": "added",
3724 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3725 })),
3726 SessionIndexChange::Updated { descriptor } => Ok(json!({
3727 "kind": "updated",
3728 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3729 })),
3730 SessionIndexChange::Removed { key } => Ok(json!({
3731 "kind": "removed",
3732 "key": key,
3733 })),
3734 })
3735 .collect()
3736}
3737
3738fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
3739 use crate::{SessionPresence, SessionTurnState};
3740 match (activity.presence, activity.turn) {
3741 (SessionPresence::Persisted, _) => None,
3742 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
3743 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
3744 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
3746 (SessionPresence::Running, SessionTurnState::Unknown)
3750 if activity.evidence.native_state.is_none() =>
3751 {
3752 None
3753 }
3754 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
3755 }
3756}
3757
3758#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
3759#[serde(rename_all = "kebab-case")]
3760enum TransferFormat {
3761 ClaudeCode,
3762 Codex,
3763 #[serde(rename = "opencode", alias = "open-code")]
3764 OpenCode,
3765 Pi,
3766 Grok,
3767 Gemini,
3768 Goose,
3769 Hermes,
3773}
3774
3775impl TransferFormat {
3776 fn id(self) -> &'static str {
3777 match self {
3778 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
3779 Self::Codex => HarnessId::CODEX,
3780 Self::OpenCode => HarnessId::OPENCODE,
3781 Self::Pi => HarnessId::PI,
3782 Self::Grok => HarnessId::GROK,
3783 Self::Gemini => HarnessId::GEMINI,
3784 Self::Goose => HarnessId::GOOSE,
3785 Self::Hermes => HarnessId::HERMES,
3786 }
3787 }
3788}
3789
3790impl From<TransferFormat> for SessionFormat {
3791 fn from(value: TransferFormat) -> Self {
3792 match value {
3793 TransferFormat::ClaudeCode => Self::ClaudeCode,
3794 TransferFormat::Codex => Self::Codex,
3795 TransferFormat::OpenCode => Self::OpenCode,
3796 TransferFormat::Pi => Self::Pi,
3797 TransferFormat::Grok => Self::Grok,
3798 TransferFormat::Gemini => Self::Gemini,
3799 TransferFormat::Goose => Self::Goose,
3800 TransferFormat::Hermes => Self::Codex,
3802 }
3803 }
3804}
3805
3806#[derive(Deserialize)]
3807struct ImportSessionParams {
3808 source_harness: TransferFormat,
3809 content: String,
3810}
3811
3812#[derive(Deserialize)]
3813struct ExportSessionParams {
3814 locator: SessionLocator,
3815 target_harness: TransferFormat,
3816}
3817
3818#[derive(Deserialize)]
3819struct ReduceSessionParams {
3820 locator: SessionLocator,
3821 target_harness: TransferFormat,
3822 #[serde(default = "default_keep_last")]
3823 keep_last: usize,
3824}
3825
3826fn default_keep_last() -> usize {
3827 6
3828}
3829
3830#[derive(Deserialize)]
3831struct BranchSessionParams {
3832 locator: SessionLocator,
3833 #[serde(default)]
3834 target_harness: Option<TransferFormat>,
3835}
3836
3837#[derive(Deserialize)]
3838struct HandoffSessionParams {
3839 locator: SessionLocator,
3840 target_harness: TransferFormat,
3841 #[serde(default)]
3842 cwd: Option<PathBuf>,
3843}
3844
3845#[derive(Deserialize)]
3846struct MaterializeSessionParams {
3847 artifact: crate::native_materialize::MaterializeArtifact,
3848 cwd: PathBuf,
3849 #[serde(default)]
3851 homes: HarnessHomes,
3852}
3853
3854#[derive(Debug, Clone, Copy, Default, Deserialize)]
3857#[serde(rename_all = "snake_case")]
3858enum ResumePolicy {
3859 Default,
3860 #[default]
3861 Yolo,
3862}
3863
3864#[derive(Deserialize)]
3865struct ResumeInstructionsParams {
3866 locator: SessionLocator,
3867 #[serde(default)]
3868 cwd: Option<PathBuf>,
3869 #[serde(default)]
3870 policy: ResumePolicy,
3871}
3872
3873#[derive(Deserialize)]
3875struct WorkflowLoadParams {
3876 from: crate::workflow_doors::WorkflowHarness,
3877 home: PathBuf,
3878}
3879
3880#[derive(Deserialize)]
3883struct OrchestrationLoadParams {
3884 root: PathBuf,
3885 #[serde(default)]
3886 flavor: crate::orchestration_doors::HomeFlavor,
3887}
3888
3889#[derive(Deserialize)]
3892struct OrchestrationSaveParams {
3893 root: PathBuf,
3894 orchestration: crate::orchestration::Orchestration,
3895 #[serde(default)]
3896 vault: BTreeMap<String, String>,
3897}
3898
3899#[derive(Deserialize)]
3901struct OrchestrationCompileParams {
3902 from: crate::orchestration_doors::OrchestrationHarness,
3903 home: PathBuf,
3904}
3905
3906#[derive(Deserialize)]
3910struct OrchestrationDecompileParams {
3911 to: crate::orchestration_doors::OrchestrationHarness,
3912 orchestration: crate::orchestration::Orchestration,
3913 source: PathBuf,
3914 #[serde(default)]
3915 source_flavor: crate::orchestration_doors::SourceFlavor,
3916 dest: PathBuf,
3917 #[serde(default)]
3918 vault: BTreeMap<String, String>,
3919}
3920
3921#[derive(Deserialize)]
3924struct OrchestrationImportParams {
3925 from: crate::orchestration_doors::OrchestrationHarness,
3926 home: PathBuf,
3927 into: PathBuf,
3928}
3929
3930#[derive(Deserialize)]
3933struct OrchestrationExportParams {
3934 to: crate::orchestration_doors::OrchestrationHarness,
3935 root: PathBuf,
3936 dest: PathBuf,
3937}
3938
3939#[derive(Deserialize)]
3941struct JobsGetParams {
3942 harness: String,
3943 id: String,
3944 #[serde(default)]
3945 homes: crate::HarnessHomes,
3946}
3947
3948fn mutate_job(
3956 verb: crate::jobs_control::JobVerb,
3957 params: Value,
3958) -> std::result::Result<Value, ServiceError> {
3959 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
3960 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
3961 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
3962 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3963}
3964
3965fn mutate_skill(
3972 verb: crate::skills_control::SkillVerb,
3973 params: Value,
3974) -> std::result::Result<Value, ServiceError> {
3975 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
3976 if !crate::skills_control::supports_skill_control(&mutation.harness) {
3977 return Err(ServiceError::UnsupportedAction(format!(
3978 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
3979 mutation.harness,
3980 verb.as_str(),
3981 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
3982 )));
3983 }
3984 let outcome =
3985 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
3986 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3987}
3988
3989fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
3991 match error {
3992 crate::skills_control::SkillControlError::Unsupported(message) => {
3993 ServiceError::UnsupportedAction(message)
3994 }
3995 crate::skills_control::SkillControlError::Invalid(message) => {
3996 ServiceError::InvalidParams(message)
3997 }
3998 crate::skills_control::SkillControlError::Failed(message) => {
3999 ServiceError::Operation(message)
4000 }
4001 }
4002}
4003
4004fn mutate_profile(
4012 verb: crate::profiles_control::ProfileVerb,
4013 params: Value,
4014) -> std::result::Result<Value, ServiceError> {
4015 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4016 let outcome =
4017 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4018 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4019}
4020
4021fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4023 match error {
4024 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4025 ServiceError::UnsupportedAction(message)
4026 }
4027 crate::profiles_control::ProfileControlError::Invalid(message) => {
4028 ServiceError::InvalidParams(message)
4029 }
4030 crate::profiles_control::ProfileControlError::Failed(message) => {
4031 ServiceError::Operation(message)
4032 }
4033 }
4034}
4035
4036fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4040 match error {
4041 crate::jobs_control::JobControlError::Unsupported(message) => {
4042 ServiceError::UnsupportedAction(message)
4043 }
4044 crate::jobs_control::JobControlError::Invalid(message) => {
4045 ServiceError::InvalidParams(message)
4046 }
4047 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4048 }
4049}
4050
4051fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4056 match error {
4057 crate::SessionControlError::Unsupported(message) => {
4058 ServiceError::UnsupportedAction(message)
4059 }
4060 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4061 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4062 }
4063}
4064
4065fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4070 if crate::jobs::supports_jobs(harness) {
4071 return Ok(());
4072 }
4073 Err(ServiceError::UnsupportedAction(format!(
4074 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4075 crate::jobs::JOB_HARNESSES.join(", ")
4076 )))
4077}
4078
4079#[derive(Deserialize)]
4081struct RunsGetParams {
4082 harness: String,
4083 id: String,
4084 #[serde(default)]
4085 homes: crate::HarnessHomes,
4086}
4087
4088fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4093 if crate::runs::supports_runs(harness) {
4094 return Ok(());
4095 }
4096 Err(ServiceError::UnsupportedAction(format!(
4097 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4098 crate::runs::RUN_HARNESSES.join(", ")
4099 )))
4100}
4101
4102#[derive(Serialize)]
4103struct SessionArtifact {
4104 source_harness: HarnessId,
4105 target_harness: &'static str,
4106 session_id: Option<String>,
4107 content: String,
4108 suggested_filename: String,
4109 files: Vec<SessionArtifactFile>,
4110 fidelity: Fidelity,
4111 residue: Vec<String>,
4112}
4113
4114#[derive(Serialize)]
4115struct SessionArtifactFile {
4116 path: String,
4117 content: String,
4118 role: ArtifactFileRole,
4119}
4120
4121#[derive(Serialize)]
4122#[serde(rename_all = "snake_case")]
4123enum ArtifactFileRole {
4124 Primary,
4125 Subagent,
4126 Bundle,
4127 SourceRecovery,
4128}
4129
4130#[derive(Serialize)]
4131struct StructuredLaunch {
4132 cwd: PathBuf,
4133 program: String,
4134 arguments: Vec<String>,
4135 env: BTreeMap<String, String>,
4136}
4137
4138struct HandoffInstructions {
4139 launch: StructuredLaunch,
4140 materialize: Option<StructuredLaunch>,
4141 requires_materialization: bool,
4142 note: String,
4143}
4144
4145#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4146#[serde(rename_all = "snake_case")]
4147enum HarnessProbeLevel {
4148 #[default]
4149 Passive,
4150 Handshake,
4151}
4152
4153#[derive(Default, Deserialize)]
4154#[serde(default)]
4155struct HarnessInventoryParams {
4156 harness: Option<HarnessId>,
4157 harnesses: Vec<HarnessId>,
4158 workspace: Option<PathBuf>,
4159 probe: HarnessProbeLevel,
4160 include_sessions: bool,
4161 skip_versions: bool,
4163}
4164
4165#[derive(Deserialize)]
4166struct HarnessAuthenticationParams {
4167 harness: HarnessId,
4168}
4169
4170#[derive(Deserialize)]
4171struct BeginHarnessAuthenticationParams {
4172 harness: HarnessId,
4173 #[serde(default = "local_browser_authentication_environment")]
4174 environment: crate::HarnessAuthenticationEnvironment,
4175 #[serde(default)]
4176 method: Option<crate::HarnessAuthenticationMethodId>,
4177 #[serde(default)]
4178 cwd: Option<PathBuf>,
4179}
4180
4181fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4182 crate::HarnessAuthenticationEnvironment::LocalBrowser
4183}
4184
4185#[derive(Serialize)]
4186struct HarnessInventoryReport {
4187 probe: HarnessProbeLevel,
4188 workspace: Option<PathBuf>,
4189 harnesses: Vec<LocalHarness>,
4190}
4191
4192#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4193#[serde(rename_all = "snake_case")]
4194enum HarnessAuthState {
4195 Ready,
4196 Configured,
4197 Required,
4198 Unknown,
4199}
4200
4201#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4202#[serde(rename_all = "snake_case")]
4203enum HarnessRuntimeState {
4204 Ready,
4205 Degraded,
4206 Unavailable,
4207}
4208
4209#[derive(Serialize)]
4210struct HarnessSessionCounts {
4211 global: Option<usize>,
4212 workspace: Option<usize>,
4213}
4214
4215#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4226#[serde(rename_all = "snake_case")]
4227pub enum GatewayState {
4228 Up,
4229 Down,
4230 Unknown,
4231}
4232
4233#[derive(Debug, Clone, Serialize)]
4235pub struct GatewayHealth {
4236 pub state: GatewayState,
4237 #[serde(skip_serializing_if = "Option::is_none")]
4242 pub endpoint: Option<String>,
4243 #[serde(skip_serializing_if = "Option::is_none")]
4244 pub version: Option<String>,
4245 pub evidence: String,
4247 pub checked_at_ms: u64,
4248}
4249
4250fn openclaw_gateway_endpoint(home: &Path) -> String {
4254 let config_path = home.join(".openclaw/openclaw.json");
4255 let gateway = std::fs::read_to_string(&config_path)
4256 .ok()
4257 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4258 .and_then(|config| config.get("gateway").cloned());
4259 if let Some(url) = gateway
4260 .as_ref()
4261 .and_then(|gateway| gateway.get("url"))
4262 .and_then(serde_json::Value::as_str)
4263 {
4264 return url.to_string();
4265 }
4266 let port = gateway
4267 .as_ref()
4268 .and_then(|gateway| gateway.get("port"))
4269 .and_then(serde_json::Value::as_u64)
4270 .unwrap_or(18789);
4271 format!("ws://127.0.0.1:{port}")
4272}
4273
4274fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4281 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4282 let output = std::process::Command::new(&program)
4283 .args(["gateway", "status"])
4284 .stdin(std::process::Stdio::null())
4285 .output()
4286 .ok()?;
4287 let text = format!(
4288 "{}{}",
4289 String::from_utf8_lossy(&output.stdout),
4290 String::from_utf8_lossy(&output.stderr)
4291 );
4292 let verdict = text.lines().find_map(|line| {
4293 let l = line.trim();
4294 if l.contains("supervised by launchd (PID")
4295 || l.contains("supervised by systemd (PID")
4296 || l.contains("Gateway is running")
4297 || l.contains("process is running")
4298 {
4299 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4300 } else if l.contains("not running") || l.contains("not installed") {
4301 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4302 } else {
4303 None
4304 }
4305 });
4306 verdict
4307}
4308
4309fn gateway_health(
4310 id: &str,
4311 installed: bool,
4312 running: Option<&RunningInstance>,
4313 version: Option<&str>,
4314) -> GatewayHealth {
4315 let checked_at_ms = now_epoch_ms();
4316 let home = supercode_interchange::user_home()
4317 .map(std::path::PathBuf::into_os_string)
4318 .map(PathBuf::from);
4319 match id {
4320 HarnessId::HERMES | HarnessId::OPENCLAW => {
4321 let endpoint = (id == HarnessId::OPENCLAW)
4322 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4323 .flatten();
4324 let (state, evidence) = match running {
4325 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4326 None if !installed => (
4327 GatewayState::Unknown,
4328 format!("`{id}` is not installed; no gateway to probe"),
4329 ),
4330 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4331 Some((state, evidence)) => (state, evidence),
4334 None => (
4335 GatewayState::Down,
4336 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4337 ),
4338 },
4339 None => (
4340 GatewayState::Down,
4341 format!(
4342 "no TCP listener at {}",
4343 endpoint.as_deref().unwrap_or("the gateway endpoint")
4344 ),
4345 ),
4346 };
4347 GatewayHealth {
4348 state,
4349 endpoint,
4350 version: version.map(str::to_string),
4351 evidence,
4352 checked_at_ms,
4353 }
4354 }
4355 HarnessId::ORCHESTRATOR => {
4362 let root = crate::HarnessHomes::default().orchestrator;
4363 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4364 Some(lease) if lease.is_live() => (
4365 GatewayState::Up,
4366 format!(
4367 "`{}` names pid {} (started {}), which is live",
4368 crate::orchestrator::lock_path(&root).display(),
4369 lease.pid,
4370 lease.started_at
4371 ),
4372 ),
4373 Some(lease) => (
4374 GatewayState::Down,
4375 format!(
4376 "stale lease `{}`: pid {} is gone",
4377 crate::orchestrator::lock_path(&root).display(),
4378 lease.pid
4379 ),
4380 ),
4381 None => (
4382 GatewayState::Down,
4383 format!(
4384 "no lease at `{}`; `supercode orchestrator start` writes one",
4385 crate::orchestrator::lock_path(&root).display()
4386 ),
4387 ),
4388 };
4389 GatewayHealth {
4390 state,
4391 endpoint: None,
4392 version: version.map(str::to_string),
4393 evidence,
4394 checked_at_ms,
4395 }
4396 }
4397 _ => GatewayHealth {
4398 state: GatewayState::Unknown,
4399 endpoint: None,
4400 version: version.map(str::to_string),
4401 evidence: format!("`{id}` runs per session, not as a gateway"),
4402 checked_at_ms,
4403 },
4404 }
4405}
4406
4407#[derive(Debug, Clone, Serialize)]
4408struct RunningInstance {
4409 method: RunningInstanceMethod,
4411 evidence: String,
4413 checked_at_ms: u64,
4415}
4416
4417#[derive(Debug, Clone, Copy, Serialize)]
4418#[serde(rename_all = "snake_case")]
4419enum RunningInstanceMethod {
4420 GatewayConnect,
4423 StoreWalActivity,
4426}
4427
4428fn now_epoch_ms() -> u64 {
4429 std::time::SystemTime::now()
4430 .duration_since(std::time::UNIX_EPOCH)
4431 .map(|elapsed| elapsed.as_millis() as u64)
4432 .unwrap_or(0)
4433}
4434
4435fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4439 let config_path = home.join(".openclaw/openclaw.json");
4440 let text = std::fs::read_to_string(&config_path).ok();
4441 let gateway = text
4442 .as_deref()
4443 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4444 .and_then(|config| config.get("gateway").cloned());
4445 let address = gateway
4446 .as_ref()
4447 .and_then(|gateway| gateway.get("url"))
4448 .and_then(serde_json::Value::as_str)
4449 .and_then(|url| {
4450 url.split("://").nth(1).map(|rest| {
4451 rest.trim_end_matches('/')
4452 .split('/')
4453 .next()
4454 .unwrap_or(rest)
4455 .to_string()
4456 })
4457 })
4458 .unwrap_or_else(|| {
4459 let port = gateway
4460 .as_ref()
4461 .and_then(|gateway| gateway.get("port"))
4462 .and_then(serde_json::Value::as_u64)
4463 .unwrap_or(18789);
4464 format!("127.0.0.1:{port}")
4465 });
4466 let reachable = std::net::TcpStream::connect_timeout(
4467 &address.parse().ok()?,
4468 std::time::Duration::from_millis(400),
4469 )
4470 .is_ok();
4471 reachable.then(|| RunningInstance {
4472 method: RunningInstanceMethod::GatewayConnect,
4473 evidence: format!(
4474 "gateway endpoint {address} accepted a TCP connect (from {})",
4475 config_path.display()
4476 ),
4477 checked_at_ms: now_epoch_ms(),
4478 })
4479}
4480
4481fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4486 let wal = home.join(".hermes/state.db-wal");
4487 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4488 let age_ms = std::time::SystemTime::now()
4489 .duration_since(modified)
4490 .map(|age| age.as_millis() as u64)
4491 .unwrap_or(u64::MAX);
4492 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4493 method: RunningInstanceMethod::StoreWalActivity,
4494 evidence: format!(
4495 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4496 wal.display()
4497 ),
4498 checked_at_ms: now_epoch_ms(),
4499 })
4500}
4501
4502fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4504 let home = supercode_interchange::user_home()
4505 .map(std::path::PathBuf::into_os_string)
4506 .map(PathBuf::from)?;
4507 match id {
4508 HarnessId::OPENCLAW => probe_openclaw_running(&home),
4509 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4510 _ => None,
4511 }
4512}
4513
4514#[derive(Serialize)]
4515struct LocalHarness {
4516 id: HarnessId,
4517 display_name: String,
4518 supported: bool,
4519 installed: bool,
4520 executable: Option<String>,
4521 version: Option<String>,
4522 auth: HarnessAuthState,
4523 runtime: HarnessRuntimeState,
4524 protocol: String,
4525 capabilities: crate::RuntimeCapabilities,
4526 effective_capabilities: crate::RuntimeCapabilities,
4527 sessions: HarnessSessionCounts,
4528 #[serde(skip_serializing_if = "Option::is_none")]
4531 running: Option<RunningInstance>,
4532 gateway: GatewayHealth,
4534 reason: Option<String>,
4535 repair: Option<String>,
4536}
4537
4538#[derive(Clone, Deserialize)]
4539struct RuntimeBackendParams {
4540 harness: HarnessId,
4541 #[serde(default)]
4542 protocol: Option<String>,
4543 #[serde(default)]
4544 launch: Option<RuntimeLaunch>,
4545 #[serde(default)]
4546 base_url: Option<String>,
4547 #[serde(default)]
4548 policy: RuntimePolicy,
4549}
4550
4551#[derive(Debug, Clone, Copy, Default, Deserialize)]
4554#[serde(rename_all = "snake_case")]
4555enum RuntimePolicy {
4556 Default,
4557 #[default]
4558 Yolo,
4559}
4560
4561#[derive(Deserialize)]
4562struct RuntimeStartParams {
4563 #[serde(flatten)]
4564 backend: RuntimeBackendParams,
4565 cwd: PathBuf,
4566 #[serde(default)]
4569 mcp_servers: Vec<crate::McpServerLaunch>,
4570 #[serde(default)]
4572 approval_policy: Option<String>,
4573}
4574
4575#[derive(Deserialize)]
4576struct RuntimeAttachParams {
4577 #[serde(flatten)]
4578 backend: RuntimeBackendParams,
4579 runtime_id: String,
4580 #[serde(default)]
4581 cwd: Option<PathBuf>,
4582 #[serde(default)]
4585 mcp_servers: Vec<crate::McpServerLaunch>,
4586 #[serde(default)]
4588 approval_policy: Option<String>,
4589}
4590
4591#[derive(Deserialize)]
4592struct RuntimeConnectionParams {
4593 connection: String,
4594}
4595
4596#[derive(Deserialize)]
4597struct RuntimeInputParams {
4598 connection: String,
4599 text: String,
4600 #[serde(default)]
4601 image_urls: Vec<String>,
4602}
4603
4604const MAX_RUNTIME_IMAGES: usize = 4;
4605const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
4606const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
4607
4608fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
4609 if image_urls.len() > MAX_RUNTIME_IMAGES {
4610 return Err(ServiceError::InvalidParams(format!(
4611 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
4612 )));
4613 }
4614 let mut total = 0usize;
4615 for url in &image_urls {
4616 if !(url.starts_with("data:image/")
4617 || url.starts_with("https://")
4618 || url.starts_with("http://"))
4619 {
4620 return Err(ServiceError::InvalidParams(
4621 "runtime images must be image data URLs or HTTP(S) URLs".into(),
4622 ));
4623 }
4624 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
4625 return Err(ServiceError::InvalidParams(format!(
4626 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
4627 )));
4628 }
4629 total = total.saturating_add(url.len());
4630 }
4631 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
4632 return Err(ServiceError::InvalidParams(format!(
4633 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
4634 )));
4635 }
4636 Ok(image_urls)
4637}
4638
4639#[derive(Deserialize)]
4640struct RuntimeRespondParams {
4641 connection: String,
4642 request_id: Value,
4643 response: Value,
4644}
4645
4646fn default_reduction_store_root() -> PathBuf {
4647 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
4648 return PathBuf::from(root).join("sessions");
4649 }
4650 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
4651 return PathBuf::from(home).join(".supercode").join("sessions");
4652 }
4653 PathBuf::from(".supercode").join("sessions")
4654}
4655
4656fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
4657 let mut output = String::new();
4658 for message in messages {
4659 output.push_str(
4660 &serde_json::to_string(message)
4661 .map_err(|error| ServiceError::Operation(error.to_string()))?,
4662 );
4663 output.push('\n');
4664 }
4665 Ok(output)
4666}
4667
4668fn parse_messages_jsonl(
4669 content: &str,
4670) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
4671 content
4672 .lines()
4673 .enumerate()
4674 .filter(|(_, line)| !line.trim().is_empty())
4675 .map(|(index, line)| {
4676 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
4677 ServiceError::Operation(format!(
4678 "reduced transcript line {} is invalid: {error}",
4679 index + 1
4680 ))
4681 })
4682 })
4683 .collect()
4684}
4685
4686fn reduced_bootstrap_prompt(
4687 source: &SessionLocator,
4688 target: TransferFormat,
4689 view_jsonl: &str,
4690 sidecar_path: &Path,
4691 reduction_log_path: &Path,
4692) -> String {
4693 format!(
4694 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
4695 \n\
4696 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\
4697 \n\
4698 <supercode-reduced-session source-session=\"{source_id}\">\n\
4699 {view_jsonl}\
4700 </supercode-reduced-session>\n\
4701 \n\
4702 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
4703 source_harness = source.harness.as_str(),
4704 target_harness = target.id(),
4705 sidecar = sidecar_path.display(),
4706 log = reduction_log_path.display(),
4707 source_id = source.session_id,
4708 )
4709}
4710
4711fn session_artifact(
4712 locator: &SessionLocator,
4713 session: &Session,
4714 target: TransferFormat,
4715) -> std::result::Result<SessionArtifact, ServiceError> {
4716 session_artifact_with_id(locator, session, target, None)
4717}
4718
4719fn session_artifact_with_id(
4720 locator: &SessionLocator,
4721 session: &Session,
4722 target: TransferFormat,
4723 target_session_id: Option<&str>,
4724) -> std::result::Result<SessionArtifact, ServiceError> {
4725 let format: SessionFormat = target.into();
4726 let diagonal = format.source() == session.meta.source;
4727 crate::residue_store::store_segments(session);
4728 let has_appended_turns = session
4729 .imported_message_count
4730 .is_some_and(|imported| imported < session.messages.len());
4731 let mut restoration = None;
4732 let content = if let Some(id) = target_session_id {
4733 if diagonal && format != SessionFormat::OpenCode {
4734 session
4735 .to_jsonl_spliced(format, Some(id))
4736 .map_err(operation)?
4737 } else {
4738 let mut rewritten = session.clone();
4739 rewritten.meta.session_id = Some(id.to_string());
4740 rewritten.to_jsonl(format).map_err(operation)?
4741 }
4742 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
4743 session.raw_verbatim()
4744 } else if diagonal {
4745 session.to_jsonl_spliced(format, None).map_err(operation)?
4746 } else {
4747 match session
4750 .restore_residue(format, crate::residue_store::lookup)
4751 .map_err(operation)?
4752 {
4753 Some((content, report)) => {
4754 restoration = Some(report);
4755 content
4756 }
4757 None => session.to_jsonl(format).map_err(operation)?,
4758 }
4759 };
4760 let stem = sanitize_filename(
4761 target_session_id
4762 .or(session.meta.session_id.as_deref())
4763 .unwrap_or(&locator.session_id),
4764 );
4765 let suggested_filename = if diagonal && target == TransferFormat::Grok {
4766 "chat_history.jsonl".to_string()
4767 } else if target == TransferFormat::Goose {
4768 format!("{stem}.goose.json")
4769 } else {
4770 format!("{stem}.{}.jsonl", target.id())
4771 };
4772 let mut files = vec![SessionArtifactFile {
4773 path: suggested_filename.clone(),
4774 content: content.clone(),
4775 role: ArtifactFileRole::Primary,
4776 }];
4777 if target == TransferFormat::ClaudeCode {
4778 let bundle_stem = Path::new(&suggested_filename)
4779 .file_stem()
4780 .and_then(|stem| stem.to_str())
4781 .unwrap_or(&stem);
4782 let mut child_paths = BTreeSet::new();
4783 for (index, subagent) in session.subagents.iter().enumerate() {
4784 let agent_id = subagent
4785 .meta
4786 .agent_id
4787 .as_deref()
4788 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
4789 .map(sanitize_filename)
4790 .filter(|id| !id.is_empty())
4791 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4792 let child_has_appended_turns = subagent
4793 .imported_message_count
4794 .is_some_and(|imported| imported < subagent.messages.len());
4795 let child_content = if target_session_id.is_none()
4796 && subagent.meta.source == SessionSource::ClaudeCode
4797 && subagent.raw_is_verbatim
4798 && !child_has_appended_turns
4799 {
4800 subagent.raw_verbatim()
4801 } else if subagent.meta.source == SessionSource::ClaudeCode {
4802 subagent
4803 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
4804 .map_err(operation)?
4805 } else {
4806 let mut child = subagent.clone();
4807 if let Some(id) = target_session_id {
4808 child.meta.session_id = Some(id.to_string());
4809 }
4810 child
4811 .to_jsonl(SessionFormat::ClaudeCode)
4812 .map_err(operation)?
4813 };
4814 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
4815 if !child_paths.insert(path.clone()) {
4816 return Err(ServiceError::Operation(format!(
4817 "Claude subagent ids collide at artifact path `{path}`"
4818 )));
4819 }
4820 files.push(SessionArtifactFile {
4821 path,
4822 content: child_content,
4823 role: ArtifactFileRole::Subagent,
4824 });
4825 }
4826 }
4827 if diagonal && target == TransferFormat::Grok {
4828 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
4829 }
4830 if !diagonal || !session.raw_is_verbatim {
4831 files.push(SessionArtifactFile {
4832 path: "recovery/source.supercode.jsonl".into(),
4833 content: session.to_native_jsonl(),
4834 role: ArtifactFileRole::SourceRecovery,
4835 });
4836 for (index, subagent) in session.subagents.iter().enumerate() {
4837 let id = subagent
4838 .meta
4839 .agent_id
4840 .as_deref()
4841 .map(sanitize_filename)
4842 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4843 files.push(SessionArtifactFile {
4844 path: format!("recovery/subagents/{id}.supercode.jsonl"),
4845 content: subagent.to_native_jsonl(),
4846 role: ArtifactFileRole::SourceRecovery,
4847 });
4848 }
4849 }
4850 if !diagonal && session.meta.source == SessionSource::Grok {
4851 append_grok_bundle_files(
4852 locator,
4853 "recovery/grok/",
4854 ArtifactFileRole::SourceRecovery,
4855 &mut files,
4856 )?;
4857 }
4858 let (fidelity, residue) = if diagonal
4859 && target_session_id.is_none()
4860 && session.raw_is_verbatim
4861 && !has_appended_turns
4862 {
4863 (Fidelity::ByteLossless, Vec::new())
4864 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
4865 (
4866 Fidelity::ValueLossless,
4867 vec![if target_session_id.is_some() {
4868 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
4869 } else {
4870 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
4871 }],
4872 )
4873 } else {
4874 match restoration {
4875 Some(report) if report.rendered_messages == 0 => (
4876 Fidelity::ByteLossless,
4877 vec![format!(
4878 "restored verbatim from this conversation's {} source records in the residue store",
4879 target.id()
4880 )],
4881 ),
4882 Some(report) => (
4883 Fidelity::Semantic,
4884 vec![format!(
4885 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
4886 report.restored_messages,
4887 report.restored_messages + report.rendered_messages,
4888 report.rendered_messages,
4889 target.id()
4890 )],
4891 ),
4892 None => (
4893 Fidelity::Semantic,
4894 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
4895 ),
4896 }
4897 };
4898 Ok(SessionArtifact {
4899 source_harness: locator.harness.clone(),
4900 target_harness: target.id(),
4901 session_id: target_session_id
4902 .map(str::to_string)
4903 .or_else(|| session.meta.session_id.clone()),
4904 content,
4905 suggested_filename,
4906 files,
4907 fidelity,
4908 residue,
4909 })
4910}
4911
4912fn append_grok_bundle_files(
4913 locator: &SessionLocator,
4914 prefix: &str,
4915 role: ArtifactFileRole,
4916 files: &mut Vec<SessionArtifactFile>,
4917) -> std::result::Result<(), ServiceError> {
4918 let primary = locator.storage.path();
4919 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
4920 return Err(ServiceError::Operation(format!(
4921 "Grok bundle locator must name chat_history.jsonl, got {}",
4922 primary.display()
4923 )));
4924 }
4925 let parent = primary.parent().ok_or_else(|| {
4926 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
4927 })?;
4928 for name in ["summary.json", "updates.jsonl"] {
4929 let path = parent.join(name);
4930 let metadata = match std::fs::symlink_metadata(&path) {
4931 Ok(metadata) => metadata,
4932 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
4933 Err(error) => return Err(ServiceError::Operation(error.to_string())),
4934 };
4935 if metadata.file_type().is_symlink() || !metadata.is_file() {
4936 return Err(ServiceError::Operation(format!(
4937 "refusing non-regular Grok bundle member {}",
4938 path.display()
4939 )));
4940 }
4941 let content = std::fs::read_to_string(&path).map_err(|error| {
4942 ServiceError::Operation(format!(
4943 "Grok bundle member {} is not representable as UTF-8: {error}",
4944 path.display()
4945 ))
4946 })?;
4947 files.push(SessionArtifactFile {
4948 path: format!("{prefix}{name}"),
4949 content,
4950 role: match role {
4951 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
4952 _ => ArtifactFileRole::SourceRecovery,
4953 },
4954 });
4955 }
4956 Ok(())
4957}
4958
4959fn handoff_artifact(
4960 locator: &SessionLocator,
4961 session: &Session,
4962 target: TransferFormat,
4963) -> std::result::Result<SessionArtifact, ServiceError> {
4964 let target_session_id = target_session_id(target);
4965 session_artifact_with_id(locator, session, target, Some(&target_session_id))
4966}
4967
4968fn target_session_id(target: TransferFormat) -> String {
4969 let uuid = generated_session_id();
4970 match target {
4971 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
4972 TransferFormat::ClaudeCode
4973 | TransferFormat::Codex
4974 | TransferFormat::Pi
4975 | TransferFormat::Grok
4976 | TransferFormat::Gemini
4977 | TransferFormat::Goose
4978 | TransferFormat::Hermes => uuid,
4979 }
4980}
4981
4982fn sanitize_filename(value: &str) -> String {
4983 let value = value
4984 .chars()
4985 .map(|character| {
4986 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
4987 character
4988 } else {
4989 '-'
4990 }
4991 })
4992 .collect::<String>();
4993 let value = value.trim_matches('-');
4994 if value.is_empty() {
4995 "session".into()
4996 } else {
4997 value.chars().take(100).collect()
4998 }
4999}
5000
5001fn handoff_instructions(
5002 target: TransferFormat,
5003 session_id: &str,
5004 cwd: &Path,
5005) -> HandoffInstructions {
5006 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5007 cwd: cwd.to_path_buf(),
5008 program: program.into(),
5009 arguments,
5010 env: BTreeMap::new(),
5011 };
5012 match target {
5013 TransferFormat::ClaudeCode => HandoffInstructions {
5014 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5015 materialize: None,
5016 requires_materialization: true,
5017 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(),
5018 },
5019 TransferFormat::Hermes => HandoffInstructions {
5020 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5021 materialize: None,
5022 requires_materialization: true,
5023 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(),
5024 },
5025 TransferFormat::Codex => HandoffInstructions {
5026 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5027 materialize: None,
5028 requires_materialization: true,
5029 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5030 },
5031 TransferFormat::OpenCode => HandoffInstructions {
5032 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5033 materialize: Some(launch(
5034 "opencode",
5035 vec!["import".into(), "{artifact_path}".into()],
5036 )),
5037 requires_materialization: true,
5038 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5039 },
5040 TransferFormat::Pi => HandoffInstructions {
5041 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5042 materialize: None,
5043 requires_materialization: true,
5044 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5045 },
5046 TransferFormat::Grok => HandoffInstructions {
5047 launch: launch(
5048 "grok",
5049 vec!["--resume".into(), "{materialized_session_id}".into()],
5050 ),
5051 materialize: None,
5052 requires_materialization: true,
5053 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(),
5054 },
5055 TransferFormat::Gemini => HandoffInstructions {
5056 launch: launch(
5057 "gemini",
5058 vec!["--session-file".into(), "{artifact_path}".into()],
5059 ),
5060 materialize: None,
5061 requires_materialization: true,
5062 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(),
5063 },
5064 TransferFormat::Goose => HandoffInstructions {
5065 launch: launch(
5066 "goose",
5067 vec![
5068 "session".into(),
5069 "--resume".into(),
5070 "--session-id".into(),
5071 "{imported_session_id}".into(),
5072 ],
5073 ),
5074 materialize: Some(launch(
5075 "goose",
5076 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5077 )),
5078 requires_materialization: true,
5079 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(),
5080 },
5081 }
5082}
5083
5084fn resume_launch(
5085 harness: &str,
5086 session_id: &str,
5087 cwd: &Path,
5088 policy: ResumePolicy,
5089) -> std::result::Result<StructuredLaunch, ServiceError> {
5090 let mut arguments = Vec::new();
5091 let program = match harness {
5092 HarnessId::GROK => {
5093 if matches!(policy, ResumePolicy::Yolo) {
5094 if crate::support::self_sandbox_supported() {
5095 arguments.extend(["--sandbox".into(), "workspace".into()]);
5096 }
5097 arguments.push("--always-approve".into());
5098 }
5099 arguments.extend(["--resume".into(), session_id.into()]);
5100 "grok"
5101 }
5102 HarnessId::CODEX => {
5103 arguments.extend(crate::startup_prompts::startup_arguments(
5104 harness,
5105 Some(cwd),
5106 &[],
5107 matches!(policy, ResumePolicy::Yolo),
5108 ));
5109 arguments.extend(["resume".into(), session_id.into()]);
5110 "codex"
5111 }
5112 HarnessId::CLAUDE_CODE => {
5113 arguments.extend(crate::startup_prompts::startup_arguments(
5114 harness,
5115 Some(cwd),
5116 &[],
5117 matches!(policy, ResumePolicy::Yolo),
5118 ));
5119 arguments.extend(["--resume".into(), session_id.into()]);
5120 "claude"
5121 }
5122 HarnessId::GEMINI => {
5123 arguments.extend(crate::startup_prompts::startup_arguments(
5124 harness,
5125 Some(cwd),
5126 &[],
5127 matches!(policy, ResumePolicy::Yolo),
5128 ));
5129 arguments.extend(["--resume".into(), session_id.into()]);
5130 "gemini"
5131 }
5132 HarnessId::GOOSE => {
5133 arguments.extend([
5134 "session".into(),
5135 "--resume".into(),
5136 "--session-id".into(),
5137 session_id.into(),
5138 ]);
5139 "goose"
5140 }
5141 HarnessId::PI => {
5142 arguments.extend(crate::startup_prompts::startup_arguments(
5143 harness,
5144 Some(cwd),
5145 &[],
5146 matches!(policy, ResumePolicy::Yolo),
5147 ));
5148 arguments.extend(["--session".into(), session_id.into()]);
5149 "pi"
5150 }
5151 HarnessId::OPENCODE => {
5152 arguments.extend(["--session".into(), session_id.into()]);
5153 "opencode"
5154 }
5155 HarnessId::SUPERCODE => {
5156 if matches!(policy, ResumePolicy::Yolo) {
5157 arguments.push("--dangerous".into());
5158 }
5159 arguments.extend(["resume".into(), session_id.into()]);
5160 "supercode"
5161 }
5162 other => {
5163 return Err(ServiceError::InvalidParams(format!(
5164 "no structured resume launch is registered for harness `{other}`"
5165 )))
5166 }
5167 };
5168 Ok(StructuredLaunch {
5169 cwd: cwd.to_path_buf(),
5170 env: if program == "grok" {
5171 crate::support::grok_home_env()
5172 } else {
5173 BTreeMap::new()
5174 },
5175 program: program.into(),
5176 arguments,
5177 })
5178}
5179
5180fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5186 let digest = blake3::hash(address.as_bytes()).to_hex();
5187 let path = std::env::temp_dir().join(format!(
5188 "supercode-openclaw-gateway-token-{}",
5189 &digest.as_str()[..16]
5190 ));
5191 #[cfg(unix)]
5192 {
5193 use std::io::Write;
5194 use std::os::unix::fs::OpenOptionsExt;
5195 let mut file = std::fs::OpenOptions::new()
5196 .write(true)
5197 .create(true)
5198 .truncate(true)
5199 .mode(0o600)
5200 .open(&path)?;
5201 file.write_all(secret.as_bytes())?;
5202 }
5203 #[cfg(not(unix))]
5204 std::fs::write(&path, secret)?;
5205 Ok(path)
5206}
5207
5208fn open_connect_descriptor(
5214 descriptor: &crate::HarnessSupportDescriptor,
5215 home: &Path,
5216) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5217 let Some(connect) = &descriptor.runtime.connect_launch else {
5218 return Err(ServiceError::InvalidParams(format!(
5219 "harness `{}` has no registered connect-mode launch",
5220 descriptor.id.as_str()
5221 )));
5222 };
5223 let resolved = connect
5224 .resolve(home)
5225 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5226 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5227 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5228 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5229 if let Some(token) = resolved.auth {
5230 backend = backend.with_bearer(token);
5231 }
5232 Ok(Box::new(backend))
5233 }
5234 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5235 let mut env = BTreeMap::new();
5246 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5247 if let Some(token) = resolved.auth {
5248 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5249 .map_err(|error| {
5250 ServiceError::UnsupportedAction(format!(
5251 "could not stage the gateway credential for the bridge: {error}"
5252 ))
5253 })?;
5254 arguments.push("--token-file".into());
5255 arguments.push(token_path.to_string_lossy().into_owned());
5256 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5257 }
5258 let program = descriptor
5263 .runtime
5264 .default_launch
5265 .as_ref()
5266 .map(|launch| launch.program.clone())
5267 .unwrap_or_else(|| "openclaw".into());
5268 let launch = RuntimeLaunch {
5269 program,
5270 arguments,
5271 env,
5272 };
5273 Ok(Box::new(
5274 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5275 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5276 ))
5277 }
5278 _ => Err(ServiceError::UnsupportedAction(format!(
5279 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5280 descriptor.id.as_str(),
5281 connect.protocol
5282 ))),
5283 }
5284}
5285
5286fn registry_connect_descriptor(
5289 params: &RuntimeBackendParams,
5290) -> Option<crate::HarnessSupportDescriptor> {
5291 if params.launch.is_some() || params.base_url.is_some() {
5292 return None;
5293 }
5294 harness_support_registry()
5295 .harnesses
5296 .into_iter()
5297 .find(|descriptor| descriptor.id == params.harness)
5298 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5299}
5300
5301fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5302 supercode_interchange::user_home()
5303 .map(std::path::PathBuf::into_os_string)
5304 .map(PathBuf::from)
5305 .ok_or_else(|| {
5306 ServiceError::UnsupportedAction(
5307 "connect-mode launches need HOME to locate the harness config".into(),
5308 )
5309 })
5310}
5311
5312pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5315 "harness.v1.runtimes.start",
5316 "harness.v1.runtimes.resume",
5317 "harness.v1.runtimes.attach",
5318 "harness.v1.runtimes.attach_existing",
5319];
5320
5321pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5326
5327pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5333
5334pub const DETACHED_METHODS: &[&str] = &[
5340 "harness.v1.harnesses.list",
5341 "harness.v1.harnesses.probe",
5342 "harness.v1.sessions.message",
5343 "harness.v1.sessions.new",
5344 "harness.v1.sessions.reset",
5345 "harness.v1.sessions.archive",
5346 "harness.v1.sessions.delete",
5347];
5348
5349pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5355
5356pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5364
5365async fn within_control_deadline<F: std::future::Future>(
5368 method: &str,
5369 call: F,
5370) -> std::result::Result<F::Output, ServiceError> {
5371 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5372 .await
5373 .map_err(|_| {
5374 ServiceError::Operation(format!(
5375 "`{method}` gave up after {}s: the runtime did not answer",
5376 RUNTIME_CONTROL_DEADLINE.as_secs()
5377 ))
5378 })
5379}
5380
5381pub struct RuntimeOpen {
5385 id: Value,
5386 method: String,
5387 params: Value,
5388}
5389
5390impl RuntimeOpen {
5391 pub async fn open(self) -> OpenedRuntime {
5395 let Self { id, method, params } = self;
5396 let outcome = open_runtime(&method, params).await;
5397 OpenedRuntime { id, outcome }
5398 }
5399}
5400
5401pub struct OpenedRuntime {
5404 id: Value,
5405 outcome: std::result::Result<OpenRuntime, ServiceError>,
5406}
5407
5408pub struct DetachedCall {
5413 id: Value,
5414 method: String,
5415 work: std::result::Result<Work, ServiceError>,
5416}
5417
5418impl DetachedCall {
5419 pub async fn run(self) -> DetachedAnswer {
5422 let Self { id, method, work } = self;
5423 match work {
5424 Ok(Work::Runtime(work)) => {
5429 let (result, returned) = work.run().await;
5430 DetachedAnswer {
5431 response: service_response(id, result),
5432 returned,
5433 }
5434 }
5435 Ok(Work::Free(work)) => {
5436 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5437 Ok(result) => result,
5438 Err(_) => Err(ServiceError::Operation(format!(
5439 "`{method}` gave up after {}s: the harness it waits on did not answer",
5440 DETACHED_CALL_DEADLINE.as_secs()
5441 ))),
5442 };
5443 DetachedAnswer {
5444 response: service_response(id, result),
5445 returned: None,
5446 }
5447 }
5448 Err(error) => DetachedAnswer {
5449 response: service_response(id, Err(error)),
5450 returned: None,
5451 },
5452 }
5453 }
5454}
5455
5456pub struct DetachedAnswer {
5460 response: Value,
5461 returned: Option<ReturnedRuntime>,
5462}
5463
5464impl DetachedAnswer {
5465 pub fn into_response(self) -> Value {
5468 self.response
5469 }
5470}
5471
5472pub struct ReturnedRuntime {
5475 connection: String,
5476 runtime: Box<dyn RuntimeConnection>,
5477}
5478
5479enum Work {
5482 Free(DetachedWork),
5483 Runtime(RuntimeWork),
5484}
5485
5486enum DetachedWork {
5489 Inventory(InventoryWork),
5493 Message(MessageSessionParams),
5495 SessionMutation {
5498 verb: crate::SessionVerb,
5499 mutation: crate::SessionMutation,
5500 },
5501}
5502
5503impl DetachedWork {
5504 async fn run(self) -> std::result::Result<Value, ServiceError> {
5505 match self {
5506 Self::Inventory(work) => run_inventory(work).await,
5507 Self::Message(params) => Ok(message_live_session(¶ms).await),
5508 Self::SessionMutation { verb, mutation } => {
5509 let outcome = run_session_mutation(verb, &mutation).await?;
5510 serde_json::to_value(outcome)
5511 .map_err(|error| ServiceError::Operation(error.to_string()))
5512 }
5513 }
5514 }
5515}
5516
5517enum RuntimeWork {
5519 Close {
5521 runtime: Box<dyn RuntimeConnection>,
5522 process_group: Option<u32>,
5523 },
5524 LiveCommand {
5527 connection: String,
5528 runtime: Box<dyn RuntimeConnection>,
5529 verb: crate::SessionVerb,
5530 mutation: crate::SessionMutation,
5531 command: &'static str,
5532 session: String,
5533 },
5534}
5535
5536type RuntimeWorkAnswer = (
5539 std::result::Result<Value, ServiceError>,
5540 Option<ReturnedRuntime>,
5541);
5542
5543impl RuntimeWork {
5544 async fn run(self) -> RuntimeWorkAnswer {
5545 match self {
5546 Self::Close {
5547 runtime,
5548 process_group,
5549 } => (close_runtime(runtime, process_group).await, None),
5550 Self::LiveCommand {
5551 connection,
5552 mut runtime,
5553 verb,
5554 mutation,
5555 command,
5556 session,
5557 } => {
5558 let result =
5559 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5560 (
5561 result,
5562 Some(ReturnedRuntime {
5563 connection,
5564 runtime,
5565 }),
5566 )
5567 }
5568 }
5569 }
5570}
5571
5572async fn close_runtime(
5575 mut runtime: Box<dyn RuntimeConnection>,
5576 process_group: Option<u32>,
5577) -> std::result::Result<Value, ServiceError> {
5578 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
5579 Ok(result) => {
5580 result.map_err(operation)?;
5581 Ok(json!({"closed": true}))
5582 }
5583 Err(deadline) => {
5584 let killed = kill_runtime_process_group(process_group);
5589 drop(runtime);
5590 Ok(json!({
5591 "closed": true,
5592 "killed": killed,
5593 "detail": error_message(deadline),
5594 }))
5595 }
5596 }
5597}
5598
5599fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
5602 mutation
5603 .session
5604 .clone()
5605 .filter(|value| !value.trim().is_empty())
5606 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
5607}
5608
5609async fn type_live_command(
5613 runtime: &mut dyn RuntimeConnection,
5614 verb: crate::SessionVerb,
5615 mutation: &crate::SessionMutation,
5616 command: &str,
5617 session: String,
5618) -> std::result::Result<Value, ServiceError> {
5619 within_control_deadline(
5620 &format!("sessions.{}", verb.as_str()),
5621 runtime.send_input(RuntimeInput {
5622 text: command.to_string(),
5623 image_urls: Vec::new(),
5624 }),
5625 )
5626 .await?
5627 .map_err(operation)?;
5628 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
5629 .map_err(session_control_error)?;
5630 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
5631}
5632
5633enum OpenRuntime {
5636 Hosted {
5639 runtime: Box<dyn RuntimeConnection>,
5640 capabilities: crate::RuntimeCapabilities,
5641 workspace: PathBuf,
5642 fresh: bool,
5644 },
5645 Joined { runtime: Box<dyn RuntimeConnection> },
5648}
5649
5650async fn open_runtime(
5655 method: &str,
5656 params: Value,
5657) -> std::result::Result<OpenRuntime, ServiceError> {
5658 match tokio::time::timeout(
5659 RUNTIME_OPEN_DEADLINE,
5660 open_runtime_unbounded(method, params),
5661 )
5662 .await
5663 {
5664 Ok(result) => result,
5665 Err(_) => Err(ServiceError::Operation(format!(
5666 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
5667 RUNTIME_OPEN_DEADLINE.as_secs()
5668 ))),
5669 }
5670}
5671
5672async fn open_runtime_unbounded(
5673 method: &str,
5674 params: Value,
5675) -> std::result::Result<OpenRuntime, ServiceError> {
5676 match method {
5677 "harness.v1.runtimes.start" => {
5678 let params = decode::<RuntimeStartParams>(params)?;
5679 let backend = runtime_backend(¶ms.backend)?;
5680 let capabilities = backend.capabilities();
5681 let workspace = params.cwd.clone();
5682 let runtime = backend
5683 .start(RuntimeStartRequest {
5684 cwd: params.cwd,
5685 launch: runtime_launch(¶ms.backend),
5686 mcp_servers: params.mcp_servers,
5687 approval_policy: params.approval_policy,
5688 })
5689 .await
5690 .map_err(operation)?;
5691 Ok(OpenRuntime::Hosted {
5692 runtime,
5693 capabilities,
5694 workspace,
5695 fresh: true,
5696 })
5697 }
5698 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
5699 let params = decode::<RuntimeAttachParams>(params)?;
5700 let backend = runtime_backend(¶ms.backend)?;
5701 let capabilities = backend.capabilities();
5702 let workspace = params
5703 .cwd
5704 .clone()
5705 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
5706 let runtime = backend
5707 .attach(RuntimeAttachRequest {
5708 runtime_id: params.runtime_id,
5709 cwd: params.cwd,
5710 launch: runtime_launch(¶ms.backend),
5711 mcp_servers: params.mcp_servers,
5712 approval_policy: params.approval_policy,
5713 })
5714 .await
5715 .map_err(operation)?;
5716 Ok(OpenRuntime::Hosted {
5717 runtime,
5718 capabilities,
5719 workspace,
5720 fresh: false,
5721 })
5722 }
5723 "harness.v1.runtimes.attach_existing" => {
5724 let params = decode::<RuntimeAttachParams>(params)?;
5725 let backend: Box<dyn RuntimeBackend> = match params
5726 .backend
5727 .base_url
5728 .as_deref()
5729 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
5730 {
5731 Some(endpoint) => {
5732 #[cfg(not(feature = "adapter-api"))]
5733 {
5734 let _ = endpoint;
5735 return Err(ServiceError::UnsupportedAction(
5736 "live HTTP attachment adapter is not compiled".into(),
5737 ));
5738 }
5739 #[cfg(feature = "adapter-api")]
5740 {
5741 let workspace = params.cwd.clone().ok_or_else(|| {
5742 ServiceError::InvalidParams(
5743 "Volter Harness live attach requires the project cwd".into(),
5744 )
5745 })?;
5746 let source = LiveRuntimeSource {
5747 harness: params.backend.harness.as_str().to_string(),
5748 session_id: params.runtime_id.clone(),
5749 workspace,
5750 };
5751 let receipt = resolve_live_runtime(&endpoint, &source)
5752 .map_err(|error| ServiceError::Operation(error.to_string()))?;
5753 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
5754 }
5755 }
5756 None => runtime_backend(¶ms.backend)?,
5757 };
5758 let capabilities = backend.capabilities();
5759 if !capabilities.attach_existing_process {
5760 return Err(ServiceError::Operation(format!(
5761 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
5762 backend.harness().as_str()
5763 )));
5764 }
5765 let runtime = backend
5766 .attach_existing(RuntimeAttachRequest {
5767 runtime_id: params.runtime_id,
5768 cwd: params.cwd,
5769 launch: runtime_launch(¶ms.backend),
5770 mcp_servers: params.mcp_servers,
5771 approval_policy: params.approval_policy,
5772 })
5773 .await
5774 .map_err(operation)?;
5775 Ok(OpenRuntime::Joined { runtime })
5776 }
5777 _ => Err(ServiceError::MethodNotFound),
5778 }
5779}
5780
5781fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
5783 match result {
5784 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
5785 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
5786 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
5787 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
5788 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
5789 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
5790 }
5791}
5792
5793fn runtime_backend(
5794 params: &RuntimeBackendParams,
5795) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5796 if let Some(descriptor) = registry_connect_descriptor(params) {
5797 return open_connect_descriptor(&descriptor, &service_home()?);
5798 }
5799 if params.protocol.as_deref() == Some("acp") {
5800 let launch = params
5801 .launch
5802 .clone()
5803 .or_else(|| {
5804 harness_support_registry()
5805 .harnesses
5806 .into_iter()
5807 .find(|harness| harness.id == params.harness)
5808 .filter(|harness| {
5809 harness.runtime.implementation == ImplementationKind::GenericProtocol
5810 && harness.runtime.protocol.starts_with("acp")
5811 })
5812 .and_then(|harness| harness.runtime.default_launch)
5813 })
5814 .ok_or_else(|| {
5815 ServiceError::InvalidParams(
5816 "an ACP runtime requires `launch` unless the harness has a registered default"
5817 .into(),
5818 )
5819 })?;
5820 let resume_session = harness_support_registry()
5821 .harnesses
5822 .into_iter()
5823 .find(|harness| harness.id == params.harness)
5824 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
5825 return Ok(Box::new(
5826 AcpRuntimeBackend::new(params.harness.clone(), launch)
5827 .with_resume_support(resume_session),
5828 ));
5829 }
5830 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
5831 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
5832 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
5833 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
5834 HarnessId::OPENCODE => match ¶ms.base_url {
5835 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
5836 None => Box::new(OpenCodeRuntimeBackend::new()),
5837 },
5838 harness => {
5839 let descriptor = harness_support_registry()
5840 .harnesses
5841 .into_iter()
5842 .find(|descriptor| descriptor.id.as_str() == harness)
5843 .filter(|descriptor| {
5844 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
5845 && descriptor.runtime.protocol.starts_with("acp")
5846 });
5847 let Some(descriptor) = descriptor else {
5848 return Err(ServiceError::InvalidParams(format!(
5849 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
5850 )));
5851 };
5852 let resume = descriptor.runtime.capabilities.resume_session;
5853 Box::new(
5854 AcpRuntimeBackend::new(
5855 descriptor.id,
5856 descriptor
5857 .runtime
5858 .default_launch
5859 .expect("generic ACP registry entry includes its launch"),
5860 )
5861 .with_resume_support(resume),
5862 )
5863 }
5864 };
5865 Ok(backend)
5866}
5867
5868fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
5869 if let Some(launch) = ¶ms.launch {
5870 return Some(launch.clone());
5871 }
5872 if !matches!(params.policy, RuntimePolicy::Yolo) {
5873 return None;
5874 }
5875 let launch = match params.harness.as_str() {
5876 HarnessId::GROK => RuntimeLaunch {
5877 program: "grok".into(),
5878 arguments: {
5879 let mut arguments: Vec<String> = Vec::new();
5880 if crate::support::self_sandbox_supported() {
5881 arguments.extend(["--sandbox".into(), "workspace".into()]);
5882 }
5883 arguments.extend([
5884 "--always-approve".into(),
5885 "agent".into(),
5886 "--no-leader".into(),
5887 "stdio".into(),
5888 ]);
5889 arguments
5890 },
5891 env: crate::support::grok_env(),
5892 },
5893 HarnessId::CODEX => RuntimeLaunch {
5894 program: "codex".into(),
5895 arguments: vec![
5896 "--dangerously-bypass-approvals-and-sandbox".into(),
5897 "--dangerously-bypass-hook-trust".into(),
5898 "app-server".into(),
5899 ],
5900 env: BTreeMap::new(),
5901 },
5902 HarnessId::CLAUDE_CODE => RuntimeLaunch {
5903 program: "claude".into(),
5904 arguments: vec![
5905 "--dangerously-skip-permissions".into(),
5906 "--print".into(),
5907 "--input-format".into(),
5908 "stream-json".into(),
5909 "--output-format".into(),
5910 "stream-json".into(),
5911 "--verbose".into(),
5912 ],
5913 env: BTreeMap::new(),
5914 },
5915 HarnessId::PI => RuntimeLaunch {
5916 program: "pi".into(),
5917 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
5918 env: BTreeMap::new(),
5919 },
5920 HarnessId::OPENCODE => RuntimeLaunch {
5921 program: "opencode".into(),
5922 arguments: vec!["serve".into()],
5923 env: BTreeMap::new(),
5924 },
5925 HarnessId::GEMINI => RuntimeLaunch {
5926 program: "gemini".into(),
5927 arguments: vec!["--acp".into(), "--yolo".into()],
5928 env: BTreeMap::new(),
5929 },
5930 HarnessId::GOOSE => RuntimeLaunch {
5931 program: "goose".into(),
5932 arguments: vec!["acp".into()],
5933 env: BTreeMap::new(),
5934 },
5935 HarnessId::SUPERCODE => RuntimeLaunch {
5936 program: "supercode".into(),
5937 arguments: vec!["acp".into(), "--dangerous".into()],
5938 env: BTreeMap::new(),
5939 },
5940 _ => return None,
5941 };
5942 Some(launch)
5943}
5944
5945struct IsolatedProbeHome {
5951 launch: RuntimeLaunch,
5952 root: PathBuf,
5953}
5954
5955impl IsolatedProbeHome {
5956 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
5957 let root = std::env::temp_dir().join(format!(
5958 "supercode-harness-probe-{harness}-{}",
5959 generated_session_id()
5960 ));
5961 std::fs::create_dir_all(&root)?;
5962 set_private_dir_permissions(&root)?;
5963
5964 if let Some(source_home) = supercode_interchange::user_home()
5965 .map(std::path::PathBuf::into_os_string)
5966 .map(PathBuf::from)
5967 {
5968 for relative in probe_auth_files(harness) {
5969 copy_probe_file(&source_home, &root, relative)?;
5970 }
5971 }
5972 if harness == HarnessId::SUPERCODE {
5977 let config_home = crate::agent::global_instructions_dir();
5978 for file in ["config.toml", "credentials.toml"] {
5979 copy_probe_path(
5980 &config_home.join(file),
5981 &root.join(".config/supercode").join(file),
5982 )?;
5983 }
5984 }
5985 configure_isolated_probe_auth(harness, &root)?;
5986
5987 let root_text = root.to_string_lossy().into_owned();
5988 for (key, value) in [
5989 ("HOME", root_text.clone()),
5990 (
5991 "XDG_CACHE_HOME",
5992 root.join(".cache").to_string_lossy().into_owned(),
5993 ),
5994 (
5995 "XDG_CONFIG_HOME",
5996 root.join(".config").to_string_lossy().into_owned(),
5997 ),
5998 (
5999 "XDG_DATA_HOME",
6000 root.join(".local/share").to_string_lossy().into_owned(),
6001 ),
6002 ] {
6003 launch.env.insert(key.into(), value);
6004 }
6005 let scoped = match harness {
6006 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6007 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6008 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6009 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6010 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6011 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6012 _ => None,
6013 };
6014 if let Some((key, value)) = scoped {
6015 launch
6016 .env
6017 .insert(key.into(), value.to_string_lossy().into_owned());
6018 }
6019 Ok(Self { launch, root })
6020 }
6021
6022 fn cleanup(&self) -> std::io::Result<()> {
6023 match std::fs::remove_dir_all(&self.root) {
6024 Ok(()) => Ok(()),
6025 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6026 Err(error) => Err(error),
6027 }
6028 }
6029}
6030
6031impl Drop for IsolatedProbeHome {
6032 fn drop(&mut self) {
6033 let _ = self.cleanup();
6034 }
6035}
6036
6037fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6038 match harness {
6039 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6040 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6044 HarnessId::CODEX => &[".codex/auth.json"],
6045 HarnessId::GEMINI => &[
6046 ".gemini/google_accounts.json",
6047 ".gemini/oauth_creds.json",
6048 ".gemini/settings.json",
6049 ],
6050 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6051 HarnessId::OPENCODE => &[
6052 ".config/opencode/auth.json",
6053 ".local/share/opencode/auth.json",
6054 ],
6055 HarnessId::PI => &[".pi/agent/auth.json"],
6056 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6061 _ => &[],
6062 }
6063}
6064
6065fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6066 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6067}
6068
6069fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6070 if !source.is_file() {
6071 return Ok(());
6072 }
6073 if let Some(parent) = destination.parent() {
6074 std::fs::create_dir_all(parent)?;
6075 set_private_dir_permissions(parent)?;
6076 }
6077 std::fs::copy(source, destination)?;
6078 set_private_file_permissions(destination)
6079}
6080
6081fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6082 if harness != HarnessId::GEMINI {
6083 return Ok(());
6084 }
6085 let oauth = probe_home.join(".gemini/oauth_creds.json");
6086 if !oauth.is_file() {
6087 return Ok(());
6088 }
6089 let settings_path = probe_home.join(".gemini/settings.json");
6090 let mut settings = std::fs::read_to_string(&settings_path)
6091 .ok()
6092 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6093 .unwrap_or_else(|| json!({}));
6094 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6095 std::fs::write(
6096 &settings_path,
6097 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6098 )?;
6099 set_private_file_permissions(&settings_path)
6100}
6101
6102#[cfg(unix)]
6103fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6104 use std::os::unix::fs::PermissionsExt;
6105 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6106}
6107
6108#[cfg(not(unix))]
6109fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6110 Ok(())
6111}
6112
6113#[cfg(unix)]
6114fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6115 use std::os::unix::fs::PermissionsExt;
6116 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6117}
6118
6119#[cfg(not(unix))]
6120fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6121 Ok(())
6122}
6123
6124fn find_executable(program: &str) -> Option<PathBuf> {
6125 let candidate = PathBuf::from(program);
6126 if candidate.components().count() > 1 {
6127 return candidate.is_file().then_some(candidate);
6128 }
6129 let path = std::env::var_os("PATH")?;
6130 for directory in std::env::split_paths(&path) {
6131 let candidate = directory.join(program);
6132 if candidate.is_file() {
6133 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6134 }
6135 #[cfg(windows)]
6136 {
6137 for extension in ["exe", "cmd", "bat"] {
6138 let candidate = directory.join(format!("{program}.{extension}"));
6139 if candidate.is_file() {
6140 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6141 }
6142 }
6143 }
6144 }
6145 None
6146}
6147
6148async fn executable_version(executable: &Path) -> Option<String> {
6149 let mut command = tokio::process::Command::new(executable);
6150 command
6151 .arg("--version")
6152 .stdin(std::process::Stdio::null())
6153 .stdout(std::process::Stdio::piped())
6154 .stderr(std::process::Stdio::piped())
6155 .kill_on_drop(true);
6156 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6157 .await
6158 .ok()?
6159 .ok()?;
6160 let stdout = String::from_utf8_lossy(&output.stdout);
6161 let stderr = String::from_utf8_lossy(&output.stderr);
6162 stdout
6163 .lines()
6164 .chain(stderr.lines())
6165 .map(str::trim)
6166 .find(|line| !line.is_empty())
6167 .map(|line| truncate_text(line, 200))
6168}
6169
6170pub(crate) fn auth_evidence(harness: &str) -> bool {
6171 let env_names: &[&str] = match harness {
6172 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6173 HarnessId::CODEX => &["OPENAI_API_KEY"],
6174 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6175 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6176 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6177 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6178 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6179 _ => &[],
6180 };
6181 if env_names
6182 .iter()
6183 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6184 {
6185 return true;
6186 }
6187 let Some(home) = supercode_interchange::user_home()
6188 .map(std::path::PathBuf::into_os_string)
6189 .map(PathBuf::from)
6190 else {
6191 return false;
6192 };
6193 let files: Vec<PathBuf> = match harness {
6194 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6195 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6196 HarnessId::OPENCODE => vec![
6197 home.join(".local/share/opencode/auth.json"),
6198 home.join(".config/opencode/auth.json"),
6199 ],
6200 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6201 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6202 HarnessId::GEMINI => vec![
6203 home.join(".gemini/oauth_creds.json"),
6204 home.join(".gemini/google_accounts.json"),
6205 ],
6206 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6207 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6208 _ => Vec::new(),
6209 };
6210 if files.into_iter().any(|path| {
6211 std::fs::metadata(path)
6212 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6213 .unwrap_or(false)
6214 }) {
6215 return true;
6216 }
6217 if harness == HarnessId::CLAUDE_CODE {
6224 return std::fs::read_to_string(home.join(".claude.json"))
6225 .map(|text| text.contains("\"oauthAccount\""))
6226 .unwrap_or(false);
6227 }
6228 false
6229}
6230
6231fn looks_like_auth_error(message: &str) -> bool {
6232 let message = message.to_ascii_lowercase();
6233 [
6234 "auth",
6235 "login",
6236 "sign in",
6237 "sign-in",
6238 "credential",
6239 "unauthorized",
6240 "forbidden",
6241 "token",
6242 ]
6243 .iter()
6244 .any(|needle| message.contains(needle))
6245}
6246
6247fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6248 crate::RuntimeCapabilities {
6249 start_session: false,
6250 resume_session: false,
6251 attach_existing_process: false,
6252 send_input: false,
6253 stream_events: false,
6254 interrupt: false,
6255 steer: false,
6256 respond_to_requests: false,
6257 }
6258}
6259
6260fn truncate_text(text: &str, max_chars: usize) -> String {
6261 let mut chars = text.chars();
6262 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6263 if chars.next().is_some() {
6264 format!("{truncated}…")
6265 } else {
6266 truncated
6267 }
6268}
6269
6270fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6277 match &handle.endpoint {
6278 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6279 crate::RuntimeEndpoint::Http { .. } => None,
6280 }
6281}
6282
6283fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6289 match process_group {
6290 #[cfg(unix)]
6291 Some(pid) => {
6292 crate::lsp::kill_process_group(pid);
6293 true
6294 }
6295 #[cfg(not(unix))]
6296 Some(_) => false,
6297 None => false,
6298 }
6299}
6300
6301fn error_message(error: ServiceError) -> String {
6302 match error {
6303 ServiceError::InvalidParams(message)
6304 | ServiceError::Operation(message)
6305 | ServiceError::UnsupportedAction(message) => message,
6306 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6307 ServiceError::Sdk(error) => error.to_string(),
6308 }
6309}
6310
6311#[derive(Debug)]
6312enum ServiceError {
6313 InvalidParams(String),
6314 MethodNotFound,
6315 UnsupportedAction(String),
6316 Operation(String),
6317 Sdk(SdkError),
6318}
6319
6320fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6321 match error {
6322 ServiceError::InvalidParams(message) => {
6323 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6324 }
6325 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6326 SdkError::unsupported(operation)
6327 }
6328 ServiceError::Operation(message) => {
6329 let code = if message.contains("already in progress") {
6330 SdkErrorCode::Busy
6331 } else if message.contains("not supported by this runtime") {
6332 SdkErrorCode::UnsupportedAction
6333 } else if message.contains("unknown runtime connection") {
6334 SdkErrorCode::NotFound
6335 } else {
6336 SdkErrorCode::Execution
6337 };
6338 SdkError::new(code, operation, message)
6339 }
6340 ServiceError::Sdk(error) => error,
6341 }
6342}
6343
6344fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6345 let error_code = error.code();
6346 let code = match error_code {
6347 SdkErrorCode::Unauthenticated => -32030,
6348 SdkErrorCode::Unauthorized => -32031,
6349 SdkErrorCode::ControllerRequired => -32032,
6350 SdkErrorCode::LeaseExpired => -32033,
6351 SdkErrorCode::InvalidArgument => -32602,
6352 SdkErrorCode::NotFound => -32004,
6353 SdkErrorCode::Busy => -32000,
6354 SdkErrorCode::UnsupportedAction => -32020,
6355 SdkErrorCode::Execution => -32002,
6356 SdkErrorCode::Transport => -32003,
6357 };
6358 json!({
6359 "jsonrpc": "2.0",
6360 "id": id,
6361 "error": {
6362 "code": code,
6363 "name": error_code,
6364 "operation": error.operation(),
6365 "message": error.to_string(),
6366 },
6367 })
6368}
6369
6370fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6371 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6372}
6373
6374fn operation(error: impl Into<crate::Error>) -> ServiceError {
6375 let error = error.into();
6376 match error {
6377 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6378 error => ServiceError::Operation(error.to_string()),
6379 }
6380}
6381
6382#[derive(Debug, Clone, Deserialize, Default)]
6386#[serde(default)]
6387struct MemoryRequest {
6388 harness: Option<String>,
6390 query: Option<String>,
6392 profile: Option<String>,
6394 session: Option<String>,
6396 full: bool,
6398 regex: bool,
6400 cwd: Option<std::path::PathBuf>,
6402 homes: crate::HarnessHomes,
6404}
6405
6406fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6409 let request = decode::<MemoryRequest>(params)?;
6410 let harness = request
6411 .harness
6412 .clone()
6413 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6414 let to_service = |error: crate::memory::MemoryError| match error {
6415 crate::memory::MemoryError::UnsupportedHarness { .. }
6416 | crate::memory::MemoryError::SessionNotScoped { .. } => {
6417 ServiceError::UnsupportedAction(error.to_string())
6418 }
6419 other => ServiceError::InvalidParams(other.to_string()),
6420 };
6421 match method {
6422 "harness.v1.memory.show" => {
6423 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6424 harness,
6425 profile: request.profile,
6426 session: request.session,
6427 full: request.full,
6428 cwd: request.cwd,
6429 homes: request.homes,
6430 })
6431 .map_err(to_service)?;
6432 Ok(json!({
6433 "schema": crate::memory::MEMORY_SCHEMA,
6434 "documents": documents,
6435 }))
6436 }
6437 "harness.v1.memory.search" => {
6438 let query = request
6439 .query
6440 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6441 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6442 harness,
6443 query,
6444 profile: request.profile,
6445 regex: request.regex,
6446 cwd: request.cwd,
6447 homes: request.homes,
6448 })
6449 .map_err(to_service)?;
6450 Ok(json!({
6451 "schema": crate::memory::MEMORY_SCHEMA,
6452 "matches": matches,
6453 }))
6454 }
6455 _ => Err(ServiceError::MethodNotFound),
6456 }
6457}
6458
6459#[derive(Debug, Clone, Deserialize)]
6463#[serde(default)]
6464struct ProfilesQuery {
6465 harness: Option<String>,
6467 name: Option<String>,
6469 homes: crate::HarnessHomes,
6471}
6472
6473impl Default for ProfilesQuery {
6474 fn default() -> Self {
6475 Self {
6476 harness: None,
6477 name: None,
6478 homes: crate::HarnessHomes::default(),
6479 }
6480 }
6481}
6482
6483fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6486 let query = decode::<ProfilesQuery>(params)?;
6487 let to_service = |error: crate::profiles::ProfileError| match error {
6488 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6489 ServiceError::UnsupportedAction(error.to_string())
6490 }
6491 crate::profiles::ProfileError::NotFound { .. } => {
6492 ServiceError::InvalidParams(error.to_string())
6493 }
6494 };
6495 match method {
6496 "harness.v1.profiles.list" => {
6497 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6498 .map_err(to_service)?;
6499 Ok(json!({
6500 "schema": crate::profiles::PROFILES_SCHEMA,
6501 "profiles": profiles,
6502 }))
6503 }
6504 "harness.v1.profiles.get" => {
6505 let harness = query
6506 .harness
6507 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6508 let name = query
6509 .name
6510 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6511 let profile =
6512 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6513 Ok(json!({
6514 "schema": crate::profiles::PROFILES_SCHEMA,
6515 "profile": profile,
6516 }))
6517 }
6518 _ => Err(ServiceError::MethodNotFound),
6519 }
6520}
6521
6522#[derive(Debug, Clone, Deserialize)]
6526#[serde(default)]
6527struct ChannelsQuery {
6528 harness: Option<String>,
6530 name: Option<String>,
6532 homes: crate::HarnessHomes,
6534}
6535
6536impl Default for ChannelsQuery {
6537 fn default() -> Self {
6538 Self {
6539 harness: None,
6540 name: None,
6541 homes: crate::HarnessHomes::default(),
6542 }
6543 }
6544}
6545
6546#[derive(Debug, Clone, Deserialize)]
6550#[serde(default)]
6551struct RoutesQuery {
6552 harness: Option<String>,
6553 profile: Option<String>,
6555 homes: crate::HarnessHomes,
6556}
6557
6558impl Default for RoutesQuery {
6559 fn default() -> Self {
6560 Self {
6561 harness: None,
6562 profile: None,
6563 homes: crate::HarnessHomes::default(),
6564 }
6565 }
6566}
6567
6568#[derive(Debug, Clone, Deserialize)]
6569#[serde(default)]
6570struct TriggersQuery {
6571 harness: Option<String>,
6572 homes: crate::HarnessHomes,
6573}
6574
6575impl Default for TriggersQuery {
6576 fn default() -> Self {
6577 Self {
6578 harness: None,
6579 homes: crate::HarnessHomes::default(),
6580 }
6581 }
6582}
6583
6584fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
6585 let query = decode::<TriggersQuery>(params)?;
6586 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
6587 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6588 Ok(json!({
6589 "schema": crate::triggers::TRIGGERS_SCHEMA,
6590 "triggers": triggers,
6591 }))
6592}
6593
6594fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
6595 let query = decode::<RoutesQuery>(params)?;
6596 let routes = crate::routes::list_routes(
6597 &query.homes,
6598 query.harness.as_deref(),
6599 query.profile.as_deref(),
6600 )
6601 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6602 Ok(json!({
6603 "schema": crate::routes::ROUTES_SCHEMA,
6604 "routes": routes,
6605 }))
6606}
6607
6608fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6609 let query = decode::<ChannelsQuery>(params)?;
6610 let to_service = |error: crate::channels::ChannelError| match error {
6611 crate::channels::ChannelError::UnsupportedHarness { .. } => {
6612 ServiceError::UnsupportedAction(error.to_string())
6613 }
6614 crate::channels::ChannelError::NotFound { .. } => {
6615 ServiceError::InvalidParams(error.to_string())
6616 }
6617 };
6618 match method {
6619 "harness.v1.channels.list" => {
6620 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
6621 .map_err(to_service)?;
6622 Ok(json!({
6623 "schema": crate::channels::CHANNELS_SCHEMA,
6624 "channels": channels,
6625 }))
6626 }
6627 "harness.v1.channels.status" => {
6628 let harness = query
6629 .harness
6630 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6631 let name = query
6632 .name
6633 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6634 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
6635 .map_err(to_service)?;
6636 Ok(json!({
6637 "schema": crate::channels::CHANNELS_SCHEMA,
6638 "channel": channel,
6639 }))
6640 }
6641 _ => Err(ServiceError::MethodNotFound),
6642 }
6643}
6644
6645fn rpc_error(id: Value, code: i64, message: &str) -> Value {
6646 json!({
6647 "jsonrpc": "2.0",
6648 "id": id,
6649 "error": {"code": code, "message": message},
6650 })
6651}