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_with_metadata, resolve_live_runtime};
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.usage",
52 "harness.v1.sessions.load",
53 "harness.v1.sessions.recorded_config",
54 "harness.v1.sessions.launch_agent",
55 "harness.v1.sessions.follow",
56 "harness.v1.sessions.unfollow",
57 "harness.v1.sessions.activity.subscribe",
58 "harness.v1.sessions.activity.unsubscribe",
59 "harness.v1.sessions.activity_under",
60 "harness.v1.sessions.index.subscribe",
61 "harness.v1.sessions.index.resize",
62 "harness.v1.sessions.index.unsubscribe",
63 "harness.v1.sessions.message",
64 "harness.v1.sessions.inbox",
65 "harness.v1.sessions.import",
66 "harness.v1.sessions.export",
67 "harness.v1.sessions.translate",
68 "harness.v1.sessions.reduce",
69 "harness.v1.sessions.branch",
70 "harness.v1.sessions.handoff",
71 "harness.v1.sessions.materialize",
72 "harness.v1.sessions.resume_instructions",
73 "harness.v1.skills.list",
74 "harness.v1.skills.install",
75 "harness.v1.skills.remove",
76 "harness.v1.memory.show",
77 "harness.v1.memory.search",
78 "harness.v1.jobs.list",
79 "harness.v1.jobs.get",
80 "harness.v1.jobs.create",
81 "harness.v1.jobs.update",
82 "harness.v1.jobs.pause",
83 "harness.v1.jobs.resume",
84 "harness.v1.jobs.run",
85 "harness.v1.jobs.delete",
86 "harness.v1.jobs.notepad",
87 "harness.v1.jobs.notepad_set",
88 "harness.v1.jobs.notepad_delete",
89 "harness.v1.sessions.new",
90 "harness.v1.sessions.reset",
91 "harness.v1.sessions.archive",
92 "harness.v1.sessions.delete",
93 "harness.v1.runs.list",
94 "harness.v1.runs.get",
95 "harness.v1.approvals.list",
96 "harness.v1.approvals.resolve",
97 "harness.v1.runtimes.capabilities",
98 "harness.v1.runtimes.start",
99 "harness.v1.runtimes.resume",
100 "harness.v1.runtimes.attach_existing",
101 "harness.v1.runtimes.attach",
102 "harness.v1.runtimes.send_input",
103 "harness.v1.runtimes.interrupt",
104 "harness.v1.runtimes.steer",
105 "harness.v1.runtimes.respond",
106 "harness.v1.runtimes.terminal_instructions",
107 "harness.v1.runtimes.acquire_control",
108 "harness.v1.runtimes.heartbeat",
109 "harness.v1.runtimes.detach",
110 "harness.v1.runtimes.close",
111 "harness.v1.profiles.list",
112 "harness.v1.profiles.get",
113 "harness.v1.profiles.create",
114 "harness.v1.profiles.delete",
115 "harness.v1.channels.list",
116 "harness.v1.routes.list",
117 "harness.v1.triggers.list",
118 "harness.v1.channels.status",
119 "harness.v1.orchestration.load",
120 "harness.v1.orchestration.save",
121 "harness.v1.orchestration.compile",
122 "harness.v1.orchestration.decompile",
123 "harness.v1.orchestration.import",
124 "harness.v1.orchestration.export",
125 "harness.v1.workflow.load",
126];
127
128pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
130pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
132pub const SESSION_ACTIVITY_EVENT_METHOD: &str = "harness.v1.sessions.activity_event";
134pub const SESSION_INDEX_EVENT_METHOD: &str = "harness.v1.sessions.index_event";
136pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
138
139pub struct HarnessSessionService {
142 catalog: HarnessCatalog,
143 followers: BTreeMap<String, SessionFollower>,
144 followed_sources: BTreeMap<String, FollowedSource>,
145 activity_subscriptions: BTreeMap<String, ActivitySubscription>,
146 index_subscriptions: BTreeMap<String, crate::session_index::SessionIndexSubscription>,
147 index_notifier: Arc<Notify>,
148 #[cfg(feature = "adapter-api")]
149 activity_monitor: crate::session_activity::SessionActivityMonitor,
150 next_subscription: u64,
151 runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
152 runtimes_in_flight: BTreeSet<String>,
157 terminal_launches: BTreeMap<String, StructuredLaunch>,
158 runtime_sequences: BTreeMap<String, u64>,
159 next_runtime: u64,
160 reduction_store_root: Option<PathBuf>,
161 approvals: crate::approvals::ApprovalRegistry,
165 subagent_approvals: Option<Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>>,
168}
169
170impl Default for HarnessSessionService {
171 fn default() -> Self {
172 Self::new()
173 }
174}
175
176impl HarnessSessionService {
177 pub fn new() -> Self {
179 Self {
180 catalog: HarnessCatalog::new(),
181 followers: BTreeMap::new(),
182 followed_sources: BTreeMap::new(),
183 activity_subscriptions: BTreeMap::new(),
184 index_subscriptions: BTreeMap::new(),
185 index_notifier: Arc::new(Notify::new()),
186 #[cfg(feature = "adapter-api")]
187 activity_monitor: Default::default(),
188 next_subscription: 1,
189 runtimes: BTreeMap::new(),
190 runtimes_in_flight: BTreeSet::new(),
191 terminal_launches: BTreeMap::new(),
192 runtime_sequences: BTreeMap::new(),
193 next_runtime: 1,
194 reduction_store_root: None,
195 approvals: crate::approvals::ApprovalRegistry::new(),
196 subagent_approvals: None,
197 }
198 }
199
200 pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
205 self.reduction_store_root = Some(root.into());
206 self
207 }
208
209 pub fn observe_subagent_approvals(
217 &mut self,
218 queue: Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>,
219 ) {
220 self.subagent_approvals = Some(queue);
221 }
222
223 pub fn approvals(&self, query: &crate::approvals::ApprovalsQuery) -> Vec<crate::ApprovalRow> {
230 let now = crate::approvals::now_ms();
231 let mut rows = self.approvals.rows(now);
232 if let Some(queue) = self.subagent_approvals.as_ref() {
233 let queued = queue
234 .lock()
235 .unwrap_or_else(std::sync::PoisonError::into_inner)
236 .clone();
237 rows.extend(crate::approvals::subagent_rows(&queued, now));
238 }
239 rows.retain(|row| query.matches(row));
240 rows.sort_by(|left, right| {
241 left.requested_at_ms
242 .cmp(&right.requested_at_ms)
243 .then_with(|| left.id.cmp(&right.id))
244 });
245 rows
246 }
247
248 async fn approvals_resolve(
258 &mut self,
259 params: Value,
260 ) -> std::result::Result<Value, ServiceError> {
261 let params = decode::<crate::approvals::ApprovalsResolveParams>(params)?;
262 if params.id.trim().is_empty() {
263 return Err(ServiceError::InvalidParams(
264 "approvals resolve requires the `id` of a listed approval row".into(),
265 ));
266 }
267 let choice = match (params.decision, params.option_id.as_deref()) {
268 (Some(_), Some(_)) => {
269 return Err(ServiceError::InvalidParams(
270 "approvals resolve takes either `decision` or `option_id`, not both".into(),
271 ))
272 }
273 (Some(decision), None) => crate::approvals::ApprovalChoice::Decision(decision),
274 (None, Some(option)) => crate::approvals::ApprovalChoice::Option(option.to_string()),
275 (None, None) => {
276 return Err(ServiceError::InvalidParams(format!(
277 "approvals resolve requires `decision` ({}) or an explicit `option_id`",
278 crate::approvals::ApprovalDecision::ALL
279 .map(|decision| decision.as_str())
280 .join(" | "),
281 )))
282 }
283 };
284 let resolution = self
285 .approvals
286 .resolution(¶ms.id, &choice)
287 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?;
288 self.runtime_call(
292 "harness.v1.runtimes.respond",
293 json!({
294 "connection": resolution.connection,
295 "request_id": resolution.request_id,
296 "response": resolution.response,
297 }),
298 )
299 .await?;
300 Ok(json!({
301 "id": params.id,
302 "decision": params.decision.map(|decision| decision.as_str()),
303 "option_id": resolution.option_id,
304 "resolved": true,
305 }))
306 }
307
308 #[cfg(feature = "adapter-api")]
311 pub fn session_index_notifier(&self) -> Arc<Notify> {
312 Arc::clone(&self.index_notifier)
313 }
314
315 #[cfg(feature = "adapter-api")]
317 pub fn handle(&mut self, request: Value) -> Value {
318 let id = request.get("id").cloned().unwrap_or(Value::Null);
319 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
320 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
321 }
322 let Some(method) = request.get("method").and_then(Value::as_str) else {
323 return rpc_error(id, -32600, "request is missing `method`");
324 };
325 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
326 match self.call(method, params) {
327 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
328 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
329 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
330 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
331 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
332 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
333 }
334 }
335
336 #[cfg(feature = "adapter-api")]
339 pub async fn handle_async(&mut self, request: Value) -> Value {
340 let method = request
341 .get("method")
342 .and_then(Value::as_str)
343 .unwrap_or_default();
344 if matches!(
345 method,
346 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
347 ) {
348 let id = request.get("id").cloned().unwrap_or(Value::Null);
349 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
350 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
351 }
352 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
353 return match self.inventory_call(method, params).await {
354 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
355 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
356 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
357 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
358 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
359 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
360 };
361 }
362 if matches!(
363 method,
364 "harness.v1.harnesses.auth.methods"
365 | "harness.v1.harnesses.auth.begin"
366 | "harness.v1.harnesses.auth.verify"
367 ) {
368 let id = request.get("id").cloned().unwrap_or(Value::Null);
369 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
370 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
371 }
372 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
373 return match self.harness_authentication_call(method, params).await {
374 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
375 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
376 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
377 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
378 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
379 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
380 };
381 }
382 if matches!(
388 method,
389 "harness.v1.sessions.new"
390 | "harness.v1.sessions.reset"
391 | "harness.v1.sessions.archive"
392 | "harness.v1.sessions.delete"
393 ) {
394 let id = request.get("id").cloned().unwrap_or(Value::Null);
395 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
396 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
397 }
398 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
399 let verb = match method {
400 "harness.v1.sessions.new" => crate::SessionVerb::New,
401 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
402 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
403 _ => crate::SessionVerb::Delete,
404 };
405 return match self.mutate_session(verb, params).await {
406 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
407 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
408 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
409 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
410 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
411 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
412 };
413 }
414 if method == "harness.v1.sessions.message" {
415 let id = request.get("id").cloned().unwrap_or(Value::Null);
416 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
417 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
418 }
419 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
420 return match self.message_call(params).await {
421 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
422 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
423 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
424 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
425 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
426 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
427 };
428 }
429 if matches!(
430 method,
431 "harness.v1.harnesses.settings" | "harness.v1.harnesses.configure"
432 ) {
433 let id = request.get("id").cloned().unwrap_or(Value::Null);
434 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
435 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
436 }
437 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
438 return match self.harness_settings_call(method, params) {
439 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
440 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
441 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
442 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
443 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
444 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
445 };
446 }
447 if method == "harness.v1.sessions.activity_under" {
448 let id = request.get("id").cloned().unwrap_or(Value::Null);
449 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
450 return match activity_under_call(params).await {
451 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
452 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
453 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
454 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
455 Err(_) => rpc_error(id, -32000, "sessions.activity_under failed"),
456 };
457 }
458 if method == "harness.v1.sessions.activity.subscribe" {
459 let id = request.get("id").cloned().unwrap_or(Value::Null);
460 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
461 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
462 }
463 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
464 return match self.subscribe_session_activity(params).await {
465 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
466 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
467 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
468 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
469 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
470 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
471 };
472 }
473 if let Some(operation) = SdkOperation::from_method(method) {
474 let id = request.get("id").cloned().unwrap_or(Value::Null);
475 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
476 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
477 }
478 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
479 return match self.execute(SdkRequest { operation, params }).await {
480 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
481 Err(error) => sdk_rpc_error(id, &error),
482 };
483 }
484 if !method.starts_with("harness.v1.runtimes.") {
485 return self.handle(request);
486 }
487 let id = request.get("id").cloned().unwrap_or(Value::Null);
488 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
489 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
490 }
491 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
492 match self.runtime_call(method, params).await {
493 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
494 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
495 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
496 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
497 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
498 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
499 }
500 }
501
502 #[cfg(feature = "adapter-api")]
505 pub fn poll(&mut self) -> Vec<Value> {
506 let mut notifications = Vec::new();
507 for (subscription, follower) in &mut self.followers {
508 match follower.poll() {
509 Ok(Some(event)) => notifications.push(json!({
510 "jsonrpc": "2.0",
511 "method": SESSION_EVENT_METHOD,
512 "params": {
513 "subscription": subscription,
514 "event": event.to_json(),
515 }
516 })),
517 Ok(None) => {}
518 Err(error) => notifications.push(json!({
519 "jsonrpc": "2.0",
520 "method": SESSION_EVENT_METHOD,
521 "params": {
522 "subscription": subscription,
523 "event": {
524 "type": "watch_error",
525 "recoverable": true,
526 "message": error.to_string(),
527 },
528 }
529 })),
530 }
531 }
532 notifications
533 }
534
535 #[cfg(feature = "adapter-api")]
546 pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
547 let registry = crate::LocalRuntimeRegistry::new();
548 let authorization = crate::RuntimeAuthorization::observer();
549 let mut notifications = Vec::new();
550 for (subscription, source) in &mut self.followed_sources {
551 let state = match registry
552 .source_state(&source.harness, &source.session_id, &authorization)
553 .await
554 {
555 Ok(Some(state)) => state,
556 Ok(None) => crate::RuntimeRegistryState::Persisted,
557 Err(_) => continue,
559 };
560 if source.reported.as_deref() == Some(state.as_str()) {
561 continue;
562 }
563 source.reported = Some(state.as_str().to_string());
564 notifications.push(json!({
565 "jsonrpc": "2.0",
566 "method": SESSION_EVENT_METHOD,
567 "params": {
568 "subscription": subscription,
569 "event": {"type": "runtime_state", "state": state.as_str()},
570 },
571 }));
572 }
573 notifications
574 }
575
576 #[cfg(feature = "adapter-api")]
580 pub async fn poll_session_activities(&mut self) -> Vec<Value> {
581 let subscriptions = self
582 .activity_subscriptions
583 .iter()
584 .map(|(id, subscription)| {
585 (
586 id.clone(),
587 subscription.locators.clone(),
588 subscription.homes.clone(),
589 )
590 })
591 .collect::<Vec<_>>();
592 let mut notifications = Vec::new();
593 for (subscription_id, locators, homes) in subscriptions {
594 let Ok(mut activities) = self.activity_monitor.resolve(&locators, &homes).await else {
595 continue;
598 };
599 pane_prompt_override(&mut activities, &homes);
600 let Some(subscription) = self.activity_subscriptions.get_mut(&subscription_id) else {
601 continue;
602 };
603 let mut changed = Vec::new();
604 for activity in activities {
605 let key = activity.key();
606 if subscription
607 .reported
608 .get(&key)
609 .is_some_and(|previous| previous.same_state(&activity))
610 {
611 continue;
612 }
613 subscription.reported.insert(key, activity.clone());
614 changed.push(activity);
615 }
616 if !changed.is_empty() {
617 notifications.push(json!({
618 "jsonrpc": "2.0",
619 "method": SESSION_ACTIVITY_EVENT_METHOD,
620 "params": {
621 "subscription": subscription_id,
622 "activities": changed,
623 },
624 }));
625 }
626 }
627 notifications
628 }
629
630 #[cfg(feature = "adapter-api")]
634 pub fn poll_session_indexes(&mut self) -> Vec<Value> {
635 let mut notifications = Vec::new();
636 for (subscription, index) in &mut self.index_subscriptions {
637 let homes = index.homes().clone();
638 match index.poll() {
639 Ok(Some(delta)) => match live_index_changes(delta.changes, &homes) {
640 Ok(changes) => notifications.push(json!({
641 "jsonrpc": "2.0",
642 "method": SESSION_INDEX_EVENT_METHOD,
643 "params": {
644 "subscription": subscription,
645 "revision": delta.revision,
646 "changes": changes,
647 },
648 })),
649 Err(error) => notifications.push(json!({
650 "jsonrpc": "2.0",
651 "method": SESSION_INDEX_EVENT_METHOD,
652 "params": {
653 "subscription": subscription,
654 "error": {"recoverable": true, "message": error_message(error)},
655 },
656 })),
657 },
658 Ok(None) => {}
659 Err(error) => notifications.push(json!({
660 "jsonrpc": "2.0",
661 "method": SESSION_INDEX_EVENT_METHOD,
662 "params": {
663 "subscription": subscription,
664 "error": {"recoverable": true, "message": error},
665 },
666 })),
667 }
668 }
669 notifications
670 }
671
672 #[cfg(feature = "adapter-api")]
673 async fn subscribe_session_activity(
674 &mut self,
675 params: Value,
676 ) -> std::result::Result<Value, ServiceError> {
677 let params = decode::<ActivitySubscribeParams>(params)?;
678 if params.locators.is_empty() {
679 return Err(ServiceError::InvalidParams(
680 "sessions.activity.subscribe requires at least one locator".into(),
681 ));
682 }
683 if params.locators.len() > 2_048 {
684 return Err(ServiceError::InvalidParams(
685 "sessions.activity.subscribe accepts at most 2048 locators".into(),
686 ));
687 }
688 let mut initial = self
689 .activity_monitor
690 .resolve(¶ms.locators, ¶ms.homes)
691 .await
692 .map_err(ServiceError::Sdk)?;
693 pane_prompt_override(&mut initial, ¶ms.homes);
694 let subscription = format!("activity-sub-{}", self.next_subscription);
695 self.next_subscription += 1;
696 let reported = initial
697 .iter()
698 .cloned()
699 .map(|activity| (activity.key(), activity))
700 .collect();
701 self.activity_subscriptions.insert(
702 subscription.clone(),
703 ActivitySubscription {
704 locators: params.locators,
705 homes: params.homes,
706 reported,
707 },
708 );
709 Ok(json!({"subscription": subscription, "initial": initial}))
710 }
711
712 #[cfg(feature = "adapter-api")]
714 pub async fn poll_runtimes(&mut self) -> Vec<Value> {
715 self.poll_sdk_events()
716 .await
717 .into_iter()
718 .map(|(connection, runtime_event)| {
719 json!({
720 "jsonrpc": "2.0",
721 "method": RUNTIME_EVENT_METHOD,
722 "params": {
723 "connection": connection,
724 "session_id": runtime_event.session_id,
725 "sequence": runtime_event.event.sequence,
726 "event": {
727 "kind": runtime_event.event.kind,
728 "payload": runtime_event.event.payload,
729 },
730 },
731 })
732 })
733 .collect()
734 }
735
736 async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
737 let mut events = Vec::new();
738 let mut closed = Vec::new();
739 let now_ms = crate::approvals::now_ms();
740 for (connection, runtime) in &mut self.runtimes {
741 let session_id = runtime.handle().runtime_id.clone();
742 let harness = runtime.handle().harness.clone();
743 for _ in 0..256 {
748 match tokio::time::timeout(Duration::ZERO, runtime.next_event()).await {
749 Ok(Ok(Some(event))) => {
750 let terminal = event.kind == "transport_closed";
751 self.approvals
755 .observe(connection, &harness, &session_id, &event, now_ms);
756 let next_sequence = self
757 .runtime_sequences
758 .entry(session_id.clone())
759 .or_insert(0);
760 let sequence = event.sequence.unwrap_or_else(|| {
761 *next_sequence = next_sequence.saturating_add(1);
762 *next_sequence
763 });
764 *next_sequence = (*next_sequence).max(sequence);
765 events.push((
766 connection.clone(),
767 SdkRuntimeEvent {
768 session_id: session_id.clone(),
769 event: SdkEvent {
770 sequence,
771 kind: event.kind,
772 payload: event.payload,
773 },
774 },
775 ));
776 if terminal {
777 closed.push(connection.clone());
778 break;
779 }
780 }
781 Ok(Ok(None)) => {
782 let sequence = self
783 .runtime_sequences
784 .entry(session_id.clone())
785 .or_insert(0);
786 *sequence = sequence.saturating_add(1);
787 events.push((
788 connection.clone(),
789 SdkRuntimeEvent {
790 session_id,
791 event: SdkEvent {
792 sequence: *sequence,
793 kind: "transport_closed".into(),
794 payload: json!({"message": "Harness runtime transport closed."}),
795 },
796 },
797 ));
798 closed.push(connection.clone());
799 break;
800 }
801 Err(_) => break,
802 Ok(Err(error)) => {
803 let sequence = self
804 .runtime_sequences
805 .entry(session_id.clone())
806 .or_insert(0);
807 *sequence = sequence.saturating_add(1);
808 events.push((
809 connection.clone(),
810 SdkRuntimeEvent {
811 session_id,
812 event: SdkEvent {
813 sequence: *sequence,
814 kind: "transport_error".into(),
815 payload: json!({"message": error.to_string(), "terminal": true}),
816 },
817 },
818 ));
819 closed.push(connection.clone());
820 break;
821 }
822 }
823 }
824 }
825 for connection in closed {
826 if let Some(runtime) = self.runtimes.remove(&connection) {
827 self.runtime_sequences.remove(&runtime.handle().runtime_id);
828 }
829 self.terminal_launches.remove(&connection);
830 self.approvals.forget(&connection);
833 }
834 events
835 }
836
837 fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
838 match method {
839 "harness.v1.capabilities" => Ok(json!({
840 "version": HARNESS_SERVICE_VERSION,
841 "sdk": self.capabilities(),
842 "methods": HARNESS_SERVICE_METHODS,
843 "notifications": [
844 SESSION_EVENT_METHOD,
845 SESSION_ACTIVITY_EVENT_METHOD,
846 SESSION_INDEX_EVENT_METHOD,
847 RUNTIME_EVENT_METHOD
848 ],
849 "harnesses": harness_support_registry()
850 .harnesses
851 .into_iter()
852 .map(|harness| harness.id)
853 .collect::<Vec<_>>(),
854 })),
855 "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
856 .map_err(|error| ServiceError::Operation(error.to_string())),
857 "harness.v1.profiles.list" | "harness.v1.profiles.get" => profiles_call(method, params),
858 "harness.v1.profiles.create" => {
864 mutate_profile(crate::profiles_control::ProfileVerb::Create, params)
865 }
866 "harness.v1.profiles.delete" => {
867 mutate_profile(crate::profiles_control::ProfileVerb::Delete, params)
868 }
869 "harness.v1.channels.list" | "harness.v1.channels.status" => {
870 channels_call(method, params)
871 }
872 "harness.v1.routes.list" => routes_call(params),
875 "harness.v1.triggers.list" => triggers_call(params),
877 "harness.v1.workflow.load" => {
887 let params = decode::<WorkflowLoadParams>(params)?;
888 let read = crate::workflow_doors::load(
889 params.from.unwrap_or_else(|| {
890 crate::workflow_doors::WorkflowHarness::of_home(¶ms.home)
891 }),
892 ¶ms.home,
893 )
894 .map_err(operation)?;
895 serde_json::to_value(read)
896 .map_err(|error| ServiceError::Operation(error.to_string()))
897 }
898 "harness.v1.orchestration.load" => {
899 let params = decode::<OrchestrationLoadParams>(params)?;
900 let read = crate::orchestration_doors::load(¶ms.root, params.flavor)
901 .map_err(operation)?;
902 serde_json::to_value(read)
903 .map_err(|error| ServiceError::Operation(error.to_string()))
904 }
905 "harness.v1.orchestration.save" => {
906 let params = decode::<OrchestrationSaveParams>(params)?;
907 let saved = crate::orchestration_doors::save(
908 ¶ms.root,
909 params.orchestration,
910 params.vault,
911 )
912 .map_err(operation)?;
913 serde_json::to_value(saved)
914 .map_err(|error| ServiceError::Operation(error.to_string()))
915 }
916 "harness.v1.orchestration.compile" => {
917 let params = decode::<OrchestrationCompileParams>(params)?;
918 let read = crate::orchestration_doors::compile(params.from, ¶ms.home)
919 .map_err(operation)?;
920 serde_json::to_value(read)
921 .map_err(|error| ServiceError::Operation(error.to_string()))
922 }
923 "harness.v1.orchestration.decompile" => {
924 let params = decode::<OrchestrationDecompileParams>(params)?;
925 let report = crate::orchestration_doors::decompile(
926 params.to,
927 params.orchestration,
928 ¶ms.source,
929 params.source_flavor,
930 ¶ms.dest,
931 params.vault,
932 )
933 .map_err(operation)?;
934 serde_json::to_value(report)
935 .map_err(|error| ServiceError::Operation(error.to_string()))
936 }
937 "harness.v1.orchestration.import" => {
942 let params = decode::<OrchestrationImportParams>(params)?;
943 let imported =
944 crate::orchestration_doors::import(params.from, ¶ms.home, ¶ms.into)
945 .map_err(operation)?;
946 serde_json::to_value(imported)
947 .map_err(|error| ServiceError::Operation(error.to_string()))
948 }
949 "harness.v1.orchestration.export" => {
950 let params = decode::<OrchestrationExportParams>(params)?;
951 let report =
952 crate::orchestration_doors::export(params.to, ¶ms.root, ¶ms.dest)
953 .map_err(operation)?;
954 serde_json::to_value(report)
955 .map_err(|error| ServiceError::Operation(error.to_string()))
956 }
957 "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
963 "harness.v1.skills.list" => {
968 let query = decode::<crate::skills::SkillsQuery>(params)?;
969 if let Some(harness) = query.harness.as_deref() {
970 if !crate::skills::SKILL_HARNESSES.contains(&harness) {
971 return Err(ServiceError::UnsupportedAction(format!(
972 "`{harness}` has no skills root Volter Harness reads"
973 )));
974 }
975 }
976 serde_json::to_value(crate::skills::list_skills(&query))
977 .map_err(|error| ServiceError::Operation(error.to_string()))
978 }
979 "harness.v1.skills.install" => {
987 mutate_skill(crate::skills_control::SkillVerb::Install, params)
988 }
989 "harness.v1.skills.remove" => {
990 mutate_skill(crate::skills_control::SkillVerb::Remove, params)
991 }
992 "harness.v1.approvals.list" => {
1000 let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
1001 if let Some(harness) = query.harness.as_deref() {
1002 if !crate::approvals::lists_approvals(harness) {
1003 return Err(ServiceError::UnsupportedAction(format!(
1004 "`{harness}` has no runtime door that carries an approval request"
1005 )));
1006 }
1007 }
1008 serde_json::to_value(self.approvals(&query))
1009 .map_err(|error| ServiceError::Operation(error.to_string()))
1010 }
1011 "harness.v1.sessions.discover" => {
1012 let query = decode::<DiscoveryQuery>(params)?;
1013 let mut page = discover_session_page(&query).map_err(operation)?;
1014 let doors = crate::mail_route::LiveSessions::read(&query.homes);
1019 if let Some(text) = query
1023 .query
1024 .as_deref()
1025 .map(str::trim)
1026 .filter(|text| !text.is_empty())
1027 {
1028 let needle = text.to_lowercase();
1029 for live in doors.all() {
1030 let name = live
1031 .name
1032 .split('@')
1033 .next()
1034 .unwrap_or(&live.name)
1035 .to_lowercase();
1036 if !name.contains(&needle)
1037 || page.sessions.iter().any(|session| {
1038 session.locator.session_id == live.address.session_id
1039 })
1040 {
1041 continue;
1042 }
1043 if query
1044 .limit
1045 .is_some_and(|limit| page.sessions.len() >= limit)
1046 {
1047 break;
1048 }
1049 let mut by_id = query.clone();
1051 by_id.query = Some(live.address.session_id.clone());
1052 by_id.cursor = None;
1053 if let Some(cwd) = &live.cwd {
1054 by_id.workspace = Some(cwd.clone());
1055 by_id.workspace_subtree = false;
1056 by_id.workspace_family = None;
1057 }
1058 if let Ok(found) = discover_session_page(&by_id) {
1059 page.sessions
1060 .extend(found.sessions.into_iter().filter(|session| {
1061 session.locator.session_id == live.address.session_id
1062 }));
1063 }
1064 }
1065 }
1066 let activities = crate::session_activity::resolve_stock_session_activities(
1067 &page
1068 .sessions
1069 .iter()
1070 .map(|session| session.locator.clone())
1071 .collect::<Vec<_>>(),
1072 &query.homes,
1073 )
1074 .into_iter()
1075 .map(|activity| (activity.key(), activity))
1076 .collect::<BTreeMap<_, _>>();
1077 let sessions = page
1078 .sessions
1079 .into_iter()
1080 .map(|session| {
1081 let mut value = live_descriptor_value(&session, &doors)?;
1082 let activity_key = (
1083 session.locator.harness.as_str().to_string(),
1084 session.locator.session_id.clone(),
1085 );
1086 if let Some(activity) = activities.get(&activity_key) {
1087 let mut activity = activity.clone();
1090 let waiting = doors.all().iter().any(|live| {
1091 live.address.harness == activity_key.0
1092 && live.address.session_id == activity_key.1
1093 && live.status == "waiting"
1094 });
1095 if waiting && activity.turn != crate::SessionTurnState::NeedsInput {
1096 activity.turn = crate::SessionTurnState::NeedsInput;
1097 activity.evidence.source = "pane_prompt".into();
1098 }
1099 let activity = &activity;
1100 value["activity"] = serde_json::to_value(activity)
1101 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1102 if let Some(status) = legacy_live_status(activity) {
1103 value["live_status"] = json!(status);
1104 }
1105 }
1106 Ok(value)
1107 })
1108 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1109 let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1110 if query.search_previews {
1113 result["receipt"] = serde_json::to_value(page.receipt)
1114 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1115 }
1116 Ok(result)
1117 }
1118 "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1119 "harness.v1.sessions.usage" => {
1120 let params = decode::<LocatorParams>(params)?;
1121 let session = load_session(¶ms.locator).map_err(operation)?;
1122 Ok(
1123 json!({"harness": params.locator.harness, "session_id": session.meta.session_id,
1124 "usage": recorded_session_usage(&session)}),
1125 )
1126 }
1127 "harness.v1.sessions.recorded_config" => {
1128 let params = decode::<LocatorParams>(params)?;
1129 recorded_config(¶ms.locator).map_err(operation)
1130 }
1131 "harness.v1.sessions.launch_agent" => {
1132 let params = decode::<LocatorParams>(params)?;
1133 Ok(json!({"environment": crate::launch_agent::read(¶ms.locator)}))
1134 }
1135 "harness.v1.sessions.load" => {
1136 let params = decode::<LoadSessionParams>(params)?;
1137 load_session_door(params)
1138 }
1139 "harness.v1.sessions.follow" => {
1140 let params = decode::<LocatorParams>(params)?;
1141 let mut follower = self
1142 .catalog
1143 .follow_read_view(
1144 ¶ms.locator,
1145 params.read_fidelity(),
1146 params.include_subagents(),
1147 params.tail_messages(),
1148 params.max_message_chars(),
1149 params.display_history(),
1150 )
1151 .map_err(operation)?;
1152 let initial = follower
1153 .poll()
1154 .map_err(operation)?
1155 .map(|event| event.to_json());
1156 let subscription = format!("sub-{}", self.next_subscription);
1157 self.next_subscription += 1;
1158 self.followers.insert(subscription.clone(), follower);
1159 self.followed_sources.insert(
1160 subscription.clone(),
1161 FollowedSource {
1162 harness: params.locator.harness.as_str().to_string(),
1163 session_id: params.locator.session_id.clone(),
1164 reported: None,
1165 },
1166 );
1167 Ok(json!({"subscription": subscription, "initial": initial}))
1168 }
1169 "harness.v1.sessions.unfollow" => {
1170 let params = decode::<UnfollowParams>(params)?;
1171 self.followed_sources.remove(¶ms.subscription);
1172 Ok(json!({
1173 "removed": self.followers.remove(¶ms.subscription).is_some()
1174 }))
1175 }
1176 "harness.v1.sessions.activity.unsubscribe" => {
1177 let params = decode::<UnfollowParams>(params)?;
1178 Ok(json!({
1179 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1180 }))
1181 }
1182 "harness.v1.sessions.index.subscribe" => {
1183 let query = decode::<DiscoveryQuery>(params)?;
1184 crate::session_index::validate_query(&query)
1185 .map_err(ServiceError::InvalidParams)?;
1186 let homes = query.homes.clone();
1187 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1188 query,
1189 Arc::clone(&self.index_notifier),
1190 )
1191 .map_err(ServiceError::Operation)?;
1192 let doors = crate::mail_route::LiveSessions::read(&homes);
1193 let initial = initial
1194 .iter()
1195 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1196 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1197 let subscription = format!("index-sub-{}", self.next_subscription);
1198 self.next_subscription += 1;
1199 self.index_subscriptions.insert(subscription.clone(), index);
1200 Ok(json!({
1201 "subscription": subscription,
1202 "revision": 1,
1203 "initial": initial,
1204 }))
1205 }
1206 "harness.v1.sessions.index.resize" => {
1207 let params = decode::<IndexResizeParams>(params)?;
1208 crate::session_index::validate_limit(params.limit)
1209 .map_err(ServiceError::InvalidParams)?;
1210 let index = self
1211 .index_subscriptions
1212 .get_mut(¶ms.subscription)
1213 .ok_or_else(|| {
1214 ServiceError::InvalidParams("unknown session index subscription".into())
1215 })?;
1216 let prepared = index
1217 .prepare_resize(params.limit)
1218 .map_err(ServiceError::Operation)?;
1219 let doors = crate::mail_route::LiveSessions::read(index.homes());
1220 let initial = prepared
1221 .page
1222 .sessions
1223 .iter()
1224 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1225 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1226 let response = json!({
1227 "subscription": params.subscription,
1228 "revision": prepared.revision,
1229 "initial": initial,
1230 "receipt": prepared.page.receipt,
1231 });
1232 index.commit_resize(prepared);
1233 Ok(response)
1234 }
1235 "harness.v1.sessions.index.unsubscribe" => {
1236 let params = decode::<UnfollowParams>(params)?;
1237 Ok(json!({
1238 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1239 }))
1240 }
1241 "harness.v1.sessions.import" => {
1242 let params = decode::<ImportSessionParams>(params)?;
1243 let session = Session::load_str(¶ms.content, params.source_harness.into())
1244 .map_err(operation)?;
1245 Ok(json!({"session": normalized_session_json(&session)}))
1246 }
1247 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1248 let params = decode::<ExportSessionParams>(params)?;
1249 let session = load_session(¶ms.locator).map_err(operation)?;
1250 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1251 if method == "harness.v1.sessions.export"
1252 && params.target_harness == TransferFormat::Hermes
1253 {
1254 let imported = crate::hermes_import::import_into_hermes(&session, None)
1256 .map_err(operation)?;
1257 return Ok(json!({"artifact": artifact, "imported": imported}));
1258 }
1259 Ok(json!({"artifact": artifact}))
1260 }
1261 "harness.v1.sessions.reduce" => {
1262 let params = decode::<ReduceSessionParams>(params)?;
1263 self.reduce_session(params)
1264 }
1265 "harness.v1.sessions.branch" => {
1266 let params = decode::<BranchSessionParams>(params)?;
1267 let session = load_session(¶ms.locator).map_err(operation)?;
1268 let storage = params.locator.storage.path().display().to_string();
1269 let bootstrap_prompt = format!(
1270 "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.",
1271 params.locator.harness.as_str(), params.locator.session_id, storage
1272 );
1273 let artifact = params
1274 .target_harness
1275 .map(|target| session_artifact(¶ms.locator, &session, target))
1276 .transpose()?;
1277 Ok(json!({
1278 "parent": params.locator,
1279 "session": normalized_session_json(&session),
1280 "bootstrap_prompt": bootstrap_prompt,
1281 "artifact": artifact,
1282 }))
1283 }
1284 "harness.v1.sessions.handoff" => {
1285 let params = decode::<HandoffSessionParams>(params)?;
1286 let session = load_session(¶ms.locator).map_err(operation)?;
1287 let cwd = params
1288 .cwd
1289 .or_else(|| session.meta.cwd.clone())
1290 .unwrap_or_else(|| PathBuf::from("."));
1291 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1292 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1293 ServiceError::Operation(
1294 "handoff artifact omitted target session identity".into(),
1295 )
1296 })?;
1297 let instructions =
1298 handoff_instructions(params.target_harness, target_session_id, &cwd);
1299 Ok(json!({
1300 "artifact": artifact,
1301 "launch": instructions.launch,
1302 "materialize": instructions.materialize,
1303 "requires_materialization": instructions.requires_materialization,
1304 "note": instructions.note,
1305 }))
1306 }
1307 "harness.v1.sessions.materialize" => {
1308 let params = decode::<MaterializeSessionParams>(params)?;
1309 for file in ¶ms.artifact.files {
1313 if file.role == "source_recovery"
1314 && file.path == "recovery/source.supercode.jsonl"
1315 {
1316 if let Ok(source) = Session::from_native_str(&file.content) {
1317 crate::residue_store::store_segments(&source);
1318 }
1319 }
1320 }
1321 let locator = crate::native_materialize::materialize_native_artifact(
1322 params.artifact,
1323 ¶ms.cwd,
1324 ¶ms.homes,
1325 )
1326 .map_err(ServiceError::Operation)?;
1327 Ok(json!({"locator": locator}))
1328 }
1329 "harness.v1.jobs.list" => {
1333 let query = decode::<crate::jobs::JobsQuery>(params)?;
1334 if let Some(harness) = query.harness.as_deref() {
1335 refuse_harness_without_jobs(harness, "jobs.list")?;
1336 }
1337 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1338 serde_json::to_value(listing)
1339 .map_err(|error| ServiceError::Operation(error.to_string()))
1340 }
1341 "harness.v1.jobs.get" => {
1342 let params = decode::<JobsGetParams>(params)?;
1343 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1344 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1345 .map_err(operation)?
1346 {
1347 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1348 None => Err(ServiceError::Operation(format!(
1349 "`{}` has no scheduled job `{}`",
1350 params.harness, params.id
1351 ))),
1352 }
1353 }
1354 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1360 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1361 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1362 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1363 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1364 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1365 "harness.v1.jobs.notepad"
1366 | "harness.v1.jobs.notepad_set"
1367 | "harness.v1.jobs.notepad_delete" => {
1368 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1369 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1370 let answer = match method {
1371 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1372 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1373 _ => crate::jobs_notepad::read(&request),
1374 }
1375 .map_err(job_control_error)?;
1376 serde_json::to_value(answer)
1377 .map_err(|error| ServiceError::Operation(error.to_string()))
1378 }
1379 "harness.v1.runs.list" => {
1383 let query = decode::<crate::runs::RunsQuery>(params)?;
1384 if let Some(harness) = query.harness.as_deref() {
1385 refuse_harness_without_runs(harness, "runs.list")?;
1386 }
1387 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1388 serde_json::to_value(listing)
1389 .map_err(|error| ServiceError::Operation(error.to_string()))
1390 }
1391 "harness.v1.runs.get" => {
1392 let params = decode::<RunsGetParams>(params)?;
1393 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1394 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1395 .map_err(operation)?
1396 {
1397 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1398 None => Err(ServiceError::Operation(format!(
1399 "`{}` has no run `{}`",
1400 params.harness, params.id
1401 ))),
1402 }
1403 }
1404 "harness.v1.sessions.resume_instructions" => {
1405 let params = decode::<ResumeInstructionsParams>(params)?;
1406 let session = load_session(¶ms.locator).map_err(operation)?;
1407 let cwd = params
1408 .cwd
1409 .or(session.meta.cwd)
1410 .unwrap_or_else(|| PathBuf::from("."));
1411 let launch = resume_launch(
1412 params.locator.harness.as_str(),
1413 ¶ms.locator.session_id,
1414 &cwd,
1415 params.policy,
1416 )?;
1417 Ok(json!({"launch": launch}))
1418 }
1419 _ => Err(ServiceError::MethodNotFound),
1420 }
1421 }
1422
1423 fn reduce_session(
1424 &self,
1425 params: ReduceSessionParams,
1426 ) -> std::result::Result<Value, ServiceError> {
1427 let session = load_session(¶ms.locator).map_err(operation)?;
1428 if session.messages.is_empty() {
1429 return Err(ServiceError::InvalidParams(
1430 "cannot reduce an empty session".into(),
1431 ));
1432 }
1433 let keep_last = params.keep_last.clamp(1, 128);
1434 let policy = reduce::ReductionPolicy {
1435 clear_turns_older_than: Some(keep_last),
1436 ..Default::default()
1437 };
1438 let (view, log) =
1439 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1440 if log.reductions.is_empty() {
1441 return Err(ServiceError::UnsupportedAction(format!(
1442 "session `{}` is already too small for a meaningful reversible reduction",
1443 params.locator.session_id
1444 )));
1445 }
1446 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1447 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1448 if reduced_tokens >= source_tokens {
1449 return Err(ServiceError::UnsupportedAction(format!(
1450 "session `{}` has no token-reducing reversible projection",
1451 params.locator.session_id
1452 )));
1453 }
1454
1455 let store_root = self
1456 .reduction_store_root
1457 .clone()
1458 .unwrap_or_else(default_reduction_store_root);
1459 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1460 let rescue_id = format!("rescue-{}", generated_session_id());
1461 let imported = session
1462 .imported_message_count
1463 .unwrap_or(session.messages.len())
1464 .min(session.messages.len());
1465 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1466 let view_jsonl = messages_jsonl(&view)?;
1467 let title = format!(
1468 "Reduced {} continuation from {}",
1469 params.target_harness.id(),
1470 params.locator.session_id
1471 );
1472
1473 store
1478 .save_sidecar(&rescue_id, &sidecar_jsonl)
1479 .map_err(operation)?;
1480 store
1481 .save_reduction_log(&rescue_id, &log)
1482 .map_err(operation)?;
1483 store
1484 .save(&rescue_id, &title, &view_jsonl)
1485 .map_err(operation)?;
1486
1487 let source_bytes = serde_json::to_vec(&session.messages)
1488 .map_err(|error| ServiceError::Operation(error.to_string()))?
1489 .len() as u64;
1490 let reduced_bytes = serde_json::to_vec(&view)
1491 .map_err(|error| ServiceError::Operation(error.to_string()))?
1492 .len() as u64;
1493 store
1494 .set_reduction_stats(
1495 &rescue_id,
1496 &title,
1497 source_bytes,
1498 reduced_bytes,
1499 log.reductions.len() as u32,
1500 )
1501 .map_err(operation)?;
1502
1503 let reloaded_sidecar = store
1507 .load_sidecar(&rescue_id)
1508 .map_err(operation)?
1509 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1510 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1511 let reloaded_log = store
1512 .load_reduction_log(&rescue_id)
1513 .map_err(operation)?
1514 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1515 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1516 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1517 let (restamped_view, restamped_log) =
1524 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1525 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1526 return Err(ServiceError::Operation(
1527 "persisted reduction view does not match its durable log and sidecar".into(),
1528 ));
1529 }
1530 if restamped_log != reloaded_log {
1531 return Err(ServiceError::Operation(
1532 "reapplying the durable reduction log changed its identity".into(),
1533 ));
1534 }
1535 let inverted =
1536 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1537 if inverted != session.messages {
1538 return Err(ServiceError::Operation(
1539 "reduction inversion did not restore the source messages byte-exactly".into(),
1540 ));
1541 }
1542
1543 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1544 let sidecar_path = store.sidecar_path(&rescue_id);
1545 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1546 let bootstrap_prompt = reduced_bootstrap_prompt(
1547 ¶ms.locator,
1548 params.target_harness,
1549 &view_jsonl,
1550 &sidecar_path,
1551 &reduction_log_path,
1552 );
1553 let mut reduced_session = session.clone();
1554 reduced_session.meta.session_id = Some(rescue_id.clone());
1555 reduced_session.messages = view;
1556
1557 Ok(json!({
1558 "session": normalized_session_json(&reduced_session),
1559 "bootstrap_prompt": bootstrap_prompt,
1560 "receipt": {
1561 "id": rescue_id,
1562 "sidecar_id": rescue_id,
1563 "source_harness": params.locator.harness,
1564 "target_harness": params.target_harness.id(),
1565 "source_tokens": source_tokens,
1566 "reduced_tokens": reduced_tokens,
1567 "ratio": ratio,
1568 "source_bytes": source_bytes,
1569 "reduced_bytes": reduced_bytes,
1570 "reductions": reloaded_log.reductions.len(),
1571 "sidecar_path": sidecar_path,
1572 "reduction_log_path": reduction_log_path,
1573 "verified": true,
1574 "reversible": true,
1575 }
1576 }))
1577 }
1578
1579 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1597 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1598 return None;
1599 }
1600 let method = request.get("method").and_then(Value::as_str)?;
1601 if !RUNTIME_OPEN_METHODS.contains(&method) {
1602 return None;
1603 }
1604 Some(RuntimeOpen {
1605 id: request.get("id").cloned().unwrap_or(Value::Null),
1606 method: method.to_string(),
1607 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1608 })
1609 }
1610
1611 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1633 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1634 return None;
1635 }
1636 let method = request.get("method").and_then(Value::as_str)?;
1637 if !DETACHED_METHODS.contains(&method) {
1638 return None;
1639 }
1640 let id = request.get("id").cloned().unwrap_or(Value::Null);
1641 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1642 let work = match method {
1643 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1644 .inventory_work(method, params)
1645 .map(DetachedWork::Inventory),
1646 "harness.v1.sessions.message" => {
1647 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1648 }
1649 "harness.v1.sessions.load" => {
1650 decode::<LoadSessionParams>(params).map(DetachedWork::Load)
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::<RuntimeCloseParams>(params)
1702 .and_then(|params| self.surrender_runtime_of(¶ms))
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 #[cfg(feature = "adapter-api")]
1809 if params
1810 .base_url
1811 .as_deref()
1812 .is_some_and(|value| LiveRuntimeEndpoint::parse(value).is_ok())
1813 {
1814 return Ok(json!({
1815 "harness": ¶ms.harness,
1816 "capabilities": SupercodeHttpRuntimeBackend::live_capabilities(),
1817 }));
1818 }
1819 let backend = runtime_backend(¶ms)?;
1820 Ok(json!({
1821 "harness": backend.harness(),
1822 "capabilities": backend.capabilities(),
1823 }))
1824 }
1825 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1826 self.register_open_runtime(open_runtime(method, params).await?)
1827 .await
1828 }
1829 "harness.v1.runtimes.send_input" => {
1830 let params = decode::<RuntimeInputParams>(params)?;
1831 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1832 let runtime = self.runtime_mut(¶ms.connection)?;
1833 let turn_id = within_control_deadline(
1834 method,
1835 runtime.send_input(RuntimeInput {
1836 text: params.text,
1837 image_urls,
1838 }),
1839 )
1840 .await?
1841 .map_err(operation)?;
1842 Ok(json!({"turn_id": turn_id}))
1843 }
1844 "harness.v1.runtimes.interrupt" => {
1845 let params = decode::<RuntimeConnectionParams>(params)?;
1846 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1847 .await?
1848 .map_err(operation)?;
1849 Ok(json!({}))
1850 }
1851 "harness.v1.runtimes.steer" => {
1852 let params = decode::<RuntimeInputParams>(params)?;
1853 if !params.image_urls.is_empty() {
1854 return Err(ServiceError::InvalidParams(
1855 "runtime steering accepts text only".into(),
1856 ));
1857 }
1858 let text = params.text.trim();
1859 if text.is_empty() || text.chars().count() > 50_000 {
1860 return Err(ServiceError::InvalidParams(
1861 "runtime steering requires 1 to 50,000 text characters".into(),
1862 ));
1863 }
1864 within_control_deadline(
1865 method,
1866 self.runtime_mut(¶ms.connection)?
1867 .steer(text.to_string()),
1868 )
1869 .await?
1870 .map_err(operation)?;
1871 Ok(json!({}))
1872 }
1873 "harness.v1.runtimes.respond" => {
1874 let params = decode::<RuntimeRespondParams>(params)?;
1875 let request_id = params.request_id.clone();
1876 within_control_deadline(
1877 method,
1878 self.runtime_mut(¶ms.connection)?
1879 .respond(params.request_id, params.response),
1880 )
1881 .await?
1882 .map_err(operation)?;
1883 self.approvals.answered(¶ms.connection, &request_id);
1885 Ok(json!({}))
1886 }
1887 "harness.v1.runtimes.acquire_control" => {
1888 let params = decode::<RuntimeConnectionParams>(params)?;
1889 let snapshot = within_control_deadline(
1890 method,
1891 self.runtime_mut(¶ms.connection)?.acquire_control(),
1892 )
1893 .await?
1894 .map_err(operation)?;
1895 serde_json::to_value(snapshot)
1896 .map_err(|error| ServiceError::Operation(error.to_string()))
1897 }
1898 "harness.v1.runtimes.heartbeat" => {
1899 let params = decode::<RuntimeConnectionParams>(params)?;
1900 let snapshot = within_control_deadline(
1901 method,
1902 self.runtime_mut(¶ms.connection)?.heartbeat(),
1903 )
1904 .await?
1905 .map_err(operation)?;
1906 serde_json::to_value(snapshot)
1907 .map_err(|error| ServiceError::Operation(error.to_string()))
1908 }
1909 "harness.v1.runtimes.detach" => {
1910 let params = decode::<RuntimeConnectionParams>(params)?;
1911 let snapshot =
1912 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1913 .await?
1914 .map_err(operation)?;
1915 serde_json::to_value(snapshot)
1916 .map_err(|error| ServiceError::Operation(error.to_string()))
1917 }
1918 "harness.v1.runtimes.terminal_instructions" => {
1919 let params = decode::<RuntimeConnectionParams>(params)?;
1920 let launch = self
1921 .terminal_launches
1922 .get(¶ms.connection)
1923 .ok_or_else(|| {
1924 ServiceError::Operation(
1925 "this runtime is not hosted for terminal attachment".into(),
1926 )
1927 })?;
1928 Ok(json!({"launch":launch}))
1929 }
1930 "harness.v1.runtimes.close" => {
1931 let params = decode::<RuntimeCloseParams>(params)?;
1932 let (runtime, process_group) = self.surrender_runtime_of(¶ms)?;
1933 close_runtime(runtime, process_group).await
1934 }
1935 _ => Err(ServiceError::MethodNotFound),
1936 }
1937 }
1938
1939 #[cfg(feature = "adapter-api")]
1941 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1942 let params = decode::<MessageSessionParams>(params)?;
1943 Ok(message_live_session(¶ms).await)
1944 }
1945
1946 #[cfg(feature = "adapter-api")]
1947 fn harness_settings_call(
1948 &self,
1949 method: &str,
1950 params: Value,
1951 ) -> std::result::Result<Value, ServiceError> {
1952 let homes = crate::HarnessHomes::default();
1953 match method {
1954 "harness.v1.harnesses.settings" => {
1955 let params = decode::<HarnessSettingsParams>(params)?;
1956 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1957 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1958 serde_json::to_value(report)
1959 .map_err(|error| ServiceError::Operation(error.to_string()))
1960 }
1961 "harness.v1.harnesses.configure" => {
1962 let params = decode::<ConfigureHarnessParams>(params)?;
1963 let report = crate::configure_harness_interop_settings(
1964 &homes,
1965 ¶ms.harness,
1966 ¶ms.changes,
1967 params.expected_revision.as_deref(),
1968 )
1969 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1970 serde_json::to_value(report)
1971 .map_err(|error| ServiceError::Operation(error.to_string()))
1972 }
1973 _ => Err(ServiceError::MethodNotFound),
1974 }
1975 }
1976
1977 fn insert_runtime(
1978 &mut self,
1979 runtime: Box<dyn RuntimeConnection>,
1980 ) -> std::result::Result<Value, ServiceError> {
1981 let connection = format!("runtime-{}", self.next_runtime);
1982 self.next_runtime += 1;
1983 let handle = runtime.handle().clone();
1984 self.runtime_sequences
1985 .entry(handle.runtime_id.clone())
1986 .or_insert(0);
1987 self.runtimes.insert(connection.clone(), runtime);
1988 Ok(json!({"connection": connection, "handle": handle}))
1989 }
1990
1991 #[cfg(feature = "adapter-api")]
1992 async fn insert_hosted_runtime(
1993 &mut self,
1994 runtime: Box<dyn RuntimeConnection>,
1995 capabilities: crate::RuntimeCapabilities,
1996 workspace: PathBuf,
1997 fresh: bool,
1998 ) -> std::result::Result<Value, ServiceError> {
1999 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
2000 let token: std::sync::Arc<str> = crate::server::generate_token().into();
2001 let server = crate::server::run_frontend_http(
2002 host.clone(),
2003 host.frontend_sender(),
2004 "127.0.0.1:0",
2005 token.clone(),
2006 connection.handle().runtime_id.clone(),
2007 )
2008 .await
2009 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2010 let source = LiveRuntimeSource {
2011 harness: connection.handle().harness.as_str().to_string(),
2012 session_id: connection.handle().runtime_id.clone(),
2013 workspace: workspace.clone(),
2014 };
2015 let registration = register_live_runtime_with_metadata(
2016 connection.handle().runtime_id.clone(),
2017 source.clone(),
2018 format!("http://{}", server.address()),
2019 token.to_string(),
2020 crate::LiveRuntimeMetadata {
2021 child_pid: match &connection.handle().endpoint {
2022 crate::RuntimeEndpoint::LocalProcess { pid, command, .. }
2023 if connection.handle().harness.as_str() != "claude-code"
2024 || command.windows(2).any(|args| {
2025 args[0] == "--name"
2026 && args[1].starts_with(crate::claude_relay::RELAY_NAME_PREFIX)
2027 }) =>
2028 {
2029 *pid
2030 }
2031 _ => None,
2032 },
2033 endpoint_capabilities: vec!["http".into(), "acp".into()],
2034 ..Default::default()
2035 },
2036 )
2037 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2038 let endpoint = registration.endpoint().to_string();
2039 let launch = StructuredLaunch {
2040 cwd: workspace,
2041 program: std::env::current_exe()
2045 .ok()
2046 .map(|path| path.to_string_lossy().into_owned())
2047 .unwrap_or_else(|| "supercode".into()),
2048 arguments: vec![
2049 "open".into(),
2050 endpoint,
2051 "--harness".into(),
2052 source.harness,
2053 "--session".into(),
2054 source.session_id,
2055 ],
2056 env: BTreeMap::new(),
2057 };
2058 host.retain_registration(registration);
2059 let lease = HostedRuntimeLease {
2060 connection,
2061 _host: host,
2062 _server: server,
2063 };
2064 let opened = self.insert_runtime(Box::new(lease))?;
2065 let connection_id = opened["connection"]
2066 .as_str()
2067 .expect("insert_runtime returns a connection id")
2068 .to_string();
2069 self.terminal_launches.insert(connection_id, launch);
2070 Ok(opened)
2071 }
2072
2073 #[cfg(not(feature = "adapter-api"))]
2074 async fn insert_hosted_runtime(
2075 &mut self,
2076 runtime: Box<dyn RuntimeConnection>,
2077 _capabilities: crate::RuntimeCapabilities,
2078 _workspace: PathBuf,
2079 _fresh: bool,
2080 ) -> std::result::Result<Value, ServiceError> {
2081 self.insert_runtime(runtime)
2082 }
2083
2084 fn runtime_mut(
2085 &mut self,
2086 connection: &str,
2087 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2088 if self.runtimes_in_flight.contains(connection) {
2089 return Err(self.lent_out(connection));
2090 }
2091 self.runtimes.get_mut(connection).ok_or_else(|| {
2092 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2093 })
2094 }
2095
2096 fn lent_out(&self, connection: &str) -> ServiceError {
2100 ServiceError::Operation(format!(
2101 "runtime connection `{connection}`: a harness turn is already in progress"
2102 ))
2103 }
2104
2105 fn lend_runtime(
2108 &mut self,
2109 connection: &str,
2110 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2111 if self.runtimes_in_flight.contains(connection) {
2112 return Err(self.lent_out(connection));
2113 }
2114 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2115 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2116 })?;
2117 self.runtimes_in_flight.insert(connection.to_string());
2118 Ok(runtime)
2119 }
2120
2121 fn surrender_runtime_of(
2135 &mut self,
2136 params: &RuntimeCloseParams,
2137 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2138 if params.connection.is_empty() {
2139 let expected = params.runtime_id.as_deref().ok_or_else(|| {
2140 ServiceError::InvalidParams("close requires connection or runtime_id".into())
2141 })?;
2142 let connection = self
2143 .runtimes
2144 .iter()
2145 .find(|(_, runtime)| runtime.handle().runtime_id == expected)
2146 .map(|(connection, _)| connection.clone())
2147 .ok_or_else(|| {
2148 ServiceError::InvalidParams(format!("unknown runtime `{expected}`"))
2149 })?;
2150 return self.surrender_runtime(&connection);
2151 }
2152 if let Some(expected) = params.runtime_id.as_deref() {
2153 let held = self
2154 .runtimes
2155 .get(¶ms.connection)
2156 .map(|runtime| runtime.handle().runtime_id.clone());
2157 if held.as_deref() != Some(expected) {
2158 return Err(ServiceError::InvalidParams(format!(
2159 "runtime connection `{}` does not hold runtime `{expected}`",
2160 params.connection
2161 )));
2162 }
2163 }
2164 self.surrender_runtime(¶ms.connection)
2165 }
2166
2167 fn surrender_runtime(
2168 &mut self,
2169 connection: &str,
2170 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2171 if self.runtimes_in_flight.contains(connection) {
2172 return Err(self.lent_out(connection));
2173 }
2174 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2175 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2176 })?;
2177 let process_group = runtime_process_group(runtime.handle());
2178 let runtime_id = runtime.handle().runtime_id.clone();
2179 self.terminal_launches.remove(connection);
2180 self.runtime_sequences.remove(&runtime_id);
2181 self.approvals.forget(connection);
2182 Ok((runtime, process_group))
2183 }
2184
2185 pub fn kill_all_runtime_groups(&self) -> usize {
2196 self.runtimes
2197 .values()
2198 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2199 .count()
2200 }
2201
2202 async fn mutate_session(
2213 &mut self,
2214 verb: crate::SessionVerb,
2215 params: Value,
2216 ) -> std::result::Result<Value, ServiceError> {
2217 let mutation = decode::<crate::SessionMutation>(params)?;
2218 let door = crate::sessions_control::door(&mutation.harness, verb)
2219 .map_err(session_control_error)?;
2220 let outcome = match door {
2221 #[cfg(not(feature = "adapter-api"))]
2225 crate::SessionDoor::Live(command) => {
2226 return Err(ServiceError::Operation(format!(
2227 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2228 session, which needs this build's `adapter-api` feature",
2229 mutation.harness,
2230 verb.as_str()
2231 )));
2232 }
2233 #[cfg(feature = "adapter-api")]
2234 crate::SessionDoor::Live(command) => {
2235 let connection = mutation
2236 .connection
2237 .clone()
2238 .filter(|value| !value.trim().is_empty())
2239 .ok_or_else(|| {
2240 ServiceError::InvalidParams(format!(
2241 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2242 driven session: pass the `connection` of an open runtime \
2243 (`harness.v1.runtimes.start`)",
2244 mutation.harness,
2245 verb.as_str()
2246 ))
2247 })?;
2248 let runtime = self.runtime_mut(&connection)?;
2249 let session = live_session_name(runtime.as_ref(), &mutation);
2250 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2256 .await;
2257 }
2258 _ => run_session_mutation(verb, &mutation).await?,
2259 };
2260 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2261 }
2262
2263 async fn inventory_call(
2267 &self,
2268 method: &str,
2269 params: Value,
2270 ) -> std::result::Result<Value, ServiceError> {
2271 run_inventory(self.inventory_work(method, params)?).await
2272 }
2273
2274 fn inventory_work(
2280 &self,
2281 method: &str,
2282 params: Value,
2283 ) -> std::result::Result<InventoryWork, ServiceError> {
2284 let mut params = decode::<HarnessInventoryParams>(params)?;
2285 if method == "harness.v1.harnesses.probe" {
2286 let harness = params.harness.take().ok_or_else(|| {
2287 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2288 })?;
2289 params.harnesses = vec![harness];
2290 }
2291 let selected = params
2292 .harnesses
2293 .iter()
2294 .map(HarnessId::as_str)
2295 .collect::<std::collections::BTreeSet<_>>();
2296 let supported = harness_support_registry()
2297 .harnesses
2298 .into_iter()
2299 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2300 .collect::<Vec<_>>();
2301 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2302 let known = supported
2303 .iter()
2304 .map(|harness| harness.id.as_str())
2305 .collect::<std::collections::BTreeSet<_>>();
2306 let missing = params
2307 .harnesses
2308 .iter()
2309 .filter(|id| !known.contains(id.as_str()))
2310 .map(HarnessId::as_str)
2311 .collect::<Vec<_>>();
2312 return Err(ServiceError::InvalidParams(format!(
2313 "unknown harness(es): {}",
2314 missing.join(", ")
2315 )));
2316 }
2317 let global_counts = params
2318 .include_sessions
2319 .then(|| self.session_counts(None, ¶ms.harnesses));
2320 let workspace_counts = params
2321 .include_sessions
2322 .then(|| {
2323 params
2324 .workspace
2325 .as_deref()
2326 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2327 })
2328 .flatten();
2329 Ok(InventoryWork {
2330 params,
2331 supported,
2332 global_counts,
2333 workspace_counts,
2334 })
2335 }
2336
2337 #[cfg(feature = "adapter-api")]
2338 async fn harness_authentication_call(
2339 &self,
2340 method: &str,
2341 params: Value,
2342 ) -> std::result::Result<Value, ServiceError> {
2343 match method {
2344 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2345 let params = decode::<HarnessAuthenticationParams>(params)?;
2346 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2347 .map_err(|error| ServiceError::Operation(error.to_string()))
2348 }
2349 "harness.v1.harnesses.auth.begin" => {
2350 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2351 let cwd = params
2352 .cwd
2353 .or_else(|| std::env::current_dir().ok())
2354 .unwrap_or_else(|| PathBuf::from("."));
2355 let plan = crate::harness_authentication_plan(
2356 ¶ms.harness,
2357 params.environment,
2358 params.method,
2359 &cwd,
2360 )
2361 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2362 serde_json::to_value(plan)
2363 .map_err(|error| ServiceError::Operation(error.to_string()))
2364 }
2365 _ => Err(ServiceError::MethodNotFound),
2366 }
2367 }
2368
2369 fn session_counts(
2370 &self,
2371 workspace: Option<&Path>,
2372 harnesses: &[HarnessId],
2373 ) -> BTreeMap<String, usize> {
2374 let mut counts = BTreeMap::new();
2375 for session in self
2376 .catalog
2377 .discover(&DiscoveryQuery {
2378 workspace: workspace.map(Path::to_path_buf),
2379 harnesses: harnesses.to_vec(),
2380 ..DiscoveryQuery::default()
2381 })
2382 .unwrap_or_default()
2383 {
2384 *counts
2385 .entry(session.locator.harness.as_str().to_string())
2386 .or_insert(0) += 1;
2387 }
2388 counts
2389 }
2390}
2391
2392#[async_trait::async_trait]
2393impl SdkService for HarnessSessionService {
2394 fn capabilities(&self) -> SdkCapabilities {
2395 SdkCapabilities::default()
2396 }
2397
2398 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2399 if request.operation == SdkOperation::Events {
2400 let events = self
2401 .poll_sdk_events()
2402 .await
2403 .into_iter()
2404 .map(|(_, event)| event)
2405 .collect::<Vec<_>>();
2406 return serde_json::to_value(events).map_err(|error| {
2407 SdkError::new(
2408 SdkErrorCode::Execution,
2409 request.operation,
2410 error.to_string(),
2411 )
2412 });
2413 }
2414 if self.runtimes.is_empty()
2415 && matches!(
2416 request.operation,
2417 SdkOperation::Input
2418 | SdkOperation::Interrupt
2419 | SdkOperation::Steer
2420 | SdkOperation::Respond
2421 | SdkOperation::Close
2422 )
2423 {
2424 return Err(SdkError::unsupported(request.operation));
2425 }
2426 let method = request
2427 .operation
2428 .method()
2429 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2430 let result = match request.operation {
2431 SdkOperation::Discover
2432 | SdkOperation::Load
2433 | SdkOperation::Export
2434 | SdkOperation::ProfilesList
2435 | SdkOperation::ProfilesGet
2436 | SdkOperation::ProfilesCreate
2437 | SdkOperation::ProfilesDelete
2438 | SdkOperation::SkillsList
2439 | SdkOperation::SkillsInstall
2440 | SdkOperation::SkillsRemove
2441 | SdkOperation::ChannelsList
2442 | SdkOperation::RoutesList
2443 | SdkOperation::TriggersList
2444 | SdkOperation::ChannelsStatus
2445 | SdkOperation::MemoryShow
2446 | SdkOperation::MemorySearch
2447 | SdkOperation::JobsList
2448 | SdkOperation::JobsGet
2449 | SdkOperation::JobsCreate
2450 | SdkOperation::JobsUpdate
2451 | SdkOperation::JobsPause
2452 | SdkOperation::JobsResume
2453 | SdkOperation::JobsRun
2454 | SdkOperation::JobsDelete
2455 | SdkOperation::JobsNotepad
2456 | SdkOperation::JobsNotepadSet
2457 | SdkOperation::JobsNotepadDelete
2458 | SdkOperation::RunsList
2459 | SdkOperation::RunsGet
2460 | SdkOperation::ApprovalsList
2461 | SdkOperation::OrchestrationLoad
2462 | SdkOperation::OrchestrationSave
2463 | SdkOperation::OrchestrationCompile
2464 | SdkOperation::OrchestrationDecompile
2465 | SdkOperation::OrchestrationImport
2466 | SdkOperation::OrchestrationExport
2467 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2468 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2471 SdkOperation::Start
2472 | SdkOperation::Resume
2473 | SdkOperation::Input
2474 | SdkOperation::Interrupt
2475 | SdkOperation::Steer
2476 | SdkOperation::Respond
2477 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2478 SdkOperation::SessionsNew => {
2483 self.mutate_session(crate::SessionVerb::New, request.params)
2484 .await
2485 }
2486 SdkOperation::SessionsReset => {
2487 self.mutate_session(crate::SessionVerb::Reset, request.params)
2488 .await
2489 }
2490 SdkOperation::SessionsArchive => {
2491 self.mutate_session(crate::SessionVerb::Archive, request.params)
2492 .await
2493 }
2494 SdkOperation::SessionsDelete => {
2495 self.mutate_session(crate::SessionVerb::Delete, request.params)
2496 .await
2497 }
2498 SdkOperation::Events => unreachable!("handled before method dispatch"),
2499 };
2500 result.map_err(|error| sdk_error(request.operation, error))
2501 }
2502
2503 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2504 Ok(self
2505 .poll_sdk_events()
2506 .await
2507 .into_iter()
2508 .map(|(_, event)| event)
2509 .collect())
2510 }
2511}
2512
2513#[cfg(feature = "adapter-api")]
2514struct HostedRuntimeLease {
2515 connection: HostedHarnessConnection,
2516 _host: std::sync::Arc<HostedHarnessRuntime>,
2517 _server: crate::server::FrontendHttpServer,
2518}
2519
2520#[async_trait::async_trait]
2521#[cfg(feature = "adapter-api")]
2522impl RuntimeConnection for HostedRuntimeLease {
2523 fn handle(&self) -> &crate::RuntimeHandle {
2524 self.connection.handle()
2525 }
2526
2527 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2528 self.connection.send_input(input).await
2529 }
2530
2531 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2532 self.connection.next_event().await
2533 }
2534
2535 async fn interrupt(&mut self) -> crate::Result<()> {
2536 self.connection.interrupt().await
2537 }
2538
2539 async fn steer(&mut self, text: String) -> crate::Result<()> {
2542 self.connection.steer(text).await
2543 }
2544
2545 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2546 self.connection.respond(request_id, response).await
2547 }
2548
2549 async fn close(&mut self) -> crate::Result<()> {
2550 self.connection.close().await
2551 }
2552}
2553
2554struct InventoryWork {
2557 params: HarnessInventoryParams,
2558 supported: Vec<crate::HarnessSupportDescriptor>,
2559 global_counts: Option<BTreeMap<String, usize>>,
2560 workspace_counts: Option<BTreeMap<String, usize>>,
2561}
2562
2563async fn run_session_mutation(
2570 verb: crate::SessionVerb,
2571 mutation: &crate::SessionMutation,
2572) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2573 let door =
2580 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2581 if let crate::SessionDoor::Http = door {
2582 return crate::sessions_control::mutate(verb, mutation)
2583 .await
2584 .map_err(session_control_error);
2585 }
2586 let mutation = mutation.clone();
2587 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2588 .await
2589 .map_err(|error| {
2590 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2591 })?
2592 .map_err(session_control_error)
2593}
2594
2595async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2598 let InventoryWork {
2599 params,
2600 supported,
2601 global_counts,
2602 workspace_counts,
2603 } = work;
2604 let probes = supported.into_iter().map(|descriptor| {
2605 let global = global_counts
2606 .as_ref()
2607 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2608 let workspace = workspace_counts
2609 .as_ref()
2610 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2611 probe_harness(descriptor, ¶ms, global, workspace)
2612 });
2613 let harnesses = futures::future::join_all(probes).await;
2614 serde_json::to_value(HarnessInventoryReport {
2615 probe: params.probe,
2616 workspace: params.workspace,
2617 harnesses,
2618 })
2619 .map_err(|error| ServiceError::Operation(error.to_string()))
2620}
2621
2622async fn probe_harness(
2623 descriptor: crate::HarnessSupportDescriptor,
2624 params: &HarnessInventoryParams,
2625 global: Option<usize>,
2626 workspace: Option<usize>,
2627) -> LocalHarness {
2628 let launch = descriptor.runtime.default_launch.as_ref();
2629 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2634 .then(crate::orchestrator::daemon_entry)
2635 .and_then(Result::ok);
2636 let executable = match &orchestrator_entry {
2637 Some(entry) => Some(entry.clone()),
2638 None => launch.and_then(|launch| find_executable(&launch.program)),
2639 };
2640 let installed = executable.is_some();
2641 let version = if params.skip_versions || orchestrator_entry.is_some() {
2642 None
2645 } else {
2646 match executable.as_deref() {
2647 Some(path) => executable_version(path).await,
2648 None => None,
2649 }
2650 };
2651 let configured = auth_evidence(descriptor.id.as_str());
2652 let mut auth = if configured {
2653 HarnessAuthState::Configured
2654 } else if matches!(
2655 descriptor.id.as_str(),
2656 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2657 ) {
2658 HarnessAuthState::Required
2663 } else {
2664 HarnessAuthState::Unknown
2665 };
2666 let mut runtime = if installed {
2667 HarnessRuntimeState::Degraded
2668 } else {
2669 HarnessRuntimeState::Unavailable
2670 };
2671 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2672 let mut reason = (!installed).then(|| {
2673 if is_orchestrator {
2674 format!(
2675 "{} is supported but its daemon entry `{}` was not found",
2676 descriptor.display_name,
2677 crate::orchestrator::DAEMON_ENTRY
2678 )
2679 } else {
2680 format!(
2681 "{} is supported but `{}` was not found on PATH",
2682 descriptor.display_name,
2683 launch
2684 .map(|launch| launch.program.as_str())
2685 .unwrap_or("executable")
2686 )
2687 }
2688 });
2689 let mut repair = (!installed).then(|| {
2690 if is_orchestrator {
2691 format!(
2692 "Install the `supercode-orchestrator` package so `{}` resolves.",
2693 crate::orchestrator::DAEMON_ENTRY
2694 )
2695 } else {
2696 format!(
2697 "Install {} and ensure `{}` is on PATH.",
2698 descriptor.display_name,
2699 launch
2700 .map(|launch| launch.program.as_str())
2701 .unwrap_or("its executable")
2702 )
2703 }
2704 });
2705
2706 if installed && params.probe == HarnessProbeLevel::Handshake {
2707 let backend_params = RuntimeBackendParams {
2708 harness: descriptor.id.clone(),
2709 protocol: None,
2710 launch: None,
2711 base_url: None,
2712 policy: RuntimePolicy::Default,
2713 };
2714 match runtime_backend(&backend_params) {
2715 Ok(backend) => {
2716 let cwd = params
2717 .workspace
2718 .clone()
2719 .or_else(|| std::env::current_dir().ok())
2720 .unwrap_or_else(|| PathBuf::from("."));
2721 let isolated = descriptor
2722 .runtime
2723 .default_launch
2724 .clone()
2725 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2726 let Some(isolated) = isolated else {
2727 reason = Some(
2728 "No-prompt runtime handshake could not create its isolated harness home."
2729 .into(),
2730 );
2731 repair = Some(
2732 "Check temporary-directory permissions, then run the handshake probe again."
2733 .into(),
2734 );
2735 let running = probe_running_instance(descriptor.id.as_str());
2736 return LocalHarness {
2737 gateway: gateway_health(
2738 descriptor.id.as_str(),
2739 installed,
2740 running.as_ref(),
2741 version.as_deref(),
2742 ),
2743 id: descriptor.id,
2744 display_name: descriptor.display_name,
2745 supported: true,
2746 installed,
2747 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2748 version,
2749 auth,
2750 runtime,
2751 protocol: descriptor.runtime.protocol,
2752 capabilities: descriptor.runtime.capabilities.clone(),
2753 effective_capabilities: descriptor.runtime.capabilities,
2754 sessions: HarnessSessionCounts { global, workspace },
2755 running,
2756 reason,
2757 repair,
2758 };
2759 };
2760 match tokio::time::timeout(
2761 Duration::from_secs(30),
2762 backend.start(RuntimeStartRequest {
2763 cwd,
2764 launch: Some(isolated.launch.clone()),
2765 mcp_servers: Vec::new(),
2766 approval_policy: None,
2767 }),
2768 )
2769 .await
2770 {
2771 Ok(Ok(mut connection)) => {
2772 match stabilize_handshake(connection.as_mut()).await {
2773 Ok(()) => {
2774 auth = HarnessAuthState::Ready;
2775 runtime = HarnessRuntimeState::Ready;
2776 reason = Some(
2777 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2778 .into(),
2779 );
2780 repair = None;
2781 }
2782 Err(message) => {
2783 auth = if looks_like_auth_error(&message) {
2784 HarnessAuthState::Required
2785 } else if configured {
2786 HarnessAuthState::Configured
2787 } else {
2788 HarnessAuthState::Unknown
2789 };
2790 reason = Some(format!(
2791 "No-prompt runtime handshake became unhealthy during startup: {message}"
2792 ));
2793 repair = Some(if auth == HarnessAuthState::Required {
2794 format!(
2795 "Run `{}` interactively once and complete sign-in, then probe again.",
2796 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2797 )
2798 } else {
2799 "Run the harness directly to inspect its startup failure, then probe again."
2800 .into()
2801 });
2802 }
2803 }
2804 let _ =
2805 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2806 }
2807 Ok(Err(error)) => {
2808 let message = truncate_text(&error.to_string(), 500);
2809 auth = if looks_like_auth_error(&message) {
2810 HarnessAuthState::Required
2811 } else if configured {
2812 HarnessAuthState::Configured
2813 } else {
2814 HarnessAuthState::Unknown
2815 };
2816 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2817 repair = Some(if auth == HarnessAuthState::Required {
2818 format!(
2819 "Run `{}` interactively once and complete sign-in, then probe again.",
2820 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2821 )
2822 } else {
2823 "Check the harness installation and run the handshake probe again."
2824 .into()
2825 });
2826 }
2827 Err(_) => {
2828 reason =
2829 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2830 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2831 }
2832 }
2833 let _ = isolated.cleanup();
2842 tokio::time::sleep(Duration::from_millis(250)).await;
2843 if let Err(error) = isolated.cleanup() {
2844 auth = if configured {
2845 HarnessAuthState::Configured
2846 } else {
2847 HarnessAuthState::Unknown
2848 };
2849 runtime = HarnessRuntimeState::Degraded;
2850 reason = Some(format!(
2851 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2852 ));
2853 repair = Some(
2854 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2855 .into(),
2856 );
2857 }
2858 }
2859 Err(error) => {
2860 reason = Some(error_message(error));
2861 }
2862 }
2863 } else if installed && configured {
2864 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2865 } else if installed && auth == HarnessAuthState::Required {
2866 reason = Some("Executable found, but no native authentication evidence is present.".into());
2867 repair = Some(format!(
2868 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2869 descriptor.id.as_str()
2870 ));
2871 } else if installed {
2872 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2873 repair = Some(format!(
2874 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2875 launch
2876 .map(|launch| launch.program.as_str())
2877 .unwrap_or("the harness")
2878 ));
2879 }
2880
2881 let effective_capabilities = if installed {
2882 descriptor.runtime.capabilities.clone()
2883 } else {
2884 unavailable_capabilities()
2885 };
2886 let running = probe_running_instance(descriptor.id.as_str());
2887 LocalHarness {
2888 gateway: gateway_health(
2889 descriptor.id.as_str(),
2890 installed,
2891 running.as_ref(),
2892 version.as_deref(),
2893 ),
2894 id: descriptor.id,
2895 display_name: descriptor.display_name,
2896 supported: true,
2897 installed,
2898 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2899 version,
2900 auth,
2901 runtime,
2902 protocol: descriptor.runtime.protocol,
2903 capabilities: descriptor.runtime.capabilities,
2904 effective_capabilities,
2905 sessions: HarnessSessionCounts { global, workspace },
2906 running,
2907 reason,
2908 repair,
2909 }
2910}
2911
2912async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2913 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2914 loop {
2915 let now = tokio::time::Instant::now();
2916 if now >= deadline {
2917 return Ok(());
2918 }
2919 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2920 Err(_) => return Ok(()),
2921 Ok(Ok(Some(event))) => {
2922 if let Some(message) = handshake_event_failure(&event) {
2923 return Err(truncate_text(&message, 500));
2924 }
2925 }
2926 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2927 Ok(Err(error)) => return Err(error.to_string()),
2928 }
2929 }
2930}
2931
2932fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2933 let detail = event
2934 .payload
2935 .get("message")
2936 .or_else(|| event.payload.get("line"))
2937 .and_then(Value::as_str)
2938 .unwrap_or(event.kind.as_str());
2939 match event.kind.as_str() {
2940 "transport_closed" => Some("runtime transport closed during startup".into()),
2941 "transport_error" => Some(format!("runtime transport error: {detail}")),
2942 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2943 _ => None,
2948 }
2949}
2950
2951fn indexed_claude_window(
2952 locator: &SessionLocator,
2953 options: &SessionLoadOptions,
2954) -> std::result::Result<Option<Value>, ServiceError> {
2955 use supercode_interchange::session::ClaudeReadIndex;
2956 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2959 || options.include_subagents != Some(false)
2960 {
2961 return Ok(None);
2962 }
2963 let crate::StorageLocator::File { path } = &locator.storage else {
2964 return Ok(None);
2965 };
2966 if !ClaudeReadIndex::supports(path)
2967 .map_err(|error| ServiceError::Operation(error.to_string()))?
2968 {
2969 return Ok(None);
2970 }
2971 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2972 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2973 let total = index.len();
2974 let (offset, end) = projected_message_window(total, options);
2975 let session = index
2976 .read_messages(offset..end)
2977 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2978 let summary = index
2979 .read_summary()
2980 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2981 let selected_options = SessionLoadOptions {
2982 message_offset: None,
2983 message_limit: None,
2984 message_tail: None,
2985 ..options.clone()
2986 };
2987 let mut selected = projected_session_json(&session, &selected_options);
2988 selected["raw_record_count"] = json!(index.raw_record_count());
2989 Ok(Some(json!({
2990 "session": selected,
2991 "summary": projected_session_summary(&summary, options),
2992 "window": {
2993 "has_more": offset > 0 || end < total, "has_newer": end < total,
2994 "has_older": offset > 0, "newer_items": index.item_count(end..total),
2995 "offset": offset, "older_items": index.item_count(0..offset),
2996 "returned": end - offset, "total_messages": total,
2997 }
2998 })))
2999}
3000
3001fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
3002 let total_messages = session.messages.len();
3003 let (offset, end) = projected_message_window(total_messages, options);
3004 json!({
3005 "session": projected_session_json(session, options),
3006 "summary": projected_session_summary(session, options),
3007 "window": {
3008 "has_more": offset > 0 || end < total_messages,
3009 "has_newer": end < total_messages,
3010 "has_older": offset > 0,
3011 "newer_items": normalized_item_count(&session.messages[end..]),
3012 "offset": offset,
3013 "older_items": normalized_item_count(&session.messages[..offset]),
3014 "returned": end.saturating_sub(offset),
3015 "total_messages": total_messages,
3016 }
3017 })
3018}
3019
3020fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
3021 messages
3022 .iter()
3023 .map(|message| {
3024 let conversation = usize::from(
3025 matches!(message.role, Role::Assistant | Role::User)
3026 && message_has_content(message),
3027 );
3028 let tool_result =
3029 usize::from(message.role == Role::Tool && message_has_content(message));
3030 conversation + tool_result + message.tool_calls().len()
3031 })
3032 .sum()
3033}
3034
3035fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
3036 let mut conversational = session.messages.iter().filter(|message| {
3037 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
3038 });
3039 let first_message = conversational.clone().next();
3040 let last_message = conversational.next_back();
3041 let mut assistant = session
3042 .messages
3043 .iter()
3044 .filter(|message| message.role == Role::Assistant && message_has_content(message));
3045 let first_assistant_message = assistant.clone().next();
3046 let last_assistant_message = assistant.next_back();
3047 let end_of_turn = session
3048 .messages
3049 .iter()
3050 .rev()
3051 .find(|message| message.role != Role::System)
3052 .is_some_and(|message| {
3053 message.role == Role::Assistant
3054 && message_has_content(message)
3055 && message.tool_calls().is_empty()
3056 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
3058 });
3059 let project = |message: Option<&crate::ChatMessage>| {
3060 message.map(|message| project_inline_media(message_json(message), options))
3061 };
3062 json!({
3063 "end_of_turn": end_of_turn,
3064 "first_assistant_message": project(first_assistant_message),
3065 "first_message": project(first_message),
3066 "last_assistant_message": project(last_assistant_message),
3067 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3068 "last_message": project(last_message),
3069 })
3070}
3071
3072fn message_has_content(message: &crate::ChatMessage) -> bool {
3073 message
3074 .content
3075 .as_deref()
3076 .is_some_and(|content| !content.trim().is_empty())
3077 || message
3078 .content_parts
3079 .as_ref()
3080 .is_some_and(|parts| !parts.is_empty())
3081}
3082
3083fn message_text(message: &crate::ChatMessage) -> String {
3084 if let Some(content) = &message.content {
3085 return content.clone();
3086 }
3087 message
3088 .content_parts
3089 .as_ref()
3090 .into_iter()
3091 .flatten()
3092 .filter_map(|part| part.get("text").and_then(Value::as_str))
3093 .collect::<Vec<_>>()
3094 .join("\n")
3095}
3096
3097fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3098 let (offset, end) = projected_message_window(session.messages.len(), options);
3099 let messages = session.messages[offset..end]
3100 .iter()
3101 .map(|message| project_inline_media(message_json(message), options))
3102 .collect::<Vec<_>>();
3103 let subagents = if options.include_subagents.unwrap_or(true) {
3104 let subagent_options = SessionLoadOptions {
3109 message_limit: None,
3110 message_offset: None,
3111 message_tail: None,
3112 ..options.clone()
3113 };
3114 session
3115 .subagents
3116 .iter()
3117 .map(|subagent| projected_session_json(subagent, &subagent_options))
3118 .collect::<Vec<_>>()
3119 } else {
3120 Vec::new()
3121 };
3122 json!({
3123 "source": match session.meta.source {
3124 SessionSource::ClaudeCode => "claude_code",
3125 SessionSource::Codex => "codex",
3126 SessionSource::Gemini => "gemini",
3127 SessionSource::Goose => "goose",
3128 SessionSource::Grok => "grok",
3129 SessionSource::Native => "native",
3130 SessionSource::OpenClaw => "openclaw",
3131 SessionSource::Hermes => "hermes",
3132 SessionSource::OpenCode => "opencode",
3133 SessionSource::Pi => "pi",
3134 },
3135 "session_id": session.meta.session_id,
3136 "ended_at": session.meta.ended_at,
3137 "end_reason": session.meta.end_reason,
3138 "model": session.meta.model,
3139 "resolved_profile": recorded_session_profile(session),
3140 "resolved_profile_history": recorded_session_profile_history(session),
3141 "cwd": session.meta.cwd,
3142 "system_prompt": session.meta.system_prompt,
3143 "agent_id": session.meta.agent_id,
3144 "parent_tool_use_id": session.meta.parent_tool_use_id,
3145 "lineage": session.meta.lineage,
3146 "messages": messages,
3147 "subagents": subagents,
3148 "raw_record_count": session.raw.len(),
3149 "parse_error_lines": session.parse_error_lines,
3150 })
3151}
3152
3153pub fn recorded_config(locator: &SessionLocator) -> crate::Result<Value> {
3161 let session = load_session(locator)?;
3162 let mut config = json!({"harness": locator.harness, "session_id": session.meta.session_id,
3163 "cwd": session.meta.cwd, "model": session.meta.model,
3164 "approval_policy": null, "sandbox_policy": null, "permission_mode": null});
3165 for header in &session.meta.codex_headers {
3166 let payload = header.get("payload").unwrap_or(header);
3167 for key in ["cwd", "model", "approval_policy", "sandbox_policy"] {
3168 if let Some(value) = payload.get(key) {
3169 config[key] = value.clone();
3170 }
3171 }
3172 }
3173 if let Ok(manifest) = crate::claude_runtime_state::ClaudeRuntimeManifest::from_session(&session)
3175 {
3176 config["permission_mode"] = json!(manifest.posture.permission_mode);
3177 }
3178 Ok(config)
3179}
3180
3181fn recorded_session_usage(session: &Session) -> Vec<Value> {
3182 let mut days: std::collections::BTreeMap<(String, String), [u64; 4]> =
3183 std::collections::BTreeMap::new();
3184 let mut add =
3185 |at: Option<&str>, model: &str, input: u64, output: u64, read: u64, write: u64| {
3186 let day = at
3187 .and_then(|t| t.get(..10))
3188 .unwrap_or("unknown")
3189 .to_string();
3190 let entry = days.entry((day, model.to_string())).or_insert([0; 4]);
3191 entry[0] += input;
3192 entry[1] += output;
3193 entry[2] += read;
3194 entry[3] += write;
3195 };
3196 let records = session
3197 .raw
3198 .iter()
3199 .filter_map(|line| serde_json::from_str::<Value>(line).ok());
3200 let n = |v: Option<&Value>| v.and_then(Value::as_u64).unwrap_or(0);
3201 match session.meta.source {
3202 SessionSource::ClaudeCode => {
3203 let mut responses: std::collections::BTreeMap<String, Value> =
3204 std::collections::BTreeMap::new();
3205 for record in records {
3206 let Some(message) = record.get("message") else {
3207 continue;
3208 };
3209 let (Some(id), Some(_)) = (
3210 message.get("id").and_then(Value::as_str),
3211 message.get("usage"),
3212 ) else {
3213 continue;
3214 };
3215 if message.get("model").and_then(Value::as_str) == Some("<synthetic>") {
3216 continue;
3217 }
3218 responses.insert(id.to_string(), record.clone());
3219 }
3220 for record in responses.values() {
3221 let message = &record["message"];
3222 let usage = &message["usage"];
3223 add(
3224 record.get("timestamp").and_then(Value::as_str),
3225 message
3226 .get("model")
3227 .and_then(Value::as_str)
3228 .unwrap_or("unknown"),
3229 n(usage.get("input_tokens")),
3230 n(usage.get("output_tokens")),
3231 n(usage.get("cache_read_input_tokens")),
3232 n(usage.get("cache_creation_input_tokens")),
3233 );
3234 }
3235 }
3236 SessionSource::Codex => {
3237 let mut model = session
3238 .meta
3239 .model
3240 .clone()
3241 .unwrap_or_else(|| "unknown".into());
3242 let mut previous = [0u64; 3];
3243 for record in records {
3244 let payload = record.get("payload").unwrap_or(&Value::Null);
3245 if record.get("type").and_then(Value::as_str) == Some("turn_context") {
3246 if let Some(m) = payload.get("model").and_then(Value::as_str) {
3247 model = m.to_string();
3248 }
3249 continue;
3250 }
3251 if payload.get("type").and_then(Value::as_str) != Some("token_count") {
3252 continue;
3253 }
3254 let Some(total) = payload.get("info").and_then(|i| i.get("total_token_usage"))
3255 else {
3256 continue;
3257 };
3258 let now = [
3259 n(total.get("input_tokens")),
3260 n(total.get("output_tokens")),
3261 n(total.get("cached_input_tokens")),
3262 ];
3263 if now == previous {
3264 continue;
3265 }
3266 add(
3267 record.get("timestamp").and_then(Value::as_str),
3268 &model,
3269 now[0].saturating_sub(previous[0]),
3270 now[1].saturating_sub(previous[1]),
3271 now[2].saturating_sub(previous[2]),
3272 0,
3273 );
3274 previous = now;
3275 }
3276 }
3277 _ => {}
3278 }
3279 days.into_iter()
3280 .map(|((day, model), [input, output, read, write])| {
3281 let mut usage = json!({"at": format!("{day}T00:00:00Z"), "model": model,
3282 "input_tokens": input, "output_tokens": output,
3283 "cache_read_tokens": read, "cache_write_tokens": write});
3284 if let Some(price) = crate::pricing::built_in(&model) {
3287 usage["cost"] = json!({"amount": price.cost_usd(input + read + write, output),
3288 "currency": "USD", "source": "rate_table"});
3289 }
3290 usage
3291 })
3292 .collect()
3293}
3294
3295fn recorded_session_profile(session: &Session) -> Value {
3298 let harness = match session.meta.source {
3299 SessionSource::ClaudeCode => "claude-code",
3300 SessionSource::Codex => "codex",
3301 SessionSource::Gemini => "gemini",
3302 SessionSource::Goose => "goose",
3303 SessionSource::Grok => "grok",
3304 SessionSource::Native => "native",
3305 SessionSource::OpenClaw => "openclaw",
3306 SessionSource::Hermes => "hermes",
3307 SessionSource::OpenCode => "opencode",
3308 SessionSource::Pi => "pi",
3309 };
3310 let mut model = session
3311 .messages
3312 .iter()
3313 .rev()
3314 .find_map(|message| message.metadata.get("model").cloned())
3315 .or_else(|| session.meta.model.clone());
3316 let mut cwd = session
3317 .meta
3318 .cwd
3319 .as_ref()
3320 .map(|path| path.to_string_lossy().into_owned());
3321 let mut version: Option<String> = None;
3322 let mut approval = Value::Null;
3323 let mut sandbox = Value::Null;
3324 for header in &session.meta.codex_headers {
3325 let payload = header.get("payload").unwrap_or(header);
3326 if let Some(value) = payload.get("cli_version").and_then(Value::as_str) {
3327 version = Some(value.to_owned());
3328 }
3329 if let Some(value) = payload.get("model").and_then(Value::as_str) {
3330 model = Some(value.to_owned());
3331 }
3332 if let Some(value) = payload.get("cwd").and_then(Value::as_str) {
3333 cwd = Some(value.to_owned());
3334 }
3335 if let Some(value) = payload.get("approval_policy") {
3336 approval = value.clone();
3337 }
3338 if let Some(value) = payload.get("sandbox_policy") {
3339 sandbox = value.clone();
3340 }
3341 }
3342 if session.meta.source == SessionSource::ClaudeCode {
3343 version = session
3344 .raw
3345 .iter()
3346 .rev()
3347 .filter_map(|line| serde_json::from_str::<Value>(line).ok())
3348 .find_map(|record| {
3349 record
3350 .get("version")
3351 .and_then(Value::as_str)
3352 .map(str::to_owned)
3353 });
3354 }
3355 json!({"worker":{"harness":harness,"model":model,"cwd":cwd,"binary":{"version":version}},
3356 "permissions":{"native":{"approval":approval,"sandbox":sandbox}},"source":"native-session"})
3357}
3358
3359fn recorded_session_profile_history(session: &Session) -> Vec<Value> {
3362 let mut profile = json!({"worker":{"harness":if session.meta.source == SessionSource::Codex {"codex"} else {"claude-code"},
3363 "model":null,"cwd":null,"binary":{"version":null}},
3364 "permissions":{"native":{"approval":null,"sandbox":null}},"source":"native-session"});
3365 let records: Box<dyn Iterator<Item = (usize, Value)> + '_> =
3366 if session.meta.source == SessionSource::Codex {
3367 Box::new(session.meta.codex_headers.iter().cloned().enumerate())
3368 } else if session.meta.source == SessionSource::ClaudeCode {
3369 Box::new(session.raw.iter().enumerate().filter_map(|(index, line)| {
3370 serde_json::from_str::<Value>(line)
3371 .ok()
3372 .map(|record| (index, record))
3373 }))
3374 } else {
3375 return Vec::new();
3376 };
3377 let mut history = Vec::new();
3378 for (index, record) in records {
3379 let before = profile.clone();
3380 let payload = record.get("payload").unwrap_or(&record);
3381 for (native, target) in [("cwd", "cwd"), ("model", "model")] {
3382 if let Some(value) = payload.get(native).and_then(Value::as_str) {
3383 profile["worker"][target] = json!(value);
3384 }
3385 }
3386 let version = payload
3387 .get("cli_version")
3388 .or_else(|| record.get("version"))
3389 .and_then(Value::as_str);
3390 if let Some(value) = version {
3391 profile["worker"]["binary"]["version"] = json!(value);
3392 }
3393 if let Some(value) = record
3394 .get("message")
3395 .and_then(|m| m.get("model"))
3396 .and_then(Value::as_str)
3397 {
3398 profile["worker"]["model"] = json!(value);
3399 }
3400 for (native, target) in [
3401 ("approval_policy", "approval"),
3402 ("sandbox_policy", "sandbox"),
3403 ] {
3404 if let Some(value) = payload.get(native) {
3405 profile["permissions"]["native"][target] = value.clone();
3406 }
3407 }
3408 if profile != before {
3409 history.push(json!({"native_key":format!("profile-record:{index}"),"at":record.get("timestamp").and_then(Value::as_str),"profile":profile}));
3410 }
3411 }
3412 history
3413}
3414
3415fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3416 if let Some(tail) = options.message_tail {
3417 return (total.saturating_sub(tail), total);
3418 }
3419 let offset = options.message_offset.unwrap_or(0).min(total);
3420 let end = options
3421 .message_limit
3422 .map(|limit| offset.saturating_add(limit).min(total))
3423 .unwrap_or(total);
3424 (offset, end)
3425}
3426
3427fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3428 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3429 return message;
3430 };
3431 for part in parts {
3432 let Some(url) = part
3433 .get("image_url")
3434 .and_then(|image| image.get("url"))
3435 .and_then(Value::as_str)
3436 else {
3437 continue;
3438 };
3439 let Some(rest) = url.strip_prefix("data:") else {
3440 continue;
3441 };
3442 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3443 continue;
3444 };
3445 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3446 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3447 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3448 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3449 || options
3450 .max_inline_media_bytes
3451 .is_some_and(|limit| decoded_bytes > limit);
3452 if should_elide {
3453 *part = json!({
3454 "type": "media_reference",
3455 "media_type": media_type,
3456 "encoding": "base64",
3457 "encoded_bytes": encoded.len(),
3458 "decoded_bytes": decoded_bytes,
3459 "omitted": true,
3460 });
3461 }
3462 }
3463 message
3464}
3465
3466#[derive(Deserialize)]
3467struct LocatorParams {
3468 locator: SessionLocator,
3469 #[serde(default)]
3482 fidelity: Option<Fidelity>,
3483 #[serde(default)]
3486 view: Option<SessionReadView>,
3487}
3488
3489#[derive(Deserialize)]
3490struct SessionReadView {
3491 #[serde(default)]
3494 tail_messages: Option<usize>,
3495 #[serde(default)]
3498 include_subagents: bool,
3499 #[serde(default)]
3501 display_history: bool,
3502 #[serde(default)]
3505 max_message_chars: Option<usize>,
3506}
3507
3508impl LocatorParams {
3509 fn read_fidelity(&self) -> Fidelity {
3510 self.fidelity.unwrap_or(Fidelity::Semantic)
3511 }
3512
3513 fn include_subagents(&self) -> bool {
3514 self.view
3515 .as_ref()
3516 .map(|view| view.include_subagents)
3517 .unwrap_or(true)
3518 }
3519
3520 fn tail_messages(&self) -> Option<usize> {
3521 self.view
3522 .as_ref()
3523 .and_then(|view| view.tail_messages)
3524 .map(|limit| limit.clamp(1, 5_000))
3525 }
3526
3527 fn display_history(&self) -> bool {
3528 self.view.as_ref().is_some_and(|view| view.display_history)
3529 }
3530
3531 fn max_message_chars(&self) -> Option<usize> {
3532 self.view
3533 .as_ref()
3534 .and_then(|view| view.max_message_chars)
3535 .map(|limit| limit.clamp(256, 64_000))
3536 }
3537
3538 fn bound_session(&self, session: &mut Session) {
3539 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3540 }
3541}
3542
3543#[derive(Debug, Clone, Copy, Default, Deserialize)]
3544#[serde(rename_all = "snake_case")]
3545enum InlineMediaMode {
3546 #[default]
3547 Full,
3548 Metadata,
3549}
3550
3551#[derive(Debug, Clone, Default, Deserialize)]
3552#[serde(default)]
3553struct SessionLoadOptions {
3554 include_subagents: Option<bool>,
3555 inline_media: InlineMediaMode,
3556 max_inline_media_bytes: Option<usize>,
3557 message_limit: Option<usize>,
3558 message_offset: Option<usize>,
3559 message_tail: Option<usize>,
3560}
3561
3562impl SessionLoadOptions {
3563 fn validate(&self) -> std::result::Result<(), ServiceError> {
3564 if self.message_tail.is_some()
3565 && (self.message_limit.is_some() || self.message_offset.is_some())
3566 {
3567 return Err(ServiceError::InvalidParams(
3568 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3569 .into(),
3570 ));
3571 }
3572 Ok(())
3573 }
3574}
3575
3576fn load_session_door(params: LoadSessionParams) -> std::result::Result<Value, ServiceError> {
3579 if let Some(options) = ¶ms.options {
3580 options.validate()?;
3581 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
3582 return Ok(result);
3583 }
3584 return load_session(¶ms.read.locator)
3585 .map(|session| projected_session_result(&session, options))
3586 .map_err(operation);
3587 }
3588 let mut session = if params.read.display_history() {
3589 HarnessCatalog::new()
3590 .load_display_view(
3591 ¶ms.read.locator,
3592 params.read.read_fidelity(),
3593 params.read.tail_messages().unwrap_or(500),
3594 )
3595 .map_err(crate::Error::from)
3596 } else if params.read.include_subagents() {
3597 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
3598 } else {
3599 HarnessCatalog::new()
3600 .load_parent_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
3601 .map_err(crate::Error::from)
3602 }
3603 .map_err(operation)?;
3604 params.read.bound_session(&mut session);
3605 let mut value = normalized_session_json(&session);
3608 value["resolved_profile"] = recorded_session_profile(&session);
3609 value["resolved_profile_history"] = Value::Array(recorded_session_profile_history(&session));
3610 Ok(json!({"session": value}))
3611}
3612
3613#[derive(Deserialize)]
3614struct LoadSessionParams {
3615 #[serde(flatten)]
3616 read: LocatorParams,
3617 #[serde(default)]
3618 options: Option<SessionLoadOptions>,
3619}
3620
3621#[derive(Deserialize)]
3622struct UnfollowParams {
3623 subscription: String,
3624}
3625
3626#[derive(Debug, Deserialize)]
3627#[serde(deny_unknown_fields)]
3628struct IndexResizeParams {
3629 subscription: String,
3630 limit: usize,
3631}
3632
3633#[cfg(feature = "adapter-api")]
3636#[derive(Default)]
3637struct PanePrompts {
3638 waiting: std::collections::HashSet<String>,
3639 read_at: Option<std::time::Instant>,
3640 reading: bool,
3641}
3642
3643#[cfg(feature = "adapter-api")]
3644static PANE_PROMPTS: std::sync::Mutex<Option<PanePrompts>> = std::sync::Mutex::new(None);
3645
3646#[cfg(feature = "adapter-api")]
3648const PANE_PROMPT_REFRESH: std::time::Duration = std::time::Duration::from_secs(3);
3649
3650#[cfg(feature = "adapter-api")]
3657fn pane_prompt_override(activities: &mut [crate::SessionActivity], homes: &crate::HarnessHomes) {
3658 let running_codex = |activity: &crate::SessionActivity| {
3659 activity.harness.as_str() == "codex"
3660 && activity.presence == crate::session_activity::SessionPresence::Running
3661 && matches!(
3664 activity.turn,
3665 crate::SessionTurnState::Working | crate::SessionTurnState::Unknown
3666 )
3667 };
3668 if !activities.iter().any(running_codex) {
3669 return;
3670 }
3671 let waiting = {
3672 let Ok(mut guard) = PANE_PROMPTS.lock() else {
3673 return;
3674 };
3675 let prompts = guard.get_or_insert_with(PanePrompts::default);
3676 let stale = prompts
3677 .read_at
3678 .is_none_or(|at| at.elapsed() >= PANE_PROMPT_REFRESH);
3679 if stale && !prompts.reading {
3680 prompts.reading = true;
3681 let homes = homes.clone();
3682 std::thread::spawn(move || {
3683 let doors = crate::mail_route::LiveSessions::read(&homes);
3684 let waiting = doors
3685 .all()
3686 .iter()
3687 .filter(|live| live.address.harness == "codex" && live.status == "waiting")
3688 .map(|live| live.address.session_id.clone())
3689 .collect();
3690 if let Ok(mut guard) = PANE_PROMPTS.lock() {
3691 let prompts = guard.get_or_insert_with(PanePrompts::default);
3692 prompts.waiting = waiting;
3693 prompts.read_at = Some(std::time::Instant::now());
3694 prompts.reading = false;
3695 }
3696 });
3697 }
3698 prompts.waiting.clone()
3699 };
3700 for activity in activities
3701 .iter_mut()
3702 .filter(|activity| running_codex(activity))
3703 {
3704 if waiting.contains(&activity.session_id) {
3705 activity.turn = crate::SessionTurnState::NeedsInput;
3706 activity.evidence.source = "pane_prompt".into();
3707 }
3708 }
3709}
3710
3711#[derive(Deserialize)]
3712struct ActivitySubscribeParams {
3713 locators: Vec<SessionLocator>,
3714 #[serde(default)]
3715 homes: crate::HarnessHomes,
3716}
3717
3718#[derive(Debug, Clone, Deserialize)]
3722struct MessageTarget {
3723 harness: HarnessId,
3724 session_id: String,
3725 #[serde(default)]
3726 #[allow(dead_code)]
3727 storage: Option<supercode_interchange::catalog::StorageLocator>,
3728}
3729
3730#[derive(Deserialize)]
3731struct MessageSessionParams {
3732 locator: MessageTarget,
3733 text: String,
3734 #[serde(default)]
3735 subject: Option<String>,
3736 #[serde(default)]
3738 idempotency_key: Option<String>,
3739 #[serde(default)]
3741 channel: bool,
3742 #[serde(default)]
3746 from_name: Option<String>,
3747 #[serde(default)]
3749 voice_for: Option<crate::mailbox::MailAddress>,
3750 #[serde(default)]
3752 sender_name: Option<String>,
3753 #[serde(default)]
3755 in_reply_to: Option<String>,
3756 #[serde(default)]
3759 marker: Option<String>,
3760 #[serde(default)]
3763 reply_marker: Option<String>,
3764 #[serde(default)]
3766 surface: Option<String>,
3767 #[serde(default)]
3769 notify_when_idle: bool,
3770 #[serde(default)]
3775 as_user: bool,
3776 #[serde(default)]
3779 homes: crate::HarnessHomes,
3780}
3781
3782#[derive(Deserialize)]
3783#[serde(deny_unknown_fields)]
3784struct InboxParams {
3785 #[serde(default)]
3787 from_name: Option<String>,
3788 #[serde(default)]
3790 address: Option<String>,
3791 #[serde(default)]
3793 all: bool,
3794}
3795
3796#[derive(Deserialize)]
3797#[serde(deny_unknown_fields)]
3798struct HarnessSettingsParams {
3799 harness: String,
3800}
3801
3802#[derive(Deserialize)]
3803#[serde(deny_unknown_fields)]
3804struct ConfigureHarnessParams {
3805 harness: String,
3806 #[serde(default)]
3807 changes: Vec<crate::HarnessSettingChange>,
3808 #[serde(default)]
3809 expected_revision: Option<String>,
3810}
3811
3812fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3813 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3814 Ok(report) => (
3815 serde_json::to_value(report).unwrap_or(Value::Null),
3816 Value::Null,
3817 ),
3818 Err(error) => (
3819 Value::Null,
3820 Value::String(format!(
3821 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3822 )),
3823 ),
3824 }
3825}
3826
3827async fn message_live_session(params: &MessageSessionParams) -> Value {
3841 use crate::mail_route::{Delivered, NoDoor, Refused};
3842 let (inbound_controls, inbound_controls_error) =
3843 claude_inbound_controls_or_error(¶ms.homes);
3844 let refused = |reason: &str, message: String| {
3845 json!({
3846 "delivered_to_bus": false,
3847 "refusal": {"reason": reason, "message": message},
3848 "inbound_controls": inbound_controls,
3849 "inbound_controls_error": inbound_controls_error,
3850 })
3851 };
3852 if params.text.trim().is_empty() {
3853 return refused(
3854 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3855 "refusing to deliver an empty message".into(),
3856 );
3857 }
3858 let sender = match operator_address(params.from_name.as_deref()) {
3859 Ok(sender) => sender,
3860 Err(message) => return refused("invalid_sender", message),
3861 };
3862 let receiver = match crate::mailbox::MailAddress::new(
3863 crate::mailbox::local_machine_name(),
3864 params.locator.harness.as_str(),
3865 ¶ms.locator.session_id,
3866 ) {
3867 Ok(receiver) => receiver,
3868 Err(error) => return refused("delivery_failed", error.to_string()),
3869 };
3870 if params
3871 .voice_for
3872 .as_ref()
3873 .is_some_and(|voice_for| voice_for != &receiver)
3874 {
3875 return refused(
3876 "invalid_sender",
3877 "a voice front must represent the receiving session".into(),
3878 );
3879 }
3880 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3883 && !crate::mail_agent::resumable(&receiver)
3884 && crate::runtime_mail::controlled_runtime(
3885 HarnessId::CLAUDE_CODE,
3886 ¶ms.locator.session_id,
3887 )
3888 .is_none()
3889 {
3890 if let Err(refusal) =
3891 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3892 {
3893 return refused(refusal.reason.as_str(), refusal.message);
3894 }
3895 }
3896 if params.as_user {
3897 let agent_address =
3900 (receiver.harness == crate::mail_agent::AGENT_HARNESS).then(|| receiver.clone());
3901 let receiver = if receiver.harness == crate::mail_agent::AGENT_HARNESS {
3902 match crate::mail_agent::load(&receiver.session_id) {
3903 Some(agent) => agent.main_session,
3904 None => {
3905 return refused(
3906 "delivery_failed",
3907 format!(
3908 "no agent named {} is declared on this machine",
3909 receiver.session_id
3910 ),
3911 )
3912 }
3913 }
3914 } else {
3915 receiver
3916 };
3917 return message_as_user(
3918 params,
3919 sender,
3920 receiver,
3921 agent_address,
3922 inbound_controls,
3923 inbound_controls_error,
3924 )
3925 .await;
3926 }
3927 let door = crate::mail_route::door_for(¶ms.homes, &receiver);
3930 let mut envelope = match crate::mailbox::Envelope::new(
3931 sender.clone(),
3932 params
3933 .sender_name
3934 .clone()
3935 .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3936 if params.channel {
3937 crate::mailbox::MailKind::Channel
3938 } else {
3939 crate::mailbox::MailKind::Peer
3940 },
3941 crate::mailbox::ReplyVia::Command,
3942 params.text.clone(),
3943 ) {
3944 Ok(envelope) => envelope,
3945 Err(error) => return refused("delivery_failed", error.to_string()),
3946 };
3947 if let Some(key) = ¶ms.idempotency_key {
3948 if key.is_empty() || key.len() > 256 {
3949 return refused("invalid_key", "idempotency_key needs 1–256 bytes".into());
3950 }
3951 let identity = serde_json::to_vec(&(sender.to_string(), receiver.to_string(), key))
3952 .expect("string tuple serializes");
3953 envelope.id = format!("m-{}", &blake3::hash(&identity).to_hex()[..24]);
3954 let mailbox = match crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &receiver) {
3955 Ok(mailbox) => mailbox,
3956 Err(error) => return refused("delivery_failed", error.to_string()),
3957 };
3958 match mailbox.find(&envelope.id) {
3959 Ok(Some(previous)) => {
3960 let saved = previous.envelope;
3961 let subject = params.subject.as_deref().map(|value| {
3962 value
3963 .lines()
3964 .next()
3965 .unwrap_or("")
3966 .chars()
3967 .take(200)
3968 .collect::<String>()
3969 });
3970 if saved.body != envelope.body
3971 || saved.subject != subject
3972 || saved.kind != envelope.kind
3973 || saved.from_name != envelope.from_name
3974 || saved.in_reply_to != params.in_reply_to
3975 || saved.voice_for != params.voice_for
3976 {
3977 return refused(
3978 "idempotency_conflict",
3979 "this key already names different mail".into(),
3980 );
3981 }
3982 }
3983 Ok(None) => {}
3984 Err(error) => return refused("delivery_failed", error.to_string()),
3985 }
3986 }
3987 envelope.in_reply_to = params.in_reply_to.clone();
3988 envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
3989 envelope.voice_for = params.voice_for.clone();
3990 envelope.subject = params.subject.as_deref().map(|value| {
3991 value
3992 .lines()
3993 .next()
3994 .unwrap_or("")
3995 .chars()
3996 .take(200)
3997 .collect()
3998 });
3999 let channel = crate::mail_agent::Channel {
4000 marker: params.marker.clone(),
4001 reply_marker: params.reply_marker.clone(),
4002 surface: params.surface.clone(),
4003 };
4004 match crate::mail_agent::plan(&mut envelope, &receiver, &channel) {
4005 Err(error) => return refused("delivery_failed", error.to_string()),
4006 Ok(Some(plan)) => {
4007 let caller = crate::mail_route::Caller {
4008 address: sender.clone(),
4009 name: envelope.from_name.clone(),
4010 };
4011 let outcome = crate::mail_send::deliver_planned(
4012 ¶ms.homes,
4013 &caller,
4014 &envelope,
4015 &plan,
4016 false,
4017 params.notify_when_idle,
4018 )
4019 .await;
4020 return json!({
4021 "delivered_to_bus": outcome.code == 0,
4022 "message_id": envelope.id,
4023 "reply_to": sender.to_string(),
4024 "thread": plan.thread.as_ref().map(|thread| thread.id.clone()),
4025 "target": {
4026 "harness": params.locator.harness.as_str(),
4027 "session_id": params.locator.session_id,
4028 },
4029 "delivery": {"door": "agent", "how": outcome.text},
4030 "inbound_controls": inbound_controls,
4031 "inbound_controls_error": inbound_controls_error,
4032 });
4033 }
4034 Ok(None) => {}
4035 }
4036 let door = match door {
4037 Ok(door) => door,
4038 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
4039 return refused(
4040 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
4041 format!(
4042 "no running `{}` session `{}` is reachable; its transcript is persisted only",
4043 params.locator.harness.as_str(),
4044 params.locator.session_id
4045 ),
4046 )
4047 }
4048 };
4049 let how = match crate::mail_route::deliver(
4050 &envelope,
4051 &receiver,
4052 &door,
4053 true,
4054 params.notify_when_idle,
4055 )
4056 .await
4057 {
4058 Err(detail) => {
4059 return refused(
4060 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
4061 detail,
4062 )
4063 }
4064 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
4065 Ok(Err(Refused::TooLong(bytes))) => {
4066 return refused(
4067 "too_long",
4068 format!(
4069 "the message is {bytes} bytes; the limit is {}",
4070 crate::mail_route::MAX_RELAYED_BYTES
4071 ),
4072 )
4073 }
4074 Ok(Ok(delivered)) => match delivered {
4075 Delivered::Steered => "steered",
4076 Delivered::Started => "started",
4077 Delivered::Native { busy: true } => "next_tool_call",
4078 Delivered::Native { busy: false } => "started",
4079 Delivered::Hooked => "hook",
4080 Delivered::HookWoken => "woken",
4081 Delivered::Queued => "queued",
4082 Delivered::Stored => "stored",
4083 Delivered::Operator => "filed",
4084 },
4085 };
4086 json!({
4087 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
4088 "message_id": envelope.id,
4089 "reply_to": sender.to_string(),
4090 "target": {
4091 "harness": params.locator.harness.as_str(),
4092 "session_id": params.locator.session_id,
4093 "name": match &door {
4094 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
4095 _ => None,
4096 },
4097 },
4098 "delivery": {"door": door.name(), "how": how},
4099 "inbound_controls": inbound_controls,
4100 "inbound_controls_error": inbound_controls_error,
4101 })
4102}
4103
4104async fn message_as_user(
4108 params: &MessageSessionParams,
4109 sender: crate::mailbox::MailAddress,
4110 receiver: crate::mailbox::MailAddress,
4111 agent: Option<crate::mailbox::MailAddress>,
4112 inbound_controls: Value,
4113 inbound_controls_error: Value,
4114) -> Value {
4115 let mut envelope = match crate::mailbox::Envelope::new(
4116 sender.clone(),
4117 format!("{}@{}", sender.session_id, sender.machine),
4118 crate::mailbox::MailKind::User,
4119 crate::mailbox::ReplyVia::None,
4120 params.text.clone(),
4121 ) {
4122 Ok(envelope) => envelope,
4123 Err(error) => {
4124 return json!({
4125 "delivered_to_bus": false,
4126 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4127 "inbound_controls": inbound_controls,
4128 "inbound_controls_error": inbound_controls_error,
4129 })
4130 }
4131 };
4132 let mut receiver = receiver;
4137 if let Some(agent) = &agent {
4138 envelope.in_reply_to = params.in_reply_to.clone();
4139 envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
4140 let channel = crate::mail_agent::Channel::default();
4141 let planned = match crate::mail_agent::plan(&mut envelope, agent, &channel) {
4142 Ok(planned) => planned,
4143 Err(error) => {
4144 return json!({
4145 "delivered_to_bus": false,
4146 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4147 "inbound_controls": inbound_controls,
4148 "inbound_controls_error": inbound_controls_error,
4149 })
4150 }
4151 };
4152 if let Some(planned) = planned {
4153 let typed_to = planned
4154 .recipients
4155 .iter()
4156 .find(|r| r.wake && r.address.harness != "operator")
4157 .map(|r| r.address.clone());
4158 if let Some(typed_to) = typed_to {
4159 receiver = typed_to;
4160 }
4161 if let Some(agent) = &planned.agent {
4162 let _ = crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), agent)
4164 .and_then(|mailbox| mailbox.deliver_read(&envelope));
4165 }
4166 let mut copy = envelope.clone();
4169 copy.kind = crate::mailbox::MailKind::Typed;
4170 for other in planned
4171 .recipients
4172 .iter()
4173 .filter(|r| r.address != receiver && r.address != sender)
4174 {
4175 let _ = crate::mail_agent::file_unread(&other.address, ©);
4176 }
4177 if let Some(account_manager) = crate::mail_agent::owners_account_manager() {
4180 let main = account_manager.main_session.clone();
4181 let already =
4182 main == receiver || planned.recipients.iter().any(|r| r.address == main);
4183 if account_manager.name != agent.session_id && !already {
4184 let _ = crate::mail_agent::file_unread(&main, ©);
4185 }
4186 }
4187 }
4188 }
4189 if agent.is_some()
4192 && crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4193 && crate::mail_agent::resumable(&receiver)
4194 {
4195 if let Err(error) = crate::mail_agent::resume(&receiver) {
4196 return json!({
4197 "delivered_to_bus": false,
4198 "refusal": {"reason": "no_user_door", "message": format!("{receiver} was not resumed: {error}")},
4199 "inbound_controls": inbound_controls,
4200 "inbound_controls_error": inbound_controls_error,
4201 });
4202 }
4203 let started = std::time::Instant::now();
4204 while crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4205 && started.elapsed() < crate::mail_agent::RESUME_WAIT
4206 {
4207 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
4208 }
4209 }
4210 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
4211 Ok(turn) => json!({
4212 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
4213 "message_id": envelope.id,
4214 "target": {
4215 "harness": params.locator.harness.as_str(),
4216 "session_id": params.locator.session_id,
4217 },
4218 "delivery": {
4219 "door": match turn {
4220 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
4221 _ => "pane",
4222 },
4223 "how": turn.as_str(),
4224 },
4225 "inbound_controls": inbound_controls,
4226 "inbound_controls_error": inbound_controls_error,
4227 }),
4228 Err(message) => json!({
4229 "delivered_to_bus": false,
4230 "refusal": {"reason": "no_user_door", "message": message},
4231 "inbound_controls": inbound_controls,
4232 "inbound_controls_error": inbound_controls_error,
4233 }),
4234 }
4235}
4236
4237#[derive(Deserialize)]
4238#[serde(deny_unknown_fields)]
4239struct ActivityUnderParams {
4240 pids: Vec<u32>,
4242 #[serde(default)]
4243 homes: crate::HarnessHomes,
4244}
4245
4246async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
4250 let params = decode::<ActivityUnderParams>(params)?;
4251 if params.pids.len() > 1024 {
4252 return Err(ServiceError::InvalidParams(
4253 "sessions.activity_under accepts at most 1024 pids".into(),
4254 ));
4255 }
4256 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
4257 .await
4258 .map_err(ServiceError::Sdk)?;
4259 Ok(json!({
4260 "activities": found
4261 .into_iter()
4262 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
4263 .collect::<Vec<_>>(),
4264 }))
4265}
4266
4267fn operator_address(
4270 name: Option<&str>,
4271) -> std::result::Result<crate::mailbox::MailAddress, String> {
4272 let name: String = name
4273 .unwrap_or("supercode")
4274 .trim()
4275 .chars()
4276 .map(|character| {
4277 if character.is_whitespace() || character == '@' {
4278 '-'
4279 } else {
4280 character
4281 }
4282 })
4283 .collect();
4284 if name.is_empty() {
4285 return Err("from_name must not be empty".into());
4286 }
4287 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
4288 .map_err(|error| error.to_string())
4289}
4290
4291fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
4295 let address = match (¶ms.address, ¶ms.from_name) {
4296 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
4297 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
4298 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
4299 (Some(_), Some(_)) => {
4300 return Err(ServiceError::InvalidParams(
4301 "sessions.inbox takes from_name or address, not both".into(),
4302 ))
4303 }
4304 };
4305 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
4306 let mailbox =
4307 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
4308 let claimed = mailbox.claim_unread().map_err(operation)?;
4309 let mut messages: Vec<Value> = Vec::new();
4310 if params.all {
4311 for stored in mailbox.list().map_err(operation)? {
4312 if stored.state == crate::mailbox::MailState::Read {
4313 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4314 }
4315 }
4316 }
4317 for stored in &claimed {
4318 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4319 }
4320 for stored in &claimed {
4321 mailbox.acknowledge(stored).map_err(operation)?;
4322 }
4323 Ok(json!({"address": address.to_string(), "messages": messages}))
4324}
4325
4326#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4331struct FollowedSource {
4332 harness: String,
4333 session_id: String,
4334 reported: Option<String>,
4335}
4336
4337#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4338struct ActivitySubscription {
4339 locators: Vec<SessionLocator>,
4340 homes: crate::HarnessHomes,
4341 reported: BTreeMap<(String, String), crate::SessionActivity>,
4342}
4343
4344fn live_descriptor_value(
4351 session: &SessionDescriptor,
4352 doors: &crate::mail_route::LiveSessions,
4353) -> std::result::Result<Value, ServiceError> {
4354 let mut value = serde_json::to_value(session)
4355 .map_err(|error| ServiceError::Operation(error.to_string()))?;
4356 if value.get("title").is_none_or(Value::is_null) {
4358 if let Some(live) = doors.all().iter().find(|live| {
4359 live.address.harness == session.locator.harness.as_str()
4360 && live.address.session_id == session.locator.session_id
4361 }) {
4362 let name = live.name.split('@').next().unwrap_or(&live.name);
4363 if !name.is_empty() {
4364 value["title"] = json!(name);
4365 }
4366 }
4367 }
4368 if let Some(workspace) = &session.cwd {
4369 let source = LiveRuntimeSource {
4370 harness: session.locator.harness.as_str().to_string(),
4371 session_id: session.locator.session_id.clone(),
4372 workspace: workspace.clone(),
4373 };
4374 if let Some(endpoint) = discover_live_runtime(&source)
4375 .map_err(|error| ServiceError::Operation(error.to_string()))?
4376 {
4377 value["live_endpoint"] = json!(endpoint.as_str());
4378 }
4379 }
4380 if let Some(door) = doors.door(
4384 session.locator.harness.as_str(),
4385 &session.locator.session_id,
4386 ) {
4387 value["delivery"] = json!(door);
4388 }
4389 if let Some(live) = doors.all().iter().find(|live| {
4392 live.address.harness == session.locator.harness.as_str()
4393 && live.address.session_id == session.locator.session_id
4394 && live.status == "waiting"
4395 }) {
4396 if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
4397 value["pending_request"] = request;
4398 }
4399 }
4400 Ok(value)
4401}
4402
4403fn live_index_changes(
4404 changes: Vec<crate::session_index::SessionIndexChange>,
4405 homes: &HarnessHomes,
4406) -> std::result::Result<Vec<Value>, ServiceError> {
4407 use crate::session_index::SessionIndexChange;
4408 let doors = crate::mail_route::LiveSessions::read(homes);
4409 changes
4410 .into_iter()
4411 .map(|change| match change {
4412 SessionIndexChange::Added { descriptor } => Ok(json!({
4413 "kind": "added",
4414 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4415 })),
4416 SessionIndexChange::Updated { descriptor } => Ok(json!({
4417 "kind": "updated",
4418 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4419 })),
4420 SessionIndexChange::Removed { key } => Ok(json!({
4421 "kind": "removed",
4422 "key": key,
4423 })),
4424 })
4425 .collect()
4426}
4427
4428fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
4429 use crate::{SessionPresence, SessionTurnState};
4430 match (activity.presence, activity.turn) {
4431 (SessionPresence::Persisted, _) => None,
4432 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
4433 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
4434 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
4436 (SessionPresence::Running, SessionTurnState::Unknown)
4440 if activity.evidence.native_state.is_none() =>
4441 {
4442 None
4443 }
4444 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
4445 }
4446}
4447
4448#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
4449#[serde(rename_all = "kebab-case")]
4450enum TransferFormat {
4451 ClaudeCode,
4452 Codex,
4453 #[serde(rename = "opencode", alias = "open-code")]
4454 OpenCode,
4455 Pi,
4456 Grok,
4457 Gemini,
4458 Goose,
4459 Hermes,
4463}
4464
4465impl TransferFormat {
4466 fn id(self) -> &'static str {
4467 match self {
4468 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
4469 Self::Codex => HarnessId::CODEX,
4470 Self::OpenCode => HarnessId::OPENCODE,
4471 Self::Pi => HarnessId::PI,
4472 Self::Grok => HarnessId::GROK,
4473 Self::Gemini => HarnessId::GEMINI,
4474 Self::Goose => HarnessId::GOOSE,
4475 Self::Hermes => HarnessId::HERMES,
4476 }
4477 }
4478}
4479
4480impl From<TransferFormat> for SessionFormat {
4481 fn from(value: TransferFormat) -> Self {
4482 match value {
4483 TransferFormat::ClaudeCode => Self::ClaudeCode,
4484 TransferFormat::Codex => Self::Codex,
4485 TransferFormat::OpenCode => Self::OpenCode,
4486 TransferFormat::Pi => Self::Pi,
4487 TransferFormat::Grok => Self::Grok,
4488 TransferFormat::Gemini => Self::Gemini,
4489 TransferFormat::Goose => Self::Goose,
4490 TransferFormat::Hermes => Self::Codex,
4492 }
4493 }
4494}
4495
4496#[derive(Deserialize)]
4497struct ImportSessionParams {
4498 source_harness: TransferFormat,
4499 content: String,
4500}
4501
4502#[derive(Deserialize)]
4503struct ExportSessionParams {
4504 locator: SessionLocator,
4505 target_harness: TransferFormat,
4506}
4507
4508#[derive(Deserialize)]
4509struct ReduceSessionParams {
4510 locator: SessionLocator,
4511 target_harness: TransferFormat,
4512 #[serde(default = "default_keep_last")]
4513 keep_last: usize,
4514}
4515
4516fn default_keep_last() -> usize {
4517 6
4518}
4519
4520#[derive(Deserialize)]
4521struct BranchSessionParams {
4522 locator: SessionLocator,
4523 #[serde(default)]
4524 target_harness: Option<TransferFormat>,
4525}
4526
4527#[derive(Deserialize)]
4528struct HandoffSessionParams {
4529 locator: SessionLocator,
4530 target_harness: TransferFormat,
4531 #[serde(default)]
4532 cwd: Option<PathBuf>,
4533}
4534
4535#[derive(Deserialize)]
4536struct MaterializeSessionParams {
4537 artifact: crate::native_materialize::MaterializeArtifact,
4538 cwd: PathBuf,
4539 #[serde(default)]
4541 homes: HarnessHomes,
4542}
4543
4544#[derive(Debug, Clone, Copy, Default, Deserialize)]
4547#[serde(rename_all = "snake_case")]
4548enum ResumePolicy {
4549 Default,
4550 #[default]
4551 Yolo,
4552}
4553
4554#[derive(Deserialize)]
4555struct ResumeInstructionsParams {
4556 locator: SessionLocator,
4557 #[serde(default)]
4558 cwd: Option<PathBuf>,
4559 #[serde(default)]
4560 policy: ResumePolicy,
4561}
4562
4563#[derive(Deserialize)]
4565struct WorkflowLoadParams {
4566 #[serde(default)]
4568 from: Option<crate::workflow_doors::WorkflowHarness>,
4569 home: PathBuf,
4570}
4571
4572#[derive(Deserialize)]
4575struct OrchestrationLoadParams {
4576 root: PathBuf,
4577 #[serde(default)]
4578 flavor: crate::orchestration_doors::HomeFlavor,
4579}
4580
4581#[derive(Deserialize)]
4584struct OrchestrationSaveParams {
4585 root: PathBuf,
4586 orchestration: crate::orchestration::Orchestration,
4587 #[serde(default)]
4588 vault: BTreeMap<String, String>,
4589}
4590
4591#[derive(Deserialize)]
4593struct OrchestrationCompileParams {
4594 from: crate::orchestration_doors::OrchestrationHarness,
4595 home: PathBuf,
4596}
4597
4598#[derive(Deserialize)]
4602struct OrchestrationDecompileParams {
4603 to: crate::orchestration_doors::OrchestrationHarness,
4604 orchestration: crate::orchestration::Orchestration,
4605 source: PathBuf,
4606 #[serde(default)]
4607 source_flavor: crate::orchestration_doors::SourceFlavor,
4608 dest: PathBuf,
4609 #[serde(default)]
4610 vault: BTreeMap<String, String>,
4611}
4612
4613#[derive(Deserialize)]
4616struct OrchestrationImportParams {
4617 from: crate::orchestration_doors::OrchestrationHarness,
4618 home: PathBuf,
4619 into: PathBuf,
4620}
4621
4622#[derive(Deserialize)]
4625struct OrchestrationExportParams {
4626 to: crate::orchestration_doors::OrchestrationHarness,
4627 root: PathBuf,
4628 dest: PathBuf,
4629}
4630
4631#[derive(Deserialize)]
4633struct JobsGetParams {
4634 harness: String,
4635 id: String,
4636 #[serde(default)]
4637 homes: crate::HarnessHomes,
4638}
4639
4640fn mutate_job(
4648 verb: crate::jobs_control::JobVerb,
4649 params: Value,
4650) -> std::result::Result<Value, ServiceError> {
4651 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4652 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4653 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4654 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4655}
4656
4657fn mutate_skill(
4664 verb: crate::skills_control::SkillVerb,
4665 params: Value,
4666) -> std::result::Result<Value, ServiceError> {
4667 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4668 if !crate::skills_control::supports_skill_control(&mutation.harness) {
4669 return Err(ServiceError::UnsupportedAction(format!(
4670 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4671 mutation.harness,
4672 verb.as_str(),
4673 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4674 )));
4675 }
4676 let outcome =
4677 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4678 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4679}
4680
4681fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4683 match error {
4684 crate::skills_control::SkillControlError::Unsupported(message) => {
4685 ServiceError::UnsupportedAction(message)
4686 }
4687 crate::skills_control::SkillControlError::Invalid(message) => {
4688 ServiceError::InvalidParams(message)
4689 }
4690 crate::skills_control::SkillControlError::Failed(message) => {
4691 ServiceError::Operation(message)
4692 }
4693 }
4694}
4695
4696fn mutate_profile(
4704 verb: crate::profiles_control::ProfileVerb,
4705 params: Value,
4706) -> std::result::Result<Value, ServiceError> {
4707 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4708 let outcome =
4709 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4710 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4711}
4712
4713fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4715 match error {
4716 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4717 ServiceError::UnsupportedAction(message)
4718 }
4719 crate::profiles_control::ProfileControlError::Invalid(message) => {
4720 ServiceError::InvalidParams(message)
4721 }
4722 crate::profiles_control::ProfileControlError::Failed(message) => {
4723 ServiceError::Operation(message)
4724 }
4725 }
4726}
4727
4728fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4732 match error {
4733 crate::jobs_control::JobControlError::Unsupported(message) => {
4734 ServiceError::UnsupportedAction(message)
4735 }
4736 crate::jobs_control::JobControlError::Invalid(message) => {
4737 ServiceError::InvalidParams(message)
4738 }
4739 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4740 }
4741}
4742
4743fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4748 match error {
4749 crate::SessionControlError::Unsupported(message) => {
4750 ServiceError::UnsupportedAction(message)
4751 }
4752 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4753 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4754 }
4755}
4756
4757fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4762 if crate::jobs::supports_jobs(harness) {
4763 return Ok(());
4764 }
4765 Err(ServiceError::UnsupportedAction(format!(
4766 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4767 crate::jobs::JOB_HARNESSES.join(", ")
4768 )))
4769}
4770
4771#[derive(Deserialize)]
4773struct RunsGetParams {
4774 harness: String,
4775 id: String,
4776 #[serde(default)]
4777 homes: crate::HarnessHomes,
4778}
4779
4780fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4785 if crate::runs::supports_runs(harness) {
4786 return Ok(());
4787 }
4788 Err(ServiceError::UnsupportedAction(format!(
4789 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4790 crate::runs::RUN_HARNESSES.join(", ")
4791 )))
4792}
4793
4794#[derive(Serialize)]
4795struct SessionArtifact {
4796 source_harness: HarnessId,
4797 target_harness: &'static str,
4798 session_id: Option<String>,
4799 content: String,
4800 suggested_filename: String,
4801 files: Vec<SessionArtifactFile>,
4802 fidelity: Fidelity,
4803 residue: Vec<String>,
4804}
4805
4806#[derive(Serialize)]
4807struct SessionArtifactFile {
4808 path: String,
4809 content: String,
4810 role: ArtifactFileRole,
4811}
4812
4813#[derive(Serialize)]
4814#[serde(rename_all = "snake_case")]
4815enum ArtifactFileRole {
4816 Primary,
4817 Subagent,
4818 Bundle,
4819 SourceRecovery,
4820}
4821
4822#[derive(Serialize)]
4823struct StructuredLaunch {
4824 cwd: PathBuf,
4825 program: String,
4826 arguments: Vec<String>,
4827 env: BTreeMap<String, String>,
4828}
4829
4830struct HandoffInstructions {
4831 launch: StructuredLaunch,
4832 materialize: Option<StructuredLaunch>,
4833 requires_materialization: bool,
4834 note: String,
4835}
4836
4837#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4838#[serde(rename_all = "snake_case")]
4839enum HarnessProbeLevel {
4840 #[default]
4841 Passive,
4842 Handshake,
4843}
4844
4845#[derive(Default, Deserialize)]
4846#[serde(default)]
4847struct HarnessInventoryParams {
4848 harness: Option<HarnessId>,
4849 harnesses: Vec<HarnessId>,
4850 workspace: Option<PathBuf>,
4851 probe: HarnessProbeLevel,
4852 include_sessions: bool,
4853 skip_versions: bool,
4855}
4856
4857#[derive(Deserialize)]
4858struct HarnessAuthenticationParams {
4859 harness: HarnessId,
4860}
4861
4862#[derive(Deserialize)]
4863struct BeginHarnessAuthenticationParams {
4864 harness: HarnessId,
4865 #[serde(default = "local_browser_authentication_environment")]
4866 environment: crate::HarnessAuthenticationEnvironment,
4867 #[serde(default)]
4868 method: Option<crate::HarnessAuthenticationMethodId>,
4869 #[serde(default)]
4870 cwd: Option<PathBuf>,
4871}
4872
4873fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4874 crate::HarnessAuthenticationEnvironment::LocalBrowser
4875}
4876
4877#[derive(Serialize)]
4878struct HarnessInventoryReport {
4879 probe: HarnessProbeLevel,
4880 workspace: Option<PathBuf>,
4881 harnesses: Vec<LocalHarness>,
4882}
4883
4884#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4885#[serde(rename_all = "snake_case")]
4886enum HarnessAuthState {
4887 Ready,
4888 Configured,
4889 Required,
4890 Unknown,
4891}
4892
4893#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4894#[serde(rename_all = "snake_case")]
4895enum HarnessRuntimeState {
4896 Ready,
4897 Degraded,
4898 Unavailable,
4899}
4900
4901#[derive(Serialize)]
4902struct HarnessSessionCounts {
4903 global: Option<usize>,
4904 workspace: Option<usize>,
4905}
4906
4907#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4918#[serde(rename_all = "snake_case")]
4919pub enum GatewayState {
4920 Up,
4921 Down,
4922 Unknown,
4923}
4924
4925#[derive(Debug, Clone, Serialize)]
4927pub struct GatewayHealth {
4928 pub state: GatewayState,
4929 #[serde(skip_serializing_if = "Option::is_none")]
4934 pub endpoint: Option<String>,
4935 #[serde(skip_serializing_if = "Option::is_none")]
4936 pub version: Option<String>,
4937 pub evidence: String,
4939 pub checked_at_ms: u64,
4940}
4941
4942fn openclaw_gateway_endpoint(home: &Path) -> String {
4946 let config_path = home.join(".openclaw/openclaw.json");
4947 let gateway = std::fs::read_to_string(&config_path)
4948 .ok()
4949 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4950 .and_then(|config| config.get("gateway").cloned());
4951 if let Some(url) = gateway
4952 .as_ref()
4953 .and_then(|gateway| gateway.get("url"))
4954 .and_then(serde_json::Value::as_str)
4955 {
4956 return url.to_string();
4957 }
4958 let port = gateway
4959 .as_ref()
4960 .and_then(|gateway| gateway.get("port"))
4961 .and_then(serde_json::Value::as_u64)
4962 .unwrap_or(18789);
4963 format!("ws://127.0.0.1:{port}")
4964}
4965
4966fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4973 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4974 let output = std::process::Command::new(&program)
4975 .args(["gateway", "status"])
4976 .stdin(std::process::Stdio::null())
4977 .output()
4978 .ok()?;
4979 let text = format!(
4980 "{}{}",
4981 String::from_utf8_lossy(&output.stdout),
4982 String::from_utf8_lossy(&output.stderr)
4983 );
4984 let verdict = text.lines().find_map(|line| {
4985 let l = line.trim();
4986 if l.contains("supervised by launchd (PID")
4987 || l.contains("supervised by systemd (PID")
4988 || l.contains("Gateway is running")
4989 || l.contains("process is running")
4990 {
4991 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4992 } else if l.contains("not running") || l.contains("not installed") {
4993 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4994 } else {
4995 None
4996 }
4997 });
4998 verdict
4999}
5000
5001fn gateway_health(
5002 id: &str,
5003 installed: bool,
5004 running: Option<&RunningInstance>,
5005 version: Option<&str>,
5006) -> GatewayHealth {
5007 let checked_at_ms = now_epoch_ms();
5008 let home = supercode_interchange::user_home()
5009 .map(std::path::PathBuf::into_os_string)
5010 .map(PathBuf::from);
5011 match id {
5012 HarnessId::HERMES | HarnessId::OPENCLAW => {
5013 let endpoint = (id == HarnessId::OPENCLAW)
5014 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
5015 .flatten();
5016 let (state, evidence) = match running {
5017 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
5018 None if !installed => (
5019 GatewayState::Unknown,
5020 format!("`{id}` is not installed; no gateway to probe"),
5021 ),
5022 None if id == HarnessId::HERMES => match hermes_gateway_status() {
5023 Some((state, evidence)) => (state, evidence),
5026 None => (
5027 GatewayState::Down,
5028 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
5029 ),
5030 },
5031 None => (
5032 GatewayState::Down,
5033 format!(
5034 "no TCP listener at {}",
5035 endpoint.as_deref().unwrap_or("the gateway endpoint")
5036 ),
5037 ),
5038 };
5039 GatewayHealth {
5040 state,
5041 endpoint,
5042 version: version.map(str::to_string),
5043 evidence,
5044 checked_at_ms,
5045 }
5046 }
5047 HarnessId::ORCHESTRATOR => {
5054 let root = crate::HarnessHomes::default().orchestrator;
5055 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
5056 Some(lease) if lease.is_live() => (
5057 GatewayState::Up,
5058 format!(
5059 "`{}` names pid {} (started {}), which is live",
5060 crate::orchestrator::lock_path(&root).display(),
5061 lease.pid,
5062 lease.started_at
5063 ),
5064 ),
5065 Some(lease) => (
5066 GatewayState::Down,
5067 format!(
5068 "stale lease `{}`: pid {} is gone",
5069 crate::orchestrator::lock_path(&root).display(),
5070 lease.pid
5071 ),
5072 ),
5073 None => (
5074 GatewayState::Down,
5075 format!(
5076 "no lease at `{}`; `supercode orchestrator start` writes one",
5077 crate::orchestrator::lock_path(&root).display()
5078 ),
5079 ),
5080 };
5081 GatewayHealth {
5082 state,
5083 endpoint: None,
5084 version: version.map(str::to_string),
5085 evidence,
5086 checked_at_ms,
5087 }
5088 }
5089 _ => GatewayHealth {
5090 state: GatewayState::Unknown,
5091 endpoint: None,
5092 version: version.map(str::to_string),
5093 evidence: format!("`{id}` runs per session, not as a gateway"),
5094 checked_at_ms,
5095 },
5096 }
5097}
5098
5099#[derive(Debug, Clone, Serialize)]
5100struct RunningInstance {
5101 method: RunningInstanceMethod,
5103 evidence: String,
5105 checked_at_ms: u64,
5107}
5108
5109#[derive(Debug, Clone, Copy, Serialize)]
5110#[serde(rename_all = "snake_case")]
5111enum RunningInstanceMethod {
5112 GatewayConnect,
5115 StoreWalActivity,
5118}
5119
5120fn now_epoch_ms() -> u64 {
5121 std::time::SystemTime::now()
5122 .duration_since(std::time::UNIX_EPOCH)
5123 .map(|elapsed| elapsed.as_millis() as u64)
5124 .unwrap_or(0)
5125}
5126
5127fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
5131 let config_path = home.join(".openclaw/openclaw.json");
5132 let text = std::fs::read_to_string(&config_path).ok();
5133 let gateway = text
5134 .as_deref()
5135 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
5136 .and_then(|config| config.get("gateway").cloned());
5137 let address = gateway
5138 .as_ref()
5139 .and_then(|gateway| gateway.get("url"))
5140 .and_then(serde_json::Value::as_str)
5141 .and_then(|url| {
5142 url.split("://").nth(1).map(|rest| {
5143 rest.trim_end_matches('/')
5144 .split('/')
5145 .next()
5146 .unwrap_or(rest)
5147 .to_string()
5148 })
5149 })
5150 .unwrap_or_else(|| {
5151 let port = gateway
5152 .as_ref()
5153 .and_then(|gateway| gateway.get("port"))
5154 .and_then(serde_json::Value::as_u64)
5155 .unwrap_or(18789);
5156 format!("127.0.0.1:{port}")
5157 });
5158 let reachable = std::net::TcpStream::connect_timeout(
5159 &address.parse().ok()?,
5160 std::time::Duration::from_millis(400),
5161 )
5162 .is_ok();
5163 reachable.then(|| RunningInstance {
5164 method: RunningInstanceMethod::GatewayConnect,
5165 evidence: format!(
5166 "gateway endpoint {address} accepted a TCP connect (from {})",
5167 config_path.display()
5168 ),
5169 checked_at_ms: now_epoch_ms(),
5170 })
5171}
5172
5173fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
5178 let wal = home.join(".hermes/state.db-wal");
5179 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
5180 let age_ms = std::time::SystemTime::now()
5181 .duration_since(modified)
5182 .map(|age| age.as_millis() as u64)
5183 .unwrap_or(u64::MAX);
5184 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
5185 method: RunningInstanceMethod::StoreWalActivity,
5186 evidence: format!(
5187 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
5188 wal.display()
5189 ),
5190 checked_at_ms: now_epoch_ms(),
5191 })
5192}
5193
5194fn probe_running_instance(id: &str) -> Option<RunningInstance> {
5196 let home = supercode_interchange::user_home()
5197 .map(std::path::PathBuf::into_os_string)
5198 .map(PathBuf::from)?;
5199 match id {
5200 HarnessId::OPENCLAW => probe_openclaw_running(&home),
5201 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
5202 _ => None,
5203 }
5204}
5205
5206#[derive(Serialize)]
5207struct LocalHarness {
5208 id: HarnessId,
5209 display_name: String,
5210 supported: bool,
5211 installed: bool,
5212 executable: Option<String>,
5213 version: Option<String>,
5214 auth: HarnessAuthState,
5215 runtime: HarnessRuntimeState,
5216 protocol: String,
5217 capabilities: crate::RuntimeCapabilities,
5218 effective_capabilities: crate::RuntimeCapabilities,
5219 sessions: HarnessSessionCounts,
5220 #[serde(skip_serializing_if = "Option::is_none")]
5223 running: Option<RunningInstance>,
5224 gateway: GatewayHealth,
5226 reason: Option<String>,
5227 repair: Option<String>,
5228}
5229
5230#[derive(Clone, Deserialize)]
5231struct RuntimeBackendParams {
5232 harness: HarnessId,
5233 #[serde(default)]
5234 protocol: Option<String>,
5235 #[serde(default)]
5236 launch: Option<RuntimeLaunch>,
5237 #[serde(default)]
5238 base_url: Option<String>,
5239 #[serde(default)]
5240 policy: RuntimePolicy,
5241}
5242
5243#[derive(Debug, Clone, Copy, Default, Deserialize)]
5246#[serde(rename_all = "snake_case")]
5247enum RuntimePolicy {
5248 Default,
5249 #[default]
5250 Yolo,
5251}
5252
5253#[derive(Deserialize)]
5254struct RuntimeStartParams {
5255 #[serde(flatten)]
5256 backend: RuntimeBackendParams,
5257 cwd: PathBuf,
5258 #[serde(default)]
5261 mcp_servers: Vec<crate::McpServerLaunch>,
5262 #[serde(default)]
5264 approval_policy: Option<String>,
5265}
5266
5267#[derive(Deserialize)]
5268struct RuntimeAttachParams {
5269 #[serde(flatten)]
5270 backend: RuntimeBackendParams,
5271 runtime_id: String,
5272 #[serde(default)]
5273 cwd: Option<PathBuf>,
5274 #[serde(default)]
5277 mcp_servers: Vec<crate::McpServerLaunch>,
5278 #[serde(default)]
5280 approval_policy: Option<String>,
5281}
5282
5283#[derive(Deserialize)]
5284struct RuntimeConnectionParams {
5285 connection: String,
5286}
5287
5288#[derive(Deserialize)]
5291struct RuntimeCloseParams {
5292 #[serde(default)]
5293 connection: String,
5294 #[serde(default)]
5295 runtime_id: Option<String>,
5296}
5297
5298#[derive(Deserialize)]
5299struct RuntimeInputParams {
5300 connection: String,
5301 text: String,
5302 #[serde(default)]
5303 image_urls: Vec<String>,
5304}
5305
5306const MAX_RUNTIME_IMAGES: usize = 4;
5307const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
5308const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
5309
5310fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
5311 if image_urls.len() > MAX_RUNTIME_IMAGES {
5312 return Err(ServiceError::InvalidParams(format!(
5313 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
5314 )));
5315 }
5316 let mut total = 0usize;
5317 for url in &image_urls {
5318 if !(url.starts_with("data:image/")
5319 || url.starts_with("https://")
5320 || url.starts_with("http://"))
5321 {
5322 return Err(ServiceError::InvalidParams(
5323 "runtime images must be image data URLs or HTTP(S) URLs".into(),
5324 ));
5325 }
5326 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
5327 return Err(ServiceError::InvalidParams(format!(
5328 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
5329 )));
5330 }
5331 total = total.saturating_add(url.len());
5332 }
5333 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
5334 return Err(ServiceError::InvalidParams(format!(
5335 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
5336 )));
5337 }
5338 Ok(image_urls)
5339}
5340
5341#[derive(Deserialize)]
5342struct RuntimeRespondParams {
5343 connection: String,
5344 request_id: Value,
5345 response: Value,
5346}
5347
5348fn default_reduction_store_root() -> PathBuf {
5349 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
5350 return PathBuf::from(root).join("sessions");
5351 }
5352 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
5353 return PathBuf::from(home).join(".supercode").join("sessions");
5354 }
5355 PathBuf::from(".supercode").join("sessions")
5356}
5357
5358fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
5359 let mut output = String::new();
5360 for message in messages {
5361 output.push_str(
5362 &serde_json::to_string(message)
5363 .map_err(|error| ServiceError::Operation(error.to_string()))?,
5364 );
5365 output.push('\n');
5366 }
5367 Ok(output)
5368}
5369
5370fn parse_messages_jsonl(
5371 content: &str,
5372) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
5373 content
5374 .lines()
5375 .enumerate()
5376 .filter(|(_, line)| !line.trim().is_empty())
5377 .map(|(index, line)| {
5378 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
5379 ServiceError::Operation(format!(
5380 "reduced transcript line {} is invalid: {error}",
5381 index + 1
5382 ))
5383 })
5384 })
5385 .collect()
5386}
5387
5388fn reduced_bootstrap_prompt(
5389 source: &SessionLocator,
5390 target: TransferFormat,
5391 view_jsonl: &str,
5392 sidecar_path: &Path,
5393 reduction_log_path: &Path,
5394) -> String {
5395 format!(
5396 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
5397 \n\
5398 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\
5399 \n\
5400 <supercode-reduced-session source-session=\"{source_id}\">\n\
5401 {view_jsonl}\
5402 </supercode-reduced-session>\n\
5403 \n\
5404 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
5405 source_harness = source.harness.as_str(),
5406 target_harness = target.id(),
5407 sidecar = sidecar_path.display(),
5408 log = reduction_log_path.display(),
5409 source_id = source.session_id,
5410 )
5411}
5412
5413fn session_artifact(
5414 locator: &SessionLocator,
5415 session: &Session,
5416 target: TransferFormat,
5417) -> std::result::Result<SessionArtifact, ServiceError> {
5418 session_artifact_with_id(locator, session, target, None)
5419}
5420
5421fn session_artifact_with_id(
5422 locator: &SessionLocator,
5423 session: &Session,
5424 target: TransferFormat,
5425 target_session_id: Option<&str>,
5426) -> std::result::Result<SessionArtifact, ServiceError> {
5427 let format: SessionFormat = target.into();
5428 let diagonal = format.source() == session.meta.source;
5429 crate::residue_store::store_segments(session);
5430 let has_appended_turns = session
5431 .imported_message_count
5432 .is_some_and(|imported| imported < session.messages.len());
5433 let mut restoration = None;
5434 let content = if let Some(id) = target_session_id {
5435 if diagonal && format != SessionFormat::OpenCode {
5436 session
5437 .to_jsonl_spliced(format, Some(id))
5438 .map_err(operation)?
5439 } else {
5440 let mut rewritten = session.clone();
5441 rewritten.meta.session_id = Some(id.to_string());
5442 rewritten.to_jsonl(format).map_err(operation)?
5443 }
5444 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
5445 session.raw_verbatim()
5446 } else if diagonal {
5447 session.to_jsonl_spliced(format, None).map_err(operation)?
5448 } else {
5449 match session
5452 .restore_residue(format, crate::residue_store::lookup)
5453 .map_err(operation)?
5454 {
5455 Some((content, report)) => {
5456 restoration = Some(report);
5457 content
5458 }
5459 None => session.to_jsonl(format).map_err(operation)?,
5460 }
5461 };
5462 let stem = sanitize_filename(
5463 target_session_id
5464 .or(session.meta.session_id.as_deref())
5465 .unwrap_or(&locator.session_id),
5466 );
5467 let suggested_filename = if diagonal && target == TransferFormat::Grok {
5468 "chat_history.jsonl".to_string()
5469 } else if target == TransferFormat::Goose {
5470 format!("{stem}.goose.json")
5471 } else {
5472 format!("{stem}.{}.jsonl", target.id())
5473 };
5474 let mut files = vec![SessionArtifactFile {
5475 path: suggested_filename.clone(),
5476 content: content.clone(),
5477 role: ArtifactFileRole::Primary,
5478 }];
5479 if target == TransferFormat::ClaudeCode {
5480 let bundle_stem = Path::new(&suggested_filename)
5481 .file_stem()
5482 .and_then(|stem| stem.to_str())
5483 .unwrap_or(&stem);
5484 let mut child_paths = BTreeSet::new();
5485 for (index, subagent) in session.subagents.iter().enumerate() {
5486 let agent_id = subagent
5487 .meta
5488 .agent_id
5489 .as_deref()
5490 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
5491 .map(sanitize_filename)
5492 .filter(|id| !id.is_empty())
5493 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5494 let child_has_appended_turns = subagent
5495 .imported_message_count
5496 .is_some_and(|imported| imported < subagent.messages.len());
5497 let child_content = if target_session_id.is_none()
5498 && subagent.meta.source == SessionSource::ClaudeCode
5499 && subagent.raw_is_verbatim
5500 && !child_has_appended_turns
5501 {
5502 subagent.raw_verbatim()
5503 } else if subagent.meta.source == SessionSource::ClaudeCode {
5504 subagent
5505 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
5506 .map_err(operation)?
5507 } else {
5508 let mut child = subagent.clone();
5509 if let Some(id) = target_session_id {
5510 child.meta.session_id = Some(id.to_string());
5511 }
5512 child
5513 .to_jsonl(SessionFormat::ClaudeCode)
5514 .map_err(operation)?
5515 };
5516 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
5517 if !child_paths.insert(path.clone()) {
5518 return Err(ServiceError::Operation(format!(
5519 "Claude subagent ids collide at artifact path `{path}`"
5520 )));
5521 }
5522 files.push(SessionArtifactFile {
5523 path,
5524 content: child_content,
5525 role: ArtifactFileRole::Subagent,
5526 });
5527 }
5528 }
5529 if diagonal && target == TransferFormat::Grok {
5530 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
5531 }
5532 if !diagonal || !session.raw_is_verbatim {
5533 files.push(SessionArtifactFile {
5534 path: "recovery/source.supercode.jsonl".into(),
5535 content: session.to_native_jsonl(),
5536 role: ArtifactFileRole::SourceRecovery,
5537 });
5538 for (index, subagent) in session.subagents.iter().enumerate() {
5539 let id = subagent
5540 .meta
5541 .agent_id
5542 .as_deref()
5543 .map(sanitize_filename)
5544 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5545 files.push(SessionArtifactFile {
5546 path: format!("recovery/subagents/{id}.supercode.jsonl"),
5547 content: subagent.to_native_jsonl(),
5548 role: ArtifactFileRole::SourceRecovery,
5549 });
5550 }
5551 }
5552 if !diagonal && session.meta.source == SessionSource::Grok {
5553 append_grok_bundle_files(
5554 locator,
5555 "recovery/grok/",
5556 ArtifactFileRole::SourceRecovery,
5557 &mut files,
5558 )?;
5559 }
5560 let (fidelity, residue) = if diagonal
5561 && target_session_id.is_none()
5562 && session.raw_is_verbatim
5563 && !has_appended_turns
5564 {
5565 (Fidelity::ByteLossless, Vec::new())
5566 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
5567 (
5568 Fidelity::ValueLossless,
5569 vec![if target_session_id.is_some() {
5570 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
5571 } else {
5572 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
5573 }],
5574 )
5575 } else {
5576 match restoration {
5577 Some(report) if report.rendered_messages == 0 => (
5578 Fidelity::ByteLossless,
5579 vec![format!(
5580 "restored verbatim from this conversation's {} source records in the residue store",
5581 target.id()
5582 )],
5583 ),
5584 Some(report) => (
5585 Fidelity::Semantic,
5586 vec![format!(
5587 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
5588 report.restored_messages,
5589 report.restored_messages + report.rendered_messages,
5590 report.rendered_messages,
5591 target.id()
5592 )],
5593 ),
5594 None => (
5595 Fidelity::Semantic,
5596 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
5597 ),
5598 }
5599 };
5600 Ok(SessionArtifact {
5601 source_harness: locator.harness.clone(),
5602 target_harness: target.id(),
5603 session_id: target_session_id
5604 .map(str::to_string)
5605 .or_else(|| session.meta.session_id.clone()),
5606 content,
5607 suggested_filename,
5608 files,
5609 fidelity,
5610 residue,
5611 })
5612}
5613
5614fn append_grok_bundle_files(
5615 locator: &SessionLocator,
5616 prefix: &str,
5617 role: ArtifactFileRole,
5618 files: &mut Vec<SessionArtifactFile>,
5619) -> std::result::Result<(), ServiceError> {
5620 let primary = locator.storage.path();
5621 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5622 return Err(ServiceError::Operation(format!(
5623 "Grok bundle locator must name chat_history.jsonl, got {}",
5624 primary.display()
5625 )));
5626 }
5627 let parent = primary.parent().ok_or_else(|| {
5628 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5629 })?;
5630 for name in ["summary.json", "updates.jsonl"] {
5631 let path = parent.join(name);
5632 let metadata = match std::fs::symlink_metadata(&path) {
5633 Ok(metadata) => metadata,
5634 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5635 Err(error) => return Err(ServiceError::Operation(error.to_string())),
5636 };
5637 if metadata.file_type().is_symlink() || !metadata.is_file() {
5638 return Err(ServiceError::Operation(format!(
5639 "refusing non-regular Grok bundle member {}",
5640 path.display()
5641 )));
5642 }
5643 let content = std::fs::read_to_string(&path).map_err(|error| {
5644 ServiceError::Operation(format!(
5645 "Grok bundle member {} is not representable as UTF-8: {error}",
5646 path.display()
5647 ))
5648 })?;
5649 files.push(SessionArtifactFile {
5650 path: format!("{prefix}{name}"),
5651 content,
5652 role: match role {
5653 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5654 _ => ArtifactFileRole::SourceRecovery,
5655 },
5656 });
5657 }
5658 Ok(())
5659}
5660
5661fn handoff_artifact(
5662 locator: &SessionLocator,
5663 session: &Session,
5664 target: TransferFormat,
5665) -> std::result::Result<SessionArtifact, ServiceError> {
5666 let target_session_id = target_session_id(target);
5667 session_artifact_with_id(locator, session, target, Some(&target_session_id))
5668}
5669
5670fn target_session_id(target: TransferFormat) -> String {
5671 let uuid = generated_session_id();
5672 match target {
5673 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5674 TransferFormat::ClaudeCode
5675 | TransferFormat::Codex
5676 | TransferFormat::Pi
5677 | TransferFormat::Grok
5678 | TransferFormat::Gemini
5679 | TransferFormat::Goose
5680 | TransferFormat::Hermes => uuid,
5681 }
5682}
5683
5684fn sanitize_filename(value: &str) -> String {
5685 let value = value
5686 .chars()
5687 .map(|character| {
5688 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5689 character
5690 } else {
5691 '-'
5692 }
5693 })
5694 .collect::<String>();
5695 let value = value.trim_matches('-');
5696 if value.is_empty() {
5697 "session".into()
5698 } else {
5699 value.chars().take(100).collect()
5700 }
5701}
5702
5703fn handoff_instructions(
5704 target: TransferFormat,
5705 session_id: &str,
5706 cwd: &Path,
5707) -> HandoffInstructions {
5708 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5709 cwd: cwd.to_path_buf(),
5710 program: program.into(),
5711 arguments,
5712 env: BTreeMap::new(),
5713 };
5714 match target {
5715 TransferFormat::ClaudeCode => HandoffInstructions {
5716 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5717 materialize: None,
5718 requires_materialization: true,
5719 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(),
5720 },
5721 TransferFormat::Hermes => HandoffInstructions {
5722 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5723 materialize: None,
5724 requires_materialization: true,
5725 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(),
5726 },
5727 TransferFormat::Codex => HandoffInstructions {
5728 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5729 materialize: None,
5730 requires_materialization: true,
5731 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5732 },
5733 TransferFormat::OpenCode => HandoffInstructions {
5734 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5735 materialize: Some(launch(
5736 "opencode",
5737 vec!["import".into(), "{artifact_path}".into()],
5738 )),
5739 requires_materialization: true,
5740 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5741 },
5742 TransferFormat::Pi => HandoffInstructions {
5743 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5744 materialize: None,
5745 requires_materialization: true,
5746 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5747 },
5748 TransferFormat::Grok => HandoffInstructions {
5749 launch: launch(
5750 "grok",
5751 vec!["--resume".into(), "{materialized_session_id}".into()],
5752 ),
5753 materialize: None,
5754 requires_materialization: true,
5755 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(),
5756 },
5757 TransferFormat::Gemini => HandoffInstructions {
5758 launch: launch(
5759 "gemini",
5760 vec!["--session-file".into(), "{artifact_path}".into()],
5761 ),
5762 materialize: None,
5763 requires_materialization: true,
5764 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(),
5765 },
5766 TransferFormat::Goose => HandoffInstructions {
5767 launch: launch(
5768 "goose",
5769 vec![
5770 "session".into(),
5771 "--resume".into(),
5772 "--session-id".into(),
5773 "{imported_session_id}".into(),
5774 ],
5775 ),
5776 materialize: Some(launch(
5777 "goose",
5778 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5779 )),
5780 requires_materialization: true,
5781 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(),
5782 },
5783 }
5784}
5785
5786fn resume_launch(
5787 harness: &str,
5788 session_id: &str,
5789 cwd: &Path,
5790 policy: ResumePolicy,
5791) -> std::result::Result<StructuredLaunch, ServiceError> {
5792 let mut arguments = Vec::new();
5793 let program = match harness {
5794 HarnessId::GROK => {
5795 if matches!(policy, ResumePolicy::Yolo) {
5796 if crate::support::self_sandbox_supported() {
5797 arguments.extend(["--sandbox".into(), "workspace".into()]);
5798 }
5799 arguments.push("--always-approve".into());
5800 }
5801 arguments.extend(["--resume".into(), session_id.into()]);
5802 "grok"
5803 }
5804 HarnessId::CODEX => {
5805 arguments.extend(crate::startup_prompts::startup_arguments(
5806 harness,
5807 Some(cwd),
5808 &[],
5809 matches!(policy, ResumePolicy::Yolo),
5810 ));
5811 arguments.extend(["resume".into(), session_id.into()]);
5812 "codex"
5813 }
5814 HarnessId::CLAUDE_CODE => {
5815 arguments.extend(crate::startup_prompts::startup_arguments(
5816 harness,
5817 Some(cwd),
5818 &[],
5819 matches!(policy, ResumePolicy::Yolo),
5820 ));
5821 arguments.extend(["--resume".into(), session_id.into()]);
5822 "claude"
5823 }
5824 HarnessId::GEMINI => {
5825 arguments.extend(crate::startup_prompts::startup_arguments(
5826 harness,
5827 Some(cwd),
5828 &[],
5829 matches!(policy, ResumePolicy::Yolo),
5830 ));
5831 arguments.extend(["--resume".into(), session_id.into()]);
5832 "gemini"
5833 }
5834 HarnessId::GOOSE => {
5835 arguments.extend([
5836 "session".into(),
5837 "--resume".into(),
5838 "--session-id".into(),
5839 session_id.into(),
5840 ]);
5841 "goose"
5842 }
5843 HarnessId::PI => {
5844 arguments.extend(crate::startup_prompts::startup_arguments(
5845 harness,
5846 Some(cwd),
5847 &[],
5848 matches!(policy, ResumePolicy::Yolo),
5849 ));
5850 arguments.extend(["--session".into(), session_id.into()]);
5851 "pi"
5852 }
5853 HarnessId::OPENCODE => {
5854 arguments.extend(["--session".into(), session_id.into()]);
5855 "opencode"
5856 }
5857 HarnessId::SUPERCODE => {
5858 if matches!(policy, ResumePolicy::Yolo) {
5859 arguments.push("--dangerous".into());
5860 }
5861 arguments.extend(["resume".into(), session_id.into()]);
5862 "supercode"
5863 }
5864 other => {
5865 return Err(ServiceError::InvalidParams(format!(
5866 "no structured resume launch is registered for harness `{other}`"
5867 )))
5868 }
5869 };
5870 Ok(StructuredLaunch {
5871 cwd: cwd.to_path_buf(),
5872 env: if program == "grok" {
5873 crate::support::grok_home_env()
5874 } else {
5875 BTreeMap::new()
5876 },
5877 program: program.into(),
5878 arguments,
5879 })
5880}
5881
5882fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5888 let digest = blake3::hash(address.as_bytes()).to_hex();
5889 let path = std::env::temp_dir().join(format!(
5890 "supercode-openclaw-gateway-token-{}",
5891 &digest.as_str()[..16]
5892 ));
5893 #[cfg(unix)]
5894 {
5895 use std::io::Write;
5896 use std::os::unix::fs::OpenOptionsExt;
5897 let mut file = std::fs::OpenOptions::new()
5898 .write(true)
5899 .create(true)
5900 .truncate(true)
5901 .mode(0o600)
5902 .open(&path)?;
5903 file.write_all(secret.as_bytes())?;
5904 }
5905 #[cfg(not(unix))]
5906 std::fs::write(&path, secret)?;
5907 Ok(path)
5908}
5909
5910fn open_connect_descriptor(
5916 descriptor: &crate::HarnessSupportDescriptor,
5917 home: &Path,
5918) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5919 let Some(connect) = &descriptor.runtime.connect_launch else {
5920 return Err(ServiceError::InvalidParams(format!(
5921 "harness `{}` has no registered connect-mode launch",
5922 descriptor.id.as_str()
5923 )));
5924 };
5925 let resolved = connect
5926 .resolve(home)
5927 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5928 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5929 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5930 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5931 if let Some(token) = resolved.auth {
5932 backend = backend.with_bearer(token);
5933 }
5934 Ok(Box::new(backend))
5935 }
5936 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5937 let mut env = BTreeMap::new();
5948 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5949 if let Some(token) = resolved.auth {
5950 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5951 .map_err(|error| {
5952 ServiceError::UnsupportedAction(format!(
5953 "could not stage the gateway credential for the bridge: {error}"
5954 ))
5955 })?;
5956 arguments.push("--token-file".into());
5957 arguments.push(token_path.to_string_lossy().into_owned());
5958 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5959 }
5960 let program = descriptor
5965 .runtime
5966 .default_launch
5967 .as_ref()
5968 .map(|launch| launch.program.clone())
5969 .unwrap_or_else(|| "openclaw".into());
5970 let launch = RuntimeLaunch {
5971 program,
5972 arguments,
5973 env,
5974 };
5975 Ok(Box::new(
5976 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5977 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5978 ))
5979 }
5980 _ => Err(ServiceError::UnsupportedAction(format!(
5981 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5982 descriptor.id.as_str(),
5983 connect.protocol
5984 ))),
5985 }
5986}
5987
5988fn registry_connect_descriptor(
5991 params: &RuntimeBackendParams,
5992) -> Option<crate::HarnessSupportDescriptor> {
5993 if params.launch.is_some() || params.base_url.is_some() {
5994 return None;
5995 }
5996 harness_support_registry()
5997 .harnesses
5998 .into_iter()
5999 .find(|descriptor| descriptor.id == params.harness)
6000 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
6001}
6002
6003fn service_home() -> std::result::Result<PathBuf, ServiceError> {
6004 supercode_interchange::user_home()
6005 .map(std::path::PathBuf::into_os_string)
6006 .map(PathBuf::from)
6007 .ok_or_else(|| {
6008 ServiceError::UnsupportedAction(
6009 "connect-mode launches need HOME to locate the harness config".into(),
6010 )
6011 })
6012}
6013
6014pub const RUNTIME_OPEN_METHODS: &[&str] = &[
6017 "harness.v1.runtimes.start",
6018 "harness.v1.runtimes.resume",
6019 "harness.v1.runtimes.attach",
6020 "harness.v1.runtimes.attach_existing",
6021];
6022
6023pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
6028
6029pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
6035
6036pub const DETACHED_METHODS: &[&str] = &[
6042 "harness.v1.harnesses.list",
6043 "harness.v1.sessions.load",
6044 "harness.v1.harnesses.probe",
6045 "harness.v1.sessions.message",
6046 "harness.v1.sessions.new",
6047 "harness.v1.sessions.reset",
6048 "harness.v1.sessions.archive",
6049 "harness.v1.sessions.delete",
6050];
6051
6052pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
6058
6059pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
6067
6068async fn within_control_deadline<F: std::future::Future>(
6071 method: &str,
6072 call: F,
6073) -> std::result::Result<F::Output, ServiceError> {
6074 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
6075 .await
6076 .map_err(|_| {
6077 ServiceError::Operation(format!(
6078 "`{method}` gave up after {}s: the runtime did not answer",
6079 RUNTIME_CONTROL_DEADLINE.as_secs()
6080 ))
6081 })
6082}
6083
6084pub struct RuntimeOpen {
6088 id: Value,
6089 method: String,
6090 params: Value,
6091}
6092
6093impl RuntimeOpen {
6094 pub async fn open(self) -> OpenedRuntime {
6098 let Self { id, method, params } = self;
6099 let outcome = open_runtime(&method, params).await;
6100 OpenedRuntime { id, outcome }
6101 }
6102}
6103
6104pub struct OpenedRuntime {
6107 id: Value,
6108 outcome: std::result::Result<OpenRuntime, ServiceError>,
6109}
6110
6111pub struct DetachedCall {
6116 id: Value,
6117 method: String,
6118 work: std::result::Result<Work, ServiceError>,
6119}
6120
6121impl DetachedCall {
6122 pub async fn run(self) -> DetachedAnswer {
6125 let Self { id, method, work } = self;
6126 match work {
6127 Ok(Work::Runtime(work)) => {
6132 let (result, returned) = work.run().await;
6133 DetachedAnswer {
6134 response: service_response(id, result),
6135 returned,
6136 }
6137 }
6138 Ok(Work::Free(work)) => {
6139 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
6140 Ok(result) => result,
6141 Err(_) => Err(ServiceError::Operation(format!(
6142 "`{method}` gave up after {}s: the harness it waits on did not answer",
6143 DETACHED_CALL_DEADLINE.as_secs()
6144 ))),
6145 };
6146 DetachedAnswer {
6147 response: service_response(id, result),
6148 returned: None,
6149 }
6150 }
6151 Err(error) => DetachedAnswer {
6152 response: service_response(id, Err(error)),
6153 returned: None,
6154 },
6155 }
6156 }
6157}
6158
6159pub struct DetachedAnswer {
6163 response: Value,
6164 returned: Option<ReturnedRuntime>,
6165}
6166
6167impl DetachedAnswer {
6168 pub fn into_response(self) -> Value {
6171 self.response
6172 }
6173}
6174
6175pub struct ReturnedRuntime {
6178 connection: String,
6179 runtime: Box<dyn RuntimeConnection>,
6180}
6181
6182enum Work {
6185 Free(DetachedWork),
6186 Runtime(RuntimeWork),
6187}
6188
6189enum DetachedWork {
6192 Inventory(InventoryWork),
6196 Message(MessageSessionParams),
6198 Load(LoadSessionParams),
6200 SessionMutation {
6203 verb: crate::SessionVerb,
6204 mutation: crate::SessionMutation,
6205 },
6206}
6207
6208impl DetachedWork {
6209 async fn run(self) -> std::result::Result<Value, ServiceError> {
6210 match self {
6211 Self::Inventory(work) => run_inventory(work).await,
6212 Self::Message(params) => Ok(message_live_session(¶ms).await),
6213 Self::Load(params) => tokio::task::spawn_blocking(move || load_session_door(params))
6214 .await
6215 .map_err(|error| {
6216 ServiceError::Operation(format!("the session read stopped: {error}"))
6217 })?,
6218 Self::SessionMutation { verb, mutation } => {
6219 let outcome = run_session_mutation(verb, &mutation).await?;
6220 serde_json::to_value(outcome)
6221 .map_err(|error| ServiceError::Operation(error.to_string()))
6222 }
6223 }
6224 }
6225}
6226
6227enum RuntimeWork {
6229 Close {
6231 runtime: Box<dyn RuntimeConnection>,
6232 process_group: Option<u32>,
6233 },
6234 LiveCommand {
6237 connection: String,
6238 runtime: Box<dyn RuntimeConnection>,
6239 verb: crate::SessionVerb,
6240 mutation: crate::SessionMutation,
6241 command: &'static str,
6242 session: String,
6243 },
6244}
6245
6246type RuntimeWorkAnswer = (
6249 std::result::Result<Value, ServiceError>,
6250 Option<ReturnedRuntime>,
6251);
6252
6253impl RuntimeWork {
6254 async fn run(self) -> RuntimeWorkAnswer {
6255 match self {
6256 Self::Close {
6257 runtime,
6258 process_group,
6259 } => (close_runtime(runtime, process_group).await, None),
6260 Self::LiveCommand {
6261 connection,
6262 mut runtime,
6263 verb,
6264 mutation,
6265 command,
6266 session,
6267 } => {
6268 let result =
6269 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
6270 (
6271 result,
6272 Some(ReturnedRuntime {
6273 connection,
6274 runtime,
6275 }),
6276 )
6277 }
6278 }
6279 }
6280}
6281
6282async fn close_runtime(
6285 mut runtime: Box<dyn RuntimeConnection>,
6286 process_group: Option<u32>,
6287) -> std::result::Result<Value, ServiceError> {
6288 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
6289 Ok(result) => {
6290 result.map_err(operation)?;
6291 Ok(json!({"closed": true}))
6292 }
6293 Err(deadline) => {
6294 let killed = kill_runtime_process_group(process_group);
6299 drop(runtime);
6300 Ok(json!({
6301 "closed": true,
6302 "killed": killed,
6303 "detail": error_message(deadline),
6304 }))
6305 }
6306 }
6307}
6308
6309fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
6312 mutation
6313 .session
6314 .clone()
6315 .filter(|value| !value.trim().is_empty())
6316 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
6317}
6318
6319async fn type_live_command(
6323 runtime: &mut dyn RuntimeConnection,
6324 verb: crate::SessionVerb,
6325 mutation: &crate::SessionMutation,
6326 command: &str,
6327 session: String,
6328) -> std::result::Result<Value, ServiceError> {
6329 within_control_deadline(
6330 &format!("sessions.{}", verb.as_str()),
6331 runtime.send_input(RuntimeInput {
6332 text: command.to_string(),
6333 image_urls: Vec::new(),
6334 }),
6335 )
6336 .await?
6337 .map_err(operation)?;
6338 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
6339 .map_err(session_control_error)?;
6340 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
6341}
6342
6343enum OpenRuntime {
6346 Hosted {
6349 runtime: Box<dyn RuntimeConnection>,
6350 capabilities: crate::RuntimeCapabilities,
6351 workspace: PathBuf,
6352 fresh: bool,
6354 },
6355 Joined { runtime: Box<dyn RuntimeConnection> },
6358}
6359
6360async fn open_runtime(
6365 method: &str,
6366 params: Value,
6367) -> std::result::Result<OpenRuntime, ServiceError> {
6368 match tokio::time::timeout(
6369 RUNTIME_OPEN_DEADLINE,
6370 open_runtime_unbounded(method, params),
6371 )
6372 .await
6373 {
6374 Ok(result) => result,
6375 Err(_) => Err(ServiceError::Operation(format!(
6376 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
6377 RUNTIME_OPEN_DEADLINE.as_secs()
6378 ))),
6379 }
6380}
6381
6382async fn open_runtime_unbounded(
6383 method: &str,
6384 params: Value,
6385) -> std::result::Result<OpenRuntime, ServiceError> {
6386 match method {
6387 "harness.v1.runtimes.start" => {
6388 let params = decode::<RuntimeStartParams>(params)?;
6389 let backend = runtime_backend(¶ms.backend)?;
6390 let capabilities = backend.capabilities();
6391 let workspace = params.cwd.clone();
6392 let runtime = backend
6393 .start(RuntimeStartRequest {
6394 cwd: params.cwd,
6395 launch: runtime_launch(¶ms.backend),
6396 mcp_servers: params.mcp_servers,
6397 approval_policy: params.approval_policy,
6398 })
6399 .await
6400 .map_err(operation)?;
6401 Ok(OpenRuntime::Hosted {
6402 runtime,
6403 capabilities,
6404 workspace,
6405 fresh: true,
6406 })
6407 }
6408 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
6409 let params = decode::<RuntimeAttachParams>(params)?;
6410 let backend = runtime_backend(¶ms.backend)?;
6411 let capabilities = backend.capabilities();
6412 let workspace = params
6413 .cwd
6414 .clone()
6415 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
6416 let runtime = backend
6417 .attach(RuntimeAttachRequest {
6418 runtime_id: params.runtime_id,
6419 cwd: params.cwd,
6420 launch: runtime_launch(¶ms.backend),
6421 mcp_servers: params.mcp_servers,
6422 approval_policy: params.approval_policy,
6423 })
6424 .await
6425 .map_err(operation)?;
6426 Ok(OpenRuntime::Hosted {
6427 runtime,
6428 capabilities,
6429 workspace,
6430 fresh: false,
6431 })
6432 }
6433 "harness.v1.runtimes.attach_existing" => {
6434 let params = decode::<RuntimeAttachParams>(params)?;
6435 let backend: Box<dyn RuntimeBackend> = match params
6436 .backend
6437 .base_url
6438 .as_deref()
6439 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
6440 {
6441 Some(endpoint) => {
6442 #[cfg(not(feature = "adapter-api"))]
6443 {
6444 let _ = endpoint;
6445 return Err(ServiceError::UnsupportedAction(
6446 "live HTTP attachment adapter is not compiled".into(),
6447 ));
6448 }
6449 #[cfg(feature = "adapter-api")]
6450 {
6451 let workspace = params.cwd.clone().ok_or_else(|| {
6452 ServiceError::InvalidParams(
6453 "Volter Harness live attach requires the project cwd".into(),
6454 )
6455 })?;
6456 let source = LiveRuntimeSource {
6457 harness: params.backend.harness.as_str().to_string(),
6458 session_id: params.runtime_id.clone(),
6459 workspace,
6460 };
6461 let receipt = resolve_live_runtime(&endpoint, &source)
6462 .map_err(|error| ServiceError::Operation(error.to_string()))?;
6463 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
6464 }
6465 }
6466 None => runtime_backend(¶ms.backend)?,
6467 };
6468 let capabilities = backend.capabilities();
6469 if !capabilities.attach_existing_process {
6470 return Err(ServiceError::Operation(format!(
6471 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
6472 backend.harness().as_str()
6473 )));
6474 }
6475 let runtime = backend
6476 .attach_existing(RuntimeAttachRequest {
6477 runtime_id: params.runtime_id,
6478 cwd: params.cwd,
6479 launch: runtime_launch(¶ms.backend),
6480 mcp_servers: params.mcp_servers,
6481 approval_policy: params.approval_policy,
6482 })
6483 .await
6484 .map_err(operation)?;
6485 Ok(OpenRuntime::Joined { runtime })
6486 }
6487 _ => Err(ServiceError::MethodNotFound),
6488 }
6489}
6490
6491fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
6493 match result {
6494 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
6495 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
6496 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
6497 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
6498 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
6499 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
6500 }
6501}
6502
6503fn runtime_backend(
6504 params: &RuntimeBackendParams,
6505) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
6506 if let Some(descriptor) = registry_connect_descriptor(params) {
6507 return open_connect_descriptor(&descriptor, &service_home()?);
6508 }
6509 if params.protocol.as_deref() == Some("acp") {
6510 let launch = params
6511 .launch
6512 .clone()
6513 .or_else(|| {
6514 harness_support_registry()
6515 .harnesses
6516 .into_iter()
6517 .find(|harness| harness.id == params.harness)
6518 .filter(|harness| {
6519 harness.runtime.implementation == ImplementationKind::GenericProtocol
6520 && harness.runtime.protocol.starts_with("acp")
6521 })
6522 .and_then(|harness| harness.runtime.default_launch)
6523 })
6524 .ok_or_else(|| {
6525 ServiceError::InvalidParams(
6526 "an ACP runtime requires `launch` unless the harness has a registered default"
6527 .into(),
6528 )
6529 })?;
6530 let resume_session = harness_support_registry()
6531 .harnesses
6532 .into_iter()
6533 .find(|harness| harness.id == params.harness)
6534 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
6535 return Ok(Box::new(
6536 AcpRuntimeBackend::new(params.harness.clone(), launch)
6537 .with_resume_support(resume_session),
6538 ));
6539 }
6540 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
6541 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
6542 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
6543 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
6544 HarnessId::OPENCODE => match ¶ms.base_url {
6545 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
6546 None => Box::new(OpenCodeRuntimeBackend::new()),
6547 },
6548 harness => {
6549 let descriptor = harness_support_registry()
6550 .harnesses
6551 .into_iter()
6552 .find(|descriptor| descriptor.id.as_str() == harness)
6553 .filter(|descriptor| {
6554 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
6555 && descriptor.runtime.protocol.starts_with("acp")
6556 });
6557 let Some(descriptor) = descriptor else {
6558 return Err(ServiceError::InvalidParams(format!(
6559 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
6560 )));
6561 };
6562 let resume = descriptor.runtime.capabilities.resume_session;
6563 Box::new(
6564 AcpRuntimeBackend::new(
6565 descriptor.id,
6566 descriptor
6567 .runtime
6568 .default_launch
6569 .expect("generic ACP registry entry includes its launch"),
6570 )
6571 .with_resume_support(resume),
6572 )
6573 }
6574 };
6575 Ok(backend)
6576}
6577
6578fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
6579 if let Some(launch) = ¶ms.launch {
6580 return Some(launch.clone());
6581 }
6582 if !matches!(params.policy, RuntimePolicy::Yolo) {
6583 return None;
6584 }
6585 let launch = match params.harness.as_str() {
6586 HarnessId::GROK => RuntimeLaunch {
6587 program: "grok".into(),
6588 arguments: {
6589 let mut arguments: Vec<String> = Vec::new();
6590 if crate::support::self_sandbox_supported() {
6591 arguments.extend(["--sandbox".into(), "workspace".into()]);
6592 }
6593 arguments.extend([
6594 "--always-approve".into(),
6595 "agent".into(),
6596 "--no-leader".into(),
6597 "stdio".into(),
6598 ]);
6599 arguments
6600 },
6601 env: crate::support::grok_env(),
6602 },
6603 HarnessId::CODEX => RuntimeLaunch {
6604 program: "codex".into(),
6605 arguments: vec![
6606 "--dangerously-bypass-approvals-and-sandbox".into(),
6607 "--dangerously-bypass-hook-trust".into(),
6608 "app-server".into(),
6609 ],
6610 env: BTreeMap::new(),
6611 },
6612 HarnessId::CLAUDE_CODE => RuntimeLaunch {
6613 program: "claude".into(),
6614 arguments: vec![
6615 "--dangerously-skip-permissions".into(),
6616 "--print".into(),
6617 "--input-format".into(),
6618 "stream-json".into(),
6619 "--output-format".into(),
6620 "stream-json".into(),
6621 "--verbose".into(),
6622 ],
6623 env: BTreeMap::new(),
6624 },
6625 HarnessId::PI => RuntimeLaunch {
6626 program: "pi".into(),
6627 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
6628 env: BTreeMap::new(),
6629 },
6630 HarnessId::OPENCODE => RuntimeLaunch {
6631 program: "opencode".into(),
6632 arguments: vec!["serve".into()],
6633 env: BTreeMap::new(),
6634 },
6635 HarnessId::GEMINI => RuntimeLaunch {
6636 program: "gemini".into(),
6637 arguments: vec!["--acp".into(), "--yolo".into()],
6638 env: BTreeMap::new(),
6639 },
6640 HarnessId::GOOSE => RuntimeLaunch {
6641 program: "goose".into(),
6642 arguments: vec!["acp".into()],
6643 env: BTreeMap::new(),
6644 },
6645 HarnessId::SUPERCODE => RuntimeLaunch {
6646 program: "supercode".into(),
6647 arguments: vec!["acp".into(), "--dangerous".into()],
6648 env: BTreeMap::new(),
6649 },
6650 _ => return None,
6651 };
6652 Some(launch)
6653}
6654
6655struct IsolatedProbeHome {
6661 launch: RuntimeLaunch,
6662 root: PathBuf,
6663}
6664
6665impl IsolatedProbeHome {
6666 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6667 let root = std::env::temp_dir().join(format!(
6668 "supercode-harness-probe-{harness}-{}",
6669 generated_session_id()
6670 ));
6671 std::fs::create_dir_all(&root)?;
6672 set_private_dir_permissions(&root)?;
6673
6674 if let Some(source_home) = supercode_interchange::user_home()
6675 .map(std::path::PathBuf::into_os_string)
6676 .map(PathBuf::from)
6677 {
6678 for relative in probe_auth_files(harness) {
6679 copy_probe_file(&source_home, &root, relative)?;
6680 }
6681 }
6682 if harness == HarnessId::SUPERCODE {
6687 let config_home = crate::agent::global_instructions_dir();
6688 for file in ["config.toml", "credentials.toml"] {
6689 copy_probe_path(
6690 &config_home.join(file),
6691 &root.join(".config/supercode").join(file),
6692 )?;
6693 }
6694 }
6695 configure_isolated_probe_auth(harness, &root)?;
6696
6697 let root_text = root.to_string_lossy().into_owned();
6698 for (key, value) in [
6699 ("HOME", root_text.clone()),
6700 (
6701 "XDG_CACHE_HOME",
6702 root.join(".cache").to_string_lossy().into_owned(),
6703 ),
6704 (
6705 "XDG_CONFIG_HOME",
6706 root.join(".config").to_string_lossy().into_owned(),
6707 ),
6708 (
6709 "XDG_DATA_HOME",
6710 root.join(".local/share").to_string_lossy().into_owned(),
6711 ),
6712 ] {
6713 launch.env.insert(key.into(), value);
6714 }
6715 let scoped = match harness {
6716 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6717 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6718 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6719 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6720 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6721 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6722 _ => None,
6723 };
6724 if let Some((key, value)) = scoped {
6725 launch
6726 .env
6727 .insert(key.into(), value.to_string_lossy().into_owned());
6728 }
6729 Ok(Self { launch, root })
6730 }
6731
6732 fn cleanup(&self) -> std::io::Result<()> {
6733 match std::fs::remove_dir_all(&self.root) {
6734 Ok(()) => Ok(()),
6735 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6736 Err(error) => Err(error),
6737 }
6738 }
6739}
6740
6741impl Drop for IsolatedProbeHome {
6742 fn drop(&mut self) {
6743 let _ = self.cleanup();
6744 }
6745}
6746
6747fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6748 match harness {
6749 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6750 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6754 HarnessId::CODEX => &[".codex/auth.json"],
6755 HarnessId::GEMINI => &[
6756 ".gemini/google_accounts.json",
6757 ".gemini/oauth_creds.json",
6758 ".gemini/settings.json",
6759 ],
6760 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6761 HarnessId::OPENCODE => &[
6762 ".config/opencode/auth.json",
6763 ".local/share/opencode/auth.json",
6764 ],
6765 HarnessId::PI => &[".pi/agent/auth.json"],
6766 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6771 _ => &[],
6772 }
6773}
6774
6775fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6776 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6777}
6778
6779fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6780 if !source.is_file() {
6781 return Ok(());
6782 }
6783 if let Some(parent) = destination.parent() {
6784 std::fs::create_dir_all(parent)?;
6785 set_private_dir_permissions(parent)?;
6786 }
6787 std::fs::copy(source, destination)?;
6788 set_private_file_permissions(destination)
6789}
6790
6791fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6792 if harness != HarnessId::GEMINI {
6793 return Ok(());
6794 }
6795 let oauth = probe_home.join(".gemini/oauth_creds.json");
6796 if !oauth.is_file() {
6797 return Ok(());
6798 }
6799 let settings_path = probe_home.join(".gemini/settings.json");
6800 let mut settings = std::fs::read_to_string(&settings_path)
6801 .ok()
6802 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6803 .unwrap_or_else(|| json!({}));
6804 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6805 std::fs::write(
6806 &settings_path,
6807 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6808 )?;
6809 set_private_file_permissions(&settings_path)
6810}
6811
6812#[cfg(unix)]
6813fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6814 use std::os::unix::fs::PermissionsExt;
6815 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6816}
6817
6818#[cfg(not(unix))]
6819fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6820 Ok(())
6821}
6822
6823#[cfg(unix)]
6824fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6825 use std::os::unix::fs::PermissionsExt;
6826 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6827}
6828
6829#[cfg(not(unix))]
6830fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6831 Ok(())
6832}
6833
6834fn find_executable(program: &str) -> Option<PathBuf> {
6835 let candidate = PathBuf::from(program);
6836 if candidate.components().count() > 1 {
6837 return candidate.is_file().then_some(candidate);
6838 }
6839 let path = std::env::var_os("PATH")?;
6840 for directory in std::env::split_paths(&path) {
6841 let candidate = directory.join(program);
6842 if candidate.is_file() {
6843 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6844 }
6845 #[cfg(windows)]
6846 {
6847 for extension in ["exe", "cmd", "bat"] {
6848 let candidate = directory.join(format!("{program}.{extension}"));
6849 if candidate.is_file() {
6850 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6851 }
6852 }
6853 }
6854 }
6855 None
6856}
6857
6858async fn executable_version(executable: &Path) -> Option<String> {
6859 let mut command = tokio::process::Command::new(executable);
6860 command
6861 .arg("--version")
6862 .stdin(std::process::Stdio::null())
6863 .stdout(std::process::Stdio::piped())
6864 .stderr(std::process::Stdio::piped())
6865 .kill_on_drop(true);
6866 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6867 .await
6868 .ok()?
6869 .ok()?;
6870 let stdout = String::from_utf8_lossy(&output.stdout);
6871 let stderr = String::from_utf8_lossy(&output.stderr);
6872 stdout
6873 .lines()
6874 .chain(stderr.lines())
6875 .map(str::trim)
6876 .find(|line| !line.is_empty())
6877 .map(|line| truncate_text(line, 200))
6878}
6879
6880pub(crate) fn auth_evidence(harness: &str) -> bool {
6881 let env_names: &[&str] = match harness {
6882 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6883 HarnessId::CODEX => &["OPENAI_API_KEY"],
6884 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6885 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6886 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6887 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6888 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6889 _ => &[],
6890 };
6891 if env_names
6892 .iter()
6893 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6894 {
6895 return true;
6896 }
6897 let Some(home) = supercode_interchange::user_home()
6898 .map(std::path::PathBuf::into_os_string)
6899 .map(PathBuf::from)
6900 else {
6901 return false;
6902 };
6903 let files: Vec<PathBuf> = match harness {
6904 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6905 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6906 HarnessId::OPENCODE => vec![
6907 home.join(".local/share/opencode/auth.json"),
6908 home.join(".config/opencode/auth.json"),
6909 ],
6910 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6911 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6912 HarnessId::GEMINI => vec![
6913 home.join(".gemini/oauth_creds.json"),
6914 home.join(".gemini/google_accounts.json"),
6915 ],
6916 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6917 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6918 _ => Vec::new(),
6919 };
6920 if files.into_iter().any(|path| {
6921 std::fs::metadata(path)
6922 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6923 .unwrap_or(false)
6924 }) {
6925 return true;
6926 }
6927 if harness == HarnessId::CLAUDE_CODE {
6934 return std::fs::read_to_string(home.join(".claude.json"))
6935 .map(|text| text.contains("\"oauthAccount\""))
6936 .unwrap_or(false);
6937 }
6938 false
6939}
6940
6941fn looks_like_auth_error(message: &str) -> bool {
6942 let message = message.to_ascii_lowercase();
6943 [
6944 "auth",
6945 "login",
6946 "sign in",
6947 "sign-in",
6948 "credential",
6949 "unauthorized",
6950 "forbidden",
6951 "token",
6952 ]
6953 .iter()
6954 .any(|needle| message.contains(needle))
6955}
6956
6957fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6958 crate::RuntimeCapabilities {
6959 start_session: false,
6960 resume_session: false,
6961 attach_existing_process: false,
6962 send_input: false,
6963 stream_events: false,
6964 interrupt: false,
6965 steer: false,
6966 respond_to_requests: false,
6967 }
6968}
6969
6970fn truncate_text(text: &str, max_chars: usize) -> String {
6971 let mut chars = text.chars();
6972 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6973 if chars.next().is_some() {
6974 format!("{truncated}…")
6975 } else {
6976 truncated
6977 }
6978}
6979
6980fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6987 match &handle.endpoint {
6988 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6989 crate::RuntimeEndpoint::Http { .. } => None,
6990 }
6991}
6992
6993fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6999 match process_group {
7000 #[cfg(unix)]
7001 Some(pid) => {
7002 crate::lsp::kill_process_group(pid);
7003 true
7004 }
7005 #[cfg(not(unix))]
7006 Some(_) => false,
7007 None => false,
7008 }
7009}
7010
7011fn error_message(error: ServiceError) -> String {
7012 match error {
7013 ServiceError::InvalidParams(message)
7014 | ServiceError::Operation(message)
7015 | ServiceError::UnsupportedAction(message) => message,
7016 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
7017 ServiceError::Sdk(error) => error.to_string(),
7018 }
7019}
7020
7021#[derive(Debug)]
7022enum ServiceError {
7023 InvalidParams(String),
7024 MethodNotFound,
7025 UnsupportedAction(String),
7026 Operation(String),
7027 Sdk(SdkError),
7028}
7029
7030fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
7031 match error {
7032 ServiceError::InvalidParams(message) => {
7033 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
7034 }
7035 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
7036 SdkError::unsupported(operation)
7037 }
7038 ServiceError::Operation(message) => {
7039 let code = if message.contains("already in progress") {
7040 SdkErrorCode::Busy
7041 } else if message.contains("not supported by this runtime") {
7042 SdkErrorCode::UnsupportedAction
7043 } else if message.contains("unknown runtime connection") {
7044 SdkErrorCode::NotFound
7045 } else {
7046 SdkErrorCode::Execution
7047 };
7048 SdkError::new(code, operation, message)
7049 }
7050 ServiceError::Sdk(error) => error,
7051 }
7052}
7053
7054fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
7055 let error_code = error.code();
7056 let code = match error_code {
7057 SdkErrorCode::Unauthenticated => -32030,
7058 SdkErrorCode::Unauthorized => -32031,
7059 SdkErrorCode::ControllerRequired => -32032,
7060 SdkErrorCode::LeaseExpired => -32033,
7061 SdkErrorCode::InvalidArgument => -32602,
7062 SdkErrorCode::NotFound => -32004,
7063 SdkErrorCode::Busy => -32000,
7064 SdkErrorCode::UnsupportedAction => -32020,
7065 SdkErrorCode::Execution => -32002,
7066 SdkErrorCode::Transport => -32003,
7067 };
7068 json!({
7069 "jsonrpc": "2.0",
7070 "id": id,
7071 "error": {
7072 "code": code,
7073 "name": error_code,
7074 "operation": error.operation(),
7075 "message": error.to_string(),
7076 },
7077 })
7078}
7079
7080fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
7081 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
7082}
7083
7084fn operation(error: impl Into<crate::Error>) -> ServiceError {
7085 let error = error.into();
7086 match error {
7087 crate::Error::Sdk(error) => ServiceError::Sdk(error),
7088 error => ServiceError::Operation(error.to_string()),
7089 }
7090}
7091
7092#[derive(Debug, Clone, Deserialize, Default)]
7096#[serde(default)]
7097struct MemoryRequest {
7098 harness: Option<String>,
7100 query: Option<String>,
7102 profile: Option<String>,
7104 session: Option<String>,
7106 full: bool,
7108 regex: bool,
7110 cwd: Option<std::path::PathBuf>,
7112 homes: crate::HarnessHomes,
7114}
7115
7116fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7119 let request = decode::<MemoryRequest>(params)?;
7120 let harness = request
7121 .harness
7122 .clone()
7123 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7124 let to_service = |error: crate::memory::MemoryError| match error {
7125 crate::memory::MemoryError::UnsupportedHarness { .. }
7126 | crate::memory::MemoryError::SessionNotScoped { .. } => {
7127 ServiceError::UnsupportedAction(error.to_string())
7128 }
7129 other => ServiceError::InvalidParams(other.to_string()),
7130 };
7131 match method {
7132 "harness.v1.memory.show" => {
7133 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
7134 harness,
7135 profile: request.profile,
7136 session: request.session,
7137 full: request.full,
7138 cwd: request.cwd,
7139 homes: request.homes,
7140 })
7141 .map_err(to_service)?;
7142 Ok(json!({
7143 "schema": crate::memory::MEMORY_SCHEMA,
7144 "documents": documents,
7145 }))
7146 }
7147 "harness.v1.memory.search" => {
7148 let query = request
7149 .query
7150 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
7151 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
7152 harness,
7153 query,
7154 profile: request.profile,
7155 regex: request.regex,
7156 cwd: request.cwd,
7157 homes: request.homes,
7158 })
7159 .map_err(to_service)?;
7160 Ok(json!({
7161 "schema": crate::memory::MEMORY_SCHEMA,
7162 "matches": matches,
7163 }))
7164 }
7165 _ => Err(ServiceError::MethodNotFound),
7166 }
7167}
7168
7169#[derive(Debug, Clone, Deserialize)]
7173#[serde(default)]
7174struct ProfilesQuery {
7175 harness: Option<String>,
7177 name: Option<String>,
7179 homes: crate::HarnessHomes,
7181}
7182
7183impl Default for ProfilesQuery {
7184 fn default() -> Self {
7185 Self {
7186 harness: None,
7187 name: None,
7188 homes: crate::HarnessHomes::default(),
7189 }
7190 }
7191}
7192
7193fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7196 let query = decode::<ProfilesQuery>(params)?;
7197 let to_service = |error: crate::profiles::ProfileError| match error {
7198 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
7199 ServiceError::UnsupportedAction(error.to_string())
7200 }
7201 crate::profiles::ProfileError::NotFound { .. } => {
7202 ServiceError::InvalidParams(error.to_string())
7203 }
7204 };
7205 match method {
7206 "harness.v1.profiles.list" => {
7207 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
7208 .map_err(to_service)?;
7209 Ok(json!({
7210 "schema": crate::profiles::PROFILES_SCHEMA,
7211 "profiles": profiles,
7212 }))
7213 }
7214 "harness.v1.profiles.get" => {
7215 let harness = query
7216 .harness
7217 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7218 let name = query
7219 .name
7220 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7221 let profile =
7222 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
7223 Ok(json!({
7224 "schema": crate::profiles::PROFILES_SCHEMA,
7225 "profile": profile,
7226 }))
7227 }
7228 _ => Err(ServiceError::MethodNotFound),
7229 }
7230}
7231
7232#[derive(Debug, Clone, Deserialize)]
7236#[serde(default)]
7237struct ChannelsQuery {
7238 harness: Option<String>,
7240 name: Option<String>,
7242 homes: crate::HarnessHomes,
7244}
7245
7246impl Default for ChannelsQuery {
7247 fn default() -> Self {
7248 Self {
7249 harness: None,
7250 name: None,
7251 homes: crate::HarnessHomes::default(),
7252 }
7253 }
7254}
7255
7256#[derive(Debug, Clone, Deserialize)]
7260#[serde(default)]
7261struct RoutesQuery {
7262 harness: Option<String>,
7263 profile: Option<String>,
7265 homes: crate::HarnessHomes,
7266}
7267
7268impl Default for RoutesQuery {
7269 fn default() -> Self {
7270 Self {
7271 harness: None,
7272 profile: None,
7273 homes: crate::HarnessHomes::default(),
7274 }
7275 }
7276}
7277
7278#[derive(Debug, Clone, Deserialize)]
7279#[serde(default)]
7280struct TriggersQuery {
7281 harness: Option<String>,
7282 homes: crate::HarnessHomes,
7283}
7284
7285impl Default for TriggersQuery {
7286 fn default() -> Self {
7287 Self {
7288 harness: None,
7289 homes: crate::HarnessHomes::default(),
7290 }
7291 }
7292}
7293
7294fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
7295 let query = decode::<TriggersQuery>(params)?;
7296 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
7297 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7298 Ok(json!({
7299 "schema": crate::triggers::TRIGGERS_SCHEMA,
7300 "triggers": triggers,
7301 }))
7302}
7303
7304fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
7305 let query = decode::<RoutesQuery>(params)?;
7306 let routes = crate::routes::list_routes(
7307 &query.homes,
7308 query.harness.as_deref(),
7309 query.profile.as_deref(),
7310 )
7311 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7312 Ok(json!({
7313 "schema": crate::routes::ROUTES_SCHEMA,
7314 "routes": routes,
7315 }))
7316}
7317
7318fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7319 let query = decode::<ChannelsQuery>(params)?;
7320 let to_service = |error: crate::channels::ChannelError| match error {
7321 crate::channels::ChannelError::UnsupportedHarness { .. } => {
7322 ServiceError::UnsupportedAction(error.to_string())
7323 }
7324 crate::channels::ChannelError::NotFound { .. } => {
7325 ServiceError::InvalidParams(error.to_string())
7326 }
7327 };
7328 match method {
7329 "harness.v1.channels.list" => {
7330 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
7331 .map_err(to_service)?;
7332 Ok(json!({
7333 "schema": crate::channels::CHANNELS_SCHEMA,
7334 "channels": channels,
7335 }))
7336 }
7337 "harness.v1.channels.status" => {
7338 let harness = query
7339 .harness
7340 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7341 let name = query
7342 .name
7343 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7344 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
7345 .map_err(to_service)?;
7346 Ok(json!({
7347 "schema": crate::channels::CHANNELS_SCHEMA,
7348 "channel": channel,
7349 }))
7350 }
7351 _ => Err(ServiceError::MethodNotFound),
7352 }
7353}
7354
7355fn rpc_error(id: Value, code: i64, message: &str) -> Value {
7356 json!({
7357 "jsonrpc": "2.0",
7358 "id": id,
7359 "error": {"code": code, "message": message},
7360 })
7361}