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 =
888 crate::workflow_doors::load(params.from, ¶ms.home).map_err(operation)?;
889 serde_json::to_value(read)
890 .map_err(|error| ServiceError::Operation(error.to_string()))
891 }
892 "harness.v1.orchestration.load" => {
893 let params = decode::<OrchestrationLoadParams>(params)?;
894 let read = crate::orchestration_doors::load(¶ms.root, params.flavor)
895 .map_err(operation)?;
896 serde_json::to_value(read)
897 .map_err(|error| ServiceError::Operation(error.to_string()))
898 }
899 "harness.v1.orchestration.save" => {
900 let params = decode::<OrchestrationSaveParams>(params)?;
901 let saved = crate::orchestration_doors::save(
902 ¶ms.root,
903 params.orchestration,
904 params.vault,
905 )
906 .map_err(operation)?;
907 serde_json::to_value(saved)
908 .map_err(|error| ServiceError::Operation(error.to_string()))
909 }
910 "harness.v1.orchestration.compile" => {
911 let params = decode::<OrchestrationCompileParams>(params)?;
912 let read = crate::orchestration_doors::compile(params.from, ¶ms.home)
913 .map_err(operation)?;
914 serde_json::to_value(read)
915 .map_err(|error| ServiceError::Operation(error.to_string()))
916 }
917 "harness.v1.orchestration.decompile" => {
918 let params = decode::<OrchestrationDecompileParams>(params)?;
919 let report = crate::orchestration_doors::decompile(
920 params.to,
921 params.orchestration,
922 ¶ms.source,
923 params.source_flavor,
924 ¶ms.dest,
925 params.vault,
926 )
927 .map_err(operation)?;
928 serde_json::to_value(report)
929 .map_err(|error| ServiceError::Operation(error.to_string()))
930 }
931 "harness.v1.orchestration.import" => {
936 let params = decode::<OrchestrationImportParams>(params)?;
937 let imported =
938 crate::orchestration_doors::import(params.from, ¶ms.home, ¶ms.into)
939 .map_err(operation)?;
940 serde_json::to_value(imported)
941 .map_err(|error| ServiceError::Operation(error.to_string()))
942 }
943 "harness.v1.orchestration.export" => {
944 let params = decode::<OrchestrationExportParams>(params)?;
945 let report =
946 crate::orchestration_doors::export(params.to, ¶ms.root, ¶ms.dest)
947 .map_err(operation)?;
948 serde_json::to_value(report)
949 .map_err(|error| ServiceError::Operation(error.to_string()))
950 }
951 "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
957 "harness.v1.skills.list" => {
962 let query = decode::<crate::skills::SkillsQuery>(params)?;
963 if let Some(harness) = query.harness.as_deref() {
964 if !crate::skills::SKILL_HARNESSES.contains(&harness) {
965 return Err(ServiceError::UnsupportedAction(format!(
966 "`{harness}` has no skills root Volter Harness reads"
967 )));
968 }
969 }
970 serde_json::to_value(crate::skills::list_skills(&query))
971 .map_err(|error| ServiceError::Operation(error.to_string()))
972 }
973 "harness.v1.skills.install" => {
981 mutate_skill(crate::skills_control::SkillVerb::Install, params)
982 }
983 "harness.v1.skills.remove" => {
984 mutate_skill(crate::skills_control::SkillVerb::Remove, params)
985 }
986 "harness.v1.approvals.list" => {
994 let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
995 if let Some(harness) = query.harness.as_deref() {
996 if !crate::approvals::lists_approvals(harness) {
997 return Err(ServiceError::UnsupportedAction(format!(
998 "`{harness}` has no runtime door that carries an approval request"
999 )));
1000 }
1001 }
1002 serde_json::to_value(self.approvals(&query))
1003 .map_err(|error| ServiceError::Operation(error.to_string()))
1004 }
1005 "harness.v1.sessions.discover" => {
1006 let query = decode::<DiscoveryQuery>(params)?;
1007 let mut page = discover_session_page(&query).map_err(operation)?;
1008 let doors = crate::mail_route::LiveSessions::read(&query.homes);
1013 if let Some(text) = query
1017 .query
1018 .as_deref()
1019 .map(str::trim)
1020 .filter(|text| !text.is_empty())
1021 {
1022 let needle = text.to_lowercase();
1023 for live in doors.all() {
1024 let name = live
1025 .name
1026 .split('@')
1027 .next()
1028 .unwrap_or(&live.name)
1029 .to_lowercase();
1030 if !name.contains(&needle)
1031 || page.sessions.iter().any(|session| {
1032 session.locator.session_id == live.address.session_id
1033 })
1034 {
1035 continue;
1036 }
1037 if query
1038 .limit
1039 .is_some_and(|limit| page.sessions.len() >= limit)
1040 {
1041 break;
1042 }
1043 let mut by_id = query.clone();
1045 by_id.query = Some(live.address.session_id.clone());
1046 by_id.cursor = None;
1047 if let Some(cwd) = &live.cwd {
1048 by_id.workspace = Some(cwd.clone());
1049 by_id.workspace_subtree = false;
1050 by_id.workspace_family = None;
1051 }
1052 if let Ok(found) = discover_session_page(&by_id) {
1053 page.sessions
1054 .extend(found.sessions.into_iter().filter(|session| {
1055 session.locator.session_id == live.address.session_id
1056 }));
1057 }
1058 }
1059 }
1060 let activities = crate::session_activity::resolve_stock_session_activities(
1061 &page
1062 .sessions
1063 .iter()
1064 .map(|session| session.locator.clone())
1065 .collect::<Vec<_>>(),
1066 &query.homes,
1067 )
1068 .into_iter()
1069 .map(|activity| (activity.key(), activity))
1070 .collect::<BTreeMap<_, _>>();
1071 let sessions = page
1072 .sessions
1073 .into_iter()
1074 .map(|session| {
1075 let mut value = live_descriptor_value(&session, &doors)?;
1076 let activity_key = (
1077 session.locator.harness.as_str().to_string(),
1078 session.locator.session_id.clone(),
1079 );
1080 if let Some(activity) = activities.get(&activity_key) {
1081 let mut activity = activity.clone();
1084 let waiting = doors.all().iter().any(|live| {
1085 live.address.harness == activity_key.0
1086 && live.address.session_id == activity_key.1
1087 && live.status == "waiting"
1088 });
1089 if waiting && activity.turn != crate::SessionTurnState::NeedsInput {
1090 activity.turn = crate::SessionTurnState::NeedsInput;
1091 activity.evidence.source = "pane_prompt".into();
1092 }
1093 let activity = &activity;
1094 value["activity"] = serde_json::to_value(activity)
1095 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1096 if let Some(status) = legacy_live_status(activity) {
1097 value["live_status"] = json!(status);
1098 }
1099 }
1100 Ok(value)
1101 })
1102 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1103 let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1104 if query.search_previews {
1107 result["receipt"] = serde_json::to_value(page.receipt)
1108 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1109 }
1110 Ok(result)
1111 }
1112 "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1113 "harness.v1.sessions.usage" => {
1114 let params = decode::<LocatorParams>(params)?;
1115 let session = load_session(¶ms.locator).map_err(operation)?;
1116 Ok(
1117 json!({"harness": params.locator.harness, "session_id": session.meta.session_id,
1118 "usage": recorded_session_usage(&session)}),
1119 )
1120 }
1121 "harness.v1.sessions.recorded_config" => {
1122 let params = decode::<LocatorParams>(params)?;
1123 recorded_config(¶ms.locator).map_err(operation)
1124 }
1125 "harness.v1.sessions.load" => {
1126 let params = decode::<LoadSessionParams>(params)?;
1127 load_session_door(params)
1128 }
1129 "harness.v1.sessions.follow" => {
1130 let params = decode::<LocatorParams>(params)?;
1131 let mut follower = self
1132 .catalog
1133 .follow_read_view(
1134 ¶ms.locator,
1135 params.read_fidelity(),
1136 params.include_subagents(),
1137 params.tail_messages(),
1138 params.max_message_chars(),
1139 params.display_history(),
1140 )
1141 .map_err(operation)?;
1142 let initial = follower
1143 .poll()
1144 .map_err(operation)?
1145 .map(|event| event.to_json());
1146 let subscription = format!("sub-{}", self.next_subscription);
1147 self.next_subscription += 1;
1148 self.followers.insert(subscription.clone(), follower);
1149 self.followed_sources.insert(
1150 subscription.clone(),
1151 FollowedSource {
1152 harness: params.locator.harness.as_str().to_string(),
1153 session_id: params.locator.session_id.clone(),
1154 reported: None,
1155 },
1156 );
1157 Ok(json!({"subscription": subscription, "initial": initial}))
1158 }
1159 "harness.v1.sessions.unfollow" => {
1160 let params = decode::<UnfollowParams>(params)?;
1161 self.followed_sources.remove(¶ms.subscription);
1162 Ok(json!({
1163 "removed": self.followers.remove(¶ms.subscription).is_some()
1164 }))
1165 }
1166 "harness.v1.sessions.activity.unsubscribe" => {
1167 let params = decode::<UnfollowParams>(params)?;
1168 Ok(json!({
1169 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1170 }))
1171 }
1172 "harness.v1.sessions.index.subscribe" => {
1173 let query = decode::<DiscoveryQuery>(params)?;
1174 crate::session_index::validate_query(&query)
1175 .map_err(ServiceError::InvalidParams)?;
1176 let homes = query.homes.clone();
1177 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1178 query,
1179 Arc::clone(&self.index_notifier),
1180 )
1181 .map_err(ServiceError::Operation)?;
1182 let doors = crate::mail_route::LiveSessions::read(&homes);
1183 let initial = initial
1184 .iter()
1185 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1186 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1187 let subscription = format!("index-sub-{}", self.next_subscription);
1188 self.next_subscription += 1;
1189 self.index_subscriptions.insert(subscription.clone(), index);
1190 Ok(json!({
1191 "subscription": subscription,
1192 "revision": 1,
1193 "initial": initial,
1194 }))
1195 }
1196 "harness.v1.sessions.index.resize" => {
1197 let params = decode::<IndexResizeParams>(params)?;
1198 crate::session_index::validate_limit(params.limit)
1199 .map_err(ServiceError::InvalidParams)?;
1200 let index = self
1201 .index_subscriptions
1202 .get_mut(¶ms.subscription)
1203 .ok_or_else(|| {
1204 ServiceError::InvalidParams("unknown session index subscription".into())
1205 })?;
1206 let prepared = index
1207 .prepare_resize(params.limit)
1208 .map_err(ServiceError::Operation)?;
1209 let doors = crate::mail_route::LiveSessions::read(index.homes());
1210 let initial = prepared
1211 .page
1212 .sessions
1213 .iter()
1214 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1215 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1216 let response = json!({
1217 "subscription": params.subscription,
1218 "revision": prepared.revision,
1219 "initial": initial,
1220 "receipt": prepared.page.receipt,
1221 });
1222 index.commit_resize(prepared);
1223 Ok(response)
1224 }
1225 "harness.v1.sessions.index.unsubscribe" => {
1226 let params = decode::<UnfollowParams>(params)?;
1227 Ok(json!({
1228 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1229 }))
1230 }
1231 "harness.v1.sessions.import" => {
1232 let params = decode::<ImportSessionParams>(params)?;
1233 let session = Session::load_str(¶ms.content, params.source_harness.into())
1234 .map_err(operation)?;
1235 Ok(json!({"session": normalized_session_json(&session)}))
1236 }
1237 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1238 let params = decode::<ExportSessionParams>(params)?;
1239 let session = load_session(¶ms.locator).map_err(operation)?;
1240 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1241 if method == "harness.v1.sessions.export"
1242 && params.target_harness == TransferFormat::Hermes
1243 {
1244 let imported = crate::hermes_import::import_into_hermes(&session, None)
1246 .map_err(operation)?;
1247 return Ok(json!({"artifact": artifact, "imported": imported}));
1248 }
1249 Ok(json!({"artifact": artifact}))
1250 }
1251 "harness.v1.sessions.reduce" => {
1252 let params = decode::<ReduceSessionParams>(params)?;
1253 self.reduce_session(params)
1254 }
1255 "harness.v1.sessions.branch" => {
1256 let params = decode::<BranchSessionParams>(params)?;
1257 let session = load_session(¶ms.locator).map_err(operation)?;
1258 let storage = params.locator.storage.path().display().to_string();
1259 let bootstrap_prompt = format!(
1260 "Continue as a new branch from {} session {}. The frozen parent transcript is at {}. Read or load that parent for context, summarize the relevant state, then continue independently without mutating the parent session.",
1261 params.locator.harness.as_str(), params.locator.session_id, storage
1262 );
1263 let artifact = params
1264 .target_harness
1265 .map(|target| session_artifact(¶ms.locator, &session, target))
1266 .transpose()?;
1267 Ok(json!({
1268 "parent": params.locator,
1269 "session": normalized_session_json(&session),
1270 "bootstrap_prompt": bootstrap_prompt,
1271 "artifact": artifact,
1272 }))
1273 }
1274 "harness.v1.sessions.handoff" => {
1275 let params = decode::<HandoffSessionParams>(params)?;
1276 let session = load_session(¶ms.locator).map_err(operation)?;
1277 let cwd = params
1278 .cwd
1279 .or_else(|| session.meta.cwd.clone())
1280 .unwrap_or_else(|| PathBuf::from("."));
1281 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1282 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1283 ServiceError::Operation(
1284 "handoff artifact omitted target session identity".into(),
1285 )
1286 })?;
1287 let instructions =
1288 handoff_instructions(params.target_harness, target_session_id, &cwd);
1289 Ok(json!({
1290 "artifact": artifact,
1291 "launch": instructions.launch,
1292 "materialize": instructions.materialize,
1293 "requires_materialization": instructions.requires_materialization,
1294 "note": instructions.note,
1295 }))
1296 }
1297 "harness.v1.sessions.materialize" => {
1298 let params = decode::<MaterializeSessionParams>(params)?;
1299 for file in ¶ms.artifact.files {
1303 if file.role == "source_recovery"
1304 && file.path == "recovery/source.supercode.jsonl"
1305 {
1306 if let Ok(source) = Session::from_native_str(&file.content) {
1307 crate::residue_store::store_segments(&source);
1308 }
1309 }
1310 }
1311 let locator = crate::native_materialize::materialize_native_artifact(
1312 params.artifact,
1313 ¶ms.cwd,
1314 ¶ms.homes,
1315 )
1316 .map_err(ServiceError::Operation)?;
1317 Ok(json!({"locator": locator}))
1318 }
1319 "harness.v1.jobs.list" => {
1323 let query = decode::<crate::jobs::JobsQuery>(params)?;
1324 if let Some(harness) = query.harness.as_deref() {
1325 refuse_harness_without_jobs(harness, "jobs.list")?;
1326 }
1327 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1328 serde_json::to_value(listing)
1329 .map_err(|error| ServiceError::Operation(error.to_string()))
1330 }
1331 "harness.v1.jobs.get" => {
1332 let params = decode::<JobsGetParams>(params)?;
1333 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1334 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1335 .map_err(operation)?
1336 {
1337 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1338 None => Err(ServiceError::Operation(format!(
1339 "`{}` has no scheduled job `{}`",
1340 params.harness, params.id
1341 ))),
1342 }
1343 }
1344 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1350 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1351 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1352 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1353 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1354 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1355 "harness.v1.jobs.notepad"
1356 | "harness.v1.jobs.notepad_set"
1357 | "harness.v1.jobs.notepad_delete" => {
1358 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1359 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1360 let answer = match method {
1361 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1362 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1363 _ => crate::jobs_notepad::read(&request),
1364 }
1365 .map_err(job_control_error)?;
1366 serde_json::to_value(answer)
1367 .map_err(|error| ServiceError::Operation(error.to_string()))
1368 }
1369 "harness.v1.runs.list" => {
1373 let query = decode::<crate::runs::RunsQuery>(params)?;
1374 if let Some(harness) = query.harness.as_deref() {
1375 refuse_harness_without_runs(harness, "runs.list")?;
1376 }
1377 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1378 serde_json::to_value(listing)
1379 .map_err(|error| ServiceError::Operation(error.to_string()))
1380 }
1381 "harness.v1.runs.get" => {
1382 let params = decode::<RunsGetParams>(params)?;
1383 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1384 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1385 .map_err(operation)?
1386 {
1387 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1388 None => Err(ServiceError::Operation(format!(
1389 "`{}` has no run `{}`",
1390 params.harness, params.id
1391 ))),
1392 }
1393 }
1394 "harness.v1.sessions.resume_instructions" => {
1395 let params = decode::<ResumeInstructionsParams>(params)?;
1396 let session = load_session(¶ms.locator).map_err(operation)?;
1397 let cwd = params
1398 .cwd
1399 .or(session.meta.cwd)
1400 .unwrap_or_else(|| PathBuf::from("."));
1401 let launch = resume_launch(
1402 params.locator.harness.as_str(),
1403 ¶ms.locator.session_id,
1404 &cwd,
1405 params.policy,
1406 )?;
1407 Ok(json!({"launch": launch}))
1408 }
1409 _ => Err(ServiceError::MethodNotFound),
1410 }
1411 }
1412
1413 fn reduce_session(
1414 &self,
1415 params: ReduceSessionParams,
1416 ) -> std::result::Result<Value, ServiceError> {
1417 let session = load_session(¶ms.locator).map_err(operation)?;
1418 if session.messages.is_empty() {
1419 return Err(ServiceError::InvalidParams(
1420 "cannot reduce an empty session".into(),
1421 ));
1422 }
1423 let keep_last = params.keep_last.clamp(1, 128);
1424 let policy = reduce::ReductionPolicy {
1425 clear_turns_older_than: Some(keep_last),
1426 ..Default::default()
1427 };
1428 let (view, log) =
1429 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1430 if log.reductions.is_empty() {
1431 return Err(ServiceError::UnsupportedAction(format!(
1432 "session `{}` is already too small for a meaningful reversible reduction",
1433 params.locator.session_id
1434 )));
1435 }
1436 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1437 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1438 if reduced_tokens >= source_tokens {
1439 return Err(ServiceError::UnsupportedAction(format!(
1440 "session `{}` has no token-reducing reversible projection",
1441 params.locator.session_id
1442 )));
1443 }
1444
1445 let store_root = self
1446 .reduction_store_root
1447 .clone()
1448 .unwrap_or_else(default_reduction_store_root);
1449 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1450 let rescue_id = format!("rescue-{}", generated_session_id());
1451 let imported = session
1452 .imported_message_count
1453 .unwrap_or(session.messages.len())
1454 .min(session.messages.len());
1455 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1456 let view_jsonl = messages_jsonl(&view)?;
1457 let title = format!(
1458 "Reduced {} continuation from {}",
1459 params.target_harness.id(),
1460 params.locator.session_id
1461 );
1462
1463 store
1468 .save_sidecar(&rescue_id, &sidecar_jsonl)
1469 .map_err(operation)?;
1470 store
1471 .save_reduction_log(&rescue_id, &log)
1472 .map_err(operation)?;
1473 store
1474 .save(&rescue_id, &title, &view_jsonl)
1475 .map_err(operation)?;
1476
1477 let source_bytes = serde_json::to_vec(&session.messages)
1478 .map_err(|error| ServiceError::Operation(error.to_string()))?
1479 .len() as u64;
1480 let reduced_bytes = serde_json::to_vec(&view)
1481 .map_err(|error| ServiceError::Operation(error.to_string()))?
1482 .len() as u64;
1483 store
1484 .set_reduction_stats(
1485 &rescue_id,
1486 &title,
1487 source_bytes,
1488 reduced_bytes,
1489 log.reductions.len() as u32,
1490 )
1491 .map_err(operation)?;
1492
1493 let reloaded_sidecar = store
1497 .load_sidecar(&rescue_id)
1498 .map_err(operation)?
1499 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1500 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1501 let reloaded_log = store
1502 .load_reduction_log(&rescue_id)
1503 .map_err(operation)?
1504 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1505 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1506 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1507 let (restamped_view, restamped_log) =
1514 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1515 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1516 return Err(ServiceError::Operation(
1517 "persisted reduction view does not match its durable log and sidecar".into(),
1518 ));
1519 }
1520 if restamped_log != reloaded_log {
1521 return Err(ServiceError::Operation(
1522 "reapplying the durable reduction log changed its identity".into(),
1523 ));
1524 }
1525 let inverted =
1526 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1527 if inverted != session.messages {
1528 return Err(ServiceError::Operation(
1529 "reduction inversion did not restore the source messages byte-exactly".into(),
1530 ));
1531 }
1532
1533 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1534 let sidecar_path = store.sidecar_path(&rescue_id);
1535 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1536 let bootstrap_prompt = reduced_bootstrap_prompt(
1537 ¶ms.locator,
1538 params.target_harness,
1539 &view_jsonl,
1540 &sidecar_path,
1541 &reduction_log_path,
1542 );
1543 let mut reduced_session = session.clone();
1544 reduced_session.meta.session_id = Some(rescue_id.clone());
1545 reduced_session.messages = view;
1546
1547 Ok(json!({
1548 "session": normalized_session_json(&reduced_session),
1549 "bootstrap_prompt": bootstrap_prompt,
1550 "receipt": {
1551 "id": rescue_id,
1552 "sidecar_id": rescue_id,
1553 "source_harness": params.locator.harness,
1554 "target_harness": params.target_harness.id(),
1555 "source_tokens": source_tokens,
1556 "reduced_tokens": reduced_tokens,
1557 "ratio": ratio,
1558 "source_bytes": source_bytes,
1559 "reduced_bytes": reduced_bytes,
1560 "reductions": reloaded_log.reductions.len(),
1561 "sidecar_path": sidecar_path,
1562 "reduction_log_path": reduction_log_path,
1563 "verified": true,
1564 "reversible": true,
1565 }
1566 }))
1567 }
1568
1569 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1587 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1588 return None;
1589 }
1590 let method = request.get("method").and_then(Value::as_str)?;
1591 if !RUNTIME_OPEN_METHODS.contains(&method) {
1592 return None;
1593 }
1594 Some(RuntimeOpen {
1595 id: request.get("id").cloned().unwrap_or(Value::Null),
1596 method: method.to_string(),
1597 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1598 })
1599 }
1600
1601 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1623 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1624 return None;
1625 }
1626 let method = request.get("method").and_then(Value::as_str)?;
1627 if !DETACHED_METHODS.contains(&method) {
1628 return None;
1629 }
1630 let id = request.get("id").cloned().unwrap_or(Value::Null);
1631 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1632 let work = match method {
1633 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1634 .inventory_work(method, params)
1635 .map(DetachedWork::Inventory),
1636 "harness.v1.sessions.message" => {
1637 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1638 }
1639 "harness.v1.sessions.load" => {
1640 decode::<LoadSessionParams>(params).map(DetachedWork::Load)
1641 }
1642 _ => {
1643 let verb = match method {
1644 "harness.v1.sessions.new" => crate::SessionVerb::New,
1645 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1646 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1647 _ => crate::SessionVerb::Delete,
1648 };
1649 match decode::<crate::SessionMutation>(params) {
1650 Ok(mutation) => {
1651 match crate::sessions_control::door(&mutation.harness, verb) {
1652 Ok(crate::SessionDoor::Live(_)) => return None,
1655 Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1656 Err(error) => Err(session_control_error(error)),
1657 }
1658 }
1659 Err(error) => Err(error),
1660 }
1661 }
1662 };
1663 Some(DetachedCall {
1664 id,
1665 method: method.to_string(),
1666 work: work.map(Work::Free),
1667 })
1668 }
1669
1670 pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1684 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1685 return None;
1686 }
1687 let method = request.get("method").and_then(Value::as_str)?;
1688 let id = request.get("id").cloned().unwrap_or(Value::Null);
1689 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1690 let work = match method {
1691 "harness.v1.runtimes.close" => decode::<RuntimeCloseParams>(params)
1692 .and_then(|params| self.surrender_runtime_of(¶ms))
1693 .map(|(runtime, process_group)| {
1694 Work::Runtime(RuntimeWork::Close {
1695 runtime,
1696 process_group,
1697 })
1698 }),
1699 "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1700 let verb = if method == "harness.v1.sessions.new" {
1701 crate::SessionVerb::New
1702 } else {
1703 crate::SessionVerb::Reset
1704 };
1705 let mutation = decode::<crate::SessionMutation>(params).ok()?;
1706 let Ok(crate::SessionDoor::Live(command)) =
1710 crate::sessions_control::door(&mutation.harness, verb)
1711 else {
1712 return None;
1713 };
1714 let connection = mutation
1715 .connection
1716 .clone()
1717 .filter(|value| !value.trim().is_empty())?;
1718 self.lend_runtime(&connection).map(|runtime| {
1719 let session = live_session_name(runtime.as_ref(), &mutation);
1720 Work::Runtime(RuntimeWork::LiveCommand {
1721 connection,
1722 runtime,
1723 verb,
1724 mutation,
1725 command,
1726 session,
1727 })
1728 })
1729 }
1730 _ => return None,
1731 };
1732 Some(DetachedCall {
1733 id,
1734 method: method.to_string(),
1735 work,
1736 })
1737 }
1738
1739 pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1744 let DetachedAnswer { response, returned } = answer;
1745 if let Some(ReturnedRuntime {
1746 connection,
1747 runtime,
1748 }) = returned
1749 {
1750 self.runtimes_in_flight.remove(&connection);
1751 self.runtimes.insert(connection, runtime);
1752 }
1753 response
1754 }
1755
1756 pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1760 let OpenedRuntime { id, outcome } = opened;
1761 let result = match outcome {
1762 Ok(open) => self.register_open_runtime(open).await,
1763 Err(error) => Err(error),
1764 };
1765 service_response(id, result)
1766 }
1767
1768 async fn register_open_runtime(
1770 &mut self,
1771 open: OpenRuntime,
1772 ) -> std::result::Result<Value, ServiceError> {
1773 match open {
1774 OpenRuntime::Hosted {
1775 runtime,
1776 capabilities,
1777 workspace,
1778 fresh,
1779 } => {
1780 self.insert_hosted_runtime(runtime, capabilities, workspace, fresh)
1781 .await
1782 }
1783 OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1784 }
1785 }
1786
1787 async fn runtime_call(
1788 &mut self,
1789 method: &str,
1790 params: Value,
1791 ) -> std::result::Result<Value, ServiceError> {
1792 match method {
1793 "harness.v1.runtimes.capabilities" => {
1794 let params = decode::<RuntimeBackendParams>(params)?;
1795 #[cfg(feature = "adapter-api")]
1799 if params
1800 .base_url
1801 .as_deref()
1802 .is_some_and(|value| LiveRuntimeEndpoint::parse(value).is_ok())
1803 {
1804 return Ok(json!({
1805 "harness": ¶ms.harness,
1806 "capabilities": SupercodeHttpRuntimeBackend::live_capabilities(),
1807 }));
1808 }
1809 let backend = runtime_backend(¶ms)?;
1810 Ok(json!({
1811 "harness": backend.harness(),
1812 "capabilities": backend.capabilities(),
1813 }))
1814 }
1815 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1816 self.register_open_runtime(open_runtime(method, params).await?)
1817 .await
1818 }
1819 "harness.v1.runtimes.send_input" => {
1820 let params = decode::<RuntimeInputParams>(params)?;
1821 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1822 let runtime = self.runtime_mut(¶ms.connection)?;
1823 let turn_id = within_control_deadline(
1824 method,
1825 runtime.send_input(RuntimeInput {
1826 text: params.text,
1827 image_urls,
1828 }),
1829 )
1830 .await?
1831 .map_err(operation)?;
1832 Ok(json!({"turn_id": turn_id}))
1833 }
1834 "harness.v1.runtimes.interrupt" => {
1835 let params = decode::<RuntimeConnectionParams>(params)?;
1836 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1837 .await?
1838 .map_err(operation)?;
1839 Ok(json!({}))
1840 }
1841 "harness.v1.runtimes.steer" => {
1842 let params = decode::<RuntimeInputParams>(params)?;
1843 if !params.image_urls.is_empty() {
1844 return Err(ServiceError::InvalidParams(
1845 "runtime steering accepts text only".into(),
1846 ));
1847 }
1848 let text = params.text.trim();
1849 if text.is_empty() || text.chars().count() > 50_000 {
1850 return Err(ServiceError::InvalidParams(
1851 "runtime steering requires 1 to 50,000 text characters".into(),
1852 ));
1853 }
1854 within_control_deadline(
1855 method,
1856 self.runtime_mut(¶ms.connection)?
1857 .steer(text.to_string()),
1858 )
1859 .await?
1860 .map_err(operation)?;
1861 Ok(json!({}))
1862 }
1863 "harness.v1.runtimes.respond" => {
1864 let params = decode::<RuntimeRespondParams>(params)?;
1865 let request_id = params.request_id.clone();
1866 within_control_deadline(
1867 method,
1868 self.runtime_mut(¶ms.connection)?
1869 .respond(params.request_id, params.response),
1870 )
1871 .await?
1872 .map_err(operation)?;
1873 self.approvals.answered(¶ms.connection, &request_id);
1875 Ok(json!({}))
1876 }
1877 "harness.v1.runtimes.acquire_control" => {
1878 let params = decode::<RuntimeConnectionParams>(params)?;
1879 let snapshot = within_control_deadline(
1880 method,
1881 self.runtime_mut(¶ms.connection)?.acquire_control(),
1882 )
1883 .await?
1884 .map_err(operation)?;
1885 serde_json::to_value(snapshot)
1886 .map_err(|error| ServiceError::Operation(error.to_string()))
1887 }
1888 "harness.v1.runtimes.heartbeat" => {
1889 let params = decode::<RuntimeConnectionParams>(params)?;
1890 let snapshot = within_control_deadline(
1891 method,
1892 self.runtime_mut(¶ms.connection)?.heartbeat(),
1893 )
1894 .await?
1895 .map_err(operation)?;
1896 serde_json::to_value(snapshot)
1897 .map_err(|error| ServiceError::Operation(error.to_string()))
1898 }
1899 "harness.v1.runtimes.detach" => {
1900 let params = decode::<RuntimeConnectionParams>(params)?;
1901 let snapshot =
1902 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1903 .await?
1904 .map_err(operation)?;
1905 serde_json::to_value(snapshot)
1906 .map_err(|error| ServiceError::Operation(error.to_string()))
1907 }
1908 "harness.v1.runtimes.terminal_instructions" => {
1909 let params = decode::<RuntimeConnectionParams>(params)?;
1910 let launch = self
1911 .terminal_launches
1912 .get(¶ms.connection)
1913 .ok_or_else(|| {
1914 ServiceError::Operation(
1915 "this runtime is not hosted for terminal attachment".into(),
1916 )
1917 })?;
1918 Ok(json!({"launch":launch}))
1919 }
1920 "harness.v1.runtimes.close" => {
1921 let params = decode::<RuntimeCloseParams>(params)?;
1922 let (runtime, process_group) = self.surrender_runtime_of(¶ms)?;
1923 close_runtime(runtime, process_group).await
1924 }
1925 _ => Err(ServiceError::MethodNotFound),
1926 }
1927 }
1928
1929 #[cfg(feature = "adapter-api")]
1931 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1932 let params = decode::<MessageSessionParams>(params)?;
1933 Ok(message_live_session(¶ms).await)
1934 }
1935
1936 #[cfg(feature = "adapter-api")]
1937 fn harness_settings_call(
1938 &self,
1939 method: &str,
1940 params: Value,
1941 ) -> std::result::Result<Value, ServiceError> {
1942 let homes = crate::HarnessHomes::default();
1943 match method {
1944 "harness.v1.harnesses.settings" => {
1945 let params = decode::<HarnessSettingsParams>(params)?;
1946 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1947 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1948 serde_json::to_value(report)
1949 .map_err(|error| ServiceError::Operation(error.to_string()))
1950 }
1951 "harness.v1.harnesses.configure" => {
1952 let params = decode::<ConfigureHarnessParams>(params)?;
1953 let report = crate::configure_harness_interop_settings(
1954 &homes,
1955 ¶ms.harness,
1956 ¶ms.changes,
1957 params.expected_revision.as_deref(),
1958 )
1959 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1960 serde_json::to_value(report)
1961 .map_err(|error| ServiceError::Operation(error.to_string()))
1962 }
1963 _ => Err(ServiceError::MethodNotFound),
1964 }
1965 }
1966
1967 fn insert_runtime(
1968 &mut self,
1969 runtime: Box<dyn RuntimeConnection>,
1970 ) -> std::result::Result<Value, ServiceError> {
1971 let connection = format!("runtime-{}", self.next_runtime);
1972 self.next_runtime += 1;
1973 let handle = runtime.handle().clone();
1974 self.runtime_sequences
1975 .entry(handle.runtime_id.clone())
1976 .or_insert(0);
1977 self.runtimes.insert(connection.clone(), runtime);
1978 Ok(json!({"connection": connection, "handle": handle}))
1979 }
1980
1981 #[cfg(feature = "adapter-api")]
1982 async fn insert_hosted_runtime(
1983 &mut self,
1984 runtime: Box<dyn RuntimeConnection>,
1985 capabilities: crate::RuntimeCapabilities,
1986 workspace: PathBuf,
1987 fresh: bool,
1988 ) -> std::result::Result<Value, ServiceError> {
1989 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
1990 let token: std::sync::Arc<str> = crate::server::generate_token().into();
1991 let server = crate::server::run_frontend_http(
1992 host.clone(),
1993 host.frontend_sender(),
1994 "127.0.0.1:0",
1995 token.clone(),
1996 connection.handle().runtime_id.clone(),
1997 )
1998 .await
1999 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2000 let source = LiveRuntimeSource {
2001 harness: connection.handle().harness.as_str().to_string(),
2002 session_id: connection.handle().runtime_id.clone(),
2003 workspace: workspace.clone(),
2004 };
2005 let registration = register_live_runtime_with_metadata(
2006 connection.handle().runtime_id.clone(),
2007 source.clone(),
2008 format!("http://{}", server.address()),
2009 token.to_string(),
2010 crate::LiveRuntimeMetadata {
2011 child_pid: match &connection.handle().endpoint {
2012 crate::RuntimeEndpoint::LocalProcess { pid, command, .. }
2013 if connection.handle().harness.as_str() != "claude-code"
2014 || command.windows(2).any(|args| {
2015 args[0] == "--name"
2016 && args[1].starts_with(crate::claude_relay::RELAY_NAME_PREFIX)
2017 }) =>
2018 {
2019 *pid
2020 }
2021 _ => None,
2022 },
2023 endpoint_capabilities: vec!["http".into(), "acp".into()],
2024 ..Default::default()
2025 },
2026 )
2027 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2028 let endpoint = registration.endpoint().to_string();
2029 let launch = StructuredLaunch {
2030 cwd: workspace,
2031 program: std::env::current_exe()
2035 .ok()
2036 .map(|path| path.to_string_lossy().into_owned())
2037 .unwrap_or_else(|| "supercode".into()),
2038 arguments: vec![
2039 "open".into(),
2040 endpoint,
2041 "--harness".into(),
2042 source.harness,
2043 "--session".into(),
2044 source.session_id,
2045 ],
2046 env: BTreeMap::new(),
2047 };
2048 host.retain_registration(registration);
2049 let lease = HostedRuntimeLease {
2050 connection,
2051 _host: host,
2052 _server: server,
2053 };
2054 let opened = self.insert_runtime(Box::new(lease))?;
2055 let connection_id = opened["connection"]
2056 .as_str()
2057 .expect("insert_runtime returns a connection id")
2058 .to_string();
2059 self.terminal_launches.insert(connection_id, launch);
2060 Ok(opened)
2061 }
2062
2063 #[cfg(not(feature = "adapter-api"))]
2064 async fn insert_hosted_runtime(
2065 &mut self,
2066 runtime: Box<dyn RuntimeConnection>,
2067 _capabilities: crate::RuntimeCapabilities,
2068 _workspace: PathBuf,
2069 _fresh: bool,
2070 ) -> std::result::Result<Value, ServiceError> {
2071 self.insert_runtime(runtime)
2072 }
2073
2074 fn runtime_mut(
2075 &mut self,
2076 connection: &str,
2077 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2078 if self.runtimes_in_flight.contains(connection) {
2079 return Err(self.lent_out(connection));
2080 }
2081 self.runtimes.get_mut(connection).ok_or_else(|| {
2082 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2083 })
2084 }
2085
2086 fn lent_out(&self, connection: &str) -> ServiceError {
2090 ServiceError::Operation(format!(
2091 "runtime connection `{connection}`: a harness turn is already in progress"
2092 ))
2093 }
2094
2095 fn lend_runtime(
2098 &mut self,
2099 connection: &str,
2100 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2101 if self.runtimes_in_flight.contains(connection) {
2102 return Err(self.lent_out(connection));
2103 }
2104 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2105 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2106 })?;
2107 self.runtimes_in_flight.insert(connection.to_string());
2108 Ok(runtime)
2109 }
2110
2111 fn surrender_runtime_of(
2125 &mut self,
2126 params: &RuntimeCloseParams,
2127 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2128 if params.connection.is_empty() {
2129 let expected = params.runtime_id.as_deref().ok_or_else(|| {
2130 ServiceError::InvalidParams("close requires connection or runtime_id".into())
2131 })?;
2132 let connection = self
2133 .runtimes
2134 .iter()
2135 .find(|(_, runtime)| runtime.handle().runtime_id == expected)
2136 .map(|(connection, _)| connection.clone())
2137 .ok_or_else(|| {
2138 ServiceError::InvalidParams(format!("unknown runtime `{expected}`"))
2139 })?;
2140 return self.surrender_runtime(&connection);
2141 }
2142 if let Some(expected) = params.runtime_id.as_deref() {
2143 let held = self
2144 .runtimes
2145 .get(¶ms.connection)
2146 .map(|runtime| runtime.handle().runtime_id.clone());
2147 if held.as_deref() != Some(expected) {
2148 return Err(ServiceError::InvalidParams(format!(
2149 "runtime connection `{}` does not hold runtime `{expected}`",
2150 params.connection
2151 )));
2152 }
2153 }
2154 self.surrender_runtime(¶ms.connection)
2155 }
2156
2157 fn surrender_runtime(
2158 &mut self,
2159 connection: &str,
2160 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2161 if self.runtimes_in_flight.contains(connection) {
2162 return Err(self.lent_out(connection));
2163 }
2164 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2165 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2166 })?;
2167 let process_group = runtime_process_group(runtime.handle());
2168 let runtime_id = runtime.handle().runtime_id.clone();
2169 self.terminal_launches.remove(connection);
2170 self.runtime_sequences.remove(&runtime_id);
2171 self.approvals.forget(connection);
2172 Ok((runtime, process_group))
2173 }
2174
2175 pub fn kill_all_runtime_groups(&self) -> usize {
2186 self.runtimes
2187 .values()
2188 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2189 .count()
2190 }
2191
2192 async fn mutate_session(
2203 &mut self,
2204 verb: crate::SessionVerb,
2205 params: Value,
2206 ) -> std::result::Result<Value, ServiceError> {
2207 let mutation = decode::<crate::SessionMutation>(params)?;
2208 let door = crate::sessions_control::door(&mutation.harness, verb)
2209 .map_err(session_control_error)?;
2210 let outcome = match door {
2211 #[cfg(not(feature = "adapter-api"))]
2215 crate::SessionDoor::Live(command) => {
2216 return Err(ServiceError::Operation(format!(
2217 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2218 session, which needs this build's `adapter-api` feature",
2219 mutation.harness,
2220 verb.as_str()
2221 )));
2222 }
2223 #[cfg(feature = "adapter-api")]
2224 crate::SessionDoor::Live(command) => {
2225 let connection = mutation
2226 .connection
2227 .clone()
2228 .filter(|value| !value.trim().is_empty())
2229 .ok_or_else(|| {
2230 ServiceError::InvalidParams(format!(
2231 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2232 driven session: pass the `connection` of an open runtime \
2233 (`harness.v1.runtimes.start`)",
2234 mutation.harness,
2235 verb.as_str()
2236 ))
2237 })?;
2238 let runtime = self.runtime_mut(&connection)?;
2239 let session = live_session_name(runtime.as_ref(), &mutation);
2240 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2246 .await;
2247 }
2248 _ => run_session_mutation(verb, &mutation).await?,
2249 };
2250 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2251 }
2252
2253 async fn inventory_call(
2257 &self,
2258 method: &str,
2259 params: Value,
2260 ) -> std::result::Result<Value, ServiceError> {
2261 run_inventory(self.inventory_work(method, params)?).await
2262 }
2263
2264 fn inventory_work(
2270 &self,
2271 method: &str,
2272 params: Value,
2273 ) -> std::result::Result<InventoryWork, ServiceError> {
2274 let mut params = decode::<HarnessInventoryParams>(params)?;
2275 if method == "harness.v1.harnesses.probe" {
2276 let harness = params.harness.take().ok_or_else(|| {
2277 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2278 })?;
2279 params.harnesses = vec![harness];
2280 }
2281 let selected = params
2282 .harnesses
2283 .iter()
2284 .map(HarnessId::as_str)
2285 .collect::<std::collections::BTreeSet<_>>();
2286 let supported = harness_support_registry()
2287 .harnesses
2288 .into_iter()
2289 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2290 .collect::<Vec<_>>();
2291 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2292 let known = supported
2293 .iter()
2294 .map(|harness| harness.id.as_str())
2295 .collect::<std::collections::BTreeSet<_>>();
2296 let missing = params
2297 .harnesses
2298 .iter()
2299 .filter(|id| !known.contains(id.as_str()))
2300 .map(HarnessId::as_str)
2301 .collect::<Vec<_>>();
2302 return Err(ServiceError::InvalidParams(format!(
2303 "unknown harness(es): {}",
2304 missing.join(", ")
2305 )));
2306 }
2307 let global_counts = params
2308 .include_sessions
2309 .then(|| self.session_counts(None, ¶ms.harnesses));
2310 let workspace_counts = params
2311 .include_sessions
2312 .then(|| {
2313 params
2314 .workspace
2315 .as_deref()
2316 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2317 })
2318 .flatten();
2319 Ok(InventoryWork {
2320 params,
2321 supported,
2322 global_counts,
2323 workspace_counts,
2324 })
2325 }
2326
2327 #[cfg(feature = "adapter-api")]
2328 async fn harness_authentication_call(
2329 &self,
2330 method: &str,
2331 params: Value,
2332 ) -> std::result::Result<Value, ServiceError> {
2333 match method {
2334 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2335 let params = decode::<HarnessAuthenticationParams>(params)?;
2336 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2337 .map_err(|error| ServiceError::Operation(error.to_string()))
2338 }
2339 "harness.v1.harnesses.auth.begin" => {
2340 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2341 let cwd = params
2342 .cwd
2343 .or_else(|| std::env::current_dir().ok())
2344 .unwrap_or_else(|| PathBuf::from("."));
2345 let plan = crate::harness_authentication_plan(
2346 ¶ms.harness,
2347 params.environment,
2348 params.method,
2349 &cwd,
2350 )
2351 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2352 serde_json::to_value(plan)
2353 .map_err(|error| ServiceError::Operation(error.to_string()))
2354 }
2355 _ => Err(ServiceError::MethodNotFound),
2356 }
2357 }
2358
2359 fn session_counts(
2360 &self,
2361 workspace: Option<&Path>,
2362 harnesses: &[HarnessId],
2363 ) -> BTreeMap<String, usize> {
2364 let mut counts = BTreeMap::new();
2365 for session in self
2366 .catalog
2367 .discover(&DiscoveryQuery {
2368 workspace: workspace.map(Path::to_path_buf),
2369 harnesses: harnesses.to_vec(),
2370 ..DiscoveryQuery::default()
2371 })
2372 .unwrap_or_default()
2373 {
2374 *counts
2375 .entry(session.locator.harness.as_str().to_string())
2376 .or_insert(0) += 1;
2377 }
2378 counts
2379 }
2380}
2381
2382#[async_trait::async_trait]
2383impl SdkService for HarnessSessionService {
2384 fn capabilities(&self) -> SdkCapabilities {
2385 SdkCapabilities::default()
2386 }
2387
2388 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2389 if request.operation == SdkOperation::Events {
2390 let events = self
2391 .poll_sdk_events()
2392 .await
2393 .into_iter()
2394 .map(|(_, event)| event)
2395 .collect::<Vec<_>>();
2396 return serde_json::to_value(events).map_err(|error| {
2397 SdkError::new(
2398 SdkErrorCode::Execution,
2399 request.operation,
2400 error.to_string(),
2401 )
2402 });
2403 }
2404 if self.runtimes.is_empty()
2405 && matches!(
2406 request.operation,
2407 SdkOperation::Input
2408 | SdkOperation::Interrupt
2409 | SdkOperation::Steer
2410 | SdkOperation::Respond
2411 | SdkOperation::Close
2412 )
2413 {
2414 return Err(SdkError::unsupported(request.operation));
2415 }
2416 let method = request
2417 .operation
2418 .method()
2419 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2420 let result = match request.operation {
2421 SdkOperation::Discover
2422 | SdkOperation::Load
2423 | SdkOperation::Export
2424 | SdkOperation::ProfilesList
2425 | SdkOperation::ProfilesGet
2426 | SdkOperation::ProfilesCreate
2427 | SdkOperation::ProfilesDelete
2428 | SdkOperation::SkillsList
2429 | SdkOperation::SkillsInstall
2430 | SdkOperation::SkillsRemove
2431 | SdkOperation::ChannelsList
2432 | SdkOperation::RoutesList
2433 | SdkOperation::TriggersList
2434 | SdkOperation::ChannelsStatus
2435 | SdkOperation::MemoryShow
2436 | SdkOperation::MemorySearch
2437 | SdkOperation::JobsList
2438 | SdkOperation::JobsGet
2439 | SdkOperation::JobsCreate
2440 | SdkOperation::JobsUpdate
2441 | SdkOperation::JobsPause
2442 | SdkOperation::JobsResume
2443 | SdkOperation::JobsRun
2444 | SdkOperation::JobsDelete
2445 | SdkOperation::JobsNotepad
2446 | SdkOperation::JobsNotepadSet
2447 | SdkOperation::JobsNotepadDelete
2448 | SdkOperation::RunsList
2449 | SdkOperation::RunsGet
2450 | SdkOperation::ApprovalsList
2451 | SdkOperation::OrchestrationLoad
2452 | SdkOperation::OrchestrationSave
2453 | SdkOperation::OrchestrationCompile
2454 | SdkOperation::OrchestrationDecompile
2455 | SdkOperation::OrchestrationImport
2456 | SdkOperation::OrchestrationExport
2457 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2458 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2461 SdkOperation::Start
2462 | SdkOperation::Resume
2463 | SdkOperation::Input
2464 | SdkOperation::Interrupt
2465 | SdkOperation::Steer
2466 | SdkOperation::Respond
2467 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2468 SdkOperation::SessionsNew => {
2473 self.mutate_session(crate::SessionVerb::New, request.params)
2474 .await
2475 }
2476 SdkOperation::SessionsReset => {
2477 self.mutate_session(crate::SessionVerb::Reset, request.params)
2478 .await
2479 }
2480 SdkOperation::SessionsArchive => {
2481 self.mutate_session(crate::SessionVerb::Archive, request.params)
2482 .await
2483 }
2484 SdkOperation::SessionsDelete => {
2485 self.mutate_session(crate::SessionVerb::Delete, request.params)
2486 .await
2487 }
2488 SdkOperation::Events => unreachable!("handled before method dispatch"),
2489 };
2490 result.map_err(|error| sdk_error(request.operation, error))
2491 }
2492
2493 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2494 Ok(self
2495 .poll_sdk_events()
2496 .await
2497 .into_iter()
2498 .map(|(_, event)| event)
2499 .collect())
2500 }
2501}
2502
2503#[cfg(feature = "adapter-api")]
2504struct HostedRuntimeLease {
2505 connection: HostedHarnessConnection,
2506 _host: std::sync::Arc<HostedHarnessRuntime>,
2507 _server: crate::server::FrontendHttpServer,
2508}
2509
2510#[async_trait::async_trait]
2511#[cfg(feature = "adapter-api")]
2512impl RuntimeConnection for HostedRuntimeLease {
2513 fn handle(&self) -> &crate::RuntimeHandle {
2514 self.connection.handle()
2515 }
2516
2517 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2518 self.connection.send_input(input).await
2519 }
2520
2521 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2522 self.connection.next_event().await
2523 }
2524
2525 async fn interrupt(&mut self) -> crate::Result<()> {
2526 self.connection.interrupt().await
2527 }
2528
2529 async fn steer(&mut self, text: String) -> crate::Result<()> {
2532 self.connection.steer(text).await
2533 }
2534
2535 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2536 self.connection.respond(request_id, response).await
2537 }
2538
2539 async fn close(&mut self) -> crate::Result<()> {
2540 self.connection.close().await
2541 }
2542}
2543
2544struct InventoryWork {
2547 params: HarnessInventoryParams,
2548 supported: Vec<crate::HarnessSupportDescriptor>,
2549 global_counts: Option<BTreeMap<String, usize>>,
2550 workspace_counts: Option<BTreeMap<String, usize>>,
2551}
2552
2553async fn run_session_mutation(
2560 verb: crate::SessionVerb,
2561 mutation: &crate::SessionMutation,
2562) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2563 let door =
2570 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2571 if let crate::SessionDoor::Http = door {
2572 return crate::sessions_control::mutate(verb, mutation)
2573 .await
2574 .map_err(session_control_error);
2575 }
2576 let mutation = mutation.clone();
2577 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2578 .await
2579 .map_err(|error| {
2580 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2581 })?
2582 .map_err(session_control_error)
2583}
2584
2585async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2588 let InventoryWork {
2589 params,
2590 supported,
2591 global_counts,
2592 workspace_counts,
2593 } = work;
2594 let probes = supported.into_iter().map(|descriptor| {
2595 let global = global_counts
2596 .as_ref()
2597 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2598 let workspace = workspace_counts
2599 .as_ref()
2600 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2601 probe_harness(descriptor, ¶ms, global, workspace)
2602 });
2603 let harnesses = futures::future::join_all(probes).await;
2604 serde_json::to_value(HarnessInventoryReport {
2605 probe: params.probe,
2606 workspace: params.workspace,
2607 harnesses,
2608 })
2609 .map_err(|error| ServiceError::Operation(error.to_string()))
2610}
2611
2612async fn probe_harness(
2613 descriptor: crate::HarnessSupportDescriptor,
2614 params: &HarnessInventoryParams,
2615 global: Option<usize>,
2616 workspace: Option<usize>,
2617) -> LocalHarness {
2618 let launch = descriptor.runtime.default_launch.as_ref();
2619 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2624 .then(crate::orchestrator::daemon_entry)
2625 .and_then(Result::ok);
2626 let executable = match &orchestrator_entry {
2627 Some(entry) => Some(entry.clone()),
2628 None => launch.and_then(|launch| find_executable(&launch.program)),
2629 };
2630 let installed = executable.is_some();
2631 let version = if params.skip_versions || orchestrator_entry.is_some() {
2632 None
2635 } else {
2636 match executable.as_deref() {
2637 Some(path) => executable_version(path).await,
2638 None => None,
2639 }
2640 };
2641 let configured = auth_evidence(descriptor.id.as_str());
2642 let mut auth = if configured {
2643 HarnessAuthState::Configured
2644 } else if matches!(
2645 descriptor.id.as_str(),
2646 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2647 ) {
2648 HarnessAuthState::Required
2653 } else {
2654 HarnessAuthState::Unknown
2655 };
2656 let mut runtime = if installed {
2657 HarnessRuntimeState::Degraded
2658 } else {
2659 HarnessRuntimeState::Unavailable
2660 };
2661 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2662 let mut reason = (!installed).then(|| {
2663 if is_orchestrator {
2664 format!(
2665 "{} is supported but its daemon entry `{}` was not found",
2666 descriptor.display_name,
2667 crate::orchestrator::DAEMON_ENTRY
2668 )
2669 } else {
2670 format!(
2671 "{} is supported but `{}` was not found on PATH",
2672 descriptor.display_name,
2673 launch
2674 .map(|launch| launch.program.as_str())
2675 .unwrap_or("executable")
2676 )
2677 }
2678 });
2679 let mut repair = (!installed).then(|| {
2680 if is_orchestrator {
2681 format!(
2682 "Install the `supercode-orchestrator` package so `{}` resolves.",
2683 crate::orchestrator::DAEMON_ENTRY
2684 )
2685 } else {
2686 format!(
2687 "Install {} and ensure `{}` is on PATH.",
2688 descriptor.display_name,
2689 launch
2690 .map(|launch| launch.program.as_str())
2691 .unwrap_or("its executable")
2692 )
2693 }
2694 });
2695
2696 if installed && params.probe == HarnessProbeLevel::Handshake {
2697 let backend_params = RuntimeBackendParams {
2698 harness: descriptor.id.clone(),
2699 protocol: None,
2700 launch: None,
2701 base_url: None,
2702 policy: RuntimePolicy::Default,
2703 };
2704 match runtime_backend(&backend_params) {
2705 Ok(backend) => {
2706 let cwd = params
2707 .workspace
2708 .clone()
2709 .or_else(|| std::env::current_dir().ok())
2710 .unwrap_or_else(|| PathBuf::from("."));
2711 let isolated = descriptor
2712 .runtime
2713 .default_launch
2714 .clone()
2715 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2716 let Some(isolated) = isolated else {
2717 reason = Some(
2718 "No-prompt runtime handshake could not create its isolated harness home."
2719 .into(),
2720 );
2721 repair = Some(
2722 "Check temporary-directory permissions, then run the handshake probe again."
2723 .into(),
2724 );
2725 let running = probe_running_instance(descriptor.id.as_str());
2726 return LocalHarness {
2727 gateway: gateway_health(
2728 descriptor.id.as_str(),
2729 installed,
2730 running.as_ref(),
2731 version.as_deref(),
2732 ),
2733 id: descriptor.id,
2734 display_name: descriptor.display_name,
2735 supported: true,
2736 installed,
2737 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2738 version,
2739 auth,
2740 runtime,
2741 protocol: descriptor.runtime.protocol,
2742 capabilities: descriptor.runtime.capabilities.clone(),
2743 effective_capabilities: descriptor.runtime.capabilities,
2744 sessions: HarnessSessionCounts { global, workspace },
2745 running,
2746 reason,
2747 repair,
2748 };
2749 };
2750 match tokio::time::timeout(
2751 Duration::from_secs(30),
2752 backend.start(RuntimeStartRequest {
2753 cwd,
2754 launch: Some(isolated.launch.clone()),
2755 mcp_servers: Vec::new(),
2756 approval_policy: None,
2757 }),
2758 )
2759 .await
2760 {
2761 Ok(Ok(mut connection)) => {
2762 match stabilize_handshake(connection.as_mut()).await {
2763 Ok(()) => {
2764 auth = HarnessAuthState::Ready;
2765 runtime = HarnessRuntimeState::Ready;
2766 reason = Some(
2767 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2768 .into(),
2769 );
2770 repair = None;
2771 }
2772 Err(message) => {
2773 auth = if looks_like_auth_error(&message) {
2774 HarnessAuthState::Required
2775 } else if configured {
2776 HarnessAuthState::Configured
2777 } else {
2778 HarnessAuthState::Unknown
2779 };
2780 reason = Some(format!(
2781 "No-prompt runtime handshake became unhealthy during startup: {message}"
2782 ));
2783 repair = Some(if auth == HarnessAuthState::Required {
2784 format!(
2785 "Run `{}` interactively once and complete sign-in, then probe again.",
2786 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2787 )
2788 } else {
2789 "Run the harness directly to inspect its startup failure, then probe again."
2790 .into()
2791 });
2792 }
2793 }
2794 let _ =
2795 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2796 }
2797 Ok(Err(error)) => {
2798 let message = truncate_text(&error.to_string(), 500);
2799 auth = if looks_like_auth_error(&message) {
2800 HarnessAuthState::Required
2801 } else if configured {
2802 HarnessAuthState::Configured
2803 } else {
2804 HarnessAuthState::Unknown
2805 };
2806 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2807 repair = Some(if auth == HarnessAuthState::Required {
2808 format!(
2809 "Run `{}` interactively once and complete sign-in, then probe again.",
2810 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2811 )
2812 } else {
2813 "Check the harness installation and run the handshake probe again."
2814 .into()
2815 });
2816 }
2817 Err(_) => {
2818 reason =
2819 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2820 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2821 }
2822 }
2823 let _ = isolated.cleanup();
2832 tokio::time::sleep(Duration::from_millis(250)).await;
2833 if let Err(error) = isolated.cleanup() {
2834 auth = if configured {
2835 HarnessAuthState::Configured
2836 } else {
2837 HarnessAuthState::Unknown
2838 };
2839 runtime = HarnessRuntimeState::Degraded;
2840 reason = Some(format!(
2841 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2842 ));
2843 repair = Some(
2844 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2845 .into(),
2846 );
2847 }
2848 }
2849 Err(error) => {
2850 reason = Some(error_message(error));
2851 }
2852 }
2853 } else if installed && configured {
2854 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2855 } else if installed && auth == HarnessAuthState::Required {
2856 reason = Some("Executable found, but no native authentication evidence is present.".into());
2857 repair = Some(format!(
2858 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2859 descriptor.id.as_str()
2860 ));
2861 } else if installed {
2862 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2863 repair = Some(format!(
2864 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2865 launch
2866 .map(|launch| launch.program.as_str())
2867 .unwrap_or("the harness")
2868 ));
2869 }
2870
2871 let effective_capabilities = if installed {
2872 descriptor.runtime.capabilities.clone()
2873 } else {
2874 unavailable_capabilities()
2875 };
2876 let running = probe_running_instance(descriptor.id.as_str());
2877 LocalHarness {
2878 gateway: gateway_health(
2879 descriptor.id.as_str(),
2880 installed,
2881 running.as_ref(),
2882 version.as_deref(),
2883 ),
2884 id: descriptor.id,
2885 display_name: descriptor.display_name,
2886 supported: true,
2887 installed,
2888 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2889 version,
2890 auth,
2891 runtime,
2892 protocol: descriptor.runtime.protocol,
2893 capabilities: descriptor.runtime.capabilities,
2894 effective_capabilities,
2895 sessions: HarnessSessionCounts { global, workspace },
2896 running,
2897 reason,
2898 repair,
2899 }
2900}
2901
2902async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2903 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2904 loop {
2905 let now = tokio::time::Instant::now();
2906 if now >= deadline {
2907 return Ok(());
2908 }
2909 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2910 Err(_) => return Ok(()),
2911 Ok(Ok(Some(event))) => {
2912 if let Some(message) = handshake_event_failure(&event) {
2913 return Err(truncate_text(&message, 500));
2914 }
2915 }
2916 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2917 Ok(Err(error)) => return Err(error.to_string()),
2918 }
2919 }
2920}
2921
2922fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2923 let detail = event
2924 .payload
2925 .get("message")
2926 .or_else(|| event.payload.get("line"))
2927 .and_then(Value::as_str)
2928 .unwrap_or(event.kind.as_str());
2929 match event.kind.as_str() {
2930 "transport_closed" => Some("runtime transport closed during startup".into()),
2931 "transport_error" => Some(format!("runtime transport error: {detail}")),
2932 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2933 _ => None,
2938 }
2939}
2940
2941fn indexed_claude_window(
2942 locator: &SessionLocator,
2943 options: &SessionLoadOptions,
2944) -> std::result::Result<Option<Value>, ServiceError> {
2945 use supercode_interchange::session::ClaudeReadIndex;
2946 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2949 || options.include_subagents != Some(false)
2950 {
2951 return Ok(None);
2952 }
2953 let crate::StorageLocator::File { path } = &locator.storage else {
2954 return Ok(None);
2955 };
2956 if !ClaudeReadIndex::supports(path)
2957 .map_err(|error| ServiceError::Operation(error.to_string()))?
2958 {
2959 return Ok(None);
2960 }
2961 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2962 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2963 let total = index.len();
2964 let (offset, end) = projected_message_window(total, options);
2965 let session = index
2966 .read_messages(offset..end)
2967 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2968 let summary = index
2969 .read_summary()
2970 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2971 let selected_options = SessionLoadOptions {
2972 message_offset: None,
2973 message_limit: None,
2974 message_tail: None,
2975 ..options.clone()
2976 };
2977 let mut selected = projected_session_json(&session, &selected_options);
2978 selected["raw_record_count"] = json!(index.raw_record_count());
2979 Ok(Some(json!({
2980 "session": selected,
2981 "summary": projected_session_summary(&summary, options),
2982 "window": {
2983 "has_more": offset > 0 || end < total, "has_newer": end < total,
2984 "has_older": offset > 0, "newer_items": index.item_count(end..total),
2985 "offset": offset, "older_items": index.item_count(0..offset),
2986 "returned": end - offset, "total_messages": total,
2987 }
2988 })))
2989}
2990
2991fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
2992 let total_messages = session.messages.len();
2993 let (offset, end) = projected_message_window(total_messages, options);
2994 json!({
2995 "session": projected_session_json(session, options),
2996 "summary": projected_session_summary(session, options),
2997 "window": {
2998 "has_more": offset > 0 || end < total_messages,
2999 "has_newer": end < total_messages,
3000 "has_older": offset > 0,
3001 "newer_items": normalized_item_count(&session.messages[end..]),
3002 "offset": offset,
3003 "older_items": normalized_item_count(&session.messages[..offset]),
3004 "returned": end.saturating_sub(offset),
3005 "total_messages": total_messages,
3006 }
3007 })
3008}
3009
3010fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
3011 messages
3012 .iter()
3013 .map(|message| {
3014 let conversation = usize::from(
3015 matches!(message.role, Role::Assistant | Role::User)
3016 && message_has_content(message),
3017 );
3018 let tool_result =
3019 usize::from(message.role == Role::Tool && message_has_content(message));
3020 conversation + tool_result + message.tool_calls().len()
3021 })
3022 .sum()
3023}
3024
3025fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
3026 let mut conversational = session.messages.iter().filter(|message| {
3027 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
3028 });
3029 let first_message = conversational.clone().next();
3030 let last_message = conversational.next_back();
3031 let mut assistant = session
3032 .messages
3033 .iter()
3034 .filter(|message| message.role == Role::Assistant && message_has_content(message));
3035 let first_assistant_message = assistant.clone().next();
3036 let last_assistant_message = assistant.next_back();
3037 let end_of_turn = session
3038 .messages
3039 .iter()
3040 .rev()
3041 .find(|message| message.role != Role::System)
3042 .is_some_and(|message| {
3043 message.role == Role::Assistant
3044 && message_has_content(message)
3045 && message.tool_calls().is_empty()
3046 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
3048 });
3049 let project = |message: Option<&crate::ChatMessage>| {
3050 message.map(|message| project_inline_media(message_json(message), options))
3051 };
3052 json!({
3053 "end_of_turn": end_of_turn,
3054 "first_assistant_message": project(first_assistant_message),
3055 "first_message": project(first_message),
3056 "last_assistant_message": project(last_assistant_message),
3057 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3058 "last_message": project(last_message),
3059 })
3060}
3061
3062fn message_has_content(message: &crate::ChatMessage) -> bool {
3063 message
3064 .content
3065 .as_deref()
3066 .is_some_and(|content| !content.trim().is_empty())
3067 || message
3068 .content_parts
3069 .as_ref()
3070 .is_some_and(|parts| !parts.is_empty())
3071}
3072
3073fn message_text(message: &crate::ChatMessage) -> String {
3074 if let Some(content) = &message.content {
3075 return content.clone();
3076 }
3077 message
3078 .content_parts
3079 .as_ref()
3080 .into_iter()
3081 .flatten()
3082 .filter_map(|part| part.get("text").and_then(Value::as_str))
3083 .collect::<Vec<_>>()
3084 .join("\n")
3085}
3086
3087fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3088 let (offset, end) = projected_message_window(session.messages.len(), options);
3089 let messages = session.messages[offset..end]
3090 .iter()
3091 .map(|message| project_inline_media(message_json(message), options))
3092 .collect::<Vec<_>>();
3093 let subagents = if options.include_subagents.unwrap_or(true) {
3094 let subagent_options = SessionLoadOptions {
3099 message_limit: None,
3100 message_offset: None,
3101 message_tail: None,
3102 ..options.clone()
3103 };
3104 session
3105 .subagents
3106 .iter()
3107 .map(|subagent| projected_session_json(subagent, &subagent_options))
3108 .collect::<Vec<_>>()
3109 } else {
3110 Vec::new()
3111 };
3112 json!({
3113 "source": match session.meta.source {
3114 SessionSource::ClaudeCode => "claude_code",
3115 SessionSource::Codex => "codex",
3116 SessionSource::Gemini => "gemini",
3117 SessionSource::Goose => "goose",
3118 SessionSource::Grok => "grok",
3119 SessionSource::Native => "native",
3120 SessionSource::OpenClaw => "openclaw",
3121 SessionSource::Hermes => "hermes",
3122 SessionSource::OpenCode => "opencode",
3123 SessionSource::Pi => "pi",
3124 },
3125 "session_id": session.meta.session_id,
3126 "ended_at": session.meta.ended_at,
3127 "end_reason": session.meta.end_reason,
3128 "model": session.meta.model,
3129 "resolved_profile": recorded_session_profile(session),
3130 "resolved_profile_history": recorded_session_profile_history(session),
3131 "cwd": session.meta.cwd,
3132 "system_prompt": session.meta.system_prompt,
3133 "agent_id": session.meta.agent_id,
3134 "parent_tool_use_id": session.meta.parent_tool_use_id,
3135 "lineage": session.meta.lineage,
3136 "messages": messages,
3137 "subagents": subagents,
3138 "raw_record_count": session.raw.len(),
3139 "parse_error_lines": session.parse_error_lines,
3140 })
3141}
3142
3143pub fn recorded_config(locator: &SessionLocator) -> crate::Result<Value> {
3151 let session = load_session(locator)?;
3152 let mut config = json!({"harness": locator.harness, "session_id": session.meta.session_id,
3153 "cwd": session.meta.cwd, "model": session.meta.model,
3154 "approval_policy": null, "sandbox_policy": null, "permission_mode": null});
3155 for header in &session.meta.codex_headers {
3156 let payload = header.get("payload").unwrap_or(header);
3157 for key in ["cwd", "model", "approval_policy", "sandbox_policy"] {
3158 if let Some(value) = payload.get(key) {
3159 config[key] = value.clone();
3160 }
3161 }
3162 }
3163 if let Ok(manifest) = crate::claude_runtime_state::ClaudeRuntimeManifest::from_session(&session)
3165 {
3166 config["permission_mode"] = json!(manifest.posture.permission_mode);
3167 }
3168 Ok(config)
3169}
3170
3171fn recorded_session_usage(session: &Session) -> Vec<Value> {
3172 let mut days: std::collections::BTreeMap<(String, String), [u64; 4]> =
3173 std::collections::BTreeMap::new();
3174 let mut add =
3175 |at: Option<&str>, model: &str, input: u64, output: u64, read: u64, write: u64| {
3176 let day = at
3177 .and_then(|t| t.get(..10))
3178 .unwrap_or("unknown")
3179 .to_string();
3180 let entry = days.entry((day, model.to_string())).or_insert([0; 4]);
3181 entry[0] += input;
3182 entry[1] += output;
3183 entry[2] += read;
3184 entry[3] += write;
3185 };
3186 let records = session
3187 .raw
3188 .iter()
3189 .filter_map(|line| serde_json::from_str::<Value>(line).ok());
3190 let n = |v: Option<&Value>| v.and_then(Value::as_u64).unwrap_or(0);
3191 match session.meta.source {
3192 SessionSource::ClaudeCode => {
3193 let mut responses: std::collections::BTreeMap<String, Value> =
3194 std::collections::BTreeMap::new();
3195 for record in records {
3196 let Some(message) = record.get("message") else {
3197 continue;
3198 };
3199 let (Some(id), Some(_)) = (
3200 message.get("id").and_then(Value::as_str),
3201 message.get("usage"),
3202 ) else {
3203 continue;
3204 };
3205 if message.get("model").and_then(Value::as_str) == Some("<synthetic>") {
3206 continue;
3207 }
3208 responses.insert(id.to_string(), record.clone());
3209 }
3210 for record in responses.values() {
3211 let message = &record["message"];
3212 let usage = &message["usage"];
3213 add(
3214 record.get("timestamp").and_then(Value::as_str),
3215 message
3216 .get("model")
3217 .and_then(Value::as_str)
3218 .unwrap_or("unknown"),
3219 n(usage.get("input_tokens")),
3220 n(usage.get("output_tokens")),
3221 n(usage.get("cache_read_input_tokens")),
3222 n(usage.get("cache_creation_input_tokens")),
3223 );
3224 }
3225 }
3226 SessionSource::Codex => {
3227 let mut model = session
3228 .meta
3229 .model
3230 .clone()
3231 .unwrap_or_else(|| "unknown".into());
3232 let mut previous = [0u64; 3];
3233 for record in records {
3234 let payload = record.get("payload").unwrap_or(&Value::Null);
3235 if record.get("type").and_then(Value::as_str) == Some("turn_context") {
3236 if let Some(m) = payload.get("model").and_then(Value::as_str) {
3237 model = m.to_string();
3238 }
3239 continue;
3240 }
3241 if payload.get("type").and_then(Value::as_str) != Some("token_count") {
3242 continue;
3243 }
3244 let Some(total) = payload.get("info").and_then(|i| i.get("total_token_usage"))
3245 else {
3246 continue;
3247 };
3248 let now = [
3249 n(total.get("input_tokens")),
3250 n(total.get("output_tokens")),
3251 n(total.get("cached_input_tokens")),
3252 ];
3253 if now == previous {
3254 continue;
3255 }
3256 add(
3257 record.get("timestamp").and_then(Value::as_str),
3258 &model,
3259 now[0].saturating_sub(previous[0]),
3260 now[1].saturating_sub(previous[1]),
3261 now[2].saturating_sub(previous[2]),
3262 0,
3263 );
3264 previous = now;
3265 }
3266 }
3267 _ => {}
3268 }
3269 days.into_iter()
3270 .map(|((day, model), [input, output, read, write])| {
3271 let mut usage = json!({"at": format!("{day}T00:00:00Z"), "model": model,
3272 "input_tokens": input, "output_tokens": output,
3273 "cache_read_tokens": read, "cache_write_tokens": write});
3274 if let Some(price) = crate::pricing::built_in(&model) {
3277 usage["cost"] = json!({"amount": price.cost_usd(input + read + write, output),
3278 "currency": "USD", "source": "rate_table"});
3279 }
3280 usage
3281 })
3282 .collect()
3283}
3284
3285fn recorded_session_profile(session: &Session) -> Value {
3288 let harness = match session.meta.source {
3289 SessionSource::ClaudeCode => "claude-code",
3290 SessionSource::Codex => "codex",
3291 SessionSource::Gemini => "gemini",
3292 SessionSource::Goose => "goose",
3293 SessionSource::Grok => "grok",
3294 SessionSource::Native => "native",
3295 SessionSource::OpenClaw => "openclaw",
3296 SessionSource::Hermes => "hermes",
3297 SessionSource::OpenCode => "opencode",
3298 SessionSource::Pi => "pi",
3299 };
3300 let mut model = session
3301 .messages
3302 .iter()
3303 .rev()
3304 .find_map(|message| message.metadata.get("model").cloned())
3305 .or_else(|| session.meta.model.clone());
3306 let mut cwd = session
3307 .meta
3308 .cwd
3309 .as_ref()
3310 .map(|path| path.to_string_lossy().into_owned());
3311 let mut version: Option<String> = None;
3312 let mut approval = Value::Null;
3313 let mut sandbox = Value::Null;
3314 for header in &session.meta.codex_headers {
3315 let payload = header.get("payload").unwrap_or(header);
3316 if let Some(value) = payload.get("cli_version").and_then(Value::as_str) {
3317 version = Some(value.to_owned());
3318 }
3319 if let Some(value) = payload.get("model").and_then(Value::as_str) {
3320 model = Some(value.to_owned());
3321 }
3322 if let Some(value) = payload.get("cwd").and_then(Value::as_str) {
3323 cwd = Some(value.to_owned());
3324 }
3325 if let Some(value) = payload.get("approval_policy") {
3326 approval = value.clone();
3327 }
3328 if let Some(value) = payload.get("sandbox_policy") {
3329 sandbox = value.clone();
3330 }
3331 }
3332 if session.meta.source == SessionSource::ClaudeCode {
3333 version = session
3334 .raw
3335 .iter()
3336 .rev()
3337 .filter_map(|line| serde_json::from_str::<Value>(line).ok())
3338 .find_map(|record| {
3339 record
3340 .get("version")
3341 .and_then(Value::as_str)
3342 .map(str::to_owned)
3343 });
3344 }
3345 json!({"worker":{"harness":harness,"model":model,"cwd":cwd,"binary":{"version":version}},
3346 "permissions":{"native":{"approval":approval,"sandbox":sandbox}},"source":"native-session"})
3347}
3348
3349fn recorded_session_profile_history(session: &Session) -> Vec<Value> {
3352 let mut profile = json!({"worker":{"harness":if session.meta.source == SessionSource::Codex {"codex"} else {"claude-code"},
3353 "model":null,"cwd":null,"binary":{"version":null}},
3354 "permissions":{"native":{"approval":null,"sandbox":null}},"source":"native-session"});
3355 let records: Box<dyn Iterator<Item = (usize, Value)> + '_> =
3356 if session.meta.source == SessionSource::Codex {
3357 Box::new(session.meta.codex_headers.iter().cloned().enumerate())
3358 } else if session.meta.source == SessionSource::ClaudeCode {
3359 Box::new(session.raw.iter().enumerate().filter_map(|(index, line)| {
3360 serde_json::from_str::<Value>(line)
3361 .ok()
3362 .map(|record| (index, record))
3363 }))
3364 } else {
3365 return Vec::new();
3366 };
3367 let mut history = Vec::new();
3368 for (index, record) in records {
3369 let before = profile.clone();
3370 let payload = record.get("payload").unwrap_or(&record);
3371 for (native, target) in [("cwd", "cwd"), ("model", "model")] {
3372 if let Some(value) = payload.get(native).and_then(Value::as_str) {
3373 profile["worker"][target] = json!(value);
3374 }
3375 }
3376 let version = payload
3377 .get("cli_version")
3378 .or_else(|| record.get("version"))
3379 .and_then(Value::as_str);
3380 if let Some(value) = version {
3381 profile["worker"]["binary"]["version"] = json!(value);
3382 }
3383 if let Some(value) = record
3384 .get("message")
3385 .and_then(|m| m.get("model"))
3386 .and_then(Value::as_str)
3387 {
3388 profile["worker"]["model"] = json!(value);
3389 }
3390 for (native, target) in [
3391 ("approval_policy", "approval"),
3392 ("sandbox_policy", "sandbox"),
3393 ] {
3394 if let Some(value) = payload.get(native) {
3395 profile["permissions"]["native"][target] = value.clone();
3396 }
3397 }
3398 if profile != before {
3399 history.push(json!({"native_key":format!("profile-record:{index}"),"at":record.get("timestamp").and_then(Value::as_str),"profile":profile}));
3400 }
3401 }
3402 history
3403}
3404
3405fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3406 if let Some(tail) = options.message_tail {
3407 return (total.saturating_sub(tail), total);
3408 }
3409 let offset = options.message_offset.unwrap_or(0).min(total);
3410 let end = options
3411 .message_limit
3412 .map(|limit| offset.saturating_add(limit).min(total))
3413 .unwrap_or(total);
3414 (offset, end)
3415}
3416
3417fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3418 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3419 return message;
3420 };
3421 for part in parts {
3422 let Some(url) = part
3423 .get("image_url")
3424 .and_then(|image| image.get("url"))
3425 .and_then(Value::as_str)
3426 else {
3427 continue;
3428 };
3429 let Some(rest) = url.strip_prefix("data:") else {
3430 continue;
3431 };
3432 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3433 continue;
3434 };
3435 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3436 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3437 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3438 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3439 || options
3440 .max_inline_media_bytes
3441 .is_some_and(|limit| decoded_bytes > limit);
3442 if should_elide {
3443 *part = json!({
3444 "type": "media_reference",
3445 "media_type": media_type,
3446 "encoding": "base64",
3447 "encoded_bytes": encoded.len(),
3448 "decoded_bytes": decoded_bytes,
3449 "omitted": true,
3450 });
3451 }
3452 }
3453 message
3454}
3455
3456#[derive(Deserialize)]
3457struct LocatorParams {
3458 locator: SessionLocator,
3459 #[serde(default)]
3472 fidelity: Option<Fidelity>,
3473 #[serde(default)]
3476 view: Option<SessionReadView>,
3477}
3478
3479#[derive(Deserialize)]
3480struct SessionReadView {
3481 #[serde(default)]
3484 tail_messages: Option<usize>,
3485 #[serde(default)]
3488 include_subagents: bool,
3489 #[serde(default)]
3491 display_history: bool,
3492 #[serde(default)]
3495 max_message_chars: Option<usize>,
3496}
3497
3498impl LocatorParams {
3499 fn read_fidelity(&self) -> Fidelity {
3500 self.fidelity.unwrap_or(Fidelity::Semantic)
3501 }
3502
3503 fn include_subagents(&self) -> bool {
3504 self.view
3505 .as_ref()
3506 .map(|view| view.include_subagents)
3507 .unwrap_or(true)
3508 }
3509
3510 fn tail_messages(&self) -> Option<usize> {
3511 self.view
3512 .as_ref()
3513 .and_then(|view| view.tail_messages)
3514 .map(|limit| limit.clamp(1, 5_000))
3515 }
3516
3517 fn display_history(&self) -> bool {
3518 self.view.as_ref().is_some_and(|view| view.display_history)
3519 }
3520
3521 fn max_message_chars(&self) -> Option<usize> {
3522 self.view
3523 .as_ref()
3524 .and_then(|view| view.max_message_chars)
3525 .map(|limit| limit.clamp(256, 64_000))
3526 }
3527
3528 fn bound_session(&self, session: &mut Session) {
3529 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3530 }
3531}
3532
3533#[derive(Debug, Clone, Copy, Default, Deserialize)]
3534#[serde(rename_all = "snake_case")]
3535enum InlineMediaMode {
3536 #[default]
3537 Full,
3538 Metadata,
3539}
3540
3541#[derive(Debug, Clone, Default, Deserialize)]
3542#[serde(default)]
3543struct SessionLoadOptions {
3544 include_subagents: Option<bool>,
3545 inline_media: InlineMediaMode,
3546 max_inline_media_bytes: Option<usize>,
3547 message_limit: Option<usize>,
3548 message_offset: Option<usize>,
3549 message_tail: Option<usize>,
3550}
3551
3552impl SessionLoadOptions {
3553 fn validate(&self) -> std::result::Result<(), ServiceError> {
3554 if self.message_tail.is_some()
3555 && (self.message_limit.is_some() || self.message_offset.is_some())
3556 {
3557 return Err(ServiceError::InvalidParams(
3558 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3559 .into(),
3560 ));
3561 }
3562 Ok(())
3563 }
3564}
3565
3566fn load_session_door(params: LoadSessionParams) -> std::result::Result<Value, ServiceError> {
3569 if let Some(options) = ¶ms.options {
3570 options.validate()?;
3571 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
3572 return Ok(result);
3573 }
3574 return load_session(¶ms.read.locator)
3575 .map(|session| projected_session_result(&session, options))
3576 .map_err(operation);
3577 }
3578 let mut session = if params.read.display_history() {
3579 HarnessCatalog::new()
3580 .load_display_view(
3581 ¶ms.read.locator,
3582 params.read.read_fidelity(),
3583 params.read.tail_messages().unwrap_or(500),
3584 )
3585 .map_err(crate::Error::from)
3586 } else if params.read.include_subagents() {
3587 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
3588 } else {
3589 HarnessCatalog::new()
3590 .load_parent_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
3591 .map_err(crate::Error::from)
3592 }
3593 .map_err(operation)?;
3594 params.read.bound_session(&mut session);
3595 let mut value = normalized_session_json(&session);
3598 value["resolved_profile"] = recorded_session_profile(&session);
3599 value["resolved_profile_history"] = Value::Array(recorded_session_profile_history(&session));
3600 Ok(json!({"session": value}))
3601}
3602
3603#[derive(Deserialize)]
3604struct LoadSessionParams {
3605 #[serde(flatten)]
3606 read: LocatorParams,
3607 #[serde(default)]
3608 options: Option<SessionLoadOptions>,
3609}
3610
3611#[derive(Deserialize)]
3612struct UnfollowParams {
3613 subscription: String,
3614}
3615
3616#[derive(Debug, Deserialize)]
3617#[serde(deny_unknown_fields)]
3618struct IndexResizeParams {
3619 subscription: String,
3620 limit: usize,
3621}
3622
3623#[cfg(feature = "adapter-api")]
3626#[derive(Default)]
3627struct PanePrompts {
3628 waiting: std::collections::HashSet<String>,
3629 read_at: Option<std::time::Instant>,
3630 reading: bool,
3631}
3632
3633#[cfg(feature = "adapter-api")]
3634static PANE_PROMPTS: std::sync::Mutex<Option<PanePrompts>> = std::sync::Mutex::new(None);
3635
3636#[cfg(feature = "adapter-api")]
3638const PANE_PROMPT_REFRESH: std::time::Duration = std::time::Duration::from_secs(3);
3639
3640#[cfg(feature = "adapter-api")]
3647fn pane_prompt_override(activities: &mut [crate::SessionActivity], homes: &crate::HarnessHomes) {
3648 let running_codex = |activity: &crate::SessionActivity| {
3649 activity.harness.as_str() == "codex"
3650 && activity.presence == crate::session_activity::SessionPresence::Running
3651 && matches!(
3654 activity.turn,
3655 crate::SessionTurnState::Working | crate::SessionTurnState::Unknown
3656 )
3657 };
3658 if !activities.iter().any(running_codex) {
3659 return;
3660 }
3661 let waiting = {
3662 let Ok(mut guard) = PANE_PROMPTS.lock() else {
3663 return;
3664 };
3665 let prompts = guard.get_or_insert_with(PanePrompts::default);
3666 let stale = prompts
3667 .read_at
3668 .is_none_or(|at| at.elapsed() >= PANE_PROMPT_REFRESH);
3669 if stale && !prompts.reading {
3670 prompts.reading = true;
3671 let homes = homes.clone();
3672 std::thread::spawn(move || {
3673 let doors = crate::mail_route::LiveSessions::read(&homes);
3674 let waiting = doors
3675 .all()
3676 .iter()
3677 .filter(|live| live.address.harness == "codex" && live.status == "waiting")
3678 .map(|live| live.address.session_id.clone())
3679 .collect();
3680 if let Ok(mut guard) = PANE_PROMPTS.lock() {
3681 let prompts = guard.get_or_insert_with(PanePrompts::default);
3682 prompts.waiting = waiting;
3683 prompts.read_at = Some(std::time::Instant::now());
3684 prompts.reading = false;
3685 }
3686 });
3687 }
3688 prompts.waiting.clone()
3689 };
3690 for activity in activities
3691 .iter_mut()
3692 .filter(|activity| running_codex(activity))
3693 {
3694 if waiting.contains(&activity.session_id) {
3695 activity.turn = crate::SessionTurnState::NeedsInput;
3696 activity.evidence.source = "pane_prompt".into();
3697 }
3698 }
3699}
3700
3701#[derive(Deserialize)]
3702struct ActivitySubscribeParams {
3703 locators: Vec<SessionLocator>,
3704 #[serde(default)]
3705 homes: crate::HarnessHomes,
3706}
3707
3708#[derive(Debug, Clone, Deserialize)]
3712struct MessageTarget {
3713 harness: HarnessId,
3714 session_id: String,
3715 #[serde(default)]
3716 #[allow(dead_code)]
3717 storage: Option<supercode_interchange::catalog::StorageLocator>,
3718}
3719
3720#[derive(Deserialize)]
3721struct MessageSessionParams {
3722 locator: MessageTarget,
3723 text: String,
3724 #[serde(default)]
3725 subject: Option<String>,
3726 #[serde(default)]
3728 idempotency_key: Option<String>,
3729 #[serde(default)]
3731 channel: bool,
3732 #[serde(default)]
3736 from_name: Option<String>,
3737 #[serde(default)]
3739 voice_for: Option<crate::mailbox::MailAddress>,
3740 #[serde(default)]
3742 sender_name: Option<String>,
3743 #[serde(default)]
3745 in_reply_to: Option<String>,
3746 #[serde(default)]
3749 marker: Option<String>,
3750 #[serde(default)]
3753 reply_marker: Option<String>,
3754 #[serde(default)]
3756 surface: Option<String>,
3757 #[serde(default)]
3759 notify_when_idle: bool,
3760 #[serde(default)]
3765 as_user: bool,
3766 #[serde(default)]
3769 homes: crate::HarnessHomes,
3770}
3771
3772#[derive(Deserialize)]
3773#[serde(deny_unknown_fields)]
3774struct InboxParams {
3775 #[serde(default)]
3777 from_name: Option<String>,
3778 #[serde(default)]
3780 address: Option<String>,
3781 #[serde(default)]
3783 all: bool,
3784}
3785
3786#[derive(Deserialize)]
3787#[serde(deny_unknown_fields)]
3788struct HarnessSettingsParams {
3789 harness: String,
3790}
3791
3792#[derive(Deserialize)]
3793#[serde(deny_unknown_fields)]
3794struct ConfigureHarnessParams {
3795 harness: String,
3796 #[serde(default)]
3797 changes: Vec<crate::HarnessSettingChange>,
3798 #[serde(default)]
3799 expected_revision: Option<String>,
3800}
3801
3802fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3803 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3804 Ok(report) => (
3805 serde_json::to_value(report).unwrap_or(Value::Null),
3806 Value::Null,
3807 ),
3808 Err(error) => (
3809 Value::Null,
3810 Value::String(format!(
3811 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3812 )),
3813 ),
3814 }
3815}
3816
3817async fn message_live_session(params: &MessageSessionParams) -> Value {
3831 use crate::mail_route::{Delivered, NoDoor, Refused};
3832 let (inbound_controls, inbound_controls_error) =
3833 claude_inbound_controls_or_error(¶ms.homes);
3834 let refused = |reason: &str, message: String| {
3835 json!({
3836 "delivered_to_bus": false,
3837 "refusal": {"reason": reason, "message": message},
3838 "inbound_controls": inbound_controls,
3839 "inbound_controls_error": inbound_controls_error,
3840 })
3841 };
3842 if params.text.trim().is_empty() {
3843 return refused(
3844 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3845 "refusing to deliver an empty message".into(),
3846 );
3847 }
3848 let sender = match operator_address(params.from_name.as_deref()) {
3849 Ok(sender) => sender,
3850 Err(message) => return refused("invalid_sender", message),
3851 };
3852 let receiver = match crate::mailbox::MailAddress::new(
3853 crate::mailbox::local_machine_name(),
3854 params.locator.harness.as_str(),
3855 ¶ms.locator.session_id,
3856 ) {
3857 Ok(receiver) => receiver,
3858 Err(error) => return refused("delivery_failed", error.to_string()),
3859 };
3860 if params
3861 .voice_for
3862 .as_ref()
3863 .is_some_and(|voice_for| voice_for != &receiver)
3864 {
3865 return refused(
3866 "invalid_sender",
3867 "a voice front must represent the receiving session".into(),
3868 );
3869 }
3870 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3873 && !crate::mail_agent::resumable(&receiver)
3874 && crate::runtime_mail::controlled_runtime(
3875 HarnessId::CLAUDE_CODE,
3876 ¶ms.locator.session_id,
3877 )
3878 .is_none()
3879 {
3880 if let Err(refusal) =
3881 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3882 {
3883 return refused(refusal.reason.as_str(), refusal.message);
3884 }
3885 }
3886 if params.as_user {
3887 let agent_address =
3890 (receiver.harness == crate::mail_agent::AGENT_HARNESS).then(|| receiver.clone());
3891 let receiver = if receiver.harness == crate::mail_agent::AGENT_HARNESS {
3892 match crate::mail_agent::load(&receiver.session_id) {
3893 Some(agent) => agent.main_session,
3894 None => {
3895 return refused(
3896 "delivery_failed",
3897 format!(
3898 "no agent named {} is declared on this machine",
3899 receiver.session_id
3900 ),
3901 )
3902 }
3903 }
3904 } else {
3905 receiver
3906 };
3907 return message_as_user(
3908 params,
3909 sender,
3910 receiver,
3911 agent_address,
3912 inbound_controls,
3913 inbound_controls_error,
3914 )
3915 .await;
3916 }
3917 let door = crate::mail_route::door_for(¶ms.homes, &receiver);
3920 let mut envelope = match crate::mailbox::Envelope::new(
3921 sender.clone(),
3922 params
3923 .sender_name
3924 .clone()
3925 .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3926 if params.channel {
3927 crate::mailbox::MailKind::Channel
3928 } else {
3929 crate::mailbox::MailKind::Peer
3930 },
3931 crate::mailbox::ReplyVia::Command,
3932 params.text.clone(),
3933 ) {
3934 Ok(envelope) => envelope,
3935 Err(error) => return refused("delivery_failed", error.to_string()),
3936 };
3937 if let Some(key) = ¶ms.idempotency_key {
3938 if key.is_empty() || key.len() > 256 {
3939 return refused("invalid_key", "idempotency_key needs 1–256 bytes".into());
3940 }
3941 let identity = serde_json::to_vec(&(sender.to_string(), receiver.to_string(), key))
3942 .expect("string tuple serializes");
3943 envelope.id = format!("m-{}", &blake3::hash(&identity).to_hex()[..24]);
3944 let mailbox = match crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &receiver) {
3945 Ok(mailbox) => mailbox,
3946 Err(error) => return refused("delivery_failed", error.to_string()),
3947 };
3948 match mailbox.find(&envelope.id) {
3949 Ok(Some(previous)) => {
3950 let saved = previous.envelope;
3951 let subject = params.subject.as_deref().map(|value| {
3952 value
3953 .lines()
3954 .next()
3955 .unwrap_or("")
3956 .chars()
3957 .take(200)
3958 .collect::<String>()
3959 });
3960 if saved.body != envelope.body
3961 || saved.subject != subject
3962 || saved.kind != envelope.kind
3963 || saved.from_name != envelope.from_name
3964 || saved.in_reply_to != params.in_reply_to
3965 || saved.voice_for != params.voice_for
3966 {
3967 return refused(
3968 "idempotency_conflict",
3969 "this key already names different mail".into(),
3970 );
3971 }
3972 }
3973 Ok(None) => {}
3974 Err(error) => return refused("delivery_failed", error.to_string()),
3975 }
3976 }
3977 envelope.in_reply_to = params.in_reply_to.clone();
3978 envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
3979 envelope.voice_for = params.voice_for.clone();
3980 envelope.subject = params.subject.as_deref().map(|value| {
3981 value
3982 .lines()
3983 .next()
3984 .unwrap_or("")
3985 .chars()
3986 .take(200)
3987 .collect()
3988 });
3989 let channel = crate::mail_agent::Channel {
3990 marker: params.marker.clone(),
3991 reply_marker: params.reply_marker.clone(),
3992 surface: params.surface.clone(),
3993 };
3994 match crate::mail_agent::plan(&mut envelope, &receiver, &channel) {
3995 Err(error) => return refused("delivery_failed", error.to_string()),
3996 Ok(Some(plan)) => {
3997 let caller = crate::mail_route::Caller {
3998 address: sender.clone(),
3999 name: envelope.from_name.clone(),
4000 };
4001 let outcome = crate::mail_send::deliver_planned(
4002 ¶ms.homes,
4003 &caller,
4004 &envelope,
4005 &plan,
4006 false,
4007 params.notify_when_idle,
4008 )
4009 .await;
4010 return json!({
4011 "delivered_to_bus": outcome.code == 0,
4012 "message_id": envelope.id,
4013 "reply_to": sender.to_string(),
4014 "thread": plan.thread.as_ref().map(|thread| thread.id.clone()),
4015 "target": {
4016 "harness": params.locator.harness.as_str(),
4017 "session_id": params.locator.session_id,
4018 },
4019 "delivery": {"door": "agent", "how": outcome.text},
4020 "inbound_controls": inbound_controls,
4021 "inbound_controls_error": inbound_controls_error,
4022 });
4023 }
4024 Ok(None) => {}
4025 }
4026 let door = match door {
4027 Ok(door) => door,
4028 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
4029 return refused(
4030 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
4031 format!(
4032 "no running `{}` session `{}` is reachable; its transcript is persisted only",
4033 params.locator.harness.as_str(),
4034 params.locator.session_id
4035 ),
4036 )
4037 }
4038 };
4039 let how = match crate::mail_route::deliver(
4040 &envelope,
4041 &receiver,
4042 &door,
4043 true,
4044 params.notify_when_idle,
4045 )
4046 .await
4047 {
4048 Err(detail) => {
4049 return refused(
4050 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
4051 detail,
4052 )
4053 }
4054 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
4055 Ok(Err(Refused::TooLong(bytes))) => {
4056 return refused(
4057 "too_long",
4058 format!(
4059 "the message is {bytes} bytes; the limit is {}",
4060 crate::mail_route::MAX_RELAYED_BYTES
4061 ),
4062 )
4063 }
4064 Ok(Ok(delivered)) => match delivered {
4065 Delivered::Steered => "steered",
4066 Delivered::Started => "started",
4067 Delivered::Native { busy: true } => "next_tool_call",
4068 Delivered::Native { busy: false } => "started",
4069 Delivered::Hooked => "hook",
4070 Delivered::HookWoken => "woken",
4071 Delivered::Queued => "queued",
4072 Delivered::Stored => "stored",
4073 Delivered::Operator => "filed",
4074 },
4075 };
4076 json!({
4077 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
4078 "message_id": envelope.id,
4079 "reply_to": sender.to_string(),
4080 "target": {
4081 "harness": params.locator.harness.as_str(),
4082 "session_id": params.locator.session_id,
4083 "name": match &door {
4084 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
4085 _ => None,
4086 },
4087 },
4088 "delivery": {"door": door.name(), "how": how},
4089 "inbound_controls": inbound_controls,
4090 "inbound_controls_error": inbound_controls_error,
4091 })
4092}
4093
4094async fn message_as_user(
4098 params: &MessageSessionParams,
4099 sender: crate::mailbox::MailAddress,
4100 receiver: crate::mailbox::MailAddress,
4101 agent: Option<crate::mailbox::MailAddress>,
4102 inbound_controls: Value,
4103 inbound_controls_error: Value,
4104) -> Value {
4105 let mut envelope = match crate::mailbox::Envelope::new(
4106 sender.clone(),
4107 format!("{}@{}", sender.session_id, sender.machine),
4108 crate::mailbox::MailKind::User,
4109 crate::mailbox::ReplyVia::None,
4110 params.text.clone(),
4111 ) {
4112 Ok(envelope) => envelope,
4113 Err(error) => {
4114 return json!({
4115 "delivered_to_bus": false,
4116 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4117 "inbound_controls": inbound_controls,
4118 "inbound_controls_error": inbound_controls_error,
4119 })
4120 }
4121 };
4122 let mut receiver = receiver;
4127 if let Some(agent) = &agent {
4128 envelope.in_reply_to = params.in_reply_to.clone();
4129 envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
4130 let channel = crate::mail_agent::Channel::default();
4131 let planned = match crate::mail_agent::plan(&mut envelope, agent, &channel) {
4132 Ok(planned) => planned,
4133 Err(error) => {
4134 return json!({
4135 "delivered_to_bus": false,
4136 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4137 "inbound_controls": inbound_controls,
4138 "inbound_controls_error": inbound_controls_error,
4139 })
4140 }
4141 };
4142 if let Some(planned) = planned {
4143 let typed_to = planned
4144 .recipients
4145 .iter()
4146 .find(|r| r.wake && r.address.harness != "operator")
4147 .map(|r| r.address.clone());
4148 if let Some(typed_to) = typed_to {
4149 receiver = typed_to;
4150 }
4151 if let Some(agent) = &planned.agent {
4152 let _ = crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), agent)
4154 .and_then(|mailbox| mailbox.deliver_read(&envelope));
4155 }
4156 let mut copy = envelope.clone();
4159 copy.kind = crate::mailbox::MailKind::Typed;
4160 for other in planned
4161 .recipients
4162 .iter()
4163 .filter(|r| r.address != receiver && r.address != sender)
4164 {
4165 let _ = crate::mail_agent::file_unread(&other.address, ©);
4166 }
4167 if let Some(account_manager) = crate::mail_agent::owners_account_manager() {
4170 let main = account_manager.main_session.clone();
4171 let already =
4172 main == receiver || planned.recipients.iter().any(|r| r.address == main);
4173 if account_manager.name != agent.session_id && !already {
4174 let _ = crate::mail_agent::file_unread(&main, ©);
4175 }
4176 }
4177 }
4178 }
4179 if agent.is_some()
4182 && crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4183 && crate::mail_agent::resumable(&receiver)
4184 {
4185 if let Err(error) = crate::mail_agent::resume(&receiver) {
4186 return json!({
4187 "delivered_to_bus": false,
4188 "refusal": {"reason": "no_user_door", "message": format!("{receiver} was not resumed: {error}")},
4189 "inbound_controls": inbound_controls,
4190 "inbound_controls_error": inbound_controls_error,
4191 });
4192 }
4193 let started = std::time::Instant::now();
4194 while crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4195 && started.elapsed() < crate::mail_agent::RESUME_WAIT
4196 {
4197 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
4198 }
4199 }
4200 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
4201 Ok(turn) => json!({
4202 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
4203 "message_id": envelope.id,
4204 "target": {
4205 "harness": params.locator.harness.as_str(),
4206 "session_id": params.locator.session_id,
4207 },
4208 "delivery": {
4209 "door": match turn {
4210 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
4211 _ => "pane",
4212 },
4213 "how": turn.as_str(),
4214 },
4215 "inbound_controls": inbound_controls,
4216 "inbound_controls_error": inbound_controls_error,
4217 }),
4218 Err(message) => json!({
4219 "delivered_to_bus": false,
4220 "refusal": {"reason": "no_user_door", "message": message},
4221 "inbound_controls": inbound_controls,
4222 "inbound_controls_error": inbound_controls_error,
4223 }),
4224 }
4225}
4226
4227#[derive(Deserialize)]
4228#[serde(deny_unknown_fields)]
4229struct ActivityUnderParams {
4230 pids: Vec<u32>,
4232 #[serde(default)]
4233 homes: crate::HarnessHomes,
4234}
4235
4236async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
4240 let params = decode::<ActivityUnderParams>(params)?;
4241 if params.pids.len() > 1024 {
4242 return Err(ServiceError::InvalidParams(
4243 "sessions.activity_under accepts at most 1024 pids".into(),
4244 ));
4245 }
4246 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
4247 .await
4248 .map_err(ServiceError::Sdk)?;
4249 Ok(json!({
4250 "activities": found
4251 .into_iter()
4252 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
4253 .collect::<Vec<_>>(),
4254 }))
4255}
4256
4257fn operator_address(
4260 name: Option<&str>,
4261) -> std::result::Result<crate::mailbox::MailAddress, String> {
4262 let name: String = name
4263 .unwrap_or("supercode")
4264 .trim()
4265 .chars()
4266 .map(|character| {
4267 if character.is_whitespace() || character == '@' {
4268 '-'
4269 } else {
4270 character
4271 }
4272 })
4273 .collect();
4274 if name.is_empty() {
4275 return Err("from_name must not be empty".into());
4276 }
4277 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
4278 .map_err(|error| error.to_string())
4279}
4280
4281fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
4285 let address = match (¶ms.address, ¶ms.from_name) {
4286 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
4287 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
4288 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
4289 (Some(_), Some(_)) => {
4290 return Err(ServiceError::InvalidParams(
4291 "sessions.inbox takes from_name or address, not both".into(),
4292 ))
4293 }
4294 };
4295 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
4296 let mailbox =
4297 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
4298 let claimed = mailbox.claim_unread().map_err(operation)?;
4299 let mut messages: Vec<Value> = Vec::new();
4300 if params.all {
4301 for stored in mailbox.list().map_err(operation)? {
4302 if stored.state == crate::mailbox::MailState::Read {
4303 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4304 }
4305 }
4306 }
4307 for stored in &claimed {
4308 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4309 }
4310 for stored in &claimed {
4311 mailbox.acknowledge(stored).map_err(operation)?;
4312 }
4313 Ok(json!({"address": address.to_string(), "messages": messages}))
4314}
4315
4316#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4321struct FollowedSource {
4322 harness: String,
4323 session_id: String,
4324 reported: Option<String>,
4325}
4326
4327#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4328struct ActivitySubscription {
4329 locators: Vec<SessionLocator>,
4330 homes: crate::HarnessHomes,
4331 reported: BTreeMap<(String, String), crate::SessionActivity>,
4332}
4333
4334fn live_descriptor_value(
4341 session: &SessionDescriptor,
4342 doors: &crate::mail_route::LiveSessions,
4343) -> std::result::Result<Value, ServiceError> {
4344 let mut value = serde_json::to_value(session)
4345 .map_err(|error| ServiceError::Operation(error.to_string()))?;
4346 if value.get("title").is_none_or(Value::is_null) {
4348 if let Some(live) = doors.all().iter().find(|live| {
4349 live.address.harness == session.locator.harness.as_str()
4350 && live.address.session_id == session.locator.session_id
4351 }) {
4352 let name = live.name.split('@').next().unwrap_or(&live.name);
4353 if !name.is_empty() {
4354 value["title"] = json!(name);
4355 }
4356 }
4357 }
4358 if let Some(workspace) = &session.cwd {
4359 let source = LiveRuntimeSource {
4360 harness: session.locator.harness.as_str().to_string(),
4361 session_id: session.locator.session_id.clone(),
4362 workspace: workspace.clone(),
4363 };
4364 if let Some(endpoint) = discover_live_runtime(&source)
4365 .map_err(|error| ServiceError::Operation(error.to_string()))?
4366 {
4367 value["live_endpoint"] = json!(endpoint.as_str());
4368 }
4369 }
4370 if let Some(door) = doors.door(
4374 session.locator.harness.as_str(),
4375 &session.locator.session_id,
4376 ) {
4377 value["delivery"] = json!(door);
4378 }
4379 if let Some(live) = doors.all().iter().find(|live| {
4382 live.address.harness == session.locator.harness.as_str()
4383 && live.address.session_id == session.locator.session_id
4384 && live.status == "waiting"
4385 }) {
4386 if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
4387 value["pending_request"] = request;
4388 }
4389 }
4390 Ok(value)
4391}
4392
4393fn live_index_changes(
4394 changes: Vec<crate::session_index::SessionIndexChange>,
4395 homes: &HarnessHomes,
4396) -> std::result::Result<Vec<Value>, ServiceError> {
4397 use crate::session_index::SessionIndexChange;
4398 let doors = crate::mail_route::LiveSessions::read(homes);
4399 changes
4400 .into_iter()
4401 .map(|change| match change {
4402 SessionIndexChange::Added { descriptor } => Ok(json!({
4403 "kind": "added",
4404 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4405 })),
4406 SessionIndexChange::Updated { descriptor } => Ok(json!({
4407 "kind": "updated",
4408 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4409 })),
4410 SessionIndexChange::Removed { key } => Ok(json!({
4411 "kind": "removed",
4412 "key": key,
4413 })),
4414 })
4415 .collect()
4416}
4417
4418fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
4419 use crate::{SessionPresence, SessionTurnState};
4420 match (activity.presence, activity.turn) {
4421 (SessionPresence::Persisted, _) => None,
4422 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
4423 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
4424 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
4426 (SessionPresence::Running, SessionTurnState::Unknown)
4430 if activity.evidence.native_state.is_none() =>
4431 {
4432 None
4433 }
4434 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
4435 }
4436}
4437
4438#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
4439#[serde(rename_all = "kebab-case")]
4440enum TransferFormat {
4441 ClaudeCode,
4442 Codex,
4443 #[serde(rename = "opencode", alias = "open-code")]
4444 OpenCode,
4445 Pi,
4446 Grok,
4447 Gemini,
4448 Goose,
4449 Hermes,
4453}
4454
4455impl TransferFormat {
4456 fn id(self) -> &'static str {
4457 match self {
4458 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
4459 Self::Codex => HarnessId::CODEX,
4460 Self::OpenCode => HarnessId::OPENCODE,
4461 Self::Pi => HarnessId::PI,
4462 Self::Grok => HarnessId::GROK,
4463 Self::Gemini => HarnessId::GEMINI,
4464 Self::Goose => HarnessId::GOOSE,
4465 Self::Hermes => HarnessId::HERMES,
4466 }
4467 }
4468}
4469
4470impl From<TransferFormat> for SessionFormat {
4471 fn from(value: TransferFormat) -> Self {
4472 match value {
4473 TransferFormat::ClaudeCode => Self::ClaudeCode,
4474 TransferFormat::Codex => Self::Codex,
4475 TransferFormat::OpenCode => Self::OpenCode,
4476 TransferFormat::Pi => Self::Pi,
4477 TransferFormat::Grok => Self::Grok,
4478 TransferFormat::Gemini => Self::Gemini,
4479 TransferFormat::Goose => Self::Goose,
4480 TransferFormat::Hermes => Self::Codex,
4482 }
4483 }
4484}
4485
4486#[derive(Deserialize)]
4487struct ImportSessionParams {
4488 source_harness: TransferFormat,
4489 content: String,
4490}
4491
4492#[derive(Deserialize)]
4493struct ExportSessionParams {
4494 locator: SessionLocator,
4495 target_harness: TransferFormat,
4496}
4497
4498#[derive(Deserialize)]
4499struct ReduceSessionParams {
4500 locator: SessionLocator,
4501 target_harness: TransferFormat,
4502 #[serde(default = "default_keep_last")]
4503 keep_last: usize,
4504}
4505
4506fn default_keep_last() -> usize {
4507 6
4508}
4509
4510#[derive(Deserialize)]
4511struct BranchSessionParams {
4512 locator: SessionLocator,
4513 #[serde(default)]
4514 target_harness: Option<TransferFormat>,
4515}
4516
4517#[derive(Deserialize)]
4518struct HandoffSessionParams {
4519 locator: SessionLocator,
4520 target_harness: TransferFormat,
4521 #[serde(default)]
4522 cwd: Option<PathBuf>,
4523}
4524
4525#[derive(Deserialize)]
4526struct MaterializeSessionParams {
4527 artifact: crate::native_materialize::MaterializeArtifact,
4528 cwd: PathBuf,
4529 #[serde(default)]
4531 homes: HarnessHomes,
4532}
4533
4534#[derive(Debug, Clone, Copy, Default, Deserialize)]
4537#[serde(rename_all = "snake_case")]
4538enum ResumePolicy {
4539 Default,
4540 #[default]
4541 Yolo,
4542}
4543
4544#[derive(Deserialize)]
4545struct ResumeInstructionsParams {
4546 locator: SessionLocator,
4547 #[serde(default)]
4548 cwd: Option<PathBuf>,
4549 #[serde(default)]
4550 policy: ResumePolicy,
4551}
4552
4553#[derive(Deserialize)]
4555struct WorkflowLoadParams {
4556 from: crate::workflow_doors::WorkflowHarness,
4557 home: PathBuf,
4558}
4559
4560#[derive(Deserialize)]
4563struct OrchestrationLoadParams {
4564 root: PathBuf,
4565 #[serde(default)]
4566 flavor: crate::orchestration_doors::HomeFlavor,
4567}
4568
4569#[derive(Deserialize)]
4572struct OrchestrationSaveParams {
4573 root: PathBuf,
4574 orchestration: crate::orchestration::Orchestration,
4575 #[serde(default)]
4576 vault: BTreeMap<String, String>,
4577}
4578
4579#[derive(Deserialize)]
4581struct OrchestrationCompileParams {
4582 from: crate::orchestration_doors::OrchestrationHarness,
4583 home: PathBuf,
4584}
4585
4586#[derive(Deserialize)]
4590struct OrchestrationDecompileParams {
4591 to: crate::orchestration_doors::OrchestrationHarness,
4592 orchestration: crate::orchestration::Orchestration,
4593 source: PathBuf,
4594 #[serde(default)]
4595 source_flavor: crate::orchestration_doors::SourceFlavor,
4596 dest: PathBuf,
4597 #[serde(default)]
4598 vault: BTreeMap<String, String>,
4599}
4600
4601#[derive(Deserialize)]
4604struct OrchestrationImportParams {
4605 from: crate::orchestration_doors::OrchestrationHarness,
4606 home: PathBuf,
4607 into: PathBuf,
4608}
4609
4610#[derive(Deserialize)]
4613struct OrchestrationExportParams {
4614 to: crate::orchestration_doors::OrchestrationHarness,
4615 root: PathBuf,
4616 dest: PathBuf,
4617}
4618
4619#[derive(Deserialize)]
4621struct JobsGetParams {
4622 harness: String,
4623 id: String,
4624 #[serde(default)]
4625 homes: crate::HarnessHomes,
4626}
4627
4628fn mutate_job(
4636 verb: crate::jobs_control::JobVerb,
4637 params: Value,
4638) -> std::result::Result<Value, ServiceError> {
4639 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4640 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4641 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4642 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4643}
4644
4645fn mutate_skill(
4652 verb: crate::skills_control::SkillVerb,
4653 params: Value,
4654) -> std::result::Result<Value, ServiceError> {
4655 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4656 if !crate::skills_control::supports_skill_control(&mutation.harness) {
4657 return Err(ServiceError::UnsupportedAction(format!(
4658 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4659 mutation.harness,
4660 verb.as_str(),
4661 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4662 )));
4663 }
4664 let outcome =
4665 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4666 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4667}
4668
4669fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4671 match error {
4672 crate::skills_control::SkillControlError::Unsupported(message) => {
4673 ServiceError::UnsupportedAction(message)
4674 }
4675 crate::skills_control::SkillControlError::Invalid(message) => {
4676 ServiceError::InvalidParams(message)
4677 }
4678 crate::skills_control::SkillControlError::Failed(message) => {
4679 ServiceError::Operation(message)
4680 }
4681 }
4682}
4683
4684fn mutate_profile(
4692 verb: crate::profiles_control::ProfileVerb,
4693 params: Value,
4694) -> std::result::Result<Value, ServiceError> {
4695 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4696 let outcome =
4697 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4698 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4699}
4700
4701fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4703 match error {
4704 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4705 ServiceError::UnsupportedAction(message)
4706 }
4707 crate::profiles_control::ProfileControlError::Invalid(message) => {
4708 ServiceError::InvalidParams(message)
4709 }
4710 crate::profiles_control::ProfileControlError::Failed(message) => {
4711 ServiceError::Operation(message)
4712 }
4713 }
4714}
4715
4716fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4720 match error {
4721 crate::jobs_control::JobControlError::Unsupported(message) => {
4722 ServiceError::UnsupportedAction(message)
4723 }
4724 crate::jobs_control::JobControlError::Invalid(message) => {
4725 ServiceError::InvalidParams(message)
4726 }
4727 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4728 }
4729}
4730
4731fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4736 match error {
4737 crate::SessionControlError::Unsupported(message) => {
4738 ServiceError::UnsupportedAction(message)
4739 }
4740 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4741 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4742 }
4743}
4744
4745fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4750 if crate::jobs::supports_jobs(harness) {
4751 return Ok(());
4752 }
4753 Err(ServiceError::UnsupportedAction(format!(
4754 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4755 crate::jobs::JOB_HARNESSES.join(", ")
4756 )))
4757}
4758
4759#[derive(Deserialize)]
4761struct RunsGetParams {
4762 harness: String,
4763 id: String,
4764 #[serde(default)]
4765 homes: crate::HarnessHomes,
4766}
4767
4768fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4773 if crate::runs::supports_runs(harness) {
4774 return Ok(());
4775 }
4776 Err(ServiceError::UnsupportedAction(format!(
4777 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4778 crate::runs::RUN_HARNESSES.join(", ")
4779 )))
4780}
4781
4782#[derive(Serialize)]
4783struct SessionArtifact {
4784 source_harness: HarnessId,
4785 target_harness: &'static str,
4786 session_id: Option<String>,
4787 content: String,
4788 suggested_filename: String,
4789 files: Vec<SessionArtifactFile>,
4790 fidelity: Fidelity,
4791 residue: Vec<String>,
4792}
4793
4794#[derive(Serialize)]
4795struct SessionArtifactFile {
4796 path: String,
4797 content: String,
4798 role: ArtifactFileRole,
4799}
4800
4801#[derive(Serialize)]
4802#[serde(rename_all = "snake_case")]
4803enum ArtifactFileRole {
4804 Primary,
4805 Subagent,
4806 Bundle,
4807 SourceRecovery,
4808}
4809
4810#[derive(Serialize)]
4811struct StructuredLaunch {
4812 cwd: PathBuf,
4813 program: String,
4814 arguments: Vec<String>,
4815 env: BTreeMap<String, String>,
4816}
4817
4818struct HandoffInstructions {
4819 launch: StructuredLaunch,
4820 materialize: Option<StructuredLaunch>,
4821 requires_materialization: bool,
4822 note: String,
4823}
4824
4825#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4826#[serde(rename_all = "snake_case")]
4827enum HarnessProbeLevel {
4828 #[default]
4829 Passive,
4830 Handshake,
4831}
4832
4833#[derive(Default, Deserialize)]
4834#[serde(default)]
4835struct HarnessInventoryParams {
4836 harness: Option<HarnessId>,
4837 harnesses: Vec<HarnessId>,
4838 workspace: Option<PathBuf>,
4839 probe: HarnessProbeLevel,
4840 include_sessions: bool,
4841 skip_versions: bool,
4843}
4844
4845#[derive(Deserialize)]
4846struct HarnessAuthenticationParams {
4847 harness: HarnessId,
4848}
4849
4850#[derive(Deserialize)]
4851struct BeginHarnessAuthenticationParams {
4852 harness: HarnessId,
4853 #[serde(default = "local_browser_authentication_environment")]
4854 environment: crate::HarnessAuthenticationEnvironment,
4855 #[serde(default)]
4856 method: Option<crate::HarnessAuthenticationMethodId>,
4857 #[serde(default)]
4858 cwd: Option<PathBuf>,
4859}
4860
4861fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4862 crate::HarnessAuthenticationEnvironment::LocalBrowser
4863}
4864
4865#[derive(Serialize)]
4866struct HarnessInventoryReport {
4867 probe: HarnessProbeLevel,
4868 workspace: Option<PathBuf>,
4869 harnesses: Vec<LocalHarness>,
4870}
4871
4872#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4873#[serde(rename_all = "snake_case")]
4874enum HarnessAuthState {
4875 Ready,
4876 Configured,
4877 Required,
4878 Unknown,
4879}
4880
4881#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4882#[serde(rename_all = "snake_case")]
4883enum HarnessRuntimeState {
4884 Ready,
4885 Degraded,
4886 Unavailable,
4887}
4888
4889#[derive(Serialize)]
4890struct HarnessSessionCounts {
4891 global: Option<usize>,
4892 workspace: Option<usize>,
4893}
4894
4895#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4906#[serde(rename_all = "snake_case")]
4907pub enum GatewayState {
4908 Up,
4909 Down,
4910 Unknown,
4911}
4912
4913#[derive(Debug, Clone, Serialize)]
4915pub struct GatewayHealth {
4916 pub state: GatewayState,
4917 #[serde(skip_serializing_if = "Option::is_none")]
4922 pub endpoint: Option<String>,
4923 #[serde(skip_serializing_if = "Option::is_none")]
4924 pub version: Option<String>,
4925 pub evidence: String,
4927 pub checked_at_ms: u64,
4928}
4929
4930fn openclaw_gateway_endpoint(home: &Path) -> String {
4934 let config_path = home.join(".openclaw/openclaw.json");
4935 let gateway = std::fs::read_to_string(&config_path)
4936 .ok()
4937 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4938 .and_then(|config| config.get("gateway").cloned());
4939 if let Some(url) = gateway
4940 .as_ref()
4941 .and_then(|gateway| gateway.get("url"))
4942 .and_then(serde_json::Value::as_str)
4943 {
4944 return url.to_string();
4945 }
4946 let port = gateway
4947 .as_ref()
4948 .and_then(|gateway| gateway.get("port"))
4949 .and_then(serde_json::Value::as_u64)
4950 .unwrap_or(18789);
4951 format!("ws://127.0.0.1:{port}")
4952}
4953
4954fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4961 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4962 let output = std::process::Command::new(&program)
4963 .args(["gateway", "status"])
4964 .stdin(std::process::Stdio::null())
4965 .output()
4966 .ok()?;
4967 let text = format!(
4968 "{}{}",
4969 String::from_utf8_lossy(&output.stdout),
4970 String::from_utf8_lossy(&output.stderr)
4971 );
4972 let verdict = text.lines().find_map(|line| {
4973 let l = line.trim();
4974 if l.contains("supervised by launchd (PID")
4975 || l.contains("supervised by systemd (PID")
4976 || l.contains("Gateway is running")
4977 || l.contains("process is running")
4978 {
4979 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4980 } else if l.contains("not running") || l.contains("not installed") {
4981 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4982 } else {
4983 None
4984 }
4985 });
4986 verdict
4987}
4988
4989fn gateway_health(
4990 id: &str,
4991 installed: bool,
4992 running: Option<&RunningInstance>,
4993 version: Option<&str>,
4994) -> GatewayHealth {
4995 let checked_at_ms = now_epoch_ms();
4996 let home = supercode_interchange::user_home()
4997 .map(std::path::PathBuf::into_os_string)
4998 .map(PathBuf::from);
4999 match id {
5000 HarnessId::HERMES | HarnessId::OPENCLAW => {
5001 let endpoint = (id == HarnessId::OPENCLAW)
5002 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
5003 .flatten();
5004 let (state, evidence) = match running {
5005 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
5006 None if !installed => (
5007 GatewayState::Unknown,
5008 format!("`{id}` is not installed; no gateway to probe"),
5009 ),
5010 None if id == HarnessId::HERMES => match hermes_gateway_status() {
5011 Some((state, evidence)) => (state, evidence),
5014 None => (
5015 GatewayState::Down,
5016 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
5017 ),
5018 },
5019 None => (
5020 GatewayState::Down,
5021 format!(
5022 "no TCP listener at {}",
5023 endpoint.as_deref().unwrap_or("the gateway endpoint")
5024 ),
5025 ),
5026 };
5027 GatewayHealth {
5028 state,
5029 endpoint,
5030 version: version.map(str::to_string),
5031 evidence,
5032 checked_at_ms,
5033 }
5034 }
5035 HarnessId::ORCHESTRATOR => {
5042 let root = crate::HarnessHomes::default().orchestrator;
5043 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
5044 Some(lease) if lease.is_live() => (
5045 GatewayState::Up,
5046 format!(
5047 "`{}` names pid {} (started {}), which is live",
5048 crate::orchestrator::lock_path(&root).display(),
5049 lease.pid,
5050 lease.started_at
5051 ),
5052 ),
5053 Some(lease) => (
5054 GatewayState::Down,
5055 format!(
5056 "stale lease `{}`: pid {} is gone",
5057 crate::orchestrator::lock_path(&root).display(),
5058 lease.pid
5059 ),
5060 ),
5061 None => (
5062 GatewayState::Down,
5063 format!(
5064 "no lease at `{}`; `supercode orchestrator start` writes one",
5065 crate::orchestrator::lock_path(&root).display()
5066 ),
5067 ),
5068 };
5069 GatewayHealth {
5070 state,
5071 endpoint: None,
5072 version: version.map(str::to_string),
5073 evidence,
5074 checked_at_ms,
5075 }
5076 }
5077 _ => GatewayHealth {
5078 state: GatewayState::Unknown,
5079 endpoint: None,
5080 version: version.map(str::to_string),
5081 evidence: format!("`{id}` runs per session, not as a gateway"),
5082 checked_at_ms,
5083 },
5084 }
5085}
5086
5087#[derive(Debug, Clone, Serialize)]
5088struct RunningInstance {
5089 method: RunningInstanceMethod,
5091 evidence: String,
5093 checked_at_ms: u64,
5095}
5096
5097#[derive(Debug, Clone, Copy, Serialize)]
5098#[serde(rename_all = "snake_case")]
5099enum RunningInstanceMethod {
5100 GatewayConnect,
5103 StoreWalActivity,
5106}
5107
5108fn now_epoch_ms() -> u64 {
5109 std::time::SystemTime::now()
5110 .duration_since(std::time::UNIX_EPOCH)
5111 .map(|elapsed| elapsed.as_millis() as u64)
5112 .unwrap_or(0)
5113}
5114
5115fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
5119 let config_path = home.join(".openclaw/openclaw.json");
5120 let text = std::fs::read_to_string(&config_path).ok();
5121 let gateway = text
5122 .as_deref()
5123 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
5124 .and_then(|config| config.get("gateway").cloned());
5125 let address = gateway
5126 .as_ref()
5127 .and_then(|gateway| gateway.get("url"))
5128 .and_then(serde_json::Value::as_str)
5129 .and_then(|url| {
5130 url.split("://").nth(1).map(|rest| {
5131 rest.trim_end_matches('/')
5132 .split('/')
5133 .next()
5134 .unwrap_or(rest)
5135 .to_string()
5136 })
5137 })
5138 .unwrap_or_else(|| {
5139 let port = gateway
5140 .as_ref()
5141 .and_then(|gateway| gateway.get("port"))
5142 .and_then(serde_json::Value::as_u64)
5143 .unwrap_or(18789);
5144 format!("127.0.0.1:{port}")
5145 });
5146 let reachable = std::net::TcpStream::connect_timeout(
5147 &address.parse().ok()?,
5148 std::time::Duration::from_millis(400),
5149 )
5150 .is_ok();
5151 reachable.then(|| RunningInstance {
5152 method: RunningInstanceMethod::GatewayConnect,
5153 evidence: format!(
5154 "gateway endpoint {address} accepted a TCP connect (from {})",
5155 config_path.display()
5156 ),
5157 checked_at_ms: now_epoch_ms(),
5158 })
5159}
5160
5161fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
5166 let wal = home.join(".hermes/state.db-wal");
5167 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
5168 let age_ms = std::time::SystemTime::now()
5169 .duration_since(modified)
5170 .map(|age| age.as_millis() as u64)
5171 .unwrap_or(u64::MAX);
5172 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
5173 method: RunningInstanceMethod::StoreWalActivity,
5174 evidence: format!(
5175 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
5176 wal.display()
5177 ),
5178 checked_at_ms: now_epoch_ms(),
5179 })
5180}
5181
5182fn probe_running_instance(id: &str) -> Option<RunningInstance> {
5184 let home = supercode_interchange::user_home()
5185 .map(std::path::PathBuf::into_os_string)
5186 .map(PathBuf::from)?;
5187 match id {
5188 HarnessId::OPENCLAW => probe_openclaw_running(&home),
5189 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
5190 _ => None,
5191 }
5192}
5193
5194#[derive(Serialize)]
5195struct LocalHarness {
5196 id: HarnessId,
5197 display_name: String,
5198 supported: bool,
5199 installed: bool,
5200 executable: Option<String>,
5201 version: Option<String>,
5202 auth: HarnessAuthState,
5203 runtime: HarnessRuntimeState,
5204 protocol: String,
5205 capabilities: crate::RuntimeCapabilities,
5206 effective_capabilities: crate::RuntimeCapabilities,
5207 sessions: HarnessSessionCounts,
5208 #[serde(skip_serializing_if = "Option::is_none")]
5211 running: Option<RunningInstance>,
5212 gateway: GatewayHealth,
5214 reason: Option<String>,
5215 repair: Option<String>,
5216}
5217
5218#[derive(Clone, Deserialize)]
5219struct RuntimeBackendParams {
5220 harness: HarnessId,
5221 #[serde(default)]
5222 protocol: Option<String>,
5223 #[serde(default)]
5224 launch: Option<RuntimeLaunch>,
5225 #[serde(default)]
5226 base_url: Option<String>,
5227 #[serde(default)]
5228 policy: RuntimePolicy,
5229}
5230
5231#[derive(Debug, Clone, Copy, Default, Deserialize)]
5234#[serde(rename_all = "snake_case")]
5235enum RuntimePolicy {
5236 Default,
5237 #[default]
5238 Yolo,
5239}
5240
5241#[derive(Deserialize)]
5242struct RuntimeStartParams {
5243 #[serde(flatten)]
5244 backend: RuntimeBackendParams,
5245 cwd: PathBuf,
5246 #[serde(default)]
5249 mcp_servers: Vec<crate::McpServerLaunch>,
5250 #[serde(default)]
5252 approval_policy: Option<String>,
5253}
5254
5255#[derive(Deserialize)]
5256struct RuntimeAttachParams {
5257 #[serde(flatten)]
5258 backend: RuntimeBackendParams,
5259 runtime_id: String,
5260 #[serde(default)]
5261 cwd: Option<PathBuf>,
5262 #[serde(default)]
5265 mcp_servers: Vec<crate::McpServerLaunch>,
5266 #[serde(default)]
5268 approval_policy: Option<String>,
5269}
5270
5271#[derive(Deserialize)]
5272struct RuntimeConnectionParams {
5273 connection: String,
5274}
5275
5276#[derive(Deserialize)]
5279struct RuntimeCloseParams {
5280 #[serde(default)]
5281 connection: String,
5282 #[serde(default)]
5283 runtime_id: Option<String>,
5284}
5285
5286#[derive(Deserialize)]
5287struct RuntimeInputParams {
5288 connection: String,
5289 text: String,
5290 #[serde(default)]
5291 image_urls: Vec<String>,
5292}
5293
5294const MAX_RUNTIME_IMAGES: usize = 4;
5295const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
5296const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
5297
5298fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
5299 if image_urls.len() > MAX_RUNTIME_IMAGES {
5300 return Err(ServiceError::InvalidParams(format!(
5301 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
5302 )));
5303 }
5304 let mut total = 0usize;
5305 for url in &image_urls {
5306 if !(url.starts_with("data:image/")
5307 || url.starts_with("https://")
5308 || url.starts_with("http://"))
5309 {
5310 return Err(ServiceError::InvalidParams(
5311 "runtime images must be image data URLs or HTTP(S) URLs".into(),
5312 ));
5313 }
5314 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
5315 return Err(ServiceError::InvalidParams(format!(
5316 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
5317 )));
5318 }
5319 total = total.saturating_add(url.len());
5320 }
5321 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
5322 return Err(ServiceError::InvalidParams(format!(
5323 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
5324 )));
5325 }
5326 Ok(image_urls)
5327}
5328
5329#[derive(Deserialize)]
5330struct RuntimeRespondParams {
5331 connection: String,
5332 request_id: Value,
5333 response: Value,
5334}
5335
5336fn default_reduction_store_root() -> PathBuf {
5337 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
5338 return PathBuf::from(root).join("sessions");
5339 }
5340 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
5341 return PathBuf::from(home).join(".supercode").join("sessions");
5342 }
5343 PathBuf::from(".supercode").join("sessions")
5344}
5345
5346fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
5347 let mut output = String::new();
5348 for message in messages {
5349 output.push_str(
5350 &serde_json::to_string(message)
5351 .map_err(|error| ServiceError::Operation(error.to_string()))?,
5352 );
5353 output.push('\n');
5354 }
5355 Ok(output)
5356}
5357
5358fn parse_messages_jsonl(
5359 content: &str,
5360) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
5361 content
5362 .lines()
5363 .enumerate()
5364 .filter(|(_, line)| !line.trim().is_empty())
5365 .map(|(index, line)| {
5366 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
5367 ServiceError::Operation(format!(
5368 "reduced transcript line {} is invalid: {error}",
5369 index + 1
5370 ))
5371 })
5372 })
5373 .collect()
5374}
5375
5376fn reduced_bootstrap_prompt(
5377 source: &SessionLocator,
5378 target: TransferFormat,
5379 view_jsonl: &str,
5380 sidecar_path: &Path,
5381 reduction_log_path: &Path,
5382) -> String {
5383 format!(
5384 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
5385 \n\
5386 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\
5387 \n\
5388 <supercode-reduced-session source-session=\"{source_id}\">\n\
5389 {view_jsonl}\
5390 </supercode-reduced-session>\n\
5391 \n\
5392 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
5393 source_harness = source.harness.as_str(),
5394 target_harness = target.id(),
5395 sidecar = sidecar_path.display(),
5396 log = reduction_log_path.display(),
5397 source_id = source.session_id,
5398 )
5399}
5400
5401fn session_artifact(
5402 locator: &SessionLocator,
5403 session: &Session,
5404 target: TransferFormat,
5405) -> std::result::Result<SessionArtifact, ServiceError> {
5406 session_artifact_with_id(locator, session, target, None)
5407}
5408
5409fn session_artifact_with_id(
5410 locator: &SessionLocator,
5411 session: &Session,
5412 target: TransferFormat,
5413 target_session_id: Option<&str>,
5414) -> std::result::Result<SessionArtifact, ServiceError> {
5415 let format: SessionFormat = target.into();
5416 let diagonal = format.source() == session.meta.source;
5417 crate::residue_store::store_segments(session);
5418 let has_appended_turns = session
5419 .imported_message_count
5420 .is_some_and(|imported| imported < session.messages.len());
5421 let mut restoration = None;
5422 let content = if let Some(id) = target_session_id {
5423 if diagonal && format != SessionFormat::OpenCode {
5424 session
5425 .to_jsonl_spliced(format, Some(id))
5426 .map_err(operation)?
5427 } else {
5428 let mut rewritten = session.clone();
5429 rewritten.meta.session_id = Some(id.to_string());
5430 rewritten.to_jsonl(format).map_err(operation)?
5431 }
5432 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
5433 session.raw_verbatim()
5434 } else if diagonal {
5435 session.to_jsonl_spliced(format, None).map_err(operation)?
5436 } else {
5437 match session
5440 .restore_residue(format, crate::residue_store::lookup)
5441 .map_err(operation)?
5442 {
5443 Some((content, report)) => {
5444 restoration = Some(report);
5445 content
5446 }
5447 None => session.to_jsonl(format).map_err(operation)?,
5448 }
5449 };
5450 let stem = sanitize_filename(
5451 target_session_id
5452 .or(session.meta.session_id.as_deref())
5453 .unwrap_or(&locator.session_id),
5454 );
5455 let suggested_filename = if diagonal && target == TransferFormat::Grok {
5456 "chat_history.jsonl".to_string()
5457 } else if target == TransferFormat::Goose {
5458 format!("{stem}.goose.json")
5459 } else {
5460 format!("{stem}.{}.jsonl", target.id())
5461 };
5462 let mut files = vec![SessionArtifactFile {
5463 path: suggested_filename.clone(),
5464 content: content.clone(),
5465 role: ArtifactFileRole::Primary,
5466 }];
5467 if target == TransferFormat::ClaudeCode {
5468 let bundle_stem = Path::new(&suggested_filename)
5469 .file_stem()
5470 .and_then(|stem| stem.to_str())
5471 .unwrap_or(&stem);
5472 let mut child_paths = BTreeSet::new();
5473 for (index, subagent) in session.subagents.iter().enumerate() {
5474 let agent_id = subagent
5475 .meta
5476 .agent_id
5477 .as_deref()
5478 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
5479 .map(sanitize_filename)
5480 .filter(|id| !id.is_empty())
5481 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5482 let child_has_appended_turns = subagent
5483 .imported_message_count
5484 .is_some_and(|imported| imported < subagent.messages.len());
5485 let child_content = if target_session_id.is_none()
5486 && subagent.meta.source == SessionSource::ClaudeCode
5487 && subagent.raw_is_verbatim
5488 && !child_has_appended_turns
5489 {
5490 subagent.raw_verbatim()
5491 } else if subagent.meta.source == SessionSource::ClaudeCode {
5492 subagent
5493 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
5494 .map_err(operation)?
5495 } else {
5496 let mut child = subagent.clone();
5497 if let Some(id) = target_session_id {
5498 child.meta.session_id = Some(id.to_string());
5499 }
5500 child
5501 .to_jsonl(SessionFormat::ClaudeCode)
5502 .map_err(operation)?
5503 };
5504 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
5505 if !child_paths.insert(path.clone()) {
5506 return Err(ServiceError::Operation(format!(
5507 "Claude subagent ids collide at artifact path `{path}`"
5508 )));
5509 }
5510 files.push(SessionArtifactFile {
5511 path,
5512 content: child_content,
5513 role: ArtifactFileRole::Subagent,
5514 });
5515 }
5516 }
5517 if diagonal && target == TransferFormat::Grok {
5518 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
5519 }
5520 if !diagonal || !session.raw_is_verbatim {
5521 files.push(SessionArtifactFile {
5522 path: "recovery/source.supercode.jsonl".into(),
5523 content: session.to_native_jsonl(),
5524 role: ArtifactFileRole::SourceRecovery,
5525 });
5526 for (index, subagent) in session.subagents.iter().enumerate() {
5527 let id = subagent
5528 .meta
5529 .agent_id
5530 .as_deref()
5531 .map(sanitize_filename)
5532 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5533 files.push(SessionArtifactFile {
5534 path: format!("recovery/subagents/{id}.supercode.jsonl"),
5535 content: subagent.to_native_jsonl(),
5536 role: ArtifactFileRole::SourceRecovery,
5537 });
5538 }
5539 }
5540 if !diagonal && session.meta.source == SessionSource::Grok {
5541 append_grok_bundle_files(
5542 locator,
5543 "recovery/grok/",
5544 ArtifactFileRole::SourceRecovery,
5545 &mut files,
5546 )?;
5547 }
5548 let (fidelity, residue) = if diagonal
5549 && target_session_id.is_none()
5550 && session.raw_is_verbatim
5551 && !has_appended_turns
5552 {
5553 (Fidelity::ByteLossless, Vec::new())
5554 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
5555 (
5556 Fidelity::ValueLossless,
5557 vec![if target_session_id.is_some() {
5558 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
5559 } else {
5560 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
5561 }],
5562 )
5563 } else {
5564 match restoration {
5565 Some(report) if report.rendered_messages == 0 => (
5566 Fidelity::ByteLossless,
5567 vec![format!(
5568 "restored verbatim from this conversation's {} source records in the residue store",
5569 target.id()
5570 )],
5571 ),
5572 Some(report) => (
5573 Fidelity::Semantic,
5574 vec![format!(
5575 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
5576 report.restored_messages,
5577 report.restored_messages + report.rendered_messages,
5578 report.rendered_messages,
5579 target.id()
5580 )],
5581 ),
5582 None => (
5583 Fidelity::Semantic,
5584 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
5585 ),
5586 }
5587 };
5588 Ok(SessionArtifact {
5589 source_harness: locator.harness.clone(),
5590 target_harness: target.id(),
5591 session_id: target_session_id
5592 .map(str::to_string)
5593 .or_else(|| session.meta.session_id.clone()),
5594 content,
5595 suggested_filename,
5596 files,
5597 fidelity,
5598 residue,
5599 })
5600}
5601
5602fn append_grok_bundle_files(
5603 locator: &SessionLocator,
5604 prefix: &str,
5605 role: ArtifactFileRole,
5606 files: &mut Vec<SessionArtifactFile>,
5607) -> std::result::Result<(), ServiceError> {
5608 let primary = locator.storage.path();
5609 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5610 return Err(ServiceError::Operation(format!(
5611 "Grok bundle locator must name chat_history.jsonl, got {}",
5612 primary.display()
5613 )));
5614 }
5615 let parent = primary.parent().ok_or_else(|| {
5616 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5617 })?;
5618 for name in ["summary.json", "updates.jsonl"] {
5619 let path = parent.join(name);
5620 let metadata = match std::fs::symlink_metadata(&path) {
5621 Ok(metadata) => metadata,
5622 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5623 Err(error) => return Err(ServiceError::Operation(error.to_string())),
5624 };
5625 if metadata.file_type().is_symlink() || !metadata.is_file() {
5626 return Err(ServiceError::Operation(format!(
5627 "refusing non-regular Grok bundle member {}",
5628 path.display()
5629 )));
5630 }
5631 let content = std::fs::read_to_string(&path).map_err(|error| {
5632 ServiceError::Operation(format!(
5633 "Grok bundle member {} is not representable as UTF-8: {error}",
5634 path.display()
5635 ))
5636 })?;
5637 files.push(SessionArtifactFile {
5638 path: format!("{prefix}{name}"),
5639 content,
5640 role: match role {
5641 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5642 _ => ArtifactFileRole::SourceRecovery,
5643 },
5644 });
5645 }
5646 Ok(())
5647}
5648
5649fn handoff_artifact(
5650 locator: &SessionLocator,
5651 session: &Session,
5652 target: TransferFormat,
5653) -> std::result::Result<SessionArtifact, ServiceError> {
5654 let target_session_id = target_session_id(target);
5655 session_artifact_with_id(locator, session, target, Some(&target_session_id))
5656}
5657
5658fn target_session_id(target: TransferFormat) -> String {
5659 let uuid = generated_session_id();
5660 match target {
5661 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5662 TransferFormat::ClaudeCode
5663 | TransferFormat::Codex
5664 | TransferFormat::Pi
5665 | TransferFormat::Grok
5666 | TransferFormat::Gemini
5667 | TransferFormat::Goose
5668 | TransferFormat::Hermes => uuid,
5669 }
5670}
5671
5672fn sanitize_filename(value: &str) -> String {
5673 let value = value
5674 .chars()
5675 .map(|character| {
5676 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5677 character
5678 } else {
5679 '-'
5680 }
5681 })
5682 .collect::<String>();
5683 let value = value.trim_matches('-');
5684 if value.is_empty() {
5685 "session".into()
5686 } else {
5687 value.chars().take(100).collect()
5688 }
5689}
5690
5691fn handoff_instructions(
5692 target: TransferFormat,
5693 session_id: &str,
5694 cwd: &Path,
5695) -> HandoffInstructions {
5696 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5697 cwd: cwd.to_path_buf(),
5698 program: program.into(),
5699 arguments,
5700 env: BTreeMap::new(),
5701 };
5702 match target {
5703 TransferFormat::ClaudeCode => HandoffInstructions {
5704 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5705 materialize: None,
5706 requires_materialization: true,
5707 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(),
5708 },
5709 TransferFormat::Hermes => HandoffInstructions {
5710 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5711 materialize: None,
5712 requires_materialization: true,
5713 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(),
5714 },
5715 TransferFormat::Codex => HandoffInstructions {
5716 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5717 materialize: None,
5718 requires_materialization: true,
5719 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5720 },
5721 TransferFormat::OpenCode => HandoffInstructions {
5722 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5723 materialize: Some(launch(
5724 "opencode",
5725 vec!["import".into(), "{artifact_path}".into()],
5726 )),
5727 requires_materialization: true,
5728 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5729 },
5730 TransferFormat::Pi => HandoffInstructions {
5731 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5732 materialize: None,
5733 requires_materialization: true,
5734 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5735 },
5736 TransferFormat::Grok => HandoffInstructions {
5737 launch: launch(
5738 "grok",
5739 vec!["--resume".into(), "{materialized_session_id}".into()],
5740 ),
5741 materialize: None,
5742 requires_materialization: true,
5743 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(),
5744 },
5745 TransferFormat::Gemini => HandoffInstructions {
5746 launch: launch(
5747 "gemini",
5748 vec!["--session-file".into(), "{artifact_path}".into()],
5749 ),
5750 materialize: None,
5751 requires_materialization: true,
5752 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(),
5753 },
5754 TransferFormat::Goose => HandoffInstructions {
5755 launch: launch(
5756 "goose",
5757 vec![
5758 "session".into(),
5759 "--resume".into(),
5760 "--session-id".into(),
5761 "{imported_session_id}".into(),
5762 ],
5763 ),
5764 materialize: Some(launch(
5765 "goose",
5766 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5767 )),
5768 requires_materialization: true,
5769 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(),
5770 },
5771 }
5772}
5773
5774fn resume_launch(
5775 harness: &str,
5776 session_id: &str,
5777 cwd: &Path,
5778 policy: ResumePolicy,
5779) -> std::result::Result<StructuredLaunch, ServiceError> {
5780 let mut arguments = Vec::new();
5781 let program = match harness {
5782 HarnessId::GROK => {
5783 if matches!(policy, ResumePolicy::Yolo) {
5784 if crate::support::self_sandbox_supported() {
5785 arguments.extend(["--sandbox".into(), "workspace".into()]);
5786 }
5787 arguments.push("--always-approve".into());
5788 }
5789 arguments.extend(["--resume".into(), session_id.into()]);
5790 "grok"
5791 }
5792 HarnessId::CODEX => {
5793 arguments.extend(crate::startup_prompts::startup_arguments(
5794 harness,
5795 Some(cwd),
5796 &[],
5797 matches!(policy, ResumePolicy::Yolo),
5798 ));
5799 arguments.extend(["resume".into(), session_id.into()]);
5800 "codex"
5801 }
5802 HarnessId::CLAUDE_CODE => {
5803 arguments.extend(crate::startup_prompts::startup_arguments(
5804 harness,
5805 Some(cwd),
5806 &[],
5807 matches!(policy, ResumePolicy::Yolo),
5808 ));
5809 arguments.extend(["--resume".into(), session_id.into()]);
5810 "claude"
5811 }
5812 HarnessId::GEMINI => {
5813 arguments.extend(crate::startup_prompts::startup_arguments(
5814 harness,
5815 Some(cwd),
5816 &[],
5817 matches!(policy, ResumePolicy::Yolo),
5818 ));
5819 arguments.extend(["--resume".into(), session_id.into()]);
5820 "gemini"
5821 }
5822 HarnessId::GOOSE => {
5823 arguments.extend([
5824 "session".into(),
5825 "--resume".into(),
5826 "--session-id".into(),
5827 session_id.into(),
5828 ]);
5829 "goose"
5830 }
5831 HarnessId::PI => {
5832 arguments.extend(crate::startup_prompts::startup_arguments(
5833 harness,
5834 Some(cwd),
5835 &[],
5836 matches!(policy, ResumePolicy::Yolo),
5837 ));
5838 arguments.extend(["--session".into(), session_id.into()]);
5839 "pi"
5840 }
5841 HarnessId::OPENCODE => {
5842 arguments.extend(["--session".into(), session_id.into()]);
5843 "opencode"
5844 }
5845 HarnessId::SUPERCODE => {
5846 if matches!(policy, ResumePolicy::Yolo) {
5847 arguments.push("--dangerous".into());
5848 }
5849 arguments.extend(["resume".into(), session_id.into()]);
5850 "supercode"
5851 }
5852 other => {
5853 return Err(ServiceError::InvalidParams(format!(
5854 "no structured resume launch is registered for harness `{other}`"
5855 )))
5856 }
5857 };
5858 Ok(StructuredLaunch {
5859 cwd: cwd.to_path_buf(),
5860 env: if program == "grok" {
5861 crate::support::grok_home_env()
5862 } else {
5863 BTreeMap::new()
5864 },
5865 program: program.into(),
5866 arguments,
5867 })
5868}
5869
5870fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5876 let digest = blake3::hash(address.as_bytes()).to_hex();
5877 let path = std::env::temp_dir().join(format!(
5878 "supercode-openclaw-gateway-token-{}",
5879 &digest.as_str()[..16]
5880 ));
5881 #[cfg(unix)]
5882 {
5883 use std::io::Write;
5884 use std::os::unix::fs::OpenOptionsExt;
5885 let mut file = std::fs::OpenOptions::new()
5886 .write(true)
5887 .create(true)
5888 .truncate(true)
5889 .mode(0o600)
5890 .open(&path)?;
5891 file.write_all(secret.as_bytes())?;
5892 }
5893 #[cfg(not(unix))]
5894 std::fs::write(&path, secret)?;
5895 Ok(path)
5896}
5897
5898fn open_connect_descriptor(
5904 descriptor: &crate::HarnessSupportDescriptor,
5905 home: &Path,
5906) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5907 let Some(connect) = &descriptor.runtime.connect_launch else {
5908 return Err(ServiceError::InvalidParams(format!(
5909 "harness `{}` has no registered connect-mode launch",
5910 descriptor.id.as_str()
5911 )));
5912 };
5913 let resolved = connect
5914 .resolve(home)
5915 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5916 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5917 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5918 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5919 if let Some(token) = resolved.auth {
5920 backend = backend.with_bearer(token);
5921 }
5922 Ok(Box::new(backend))
5923 }
5924 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5925 let mut env = BTreeMap::new();
5936 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5937 if let Some(token) = resolved.auth {
5938 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5939 .map_err(|error| {
5940 ServiceError::UnsupportedAction(format!(
5941 "could not stage the gateway credential for the bridge: {error}"
5942 ))
5943 })?;
5944 arguments.push("--token-file".into());
5945 arguments.push(token_path.to_string_lossy().into_owned());
5946 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5947 }
5948 let program = descriptor
5953 .runtime
5954 .default_launch
5955 .as_ref()
5956 .map(|launch| launch.program.clone())
5957 .unwrap_or_else(|| "openclaw".into());
5958 let launch = RuntimeLaunch {
5959 program,
5960 arguments,
5961 env,
5962 };
5963 Ok(Box::new(
5964 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5965 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5966 ))
5967 }
5968 _ => Err(ServiceError::UnsupportedAction(format!(
5969 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5970 descriptor.id.as_str(),
5971 connect.protocol
5972 ))),
5973 }
5974}
5975
5976fn registry_connect_descriptor(
5979 params: &RuntimeBackendParams,
5980) -> Option<crate::HarnessSupportDescriptor> {
5981 if params.launch.is_some() || params.base_url.is_some() {
5982 return None;
5983 }
5984 harness_support_registry()
5985 .harnesses
5986 .into_iter()
5987 .find(|descriptor| descriptor.id == params.harness)
5988 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5989}
5990
5991fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5992 supercode_interchange::user_home()
5993 .map(std::path::PathBuf::into_os_string)
5994 .map(PathBuf::from)
5995 .ok_or_else(|| {
5996 ServiceError::UnsupportedAction(
5997 "connect-mode launches need HOME to locate the harness config".into(),
5998 )
5999 })
6000}
6001
6002pub const RUNTIME_OPEN_METHODS: &[&str] = &[
6005 "harness.v1.runtimes.start",
6006 "harness.v1.runtimes.resume",
6007 "harness.v1.runtimes.attach",
6008 "harness.v1.runtimes.attach_existing",
6009];
6010
6011pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
6016
6017pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
6023
6024pub const DETACHED_METHODS: &[&str] = &[
6030 "harness.v1.harnesses.list",
6031 "harness.v1.sessions.load",
6032 "harness.v1.harnesses.probe",
6033 "harness.v1.sessions.message",
6034 "harness.v1.sessions.new",
6035 "harness.v1.sessions.reset",
6036 "harness.v1.sessions.archive",
6037 "harness.v1.sessions.delete",
6038];
6039
6040pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
6046
6047pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
6055
6056async fn within_control_deadline<F: std::future::Future>(
6059 method: &str,
6060 call: F,
6061) -> std::result::Result<F::Output, ServiceError> {
6062 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
6063 .await
6064 .map_err(|_| {
6065 ServiceError::Operation(format!(
6066 "`{method}` gave up after {}s: the runtime did not answer",
6067 RUNTIME_CONTROL_DEADLINE.as_secs()
6068 ))
6069 })
6070}
6071
6072pub struct RuntimeOpen {
6076 id: Value,
6077 method: String,
6078 params: Value,
6079}
6080
6081impl RuntimeOpen {
6082 pub async fn open(self) -> OpenedRuntime {
6086 let Self { id, method, params } = self;
6087 let outcome = open_runtime(&method, params).await;
6088 OpenedRuntime { id, outcome }
6089 }
6090}
6091
6092pub struct OpenedRuntime {
6095 id: Value,
6096 outcome: std::result::Result<OpenRuntime, ServiceError>,
6097}
6098
6099pub struct DetachedCall {
6104 id: Value,
6105 method: String,
6106 work: std::result::Result<Work, ServiceError>,
6107}
6108
6109impl DetachedCall {
6110 pub async fn run(self) -> DetachedAnswer {
6113 let Self { id, method, work } = self;
6114 match work {
6115 Ok(Work::Runtime(work)) => {
6120 let (result, returned) = work.run().await;
6121 DetachedAnswer {
6122 response: service_response(id, result),
6123 returned,
6124 }
6125 }
6126 Ok(Work::Free(work)) => {
6127 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
6128 Ok(result) => result,
6129 Err(_) => Err(ServiceError::Operation(format!(
6130 "`{method}` gave up after {}s: the harness it waits on did not answer",
6131 DETACHED_CALL_DEADLINE.as_secs()
6132 ))),
6133 };
6134 DetachedAnswer {
6135 response: service_response(id, result),
6136 returned: None,
6137 }
6138 }
6139 Err(error) => DetachedAnswer {
6140 response: service_response(id, Err(error)),
6141 returned: None,
6142 },
6143 }
6144 }
6145}
6146
6147pub struct DetachedAnswer {
6151 response: Value,
6152 returned: Option<ReturnedRuntime>,
6153}
6154
6155impl DetachedAnswer {
6156 pub fn into_response(self) -> Value {
6159 self.response
6160 }
6161}
6162
6163pub struct ReturnedRuntime {
6166 connection: String,
6167 runtime: Box<dyn RuntimeConnection>,
6168}
6169
6170enum Work {
6173 Free(DetachedWork),
6174 Runtime(RuntimeWork),
6175}
6176
6177enum DetachedWork {
6180 Inventory(InventoryWork),
6184 Message(MessageSessionParams),
6186 Load(LoadSessionParams),
6188 SessionMutation {
6191 verb: crate::SessionVerb,
6192 mutation: crate::SessionMutation,
6193 },
6194}
6195
6196impl DetachedWork {
6197 async fn run(self) -> std::result::Result<Value, ServiceError> {
6198 match self {
6199 Self::Inventory(work) => run_inventory(work).await,
6200 Self::Message(params) => Ok(message_live_session(¶ms).await),
6201 Self::Load(params) => tokio::task::spawn_blocking(move || load_session_door(params))
6202 .await
6203 .map_err(|error| {
6204 ServiceError::Operation(format!("the session read stopped: {error}"))
6205 })?,
6206 Self::SessionMutation { verb, mutation } => {
6207 let outcome = run_session_mutation(verb, &mutation).await?;
6208 serde_json::to_value(outcome)
6209 .map_err(|error| ServiceError::Operation(error.to_string()))
6210 }
6211 }
6212 }
6213}
6214
6215enum RuntimeWork {
6217 Close {
6219 runtime: Box<dyn RuntimeConnection>,
6220 process_group: Option<u32>,
6221 },
6222 LiveCommand {
6225 connection: String,
6226 runtime: Box<dyn RuntimeConnection>,
6227 verb: crate::SessionVerb,
6228 mutation: crate::SessionMutation,
6229 command: &'static str,
6230 session: String,
6231 },
6232}
6233
6234type RuntimeWorkAnswer = (
6237 std::result::Result<Value, ServiceError>,
6238 Option<ReturnedRuntime>,
6239);
6240
6241impl RuntimeWork {
6242 async fn run(self) -> RuntimeWorkAnswer {
6243 match self {
6244 Self::Close {
6245 runtime,
6246 process_group,
6247 } => (close_runtime(runtime, process_group).await, None),
6248 Self::LiveCommand {
6249 connection,
6250 mut runtime,
6251 verb,
6252 mutation,
6253 command,
6254 session,
6255 } => {
6256 let result =
6257 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
6258 (
6259 result,
6260 Some(ReturnedRuntime {
6261 connection,
6262 runtime,
6263 }),
6264 )
6265 }
6266 }
6267 }
6268}
6269
6270async fn close_runtime(
6273 mut runtime: Box<dyn RuntimeConnection>,
6274 process_group: Option<u32>,
6275) -> std::result::Result<Value, ServiceError> {
6276 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
6277 Ok(result) => {
6278 result.map_err(operation)?;
6279 Ok(json!({"closed": true}))
6280 }
6281 Err(deadline) => {
6282 let killed = kill_runtime_process_group(process_group);
6287 drop(runtime);
6288 Ok(json!({
6289 "closed": true,
6290 "killed": killed,
6291 "detail": error_message(deadline),
6292 }))
6293 }
6294 }
6295}
6296
6297fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
6300 mutation
6301 .session
6302 .clone()
6303 .filter(|value| !value.trim().is_empty())
6304 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
6305}
6306
6307async fn type_live_command(
6311 runtime: &mut dyn RuntimeConnection,
6312 verb: crate::SessionVerb,
6313 mutation: &crate::SessionMutation,
6314 command: &str,
6315 session: String,
6316) -> std::result::Result<Value, ServiceError> {
6317 within_control_deadline(
6318 &format!("sessions.{}", verb.as_str()),
6319 runtime.send_input(RuntimeInput {
6320 text: command.to_string(),
6321 image_urls: Vec::new(),
6322 }),
6323 )
6324 .await?
6325 .map_err(operation)?;
6326 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
6327 .map_err(session_control_error)?;
6328 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
6329}
6330
6331enum OpenRuntime {
6334 Hosted {
6337 runtime: Box<dyn RuntimeConnection>,
6338 capabilities: crate::RuntimeCapabilities,
6339 workspace: PathBuf,
6340 fresh: bool,
6342 },
6343 Joined { runtime: Box<dyn RuntimeConnection> },
6346}
6347
6348async fn open_runtime(
6353 method: &str,
6354 params: Value,
6355) -> std::result::Result<OpenRuntime, ServiceError> {
6356 match tokio::time::timeout(
6357 RUNTIME_OPEN_DEADLINE,
6358 open_runtime_unbounded(method, params),
6359 )
6360 .await
6361 {
6362 Ok(result) => result,
6363 Err(_) => Err(ServiceError::Operation(format!(
6364 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
6365 RUNTIME_OPEN_DEADLINE.as_secs()
6366 ))),
6367 }
6368}
6369
6370async fn open_runtime_unbounded(
6371 method: &str,
6372 params: Value,
6373) -> std::result::Result<OpenRuntime, ServiceError> {
6374 match method {
6375 "harness.v1.runtimes.start" => {
6376 let params = decode::<RuntimeStartParams>(params)?;
6377 let backend = runtime_backend(¶ms.backend)?;
6378 let capabilities = backend.capabilities();
6379 let workspace = params.cwd.clone();
6380 let runtime = backend
6381 .start(RuntimeStartRequest {
6382 cwd: params.cwd,
6383 launch: runtime_launch(¶ms.backend),
6384 mcp_servers: params.mcp_servers,
6385 approval_policy: params.approval_policy,
6386 })
6387 .await
6388 .map_err(operation)?;
6389 Ok(OpenRuntime::Hosted {
6390 runtime,
6391 capabilities,
6392 workspace,
6393 fresh: true,
6394 })
6395 }
6396 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
6397 let params = decode::<RuntimeAttachParams>(params)?;
6398 let backend = runtime_backend(¶ms.backend)?;
6399 let capabilities = backend.capabilities();
6400 let workspace = params
6401 .cwd
6402 .clone()
6403 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
6404 let runtime = backend
6405 .attach(RuntimeAttachRequest {
6406 runtime_id: params.runtime_id,
6407 cwd: params.cwd,
6408 launch: runtime_launch(¶ms.backend),
6409 mcp_servers: params.mcp_servers,
6410 approval_policy: params.approval_policy,
6411 })
6412 .await
6413 .map_err(operation)?;
6414 Ok(OpenRuntime::Hosted {
6415 runtime,
6416 capabilities,
6417 workspace,
6418 fresh: false,
6419 })
6420 }
6421 "harness.v1.runtimes.attach_existing" => {
6422 let params = decode::<RuntimeAttachParams>(params)?;
6423 let backend: Box<dyn RuntimeBackend> = match params
6424 .backend
6425 .base_url
6426 .as_deref()
6427 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
6428 {
6429 Some(endpoint) => {
6430 #[cfg(not(feature = "adapter-api"))]
6431 {
6432 let _ = endpoint;
6433 return Err(ServiceError::UnsupportedAction(
6434 "live HTTP attachment adapter is not compiled".into(),
6435 ));
6436 }
6437 #[cfg(feature = "adapter-api")]
6438 {
6439 let workspace = params.cwd.clone().ok_or_else(|| {
6440 ServiceError::InvalidParams(
6441 "Volter Harness live attach requires the project cwd".into(),
6442 )
6443 })?;
6444 let source = LiveRuntimeSource {
6445 harness: params.backend.harness.as_str().to_string(),
6446 session_id: params.runtime_id.clone(),
6447 workspace,
6448 };
6449 let receipt = resolve_live_runtime(&endpoint, &source)
6450 .map_err(|error| ServiceError::Operation(error.to_string()))?;
6451 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
6452 }
6453 }
6454 None => runtime_backend(¶ms.backend)?,
6455 };
6456 let capabilities = backend.capabilities();
6457 if !capabilities.attach_existing_process {
6458 return Err(ServiceError::Operation(format!(
6459 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
6460 backend.harness().as_str()
6461 )));
6462 }
6463 let runtime = backend
6464 .attach_existing(RuntimeAttachRequest {
6465 runtime_id: params.runtime_id,
6466 cwd: params.cwd,
6467 launch: runtime_launch(¶ms.backend),
6468 mcp_servers: params.mcp_servers,
6469 approval_policy: params.approval_policy,
6470 })
6471 .await
6472 .map_err(operation)?;
6473 Ok(OpenRuntime::Joined { runtime })
6474 }
6475 _ => Err(ServiceError::MethodNotFound),
6476 }
6477}
6478
6479fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
6481 match result {
6482 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
6483 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
6484 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
6485 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
6486 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
6487 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
6488 }
6489}
6490
6491fn runtime_backend(
6492 params: &RuntimeBackendParams,
6493) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
6494 if let Some(descriptor) = registry_connect_descriptor(params) {
6495 return open_connect_descriptor(&descriptor, &service_home()?);
6496 }
6497 if params.protocol.as_deref() == Some("acp") {
6498 let launch = params
6499 .launch
6500 .clone()
6501 .or_else(|| {
6502 harness_support_registry()
6503 .harnesses
6504 .into_iter()
6505 .find(|harness| harness.id == params.harness)
6506 .filter(|harness| {
6507 harness.runtime.implementation == ImplementationKind::GenericProtocol
6508 && harness.runtime.protocol.starts_with("acp")
6509 })
6510 .and_then(|harness| harness.runtime.default_launch)
6511 })
6512 .ok_or_else(|| {
6513 ServiceError::InvalidParams(
6514 "an ACP runtime requires `launch` unless the harness has a registered default"
6515 .into(),
6516 )
6517 })?;
6518 let resume_session = harness_support_registry()
6519 .harnesses
6520 .into_iter()
6521 .find(|harness| harness.id == params.harness)
6522 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
6523 return Ok(Box::new(
6524 AcpRuntimeBackend::new(params.harness.clone(), launch)
6525 .with_resume_support(resume_session),
6526 ));
6527 }
6528 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
6529 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
6530 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
6531 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
6532 HarnessId::OPENCODE => match ¶ms.base_url {
6533 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
6534 None => Box::new(OpenCodeRuntimeBackend::new()),
6535 },
6536 harness => {
6537 let descriptor = harness_support_registry()
6538 .harnesses
6539 .into_iter()
6540 .find(|descriptor| descriptor.id.as_str() == harness)
6541 .filter(|descriptor| {
6542 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
6543 && descriptor.runtime.protocol.starts_with("acp")
6544 });
6545 let Some(descriptor) = descriptor else {
6546 return Err(ServiceError::InvalidParams(format!(
6547 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
6548 )));
6549 };
6550 let resume = descriptor.runtime.capabilities.resume_session;
6551 Box::new(
6552 AcpRuntimeBackend::new(
6553 descriptor.id,
6554 descriptor
6555 .runtime
6556 .default_launch
6557 .expect("generic ACP registry entry includes its launch"),
6558 )
6559 .with_resume_support(resume),
6560 )
6561 }
6562 };
6563 Ok(backend)
6564}
6565
6566fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
6567 if let Some(launch) = ¶ms.launch {
6568 return Some(launch.clone());
6569 }
6570 if !matches!(params.policy, RuntimePolicy::Yolo) {
6571 return None;
6572 }
6573 let launch = match params.harness.as_str() {
6574 HarnessId::GROK => RuntimeLaunch {
6575 program: "grok".into(),
6576 arguments: {
6577 let mut arguments: Vec<String> = Vec::new();
6578 if crate::support::self_sandbox_supported() {
6579 arguments.extend(["--sandbox".into(), "workspace".into()]);
6580 }
6581 arguments.extend([
6582 "--always-approve".into(),
6583 "agent".into(),
6584 "--no-leader".into(),
6585 "stdio".into(),
6586 ]);
6587 arguments
6588 },
6589 env: crate::support::grok_env(),
6590 },
6591 HarnessId::CODEX => RuntimeLaunch {
6592 program: "codex".into(),
6593 arguments: vec![
6594 "--dangerously-bypass-approvals-and-sandbox".into(),
6595 "--dangerously-bypass-hook-trust".into(),
6596 "app-server".into(),
6597 ],
6598 env: BTreeMap::new(),
6599 },
6600 HarnessId::CLAUDE_CODE => RuntimeLaunch {
6601 program: "claude".into(),
6602 arguments: vec![
6603 "--dangerously-skip-permissions".into(),
6604 "--print".into(),
6605 "--input-format".into(),
6606 "stream-json".into(),
6607 "--output-format".into(),
6608 "stream-json".into(),
6609 "--verbose".into(),
6610 ],
6611 env: BTreeMap::new(),
6612 },
6613 HarnessId::PI => RuntimeLaunch {
6614 program: "pi".into(),
6615 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
6616 env: BTreeMap::new(),
6617 },
6618 HarnessId::OPENCODE => RuntimeLaunch {
6619 program: "opencode".into(),
6620 arguments: vec!["serve".into()],
6621 env: BTreeMap::new(),
6622 },
6623 HarnessId::GEMINI => RuntimeLaunch {
6624 program: "gemini".into(),
6625 arguments: vec!["--acp".into(), "--yolo".into()],
6626 env: BTreeMap::new(),
6627 },
6628 HarnessId::GOOSE => RuntimeLaunch {
6629 program: "goose".into(),
6630 arguments: vec!["acp".into()],
6631 env: BTreeMap::new(),
6632 },
6633 HarnessId::SUPERCODE => RuntimeLaunch {
6634 program: "supercode".into(),
6635 arguments: vec!["acp".into(), "--dangerous".into()],
6636 env: BTreeMap::new(),
6637 },
6638 _ => return None,
6639 };
6640 Some(launch)
6641}
6642
6643struct IsolatedProbeHome {
6649 launch: RuntimeLaunch,
6650 root: PathBuf,
6651}
6652
6653impl IsolatedProbeHome {
6654 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6655 let root = std::env::temp_dir().join(format!(
6656 "supercode-harness-probe-{harness}-{}",
6657 generated_session_id()
6658 ));
6659 std::fs::create_dir_all(&root)?;
6660 set_private_dir_permissions(&root)?;
6661
6662 if let Some(source_home) = supercode_interchange::user_home()
6663 .map(std::path::PathBuf::into_os_string)
6664 .map(PathBuf::from)
6665 {
6666 for relative in probe_auth_files(harness) {
6667 copy_probe_file(&source_home, &root, relative)?;
6668 }
6669 }
6670 if harness == HarnessId::SUPERCODE {
6675 let config_home = crate::agent::global_instructions_dir();
6676 for file in ["config.toml", "credentials.toml"] {
6677 copy_probe_path(
6678 &config_home.join(file),
6679 &root.join(".config/supercode").join(file),
6680 )?;
6681 }
6682 }
6683 configure_isolated_probe_auth(harness, &root)?;
6684
6685 let root_text = root.to_string_lossy().into_owned();
6686 for (key, value) in [
6687 ("HOME", root_text.clone()),
6688 (
6689 "XDG_CACHE_HOME",
6690 root.join(".cache").to_string_lossy().into_owned(),
6691 ),
6692 (
6693 "XDG_CONFIG_HOME",
6694 root.join(".config").to_string_lossy().into_owned(),
6695 ),
6696 (
6697 "XDG_DATA_HOME",
6698 root.join(".local/share").to_string_lossy().into_owned(),
6699 ),
6700 ] {
6701 launch.env.insert(key.into(), value);
6702 }
6703 let scoped = match harness {
6704 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6705 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6706 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6707 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6708 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6709 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6710 _ => None,
6711 };
6712 if let Some((key, value)) = scoped {
6713 launch
6714 .env
6715 .insert(key.into(), value.to_string_lossy().into_owned());
6716 }
6717 Ok(Self { launch, root })
6718 }
6719
6720 fn cleanup(&self) -> std::io::Result<()> {
6721 match std::fs::remove_dir_all(&self.root) {
6722 Ok(()) => Ok(()),
6723 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6724 Err(error) => Err(error),
6725 }
6726 }
6727}
6728
6729impl Drop for IsolatedProbeHome {
6730 fn drop(&mut self) {
6731 let _ = self.cleanup();
6732 }
6733}
6734
6735fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6736 match harness {
6737 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6738 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6742 HarnessId::CODEX => &[".codex/auth.json"],
6743 HarnessId::GEMINI => &[
6744 ".gemini/google_accounts.json",
6745 ".gemini/oauth_creds.json",
6746 ".gemini/settings.json",
6747 ],
6748 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6749 HarnessId::OPENCODE => &[
6750 ".config/opencode/auth.json",
6751 ".local/share/opencode/auth.json",
6752 ],
6753 HarnessId::PI => &[".pi/agent/auth.json"],
6754 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6759 _ => &[],
6760 }
6761}
6762
6763fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6764 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6765}
6766
6767fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6768 if !source.is_file() {
6769 return Ok(());
6770 }
6771 if let Some(parent) = destination.parent() {
6772 std::fs::create_dir_all(parent)?;
6773 set_private_dir_permissions(parent)?;
6774 }
6775 std::fs::copy(source, destination)?;
6776 set_private_file_permissions(destination)
6777}
6778
6779fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6780 if harness != HarnessId::GEMINI {
6781 return Ok(());
6782 }
6783 let oauth = probe_home.join(".gemini/oauth_creds.json");
6784 if !oauth.is_file() {
6785 return Ok(());
6786 }
6787 let settings_path = probe_home.join(".gemini/settings.json");
6788 let mut settings = std::fs::read_to_string(&settings_path)
6789 .ok()
6790 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6791 .unwrap_or_else(|| json!({}));
6792 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6793 std::fs::write(
6794 &settings_path,
6795 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6796 )?;
6797 set_private_file_permissions(&settings_path)
6798}
6799
6800#[cfg(unix)]
6801fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6802 use std::os::unix::fs::PermissionsExt;
6803 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6804}
6805
6806#[cfg(not(unix))]
6807fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6808 Ok(())
6809}
6810
6811#[cfg(unix)]
6812fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6813 use std::os::unix::fs::PermissionsExt;
6814 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6815}
6816
6817#[cfg(not(unix))]
6818fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6819 Ok(())
6820}
6821
6822fn find_executable(program: &str) -> Option<PathBuf> {
6823 let candidate = PathBuf::from(program);
6824 if candidate.components().count() > 1 {
6825 return candidate.is_file().then_some(candidate);
6826 }
6827 let path = std::env::var_os("PATH")?;
6828 for directory in std::env::split_paths(&path) {
6829 let candidate = directory.join(program);
6830 if candidate.is_file() {
6831 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6832 }
6833 #[cfg(windows)]
6834 {
6835 for extension in ["exe", "cmd", "bat"] {
6836 let candidate = directory.join(format!("{program}.{extension}"));
6837 if candidate.is_file() {
6838 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6839 }
6840 }
6841 }
6842 }
6843 None
6844}
6845
6846async fn executable_version(executable: &Path) -> Option<String> {
6847 let mut command = tokio::process::Command::new(executable);
6848 command
6849 .arg("--version")
6850 .stdin(std::process::Stdio::null())
6851 .stdout(std::process::Stdio::piped())
6852 .stderr(std::process::Stdio::piped())
6853 .kill_on_drop(true);
6854 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6855 .await
6856 .ok()?
6857 .ok()?;
6858 let stdout = String::from_utf8_lossy(&output.stdout);
6859 let stderr = String::from_utf8_lossy(&output.stderr);
6860 stdout
6861 .lines()
6862 .chain(stderr.lines())
6863 .map(str::trim)
6864 .find(|line| !line.is_empty())
6865 .map(|line| truncate_text(line, 200))
6866}
6867
6868pub(crate) fn auth_evidence(harness: &str) -> bool {
6869 let env_names: &[&str] = match harness {
6870 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6871 HarnessId::CODEX => &["OPENAI_API_KEY"],
6872 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6873 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6874 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6875 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6876 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6877 _ => &[],
6878 };
6879 if env_names
6880 .iter()
6881 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6882 {
6883 return true;
6884 }
6885 let Some(home) = supercode_interchange::user_home()
6886 .map(std::path::PathBuf::into_os_string)
6887 .map(PathBuf::from)
6888 else {
6889 return false;
6890 };
6891 let files: Vec<PathBuf> = match harness {
6892 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6893 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6894 HarnessId::OPENCODE => vec![
6895 home.join(".local/share/opencode/auth.json"),
6896 home.join(".config/opencode/auth.json"),
6897 ],
6898 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6899 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6900 HarnessId::GEMINI => vec![
6901 home.join(".gemini/oauth_creds.json"),
6902 home.join(".gemini/google_accounts.json"),
6903 ],
6904 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6905 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6906 _ => Vec::new(),
6907 };
6908 if files.into_iter().any(|path| {
6909 std::fs::metadata(path)
6910 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6911 .unwrap_or(false)
6912 }) {
6913 return true;
6914 }
6915 if harness == HarnessId::CLAUDE_CODE {
6922 return std::fs::read_to_string(home.join(".claude.json"))
6923 .map(|text| text.contains("\"oauthAccount\""))
6924 .unwrap_or(false);
6925 }
6926 false
6927}
6928
6929fn looks_like_auth_error(message: &str) -> bool {
6930 let message = message.to_ascii_lowercase();
6931 [
6932 "auth",
6933 "login",
6934 "sign in",
6935 "sign-in",
6936 "credential",
6937 "unauthorized",
6938 "forbidden",
6939 "token",
6940 ]
6941 .iter()
6942 .any(|needle| message.contains(needle))
6943}
6944
6945fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6946 crate::RuntimeCapabilities {
6947 start_session: false,
6948 resume_session: false,
6949 attach_existing_process: false,
6950 send_input: false,
6951 stream_events: false,
6952 interrupt: false,
6953 steer: false,
6954 respond_to_requests: false,
6955 }
6956}
6957
6958fn truncate_text(text: &str, max_chars: usize) -> String {
6959 let mut chars = text.chars();
6960 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6961 if chars.next().is_some() {
6962 format!("{truncated}…")
6963 } else {
6964 truncated
6965 }
6966}
6967
6968fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6975 match &handle.endpoint {
6976 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6977 crate::RuntimeEndpoint::Http { .. } => None,
6978 }
6979}
6980
6981fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6987 match process_group {
6988 #[cfg(unix)]
6989 Some(pid) => {
6990 crate::lsp::kill_process_group(pid);
6991 true
6992 }
6993 #[cfg(not(unix))]
6994 Some(_) => false,
6995 None => false,
6996 }
6997}
6998
6999fn error_message(error: ServiceError) -> String {
7000 match error {
7001 ServiceError::InvalidParams(message)
7002 | ServiceError::Operation(message)
7003 | ServiceError::UnsupportedAction(message) => message,
7004 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
7005 ServiceError::Sdk(error) => error.to_string(),
7006 }
7007}
7008
7009#[derive(Debug)]
7010enum ServiceError {
7011 InvalidParams(String),
7012 MethodNotFound,
7013 UnsupportedAction(String),
7014 Operation(String),
7015 Sdk(SdkError),
7016}
7017
7018fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
7019 match error {
7020 ServiceError::InvalidParams(message) => {
7021 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
7022 }
7023 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
7024 SdkError::unsupported(operation)
7025 }
7026 ServiceError::Operation(message) => {
7027 let code = if message.contains("already in progress") {
7028 SdkErrorCode::Busy
7029 } else if message.contains("not supported by this runtime") {
7030 SdkErrorCode::UnsupportedAction
7031 } else if message.contains("unknown runtime connection") {
7032 SdkErrorCode::NotFound
7033 } else {
7034 SdkErrorCode::Execution
7035 };
7036 SdkError::new(code, operation, message)
7037 }
7038 ServiceError::Sdk(error) => error,
7039 }
7040}
7041
7042fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
7043 let error_code = error.code();
7044 let code = match error_code {
7045 SdkErrorCode::Unauthenticated => -32030,
7046 SdkErrorCode::Unauthorized => -32031,
7047 SdkErrorCode::ControllerRequired => -32032,
7048 SdkErrorCode::LeaseExpired => -32033,
7049 SdkErrorCode::InvalidArgument => -32602,
7050 SdkErrorCode::NotFound => -32004,
7051 SdkErrorCode::Busy => -32000,
7052 SdkErrorCode::UnsupportedAction => -32020,
7053 SdkErrorCode::Execution => -32002,
7054 SdkErrorCode::Transport => -32003,
7055 };
7056 json!({
7057 "jsonrpc": "2.0",
7058 "id": id,
7059 "error": {
7060 "code": code,
7061 "name": error_code,
7062 "operation": error.operation(),
7063 "message": error.to_string(),
7064 },
7065 })
7066}
7067
7068fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
7069 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
7070}
7071
7072fn operation(error: impl Into<crate::Error>) -> ServiceError {
7073 let error = error.into();
7074 match error {
7075 crate::Error::Sdk(error) => ServiceError::Sdk(error),
7076 error => ServiceError::Operation(error.to_string()),
7077 }
7078}
7079
7080#[derive(Debug, Clone, Deserialize, Default)]
7084#[serde(default)]
7085struct MemoryRequest {
7086 harness: Option<String>,
7088 query: Option<String>,
7090 profile: Option<String>,
7092 session: Option<String>,
7094 full: bool,
7096 regex: bool,
7098 cwd: Option<std::path::PathBuf>,
7100 homes: crate::HarnessHomes,
7102}
7103
7104fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7107 let request = decode::<MemoryRequest>(params)?;
7108 let harness = request
7109 .harness
7110 .clone()
7111 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7112 let to_service = |error: crate::memory::MemoryError| match error {
7113 crate::memory::MemoryError::UnsupportedHarness { .. }
7114 | crate::memory::MemoryError::SessionNotScoped { .. } => {
7115 ServiceError::UnsupportedAction(error.to_string())
7116 }
7117 other => ServiceError::InvalidParams(other.to_string()),
7118 };
7119 match method {
7120 "harness.v1.memory.show" => {
7121 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
7122 harness,
7123 profile: request.profile,
7124 session: request.session,
7125 full: request.full,
7126 cwd: request.cwd,
7127 homes: request.homes,
7128 })
7129 .map_err(to_service)?;
7130 Ok(json!({
7131 "schema": crate::memory::MEMORY_SCHEMA,
7132 "documents": documents,
7133 }))
7134 }
7135 "harness.v1.memory.search" => {
7136 let query = request
7137 .query
7138 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
7139 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
7140 harness,
7141 query,
7142 profile: request.profile,
7143 regex: request.regex,
7144 cwd: request.cwd,
7145 homes: request.homes,
7146 })
7147 .map_err(to_service)?;
7148 Ok(json!({
7149 "schema": crate::memory::MEMORY_SCHEMA,
7150 "matches": matches,
7151 }))
7152 }
7153 _ => Err(ServiceError::MethodNotFound),
7154 }
7155}
7156
7157#[derive(Debug, Clone, Deserialize)]
7161#[serde(default)]
7162struct ProfilesQuery {
7163 harness: Option<String>,
7165 name: Option<String>,
7167 homes: crate::HarnessHomes,
7169}
7170
7171impl Default for ProfilesQuery {
7172 fn default() -> Self {
7173 Self {
7174 harness: None,
7175 name: None,
7176 homes: crate::HarnessHomes::default(),
7177 }
7178 }
7179}
7180
7181fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7184 let query = decode::<ProfilesQuery>(params)?;
7185 let to_service = |error: crate::profiles::ProfileError| match error {
7186 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
7187 ServiceError::UnsupportedAction(error.to_string())
7188 }
7189 crate::profiles::ProfileError::NotFound { .. } => {
7190 ServiceError::InvalidParams(error.to_string())
7191 }
7192 };
7193 match method {
7194 "harness.v1.profiles.list" => {
7195 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
7196 .map_err(to_service)?;
7197 Ok(json!({
7198 "schema": crate::profiles::PROFILES_SCHEMA,
7199 "profiles": profiles,
7200 }))
7201 }
7202 "harness.v1.profiles.get" => {
7203 let harness = query
7204 .harness
7205 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7206 let name = query
7207 .name
7208 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7209 let profile =
7210 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
7211 Ok(json!({
7212 "schema": crate::profiles::PROFILES_SCHEMA,
7213 "profile": profile,
7214 }))
7215 }
7216 _ => Err(ServiceError::MethodNotFound),
7217 }
7218}
7219
7220#[derive(Debug, Clone, Deserialize)]
7224#[serde(default)]
7225struct ChannelsQuery {
7226 harness: Option<String>,
7228 name: Option<String>,
7230 homes: crate::HarnessHomes,
7232}
7233
7234impl Default for ChannelsQuery {
7235 fn default() -> Self {
7236 Self {
7237 harness: None,
7238 name: None,
7239 homes: crate::HarnessHomes::default(),
7240 }
7241 }
7242}
7243
7244#[derive(Debug, Clone, Deserialize)]
7248#[serde(default)]
7249struct RoutesQuery {
7250 harness: Option<String>,
7251 profile: Option<String>,
7253 homes: crate::HarnessHomes,
7254}
7255
7256impl Default for RoutesQuery {
7257 fn default() -> Self {
7258 Self {
7259 harness: None,
7260 profile: None,
7261 homes: crate::HarnessHomes::default(),
7262 }
7263 }
7264}
7265
7266#[derive(Debug, Clone, Deserialize)]
7267#[serde(default)]
7268struct TriggersQuery {
7269 harness: Option<String>,
7270 homes: crate::HarnessHomes,
7271}
7272
7273impl Default for TriggersQuery {
7274 fn default() -> Self {
7275 Self {
7276 harness: None,
7277 homes: crate::HarnessHomes::default(),
7278 }
7279 }
7280}
7281
7282fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
7283 let query = decode::<TriggersQuery>(params)?;
7284 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
7285 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7286 Ok(json!({
7287 "schema": crate::triggers::TRIGGERS_SCHEMA,
7288 "triggers": triggers,
7289 }))
7290}
7291
7292fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
7293 let query = decode::<RoutesQuery>(params)?;
7294 let routes = crate::routes::list_routes(
7295 &query.homes,
7296 query.harness.as_deref(),
7297 query.profile.as_deref(),
7298 )
7299 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7300 Ok(json!({
7301 "schema": crate::routes::ROUTES_SCHEMA,
7302 "routes": routes,
7303 }))
7304}
7305
7306fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7307 let query = decode::<ChannelsQuery>(params)?;
7308 let to_service = |error: crate::channels::ChannelError| match error {
7309 crate::channels::ChannelError::UnsupportedHarness { .. } => {
7310 ServiceError::UnsupportedAction(error.to_string())
7311 }
7312 crate::channels::ChannelError::NotFound { .. } => {
7313 ServiceError::InvalidParams(error.to_string())
7314 }
7315 };
7316 match method {
7317 "harness.v1.channels.list" => {
7318 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
7319 .map_err(to_service)?;
7320 Ok(json!({
7321 "schema": crate::channels::CHANNELS_SCHEMA,
7322 "channels": channels,
7323 }))
7324 }
7325 "harness.v1.channels.status" => {
7326 let harness = query
7327 .harness
7328 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7329 let name = query
7330 .name
7331 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7332 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
7333 .map_err(to_service)?;
7334 Ok(json!({
7335 "schema": crate::channels::CHANNELS_SCHEMA,
7336 "channel": channel,
7337 }))
7338 }
7339 _ => Err(ServiceError::MethodNotFound),
7340 }
7341}
7342
7343fn rpc_error(id: Value, code: i64, message: &str) -> Value {
7344 json!({
7345 "jsonrpc": "2.0",
7346 "id": id,
7347 "error": {"code": code, "message": message},
7348 })
7349}