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