1use std::collections::{BTreeMap, BTreeSet};
8use std::path::{Path, PathBuf};
9use std::sync::Arc;
10use std::time::Duration;
11
12use serde::{Deserialize, Serialize};
13use serde_json::{json, Value};
14use tokio::sync::Notify;
15
16use crate::reduce;
17use crate::runtime::generated_session_id;
18#[cfg(feature = "adapter-api")]
19use crate::runtime::{HostedHarnessConnection, HostedHarnessRuntime};
20use crate::sdk::{
21 discover_session_page, load_session, load_session_with_fidelity, SdkCapabilities, SdkError,
22 SdkErrorCode, SdkEvent, SdkOperation, SdkRequest, SdkRuntimeEvent, SdkService,
23};
24use crate::Fidelity;
25#[cfg(feature = "adapter-api")]
26use crate::SupercodeHttpRuntimeBackend;
27use crate::{
28 discover_live_runtime, harness_support_registry, AcpRuntimeBackend, ClaudeCodeRuntimeBackend,
29 CodexRuntimeBackend, DiscoveryQuery, HarnessCatalog, HarnessHomes, HarnessId,
30 ImplementationKind, LiveRuntimeEndpoint, LiveRuntimeSource, OpenCodeRuntimeBackend,
31 PiRuntimeBackend, Role, RuntimeAttachRequest, RuntimeBackend, RuntimeConnection, RuntimeInput,
32 RuntimeLaunch, RuntimeStartRequest, Session, SessionDescriptor, SessionFollower, SessionFormat,
33 SessionLocator, SessionSource,
34};
35#[cfg(feature = "adapter-api")]
36use crate::{register_live_runtime, resolve_live_runtime, LiveRuntimeRegistration};
37use supercode_interchange::watch::{bound_session_view, message_json, normalized_session_json};
38
39pub const HARNESS_SERVICE_METHODS: &[&str] = &[
42 "harness.v1.support.report",
43 "harness.v1.harnesses.list",
44 "harness.v1.harnesses.probe",
45 "harness.v1.harnesses.settings",
46 "harness.v1.harnesses.configure",
47 "harness.v1.harnesses.auth.methods",
48 "harness.v1.harnesses.auth.begin",
49 "harness.v1.harnesses.auth.verify",
50 "harness.v1.sessions.discover",
51 "harness.v1.sessions.load",
52 "harness.v1.sessions.follow",
53 "harness.v1.sessions.unfollow",
54 "harness.v1.sessions.activity.subscribe",
55 "harness.v1.sessions.activity.unsubscribe",
56 "harness.v1.sessions.activity_under",
57 "harness.v1.sessions.index.subscribe",
58 "harness.v1.sessions.index.resize",
59 "harness.v1.sessions.index.unsubscribe",
60 "harness.v1.sessions.message",
61 "harness.v1.sessions.inbox",
62 "harness.v1.sessions.import",
63 "harness.v1.sessions.export",
64 "harness.v1.sessions.translate",
65 "harness.v1.sessions.reduce",
66 "harness.v1.sessions.branch",
67 "harness.v1.sessions.handoff",
68 "harness.v1.sessions.materialize",
69 "harness.v1.sessions.resume_instructions",
70 "harness.v1.skills.list",
71 "harness.v1.skills.install",
72 "harness.v1.skills.remove",
73 "harness.v1.memory.show",
74 "harness.v1.memory.search",
75 "harness.v1.jobs.list",
76 "harness.v1.jobs.get",
77 "harness.v1.jobs.create",
78 "harness.v1.jobs.update",
79 "harness.v1.jobs.pause",
80 "harness.v1.jobs.resume",
81 "harness.v1.jobs.run",
82 "harness.v1.jobs.delete",
83 "harness.v1.jobs.notepad",
84 "harness.v1.jobs.notepad_set",
85 "harness.v1.jobs.notepad_delete",
86 "harness.v1.sessions.new",
87 "harness.v1.sessions.reset",
88 "harness.v1.sessions.archive",
89 "harness.v1.sessions.delete",
90 "harness.v1.runs.list",
91 "harness.v1.runs.get",
92 "harness.v1.approvals.list",
93 "harness.v1.approvals.resolve",
94 "harness.v1.runtimes.capabilities",
95 "harness.v1.runtimes.start",
96 "harness.v1.runtimes.resume",
97 "harness.v1.runtimes.attach_existing",
98 "harness.v1.runtimes.attach",
99 "harness.v1.runtimes.send_input",
100 "harness.v1.runtimes.interrupt",
101 "harness.v1.runtimes.steer",
102 "harness.v1.runtimes.respond",
103 "harness.v1.runtimes.terminal_instructions",
104 "harness.v1.runtimes.acquire_control",
105 "harness.v1.runtimes.heartbeat",
106 "harness.v1.runtimes.detach",
107 "harness.v1.runtimes.close",
108 "harness.v1.profiles.list",
109 "harness.v1.profiles.get",
110 "harness.v1.profiles.create",
111 "harness.v1.profiles.delete",
112 "harness.v1.channels.list",
113 "harness.v1.routes.list",
114 "harness.v1.triggers.list",
115 "harness.v1.channels.status",
116 "harness.v1.orchestration.load",
117 "harness.v1.orchestration.save",
118 "harness.v1.orchestration.compile",
119 "harness.v1.orchestration.decompile",
120 "harness.v1.orchestration.import",
121 "harness.v1.orchestration.export",
122 "harness.v1.workflow.load",
123];
124
125pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
127pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
129pub const SESSION_ACTIVITY_EVENT_METHOD: &str = "harness.v1.sessions.activity_event";
131pub const SESSION_INDEX_EVENT_METHOD: &str = "harness.v1.sessions.index_event";
133pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
135
136pub struct HarnessSessionService {
139 catalog: HarnessCatalog,
140 followers: BTreeMap<String, SessionFollower>,
141 followed_sources: BTreeMap<String, FollowedSource>,
142 activity_subscriptions: BTreeMap<String, ActivitySubscription>,
143 index_subscriptions: BTreeMap<String, crate::session_index::SessionIndexSubscription>,
144 index_notifier: Arc<Notify>,
145 #[cfg(feature = "adapter-api")]
146 activity_monitor: crate::session_activity::SessionActivityMonitor,
147 next_subscription: u64,
148 runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
149 runtimes_in_flight: BTreeSet<String>,
154 terminal_launches: BTreeMap<String, StructuredLaunch>,
155 runtime_sequences: BTreeMap<String, u64>,
156 next_runtime: u64,
157 reduction_store_root: Option<PathBuf>,
158 approvals: crate::approvals::ApprovalRegistry,
162 subagent_approvals: Option<Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>>,
165}
166
167impl Default for HarnessSessionService {
168 fn default() -> Self {
169 Self::new()
170 }
171}
172
173impl HarnessSessionService {
174 pub fn new() -> Self {
176 Self {
177 catalog: HarnessCatalog::new(),
178 followers: BTreeMap::new(),
179 followed_sources: BTreeMap::new(),
180 activity_subscriptions: BTreeMap::new(),
181 index_subscriptions: BTreeMap::new(),
182 index_notifier: Arc::new(Notify::new()),
183 #[cfg(feature = "adapter-api")]
184 activity_monitor: Default::default(),
185 next_subscription: 1,
186 runtimes: BTreeMap::new(),
187 runtimes_in_flight: BTreeSet::new(),
188 terminal_launches: BTreeMap::new(),
189 runtime_sequences: BTreeMap::new(),
190 next_runtime: 1,
191 reduction_store_root: None,
192 approvals: crate::approvals::ApprovalRegistry::new(),
193 subagent_approvals: None,
194 }
195 }
196
197 pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
202 self.reduction_store_root = Some(root.into());
203 self
204 }
205
206 pub fn observe_subagent_approvals(
214 &mut self,
215 queue: Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>,
216 ) {
217 self.subagent_approvals = Some(queue);
218 }
219
220 pub fn approvals(&self, query: &crate::approvals::ApprovalsQuery) -> Vec<crate::ApprovalRow> {
227 let now = crate::approvals::now_ms();
228 let mut rows = self.approvals.rows(now);
229 if let Some(queue) = self.subagent_approvals.as_ref() {
230 let queued = queue
231 .lock()
232 .unwrap_or_else(std::sync::PoisonError::into_inner)
233 .clone();
234 rows.extend(crate::approvals::subagent_rows(&queued, now));
235 }
236 rows.retain(|row| query.matches(row));
237 rows.sort_by(|left, right| {
238 left.requested_at_ms
239 .cmp(&right.requested_at_ms)
240 .then_with(|| left.id.cmp(&right.id))
241 });
242 rows
243 }
244
245 async fn approvals_resolve(
255 &mut self,
256 params: Value,
257 ) -> std::result::Result<Value, ServiceError> {
258 let params = decode::<crate::approvals::ApprovalsResolveParams>(params)?;
259 if params.id.trim().is_empty() {
260 return Err(ServiceError::InvalidParams(
261 "approvals resolve requires the `id` of a listed approval row".into(),
262 ));
263 }
264 let choice = match (params.decision, params.option_id.as_deref()) {
265 (Some(_), Some(_)) => {
266 return Err(ServiceError::InvalidParams(
267 "approvals resolve takes either `decision` or `option_id`, not both".into(),
268 ))
269 }
270 (Some(decision), None) => crate::approvals::ApprovalChoice::Decision(decision),
271 (None, Some(option)) => crate::approvals::ApprovalChoice::Option(option.to_string()),
272 (None, None) => {
273 return Err(ServiceError::InvalidParams(format!(
274 "approvals resolve requires `decision` ({}) or an explicit `option_id`",
275 crate::approvals::ApprovalDecision::ALL
276 .map(|decision| decision.as_str())
277 .join(" | "),
278 )))
279 }
280 };
281 let resolution = self
282 .approvals
283 .resolution(¶ms.id, &choice)
284 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?;
285 self.runtime_call(
289 "harness.v1.runtimes.respond",
290 json!({
291 "connection": resolution.connection,
292 "request_id": resolution.request_id,
293 "response": resolution.response,
294 }),
295 )
296 .await?;
297 Ok(json!({
298 "id": params.id,
299 "decision": params.decision.map(|decision| decision.as_str()),
300 "option_id": resolution.option_id,
301 "resolved": true,
302 }))
303 }
304
305 #[cfg(feature = "adapter-api")]
308 pub fn session_index_notifier(&self) -> Arc<Notify> {
309 Arc::clone(&self.index_notifier)
310 }
311
312 #[cfg(feature = "adapter-api")]
314 pub fn handle(&mut self, request: Value) -> Value {
315 let id = request.get("id").cloned().unwrap_or(Value::Null);
316 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
317 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
318 }
319 let Some(method) = request.get("method").and_then(Value::as_str) else {
320 return rpc_error(id, -32600, "request is missing `method`");
321 };
322 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
323 match self.call(method, params) {
324 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
325 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
326 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
327 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
328 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
329 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
330 }
331 }
332
333 #[cfg(feature = "adapter-api")]
336 pub async fn handle_async(&mut self, request: Value) -> Value {
337 let method = request
338 .get("method")
339 .and_then(Value::as_str)
340 .unwrap_or_default();
341 if matches!(
342 method,
343 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
344 ) {
345 let id = request.get("id").cloned().unwrap_or(Value::Null);
346 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
347 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
348 }
349 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
350 return match self.inventory_call(method, params).await {
351 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
352 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
353 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
354 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
355 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
356 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
357 };
358 }
359 if matches!(
360 method,
361 "harness.v1.harnesses.auth.methods"
362 | "harness.v1.harnesses.auth.begin"
363 | "harness.v1.harnesses.auth.verify"
364 ) {
365 let id = request.get("id").cloned().unwrap_or(Value::Null);
366 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
367 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
368 }
369 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
370 return match self.harness_authentication_call(method, params).await {
371 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
372 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
373 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
374 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
375 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
376 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
377 };
378 }
379 if matches!(
385 method,
386 "harness.v1.sessions.new"
387 | "harness.v1.sessions.reset"
388 | "harness.v1.sessions.archive"
389 | "harness.v1.sessions.delete"
390 ) {
391 let id = request.get("id").cloned().unwrap_or(Value::Null);
392 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
393 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
394 }
395 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
396 let verb = match method {
397 "harness.v1.sessions.new" => crate::SessionVerb::New,
398 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
399 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
400 _ => crate::SessionVerb::Delete,
401 };
402 return match self.mutate_session(verb, params).await {
403 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
404 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
405 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
406 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
407 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
408 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
409 };
410 }
411 if method == "harness.v1.sessions.message" {
412 let id = request.get("id").cloned().unwrap_or(Value::Null);
413 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
414 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
415 }
416 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
417 return match self.message_call(params).await {
418 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
419 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
420 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
421 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
422 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
423 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
424 };
425 }
426 if matches!(
427 method,
428 "harness.v1.harnesses.settings" | "harness.v1.harnesses.configure"
429 ) {
430 let id = request.get("id").cloned().unwrap_or(Value::Null);
431 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
432 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
433 }
434 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
435 return match self.harness_settings_call(method, params) {
436 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
437 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
438 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
439 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
440 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
441 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
442 };
443 }
444 if method == "harness.v1.sessions.activity_under" {
445 let id = request.get("id").cloned().unwrap_or(Value::Null);
446 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
447 return match activity_under_call(params).await {
448 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
449 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
450 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
451 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
452 Err(_) => rpc_error(id, -32000, "sessions.activity_under failed"),
453 };
454 }
455 if method == "harness.v1.sessions.activity.subscribe" {
456 let id = request.get("id").cloned().unwrap_or(Value::Null);
457 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
458 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
459 }
460 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
461 return match self.subscribe_session_activity(params).await {
462 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
463 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
464 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
465 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
466 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
467 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
468 };
469 }
470 if let Some(operation) = SdkOperation::from_method(method) {
471 let id = request.get("id").cloned().unwrap_or(Value::Null);
472 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
473 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
474 }
475 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
476 return match self.execute(SdkRequest { operation, params }).await {
477 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
478 Err(error) => sdk_rpc_error(id, &error),
479 };
480 }
481 if !method.starts_with("harness.v1.runtimes.") {
482 return self.handle(request);
483 }
484 let id = request.get("id").cloned().unwrap_or(Value::Null);
485 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
486 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
487 }
488 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
489 match self.runtime_call(method, params).await {
490 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
491 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
492 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
493 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
494 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
495 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
496 }
497 }
498
499 #[cfg(feature = "adapter-api")]
502 pub fn poll(&mut self) -> Vec<Value> {
503 let mut notifications = Vec::new();
504 for (subscription, follower) in &mut self.followers {
505 match follower.poll() {
506 Ok(Some(event)) => notifications.push(json!({
507 "jsonrpc": "2.0",
508 "method": SESSION_EVENT_METHOD,
509 "params": {
510 "subscription": subscription,
511 "event": event.to_json(),
512 }
513 })),
514 Ok(None) => {}
515 Err(error) => notifications.push(json!({
516 "jsonrpc": "2.0",
517 "method": SESSION_EVENT_METHOD,
518 "params": {
519 "subscription": subscription,
520 "event": {
521 "type": "watch_error",
522 "recoverable": true,
523 "message": error.to_string(),
524 },
525 }
526 })),
527 }
528 }
529 notifications
530 }
531
532 #[cfg(feature = "adapter-api")]
543 pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
544 let registry = crate::LocalRuntimeRegistry::new();
545 let authorization = crate::RuntimeAuthorization::observer();
546 let mut notifications = Vec::new();
547 for (subscription, source) in &mut self.followed_sources {
548 let state = match registry
549 .source_state(&source.harness, &source.session_id, &authorization)
550 .await
551 {
552 Ok(Some(state)) => state,
553 Ok(None) => crate::RuntimeRegistryState::Persisted,
554 Err(_) => continue,
556 };
557 if source.reported.as_deref() == Some(state.as_str()) {
558 continue;
559 }
560 source.reported = Some(state.as_str().to_string());
561 notifications.push(json!({
562 "jsonrpc": "2.0",
563 "method": SESSION_EVENT_METHOD,
564 "params": {
565 "subscription": subscription,
566 "event": {"type": "runtime_state", "state": state.as_str()},
567 },
568 }));
569 }
570 notifications
571 }
572
573 #[cfg(feature = "adapter-api")]
577 pub async fn poll_session_activities(&mut self) -> Vec<Value> {
578 let subscriptions = self
579 .activity_subscriptions
580 .iter()
581 .map(|(id, subscription)| {
582 (
583 id.clone(),
584 subscription.locators.clone(),
585 subscription.homes.clone(),
586 )
587 })
588 .collect::<Vec<_>>();
589 let mut notifications = Vec::new();
590 for (subscription_id, locators, homes) in subscriptions {
591 let Ok(activities) = self.activity_monitor.resolve(&locators, &homes).await else {
592 continue;
595 };
596 let Some(subscription) = self.activity_subscriptions.get_mut(&subscription_id) else {
597 continue;
598 };
599 let mut changed = Vec::new();
600 for activity in activities {
601 let key = activity.key();
602 if subscription
603 .reported
604 .get(&key)
605 .is_some_and(|previous| previous.same_state(&activity))
606 {
607 continue;
608 }
609 subscription.reported.insert(key, activity.clone());
610 changed.push(activity);
611 }
612 if !changed.is_empty() {
613 notifications.push(json!({
614 "jsonrpc": "2.0",
615 "method": SESSION_ACTIVITY_EVENT_METHOD,
616 "params": {
617 "subscription": subscription_id,
618 "activities": changed,
619 },
620 }));
621 }
622 }
623 notifications
624 }
625
626 #[cfg(feature = "adapter-api")]
630 pub fn poll_session_indexes(&mut self) -> Vec<Value> {
631 let mut notifications = Vec::new();
632 for (subscription, index) in &mut self.index_subscriptions {
633 let homes = index.homes().clone();
634 match index.poll() {
635 Ok(Some(delta)) => match live_index_changes(delta.changes, &homes) {
636 Ok(changes) => notifications.push(json!({
637 "jsonrpc": "2.0",
638 "method": SESSION_INDEX_EVENT_METHOD,
639 "params": {
640 "subscription": subscription,
641 "revision": delta.revision,
642 "changes": changes,
643 },
644 })),
645 Err(error) => notifications.push(json!({
646 "jsonrpc": "2.0",
647 "method": SESSION_INDEX_EVENT_METHOD,
648 "params": {
649 "subscription": subscription,
650 "error": {"recoverable": true, "message": error_message(error)},
651 },
652 })),
653 },
654 Ok(None) => {}
655 Err(error) => notifications.push(json!({
656 "jsonrpc": "2.0",
657 "method": SESSION_INDEX_EVENT_METHOD,
658 "params": {
659 "subscription": subscription,
660 "error": {"recoverable": true, "message": error},
661 },
662 })),
663 }
664 }
665 notifications
666 }
667
668 #[cfg(feature = "adapter-api")]
669 async fn subscribe_session_activity(
670 &mut self,
671 params: Value,
672 ) -> std::result::Result<Value, ServiceError> {
673 let params = decode::<ActivitySubscribeParams>(params)?;
674 if params.locators.is_empty() {
675 return Err(ServiceError::InvalidParams(
676 "sessions.activity.subscribe requires at least one locator".into(),
677 ));
678 }
679 if params.locators.len() > 2_048 {
680 return Err(ServiceError::InvalidParams(
681 "sessions.activity.subscribe accepts at most 2048 locators".into(),
682 ));
683 }
684 let initial = self
685 .activity_monitor
686 .resolve(¶ms.locators, ¶ms.homes)
687 .await
688 .map_err(ServiceError::Sdk)?;
689 let subscription = format!("activity-sub-{}", self.next_subscription);
690 self.next_subscription += 1;
691 let reported = initial
692 .iter()
693 .cloned()
694 .map(|activity| (activity.key(), activity))
695 .collect();
696 self.activity_subscriptions.insert(
697 subscription.clone(),
698 ActivitySubscription {
699 locators: params.locators,
700 homes: params.homes,
701 reported,
702 },
703 );
704 Ok(json!({"subscription": subscription, "initial": initial}))
705 }
706
707 #[cfg(feature = "adapter-api")]
709 pub async fn poll_runtimes(&mut self) -> Vec<Value> {
710 self.poll_sdk_events()
711 .await
712 .into_iter()
713 .map(|(connection, runtime_event)| {
714 json!({
715 "jsonrpc": "2.0",
716 "method": RUNTIME_EVENT_METHOD,
717 "params": {
718 "connection": connection,
719 "session_id": runtime_event.session_id,
720 "sequence": runtime_event.event.sequence,
721 "event": {
722 "kind": runtime_event.event.kind,
723 "payload": runtime_event.event.payload,
724 },
725 },
726 })
727 })
728 .collect()
729 }
730
731 async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
732 let mut events = Vec::new();
733 let mut closed = Vec::new();
734 let now_ms = crate::approvals::now_ms();
735 for (connection, runtime) in &mut self.runtimes {
736 let session_id = runtime.handle().runtime_id.clone();
737 let harness = runtime.handle().harness.clone();
738 for _ in 0..256 {
743 match tokio::time::timeout(Duration::ZERO, runtime.next_event()).await {
744 Ok(Ok(Some(event))) => {
745 let terminal = event.kind == "transport_closed";
746 self.approvals
750 .observe(connection, &harness, &session_id, &event, now_ms);
751 let next_sequence = self
752 .runtime_sequences
753 .entry(session_id.clone())
754 .or_insert(0);
755 let sequence = event.sequence.unwrap_or_else(|| {
756 *next_sequence = next_sequence.saturating_add(1);
757 *next_sequence
758 });
759 *next_sequence = (*next_sequence).max(sequence);
760 events.push((
761 connection.clone(),
762 SdkRuntimeEvent {
763 session_id: session_id.clone(),
764 event: SdkEvent {
765 sequence,
766 kind: event.kind,
767 payload: event.payload,
768 },
769 },
770 ));
771 if terminal {
772 closed.push(connection.clone());
773 break;
774 }
775 }
776 Ok(Ok(None)) => {
777 let sequence = self
778 .runtime_sequences
779 .entry(session_id.clone())
780 .or_insert(0);
781 *sequence = sequence.saturating_add(1);
782 events.push((
783 connection.clone(),
784 SdkRuntimeEvent {
785 session_id,
786 event: SdkEvent {
787 sequence: *sequence,
788 kind: "transport_closed".into(),
789 payload: json!({"message": "Harness runtime transport closed."}),
790 },
791 },
792 ));
793 closed.push(connection.clone());
794 break;
795 }
796 Err(_) => break,
797 Ok(Err(error)) => {
798 let sequence = self
799 .runtime_sequences
800 .entry(session_id.clone())
801 .or_insert(0);
802 *sequence = sequence.saturating_add(1);
803 events.push((
804 connection.clone(),
805 SdkRuntimeEvent {
806 session_id,
807 event: SdkEvent {
808 sequence: *sequence,
809 kind: "transport_error".into(),
810 payload: json!({"message": error.to_string(), "terminal": true}),
811 },
812 },
813 ));
814 closed.push(connection.clone());
815 break;
816 }
817 }
818 }
819 }
820 for connection in closed {
821 if let Some(runtime) = self.runtimes.remove(&connection) {
822 self.runtime_sequences.remove(&runtime.handle().runtime_id);
823 }
824 self.terminal_launches.remove(&connection);
825 self.approvals.forget(&connection);
828 }
829 events
830 }
831
832 fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
833 match method {
834 "harness.v1.capabilities" => Ok(json!({
835 "version": HARNESS_SERVICE_VERSION,
836 "sdk": self.capabilities(),
837 "methods": HARNESS_SERVICE_METHODS,
838 "notifications": [
839 SESSION_EVENT_METHOD,
840 SESSION_ACTIVITY_EVENT_METHOD,
841 SESSION_INDEX_EVENT_METHOD,
842 RUNTIME_EVENT_METHOD
843 ],
844 "harnesses": harness_support_registry()
845 .harnesses
846 .into_iter()
847 .map(|harness| harness.id)
848 .collect::<Vec<_>>(),
849 })),
850 "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
851 .map_err(|error| ServiceError::Operation(error.to_string())),
852 "harness.v1.profiles.list" | "harness.v1.profiles.get" => profiles_call(method, params),
853 "harness.v1.profiles.create" => {
859 mutate_profile(crate::profiles_control::ProfileVerb::Create, params)
860 }
861 "harness.v1.profiles.delete" => {
862 mutate_profile(crate::profiles_control::ProfileVerb::Delete, params)
863 }
864 "harness.v1.channels.list" | "harness.v1.channels.status" => {
865 channels_call(method, params)
866 }
867 "harness.v1.routes.list" => routes_call(params),
870 "harness.v1.triggers.list" => triggers_call(params),
872 "harness.v1.workflow.load" => {
882 let params = decode::<WorkflowLoadParams>(params)?;
883 let read =
884 crate::workflow_doors::load(params.from, ¶ms.home).map_err(operation)?;
885 serde_json::to_value(read)
886 .map_err(|error| ServiceError::Operation(error.to_string()))
887 }
888 "harness.v1.orchestration.load" => {
889 let params = decode::<OrchestrationLoadParams>(params)?;
890 let read = crate::orchestration_doors::load(¶ms.root, params.flavor)
891 .map_err(operation)?;
892 serde_json::to_value(read)
893 .map_err(|error| ServiceError::Operation(error.to_string()))
894 }
895 "harness.v1.orchestration.save" => {
896 let params = decode::<OrchestrationSaveParams>(params)?;
897 let saved = crate::orchestration_doors::save(
898 ¶ms.root,
899 params.orchestration,
900 params.vault,
901 )
902 .map_err(operation)?;
903 serde_json::to_value(saved)
904 .map_err(|error| ServiceError::Operation(error.to_string()))
905 }
906 "harness.v1.orchestration.compile" => {
907 let params = decode::<OrchestrationCompileParams>(params)?;
908 let read = crate::orchestration_doors::compile(params.from, ¶ms.home)
909 .map_err(operation)?;
910 serde_json::to_value(read)
911 .map_err(|error| ServiceError::Operation(error.to_string()))
912 }
913 "harness.v1.orchestration.decompile" => {
914 let params = decode::<OrchestrationDecompileParams>(params)?;
915 let report = crate::orchestration_doors::decompile(
916 params.to,
917 params.orchestration,
918 ¶ms.source,
919 params.source_flavor,
920 ¶ms.dest,
921 params.vault,
922 )
923 .map_err(operation)?;
924 serde_json::to_value(report)
925 .map_err(|error| ServiceError::Operation(error.to_string()))
926 }
927 "harness.v1.orchestration.import" => {
932 let params = decode::<OrchestrationImportParams>(params)?;
933 let imported =
934 crate::orchestration_doors::import(params.from, ¶ms.home, ¶ms.into)
935 .map_err(operation)?;
936 serde_json::to_value(imported)
937 .map_err(|error| ServiceError::Operation(error.to_string()))
938 }
939 "harness.v1.orchestration.export" => {
940 let params = decode::<OrchestrationExportParams>(params)?;
941 let report =
942 crate::orchestration_doors::export(params.to, ¶ms.root, ¶ms.dest)
943 .map_err(operation)?;
944 serde_json::to_value(report)
945 .map_err(|error| ServiceError::Operation(error.to_string()))
946 }
947 "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
953 "harness.v1.skills.list" => {
958 let query = decode::<crate::skills::SkillsQuery>(params)?;
959 if let Some(harness) = query.harness.as_deref() {
960 if !crate::skills::SKILL_HARNESSES.contains(&harness) {
961 return Err(ServiceError::UnsupportedAction(format!(
962 "`{harness}` has no skills root Volter Harness reads"
963 )));
964 }
965 }
966 serde_json::to_value(crate::skills::list_skills(&query))
967 .map_err(|error| ServiceError::Operation(error.to_string()))
968 }
969 "harness.v1.skills.install" => {
977 mutate_skill(crate::skills_control::SkillVerb::Install, params)
978 }
979 "harness.v1.skills.remove" => {
980 mutate_skill(crate::skills_control::SkillVerb::Remove, params)
981 }
982 "harness.v1.approvals.list" => {
990 let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
991 if let Some(harness) = query.harness.as_deref() {
992 if !crate::approvals::lists_approvals(harness) {
993 return Err(ServiceError::UnsupportedAction(format!(
994 "`{harness}` has no runtime door that carries an approval request"
995 )));
996 }
997 }
998 serde_json::to_value(self.approvals(&query))
999 .map_err(|error| ServiceError::Operation(error.to_string()))
1000 }
1001 "harness.v1.sessions.discover" => {
1002 let query = decode::<DiscoveryQuery>(params)?;
1003 let mut page = discover_session_page(&query).map_err(operation)?;
1004 let doors = crate::mail_route::LiveSessions::read(&query.homes);
1009 if let Some(text) = query
1013 .query
1014 .as_deref()
1015 .map(str::trim)
1016 .filter(|text| !text.is_empty())
1017 {
1018 let needle = text.to_lowercase();
1019 for live in doors.all() {
1020 let name = live
1021 .name
1022 .split('@')
1023 .next()
1024 .unwrap_or(&live.name)
1025 .to_lowercase();
1026 if !name.contains(&needle)
1027 || page.sessions.iter().any(|session| {
1028 session.locator.session_id == live.address.session_id
1029 })
1030 {
1031 continue;
1032 }
1033 if query
1034 .limit
1035 .is_some_and(|limit| page.sessions.len() >= limit)
1036 {
1037 break;
1038 }
1039 let mut by_id = query.clone();
1041 by_id.query = Some(live.address.session_id.clone());
1042 by_id.cursor = None;
1043 if let Some(cwd) = &live.cwd {
1044 by_id.workspace = Some(cwd.clone());
1045 by_id.workspace_subtree = false;
1046 by_id.workspace_family = None;
1047 }
1048 if let Ok(found) = discover_session_page(&by_id) {
1049 page.sessions
1050 .extend(found.sessions.into_iter().filter(|session| {
1051 session.locator.session_id == live.address.session_id
1052 }));
1053 }
1054 }
1055 }
1056 let activities = crate::session_activity::resolve_stock_session_activities(
1057 &page
1058 .sessions
1059 .iter()
1060 .map(|session| session.locator.clone())
1061 .collect::<Vec<_>>(),
1062 &query.homes,
1063 )
1064 .into_iter()
1065 .map(|activity| (activity.key(), activity))
1066 .collect::<BTreeMap<_, _>>();
1067 let sessions = page
1068 .sessions
1069 .into_iter()
1070 .map(|session| {
1071 let mut value = live_descriptor_value(&session, &doors)?;
1072 let activity_key = (
1073 session.locator.harness.as_str().to_string(),
1074 session.locator.session_id.clone(),
1075 );
1076 if let Some(activity) = activities.get(&activity_key) {
1077 value["activity"] = serde_json::to_value(activity)
1078 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1079 if let Some(status) = legacy_live_status(activity) {
1080 value["live_status"] = json!(status);
1081 }
1082 }
1083 Ok(value)
1084 })
1085 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1086 let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1087 if query.search_previews {
1090 result["receipt"] = serde_json::to_value(page.receipt)
1091 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1092 }
1093 Ok(result)
1094 }
1095 "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1096 "harness.v1.sessions.load" => {
1097 let params = decode::<LoadSessionParams>(params)?;
1098 if let Some(options) = ¶ms.options {
1099 options.validate()?;
1100 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
1101 return Ok(result);
1102 }
1103 return load_session(¶ms.read.locator)
1104 .map(|session| projected_session_result(&session, options))
1105 .map_err(operation);
1106 }
1107 let mut session = if params.read.display_history() {
1108 self.catalog
1109 .load_display_view(
1110 ¶ms.read.locator,
1111 params.read.read_fidelity(),
1112 params.read.tail_messages().unwrap_or(500),
1113 )
1114 .map_err(crate::Error::from)
1115 } else if params.read.include_subagents() {
1116 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
1117 } else {
1118 self.catalog
1119 .load_parent_with_fidelity(
1120 ¶ms.read.locator,
1121 params.read.read_fidelity(),
1122 )
1123 .map_err(crate::Error::from)
1124 }
1125 .map_err(operation)?;
1126 params.read.bound_session(&mut session);
1127 Ok(json!({"session": normalized_session_json(&session)}))
1128 }
1129 "harness.v1.sessions.follow" => {
1130 let params = decode::<LocatorParams>(params)?;
1131 let mut follower = self
1132 .catalog
1133 .follow_read_view(
1134 ¶ms.locator,
1135 params.read_fidelity(),
1136 params.include_subagents(),
1137 params.tail_messages(),
1138 params.max_message_chars(),
1139 params.display_history(),
1140 )
1141 .map_err(operation)?;
1142 let initial = follower
1143 .poll()
1144 .map_err(operation)?
1145 .map(|event| event.to_json());
1146 let subscription = format!("sub-{}", self.next_subscription);
1147 self.next_subscription += 1;
1148 self.followers.insert(subscription.clone(), follower);
1149 self.followed_sources.insert(
1150 subscription.clone(),
1151 FollowedSource {
1152 harness: params.locator.harness.as_str().to_string(),
1153 session_id: params.locator.session_id.clone(),
1154 reported: None,
1155 },
1156 );
1157 Ok(json!({"subscription": subscription, "initial": initial}))
1158 }
1159 "harness.v1.sessions.unfollow" => {
1160 let params = decode::<UnfollowParams>(params)?;
1161 self.followed_sources.remove(¶ms.subscription);
1162 Ok(json!({
1163 "removed": self.followers.remove(¶ms.subscription).is_some()
1164 }))
1165 }
1166 "harness.v1.sessions.activity.unsubscribe" => {
1167 let params = decode::<UnfollowParams>(params)?;
1168 Ok(json!({
1169 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1170 }))
1171 }
1172 "harness.v1.sessions.index.subscribe" => {
1173 let query = decode::<DiscoveryQuery>(params)?;
1174 crate::session_index::validate_query(&query)
1175 .map_err(ServiceError::InvalidParams)?;
1176 let homes = query.homes.clone();
1177 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1178 query,
1179 Arc::clone(&self.index_notifier),
1180 )
1181 .map_err(ServiceError::Operation)?;
1182 let doors = crate::mail_route::LiveSessions::read(&homes);
1183 let initial = initial
1184 .iter()
1185 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1186 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1187 let subscription = format!("index-sub-{}", self.next_subscription);
1188 self.next_subscription += 1;
1189 self.index_subscriptions.insert(subscription.clone(), index);
1190 Ok(json!({
1191 "subscription": subscription,
1192 "revision": 1,
1193 "initial": initial,
1194 }))
1195 }
1196 "harness.v1.sessions.index.resize" => {
1197 let params = decode::<IndexResizeParams>(params)?;
1198 crate::session_index::validate_limit(params.limit)
1199 .map_err(ServiceError::InvalidParams)?;
1200 let index = self
1201 .index_subscriptions
1202 .get_mut(¶ms.subscription)
1203 .ok_or_else(|| {
1204 ServiceError::InvalidParams("unknown session index subscription".into())
1205 })?;
1206 let prepared = index
1207 .prepare_resize(params.limit)
1208 .map_err(ServiceError::Operation)?;
1209 let doors = crate::mail_route::LiveSessions::read(index.homes());
1210 let initial = prepared
1211 .page
1212 .sessions
1213 .iter()
1214 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1215 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1216 let response = json!({
1217 "subscription": params.subscription,
1218 "revision": prepared.revision,
1219 "initial": initial,
1220 "receipt": prepared.page.receipt,
1221 });
1222 index.commit_resize(prepared);
1223 Ok(response)
1224 }
1225 "harness.v1.sessions.index.unsubscribe" => {
1226 let params = decode::<UnfollowParams>(params)?;
1227 Ok(json!({
1228 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1229 }))
1230 }
1231 "harness.v1.sessions.import" => {
1232 let params = decode::<ImportSessionParams>(params)?;
1233 let session = Session::load_str(¶ms.content, params.source_harness.into())
1234 .map_err(operation)?;
1235 Ok(json!({"session": normalized_session_json(&session)}))
1236 }
1237 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1238 let params = decode::<ExportSessionParams>(params)?;
1239 let session = load_session(¶ms.locator).map_err(operation)?;
1240 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1241 if method == "harness.v1.sessions.export"
1242 && params.target_harness == TransferFormat::Hermes
1243 {
1244 let imported = crate::hermes_import::import_into_hermes(&session, None)
1246 .map_err(operation)?;
1247 return Ok(json!({"artifact": artifact, "imported": imported}));
1248 }
1249 Ok(json!({"artifact": artifact}))
1250 }
1251 "harness.v1.sessions.reduce" => {
1252 let params = decode::<ReduceSessionParams>(params)?;
1253 self.reduce_session(params)
1254 }
1255 "harness.v1.sessions.branch" => {
1256 let params = decode::<BranchSessionParams>(params)?;
1257 let session = load_session(¶ms.locator).map_err(operation)?;
1258 let storage = params.locator.storage.path().display().to_string();
1259 let bootstrap_prompt = format!(
1260 "Continue as a new branch from {} session {}. The frozen parent transcript is at {}. Read or load that parent for context, summarize the relevant state, then continue independently without mutating the parent session.",
1261 params.locator.harness.as_str(), params.locator.session_id, storage
1262 );
1263 let artifact = params
1264 .target_harness
1265 .map(|target| session_artifact(¶ms.locator, &session, target))
1266 .transpose()?;
1267 Ok(json!({
1268 "parent": params.locator,
1269 "session": normalized_session_json(&session),
1270 "bootstrap_prompt": bootstrap_prompt,
1271 "artifact": artifact,
1272 }))
1273 }
1274 "harness.v1.sessions.handoff" => {
1275 let params = decode::<HandoffSessionParams>(params)?;
1276 let session = load_session(¶ms.locator).map_err(operation)?;
1277 let cwd = params
1278 .cwd
1279 .or_else(|| session.meta.cwd.clone())
1280 .unwrap_or_else(|| PathBuf::from("."));
1281 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1282 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1283 ServiceError::Operation(
1284 "handoff artifact omitted target session identity".into(),
1285 )
1286 })?;
1287 let instructions =
1288 handoff_instructions(params.target_harness, target_session_id, &cwd);
1289 Ok(json!({
1290 "artifact": artifact,
1291 "launch": instructions.launch,
1292 "materialize": instructions.materialize,
1293 "requires_materialization": instructions.requires_materialization,
1294 "note": instructions.note,
1295 }))
1296 }
1297 "harness.v1.sessions.materialize" => {
1298 let params = decode::<MaterializeSessionParams>(params)?;
1299 for file in ¶ms.artifact.files {
1303 if file.role == "source_recovery"
1304 && file.path == "recovery/source.supercode.jsonl"
1305 {
1306 if let Ok(source) = Session::from_native_str(&file.content) {
1307 crate::residue_store::store_segments(&source);
1308 }
1309 }
1310 }
1311 let locator = crate::native_materialize::materialize_native_artifact(
1312 params.artifact,
1313 ¶ms.cwd,
1314 ¶ms.homes,
1315 )
1316 .map_err(ServiceError::Operation)?;
1317 Ok(json!({"locator": locator}))
1318 }
1319 "harness.v1.jobs.list" => {
1323 let query = decode::<crate::jobs::JobsQuery>(params)?;
1324 if let Some(harness) = query.harness.as_deref() {
1325 refuse_harness_without_jobs(harness, "jobs.list")?;
1326 }
1327 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1328 serde_json::to_value(listing)
1329 .map_err(|error| ServiceError::Operation(error.to_string()))
1330 }
1331 "harness.v1.jobs.get" => {
1332 let params = decode::<JobsGetParams>(params)?;
1333 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1334 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1335 .map_err(operation)?
1336 {
1337 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1338 None => Err(ServiceError::Operation(format!(
1339 "`{}` has no scheduled job `{}`",
1340 params.harness, params.id
1341 ))),
1342 }
1343 }
1344 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1350 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1351 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1352 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1353 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1354 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1355 "harness.v1.jobs.notepad"
1356 | "harness.v1.jobs.notepad_set"
1357 | "harness.v1.jobs.notepad_delete" => {
1358 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1359 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1360 let answer = match method {
1361 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1362 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1363 _ => crate::jobs_notepad::read(&request),
1364 }
1365 .map_err(job_control_error)?;
1366 serde_json::to_value(answer)
1367 .map_err(|error| ServiceError::Operation(error.to_string()))
1368 }
1369 "harness.v1.runs.list" => {
1373 let query = decode::<crate::runs::RunsQuery>(params)?;
1374 if let Some(harness) = query.harness.as_deref() {
1375 refuse_harness_without_runs(harness, "runs.list")?;
1376 }
1377 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1378 serde_json::to_value(listing)
1379 .map_err(|error| ServiceError::Operation(error.to_string()))
1380 }
1381 "harness.v1.runs.get" => {
1382 let params = decode::<RunsGetParams>(params)?;
1383 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1384 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1385 .map_err(operation)?
1386 {
1387 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1388 None => Err(ServiceError::Operation(format!(
1389 "`{}` has no run `{}`",
1390 params.harness, params.id
1391 ))),
1392 }
1393 }
1394 "harness.v1.sessions.resume_instructions" => {
1395 let params = decode::<ResumeInstructionsParams>(params)?;
1396 let session = load_session(¶ms.locator).map_err(operation)?;
1397 let cwd = params
1398 .cwd
1399 .or(session.meta.cwd)
1400 .unwrap_or_else(|| PathBuf::from("."));
1401 let launch = resume_launch(
1402 params.locator.harness.as_str(),
1403 ¶ms.locator.session_id,
1404 &cwd,
1405 params.policy,
1406 )?;
1407 Ok(json!({"launch": launch}))
1408 }
1409 _ => Err(ServiceError::MethodNotFound),
1410 }
1411 }
1412
1413 fn reduce_session(
1414 &self,
1415 params: ReduceSessionParams,
1416 ) -> std::result::Result<Value, ServiceError> {
1417 let session = load_session(¶ms.locator).map_err(operation)?;
1418 if session.messages.is_empty() {
1419 return Err(ServiceError::InvalidParams(
1420 "cannot reduce an empty session".into(),
1421 ));
1422 }
1423 let keep_last = params.keep_last.clamp(1, 128);
1424 let policy = reduce::ReductionPolicy {
1425 clear_turns_older_than: Some(keep_last),
1426 ..Default::default()
1427 };
1428 let (view, log) =
1429 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1430 if log.reductions.is_empty() {
1431 return Err(ServiceError::UnsupportedAction(format!(
1432 "session `{}` is already too small for a meaningful reversible reduction",
1433 params.locator.session_id
1434 )));
1435 }
1436 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1437 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1438 if reduced_tokens >= source_tokens {
1439 return Err(ServiceError::UnsupportedAction(format!(
1440 "session `{}` has no token-reducing reversible projection",
1441 params.locator.session_id
1442 )));
1443 }
1444
1445 let store_root = self
1446 .reduction_store_root
1447 .clone()
1448 .unwrap_or_else(default_reduction_store_root);
1449 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1450 let rescue_id = format!("rescue-{}", generated_session_id());
1451 let imported = session
1452 .imported_message_count
1453 .unwrap_or(session.messages.len())
1454 .min(session.messages.len());
1455 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1456 let view_jsonl = messages_jsonl(&view)?;
1457 let title = format!(
1458 "Reduced {} continuation from {}",
1459 params.target_harness.id(),
1460 params.locator.session_id
1461 );
1462
1463 store
1468 .save_sidecar(&rescue_id, &sidecar_jsonl)
1469 .map_err(operation)?;
1470 store
1471 .save_reduction_log(&rescue_id, &log)
1472 .map_err(operation)?;
1473 store
1474 .save(&rescue_id, &title, &view_jsonl)
1475 .map_err(operation)?;
1476
1477 let source_bytes = serde_json::to_vec(&session.messages)
1478 .map_err(|error| ServiceError::Operation(error.to_string()))?
1479 .len() as u64;
1480 let reduced_bytes = serde_json::to_vec(&view)
1481 .map_err(|error| ServiceError::Operation(error.to_string()))?
1482 .len() as u64;
1483 store
1484 .set_reduction_stats(
1485 &rescue_id,
1486 &title,
1487 source_bytes,
1488 reduced_bytes,
1489 log.reductions.len() as u32,
1490 )
1491 .map_err(operation)?;
1492
1493 let reloaded_sidecar = store
1497 .load_sidecar(&rescue_id)
1498 .map_err(operation)?
1499 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1500 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1501 let reloaded_log = store
1502 .load_reduction_log(&rescue_id)
1503 .map_err(operation)?
1504 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1505 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1506 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1507 let (restamped_view, restamped_log) =
1514 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1515 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1516 return Err(ServiceError::Operation(
1517 "persisted reduction view does not match its durable log and sidecar".into(),
1518 ));
1519 }
1520 if restamped_log != reloaded_log {
1521 return Err(ServiceError::Operation(
1522 "reapplying the durable reduction log changed its identity".into(),
1523 ));
1524 }
1525 let inverted =
1526 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1527 if inverted != session.messages {
1528 return Err(ServiceError::Operation(
1529 "reduction inversion did not restore the source messages byte-exactly".into(),
1530 ));
1531 }
1532
1533 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1534 let sidecar_path = store.sidecar_path(&rescue_id);
1535 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1536 let bootstrap_prompt = reduced_bootstrap_prompt(
1537 ¶ms.locator,
1538 params.target_harness,
1539 &view_jsonl,
1540 &sidecar_path,
1541 &reduction_log_path,
1542 );
1543 let mut reduced_session = session.clone();
1544 reduced_session.meta.session_id = Some(rescue_id.clone());
1545 reduced_session.messages = view;
1546
1547 Ok(json!({
1548 "session": normalized_session_json(&reduced_session),
1549 "bootstrap_prompt": bootstrap_prompt,
1550 "receipt": {
1551 "id": rescue_id,
1552 "sidecar_id": rescue_id,
1553 "source_harness": params.locator.harness,
1554 "target_harness": params.target_harness.id(),
1555 "source_tokens": source_tokens,
1556 "reduced_tokens": reduced_tokens,
1557 "ratio": ratio,
1558 "source_bytes": source_bytes,
1559 "reduced_bytes": reduced_bytes,
1560 "reductions": reloaded_log.reductions.len(),
1561 "sidecar_path": sidecar_path,
1562 "reduction_log_path": reduction_log_path,
1563 "verified": true,
1564 "reversible": true,
1565 }
1566 }))
1567 }
1568
1569 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1587 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1588 return None;
1589 }
1590 let method = request.get("method").and_then(Value::as_str)?;
1591 if !RUNTIME_OPEN_METHODS.contains(&method) {
1592 return None;
1593 }
1594 Some(RuntimeOpen {
1595 id: request.get("id").cloned().unwrap_or(Value::Null),
1596 method: method.to_string(),
1597 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1598 })
1599 }
1600
1601 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1623 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1624 return None;
1625 }
1626 let method = request.get("method").and_then(Value::as_str)?;
1627 if !DETACHED_METHODS.contains(&method) {
1628 return None;
1629 }
1630 let id = request.get("id").cloned().unwrap_or(Value::Null);
1631 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1632 let work = match method {
1633 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1634 .inventory_work(method, params)
1635 .map(DetachedWork::Inventory),
1636 "harness.v1.sessions.message" => {
1637 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1638 }
1639 _ => {
1640 let verb = match method {
1641 "harness.v1.sessions.new" => crate::SessionVerb::New,
1642 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1643 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1644 _ => crate::SessionVerb::Delete,
1645 };
1646 match decode::<crate::SessionMutation>(params) {
1647 Ok(mutation) => {
1648 match crate::sessions_control::door(&mutation.harness, verb) {
1649 Ok(crate::SessionDoor::Live(_)) => return None,
1652 Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1653 Err(error) => Err(session_control_error(error)),
1654 }
1655 }
1656 Err(error) => Err(error),
1657 }
1658 }
1659 };
1660 Some(DetachedCall {
1661 id,
1662 method: method.to_string(),
1663 work: work.map(Work::Free),
1664 })
1665 }
1666
1667 pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1681 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1682 return None;
1683 }
1684 let method = request.get("method").and_then(Value::as_str)?;
1685 let id = request.get("id").cloned().unwrap_or(Value::Null);
1686 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1687 let work = match method {
1688 "harness.v1.runtimes.close" => decode::<RuntimeConnectionParams>(params)
1689 .and_then(|params| self.surrender_runtime(¶ms.connection))
1690 .map(|(runtime, process_group)| {
1691 Work::Runtime(RuntimeWork::Close {
1692 runtime,
1693 process_group,
1694 })
1695 }),
1696 "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1697 let verb = if method == "harness.v1.sessions.new" {
1698 crate::SessionVerb::New
1699 } else {
1700 crate::SessionVerb::Reset
1701 };
1702 let mutation = decode::<crate::SessionMutation>(params).ok()?;
1703 let Ok(crate::SessionDoor::Live(command)) =
1707 crate::sessions_control::door(&mutation.harness, verb)
1708 else {
1709 return None;
1710 };
1711 let connection = mutation
1712 .connection
1713 .clone()
1714 .filter(|value| !value.trim().is_empty())?;
1715 self.lend_runtime(&connection).map(|runtime| {
1716 let session = live_session_name(runtime.as_ref(), &mutation);
1717 Work::Runtime(RuntimeWork::LiveCommand {
1718 connection,
1719 runtime,
1720 verb,
1721 mutation,
1722 command,
1723 session,
1724 })
1725 })
1726 }
1727 _ => return None,
1728 };
1729 Some(DetachedCall {
1730 id,
1731 method: method.to_string(),
1732 work,
1733 })
1734 }
1735
1736 pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1741 let DetachedAnswer { response, returned } = answer;
1742 if let Some(ReturnedRuntime {
1743 connection,
1744 runtime,
1745 }) = returned
1746 {
1747 self.runtimes_in_flight.remove(&connection);
1748 self.runtimes.insert(connection, runtime);
1749 }
1750 response
1751 }
1752
1753 pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1757 let OpenedRuntime { id, outcome } = opened;
1758 let result = match outcome {
1759 Ok(open) => self.register_open_runtime(open).await,
1760 Err(error) => Err(error),
1761 };
1762 service_response(id, result)
1763 }
1764
1765 async fn register_open_runtime(
1767 &mut self,
1768 open: OpenRuntime,
1769 ) -> std::result::Result<Value, ServiceError> {
1770 match open {
1771 OpenRuntime::Hosted {
1772 runtime,
1773 capabilities,
1774 workspace,
1775 fresh,
1776 } => {
1777 self.insert_hosted_runtime(runtime, capabilities, workspace, fresh)
1778 .await
1779 }
1780 OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1781 }
1782 }
1783
1784 async fn runtime_call(
1785 &mut self,
1786 method: &str,
1787 params: Value,
1788 ) -> std::result::Result<Value, ServiceError> {
1789 match method {
1790 "harness.v1.runtimes.capabilities" => {
1791 let params = decode::<RuntimeBackendParams>(params)?;
1792 let backend = runtime_backend(¶ms)?;
1793 Ok(json!({
1794 "harness": backend.harness(),
1795 "capabilities": backend.capabilities(),
1796 }))
1797 }
1798 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1799 self.register_open_runtime(open_runtime(method, params).await?)
1800 .await
1801 }
1802 "harness.v1.runtimes.send_input" => {
1803 let params = decode::<RuntimeInputParams>(params)?;
1804 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1805 let runtime = self.runtime_mut(¶ms.connection)?;
1806 let turn_id = within_control_deadline(
1807 method,
1808 runtime.send_input(RuntimeInput {
1809 text: params.text,
1810 image_urls,
1811 }),
1812 )
1813 .await?
1814 .map_err(operation)?;
1815 Ok(json!({"turn_id": turn_id}))
1816 }
1817 "harness.v1.runtimes.interrupt" => {
1818 let params = decode::<RuntimeConnectionParams>(params)?;
1819 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1820 .await?
1821 .map_err(operation)?;
1822 Ok(json!({}))
1823 }
1824 "harness.v1.runtimes.steer" => {
1825 let params = decode::<RuntimeInputParams>(params)?;
1826 if !params.image_urls.is_empty() {
1827 return Err(ServiceError::InvalidParams(
1828 "runtime steering accepts text only".into(),
1829 ));
1830 }
1831 let text = params.text.trim();
1832 if text.is_empty() || text.chars().count() > 50_000 {
1833 return Err(ServiceError::InvalidParams(
1834 "runtime steering requires 1 to 50,000 text characters".into(),
1835 ));
1836 }
1837 within_control_deadline(
1838 method,
1839 self.runtime_mut(¶ms.connection)?
1840 .steer(text.to_string()),
1841 )
1842 .await?
1843 .map_err(operation)?;
1844 Ok(json!({}))
1845 }
1846 "harness.v1.runtimes.respond" => {
1847 let params = decode::<RuntimeRespondParams>(params)?;
1848 let request_id = params.request_id.clone();
1849 within_control_deadline(
1850 method,
1851 self.runtime_mut(¶ms.connection)?
1852 .respond(params.request_id, params.response),
1853 )
1854 .await?
1855 .map_err(operation)?;
1856 self.approvals.answered(¶ms.connection, &request_id);
1858 Ok(json!({}))
1859 }
1860 "harness.v1.runtimes.acquire_control" => {
1861 let params = decode::<RuntimeConnectionParams>(params)?;
1862 let snapshot = within_control_deadline(
1863 method,
1864 self.runtime_mut(¶ms.connection)?.acquire_control(),
1865 )
1866 .await?
1867 .map_err(operation)?;
1868 serde_json::to_value(snapshot)
1869 .map_err(|error| ServiceError::Operation(error.to_string()))
1870 }
1871 "harness.v1.runtimes.heartbeat" => {
1872 let params = decode::<RuntimeConnectionParams>(params)?;
1873 let snapshot = within_control_deadline(
1874 method,
1875 self.runtime_mut(¶ms.connection)?.heartbeat(),
1876 )
1877 .await?
1878 .map_err(operation)?;
1879 serde_json::to_value(snapshot)
1880 .map_err(|error| ServiceError::Operation(error.to_string()))
1881 }
1882 "harness.v1.runtimes.detach" => {
1883 let params = decode::<RuntimeConnectionParams>(params)?;
1884 let snapshot =
1885 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1886 .await?
1887 .map_err(operation)?;
1888 serde_json::to_value(snapshot)
1889 .map_err(|error| ServiceError::Operation(error.to_string()))
1890 }
1891 "harness.v1.runtimes.terminal_instructions" => {
1892 let params = decode::<RuntimeConnectionParams>(params)?;
1893 let launch = self
1894 .terminal_launches
1895 .get(¶ms.connection)
1896 .ok_or_else(|| {
1897 ServiceError::Operation(
1898 "this runtime is not hosted for terminal attachment".into(),
1899 )
1900 })?;
1901 Ok(json!({"launch":launch}))
1902 }
1903 "harness.v1.runtimes.close" => {
1904 let params = decode::<RuntimeConnectionParams>(params)?;
1905 let (runtime, process_group) = self.surrender_runtime(¶ms.connection)?;
1906 close_runtime(runtime, process_group).await
1907 }
1908 _ => Err(ServiceError::MethodNotFound),
1909 }
1910 }
1911
1912 #[cfg(feature = "adapter-api")]
1914 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1915 let params = decode::<MessageSessionParams>(params)?;
1916 Ok(message_live_session(¶ms).await)
1917 }
1918
1919 #[cfg(feature = "adapter-api")]
1920 fn harness_settings_call(
1921 &self,
1922 method: &str,
1923 params: Value,
1924 ) -> std::result::Result<Value, ServiceError> {
1925 let homes = crate::HarnessHomes::default();
1926 match method {
1927 "harness.v1.harnesses.settings" => {
1928 let params = decode::<HarnessSettingsParams>(params)?;
1929 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1930 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1931 serde_json::to_value(report)
1932 .map_err(|error| ServiceError::Operation(error.to_string()))
1933 }
1934 "harness.v1.harnesses.configure" => {
1935 let params = decode::<ConfigureHarnessParams>(params)?;
1936 let report = crate::configure_harness_interop_settings(
1937 &homes,
1938 ¶ms.harness,
1939 ¶ms.changes,
1940 params.expected_revision.as_deref(),
1941 )
1942 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1943 serde_json::to_value(report)
1944 .map_err(|error| ServiceError::Operation(error.to_string()))
1945 }
1946 _ => Err(ServiceError::MethodNotFound),
1947 }
1948 }
1949
1950 fn insert_runtime(
1951 &mut self,
1952 runtime: Box<dyn RuntimeConnection>,
1953 ) -> std::result::Result<Value, ServiceError> {
1954 let connection = format!("runtime-{}", self.next_runtime);
1955 self.next_runtime += 1;
1956 let handle = runtime.handle().clone();
1957 self.runtime_sequences
1958 .entry(handle.runtime_id.clone())
1959 .or_insert(0);
1960 self.runtimes.insert(connection.clone(), runtime);
1961 Ok(json!({"connection": connection, "handle": handle}))
1962 }
1963
1964 #[cfg(feature = "adapter-api")]
1965 async fn insert_hosted_runtime(
1966 &mut self,
1967 runtime: Box<dyn RuntimeConnection>,
1968 capabilities: crate::RuntimeCapabilities,
1969 workspace: PathBuf,
1970 fresh: bool,
1971 ) -> std::result::Result<Value, ServiceError> {
1972 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
1973 let token: std::sync::Arc<str> = crate::server::generate_token().into();
1974 let server = crate::server::run_frontend_http(
1975 host.clone(),
1976 host.frontend_sender(),
1977 "127.0.0.1:0",
1978 token.clone(),
1979 connection.handle().runtime_id.clone(),
1980 )
1981 .await
1982 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1983 let source = LiveRuntimeSource {
1984 harness: connection.handle().harness.as_str().to_string(),
1985 session_id: connection.handle().runtime_id.clone(),
1986 workspace: workspace.clone(),
1987 };
1988 let registration = register_live_runtime(
1989 connection.handle().runtime_id.clone(),
1990 source.clone(),
1991 format!("http://{}", server.address()),
1992 token.to_string(),
1993 )
1994 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1995 let endpoint = registration.endpoint().to_string();
1996 let launch = StructuredLaunch {
1997 cwd: workspace,
1998 program: std::env::current_exe()
2002 .ok()
2003 .map(|path| path.to_string_lossy().into_owned())
2004 .unwrap_or_else(|| "supercode".into()),
2005 arguments: vec![
2006 "open".into(),
2007 endpoint,
2008 "--harness".into(),
2009 source.harness,
2010 "--session".into(),
2011 source.session_id,
2012 ],
2013 env: BTreeMap::new(),
2014 };
2015 let lease = HostedRuntimeLease {
2016 connection,
2017 _host: host,
2018 _registration: registration,
2019 _server: server,
2020 };
2021 let opened = self.insert_runtime(Box::new(lease))?;
2022 let connection_id = opened["connection"]
2023 .as_str()
2024 .expect("insert_runtime returns a connection id")
2025 .to_string();
2026 self.terminal_launches.insert(connection_id, launch);
2027 Ok(opened)
2028 }
2029
2030 #[cfg(not(feature = "adapter-api"))]
2031 async fn insert_hosted_runtime(
2032 &mut self,
2033 runtime: Box<dyn RuntimeConnection>,
2034 _capabilities: crate::RuntimeCapabilities,
2035 _workspace: PathBuf,
2036 _fresh: bool,
2037 ) -> std::result::Result<Value, ServiceError> {
2038 self.insert_runtime(runtime)
2039 }
2040
2041 fn runtime_mut(
2042 &mut self,
2043 connection: &str,
2044 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2045 if self.runtimes_in_flight.contains(connection) {
2046 return Err(self.lent_out(connection));
2047 }
2048 self.runtimes.get_mut(connection).ok_or_else(|| {
2049 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2050 })
2051 }
2052
2053 fn lent_out(&self, connection: &str) -> ServiceError {
2057 ServiceError::Operation(format!(
2058 "runtime connection `{connection}`: a harness turn is already in progress"
2059 ))
2060 }
2061
2062 fn lend_runtime(
2065 &mut self,
2066 connection: &str,
2067 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2068 if self.runtimes_in_flight.contains(connection) {
2069 return Err(self.lent_out(connection));
2070 }
2071 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2072 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2073 })?;
2074 self.runtimes_in_flight.insert(connection.to_string());
2075 Ok(runtime)
2076 }
2077
2078 fn surrender_runtime(
2090 &mut self,
2091 connection: &str,
2092 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2093 if self.runtimes_in_flight.contains(connection) {
2094 return Err(self.lent_out(connection));
2095 }
2096 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2097 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2098 })?;
2099 let process_group = runtime_process_group(runtime.handle());
2100 let runtime_id = runtime.handle().runtime_id.clone();
2101 self.terminal_launches.remove(connection);
2102 self.runtime_sequences.remove(&runtime_id);
2103 self.approvals.forget(connection);
2104 Ok((runtime, process_group))
2105 }
2106
2107 pub fn kill_all_runtime_groups(&self) -> usize {
2118 self.runtimes
2119 .values()
2120 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2121 .count()
2122 }
2123
2124 async fn mutate_session(
2135 &mut self,
2136 verb: crate::SessionVerb,
2137 params: Value,
2138 ) -> std::result::Result<Value, ServiceError> {
2139 let mutation = decode::<crate::SessionMutation>(params)?;
2140 let door = crate::sessions_control::door(&mutation.harness, verb)
2141 .map_err(session_control_error)?;
2142 let outcome = match door {
2143 #[cfg(not(feature = "adapter-api"))]
2147 crate::SessionDoor::Live(command) => {
2148 return Err(ServiceError::Operation(format!(
2149 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2150 session, which needs this build's `adapter-api` feature",
2151 mutation.harness,
2152 verb.as_str()
2153 )));
2154 }
2155 #[cfg(feature = "adapter-api")]
2156 crate::SessionDoor::Live(command) => {
2157 let connection = mutation
2158 .connection
2159 .clone()
2160 .filter(|value| !value.trim().is_empty())
2161 .ok_or_else(|| {
2162 ServiceError::InvalidParams(format!(
2163 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2164 driven session: pass the `connection` of an open runtime \
2165 (`harness.v1.runtimes.start`)",
2166 mutation.harness,
2167 verb.as_str()
2168 ))
2169 })?;
2170 let runtime = self.runtime_mut(&connection)?;
2171 let session = live_session_name(runtime.as_ref(), &mutation);
2172 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2178 .await;
2179 }
2180 _ => run_session_mutation(verb, &mutation).await?,
2181 };
2182 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2183 }
2184
2185 async fn inventory_call(
2189 &self,
2190 method: &str,
2191 params: Value,
2192 ) -> std::result::Result<Value, ServiceError> {
2193 run_inventory(self.inventory_work(method, params)?).await
2194 }
2195
2196 fn inventory_work(
2202 &self,
2203 method: &str,
2204 params: Value,
2205 ) -> std::result::Result<InventoryWork, ServiceError> {
2206 let mut params = decode::<HarnessInventoryParams>(params)?;
2207 if method == "harness.v1.harnesses.probe" {
2208 let harness = params.harness.take().ok_or_else(|| {
2209 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2210 })?;
2211 params.harnesses = vec![harness];
2212 }
2213 let selected = params
2214 .harnesses
2215 .iter()
2216 .map(HarnessId::as_str)
2217 .collect::<std::collections::BTreeSet<_>>();
2218 let supported = harness_support_registry()
2219 .harnesses
2220 .into_iter()
2221 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2222 .collect::<Vec<_>>();
2223 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2224 let known = supported
2225 .iter()
2226 .map(|harness| harness.id.as_str())
2227 .collect::<std::collections::BTreeSet<_>>();
2228 let missing = params
2229 .harnesses
2230 .iter()
2231 .filter(|id| !known.contains(id.as_str()))
2232 .map(HarnessId::as_str)
2233 .collect::<Vec<_>>();
2234 return Err(ServiceError::InvalidParams(format!(
2235 "unknown harness(es): {}",
2236 missing.join(", ")
2237 )));
2238 }
2239 let global_counts = params
2240 .include_sessions
2241 .then(|| self.session_counts(None, ¶ms.harnesses));
2242 let workspace_counts = params
2243 .include_sessions
2244 .then(|| {
2245 params
2246 .workspace
2247 .as_deref()
2248 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2249 })
2250 .flatten();
2251 Ok(InventoryWork {
2252 params,
2253 supported,
2254 global_counts,
2255 workspace_counts,
2256 })
2257 }
2258
2259 #[cfg(feature = "adapter-api")]
2260 async fn harness_authentication_call(
2261 &self,
2262 method: &str,
2263 params: Value,
2264 ) -> std::result::Result<Value, ServiceError> {
2265 match method {
2266 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2267 let params = decode::<HarnessAuthenticationParams>(params)?;
2268 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2269 .map_err(|error| ServiceError::Operation(error.to_string()))
2270 }
2271 "harness.v1.harnesses.auth.begin" => {
2272 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2273 let cwd = params
2274 .cwd
2275 .or_else(|| std::env::current_dir().ok())
2276 .unwrap_or_else(|| PathBuf::from("."));
2277 let plan = crate::harness_authentication_plan(
2278 ¶ms.harness,
2279 params.environment,
2280 params.method,
2281 &cwd,
2282 )
2283 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2284 serde_json::to_value(plan)
2285 .map_err(|error| ServiceError::Operation(error.to_string()))
2286 }
2287 _ => Err(ServiceError::MethodNotFound),
2288 }
2289 }
2290
2291 fn session_counts(
2292 &self,
2293 workspace: Option<&Path>,
2294 harnesses: &[HarnessId],
2295 ) -> BTreeMap<String, usize> {
2296 let mut counts = BTreeMap::new();
2297 for session in self
2298 .catalog
2299 .discover(&DiscoveryQuery {
2300 workspace: workspace.map(Path::to_path_buf),
2301 harnesses: harnesses.to_vec(),
2302 ..DiscoveryQuery::default()
2303 })
2304 .unwrap_or_default()
2305 {
2306 *counts
2307 .entry(session.locator.harness.as_str().to_string())
2308 .or_insert(0) += 1;
2309 }
2310 counts
2311 }
2312}
2313
2314#[async_trait::async_trait]
2315impl SdkService for HarnessSessionService {
2316 fn capabilities(&self) -> SdkCapabilities {
2317 SdkCapabilities::default()
2318 }
2319
2320 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2321 if request.operation == SdkOperation::Events {
2322 let events = self
2323 .poll_sdk_events()
2324 .await
2325 .into_iter()
2326 .map(|(_, event)| event)
2327 .collect::<Vec<_>>();
2328 return serde_json::to_value(events).map_err(|error| {
2329 SdkError::new(
2330 SdkErrorCode::Execution,
2331 request.operation,
2332 error.to_string(),
2333 )
2334 });
2335 }
2336 if self.runtimes.is_empty()
2337 && matches!(
2338 request.operation,
2339 SdkOperation::Input
2340 | SdkOperation::Interrupt
2341 | SdkOperation::Steer
2342 | SdkOperation::Respond
2343 | SdkOperation::Close
2344 )
2345 {
2346 return Err(SdkError::unsupported(request.operation));
2347 }
2348 let method = request
2349 .operation
2350 .method()
2351 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2352 let result = match request.operation {
2353 SdkOperation::Discover
2354 | SdkOperation::Load
2355 | SdkOperation::Export
2356 | SdkOperation::ProfilesList
2357 | SdkOperation::ProfilesGet
2358 | SdkOperation::ProfilesCreate
2359 | SdkOperation::ProfilesDelete
2360 | SdkOperation::SkillsList
2361 | SdkOperation::SkillsInstall
2362 | SdkOperation::SkillsRemove
2363 | SdkOperation::ChannelsList
2364 | SdkOperation::RoutesList
2365 | SdkOperation::TriggersList
2366 | SdkOperation::ChannelsStatus
2367 | SdkOperation::MemoryShow
2368 | SdkOperation::MemorySearch
2369 | SdkOperation::JobsList
2370 | SdkOperation::JobsGet
2371 | SdkOperation::JobsCreate
2372 | SdkOperation::JobsUpdate
2373 | SdkOperation::JobsPause
2374 | SdkOperation::JobsResume
2375 | SdkOperation::JobsRun
2376 | SdkOperation::JobsDelete
2377 | SdkOperation::JobsNotepad
2378 | SdkOperation::JobsNotepadSet
2379 | SdkOperation::JobsNotepadDelete
2380 | SdkOperation::RunsList
2381 | SdkOperation::RunsGet
2382 | SdkOperation::ApprovalsList
2383 | SdkOperation::OrchestrationLoad
2384 | SdkOperation::OrchestrationSave
2385 | SdkOperation::OrchestrationCompile
2386 | SdkOperation::OrchestrationDecompile
2387 | SdkOperation::OrchestrationImport
2388 | SdkOperation::OrchestrationExport
2389 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2390 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2393 SdkOperation::Start
2394 | SdkOperation::Resume
2395 | SdkOperation::Input
2396 | SdkOperation::Interrupt
2397 | SdkOperation::Steer
2398 | SdkOperation::Respond
2399 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2400 SdkOperation::SessionsNew => {
2405 self.mutate_session(crate::SessionVerb::New, request.params)
2406 .await
2407 }
2408 SdkOperation::SessionsReset => {
2409 self.mutate_session(crate::SessionVerb::Reset, request.params)
2410 .await
2411 }
2412 SdkOperation::SessionsArchive => {
2413 self.mutate_session(crate::SessionVerb::Archive, request.params)
2414 .await
2415 }
2416 SdkOperation::SessionsDelete => {
2417 self.mutate_session(crate::SessionVerb::Delete, request.params)
2418 .await
2419 }
2420 SdkOperation::Events => unreachable!("handled before method dispatch"),
2421 };
2422 result.map_err(|error| sdk_error(request.operation, error))
2423 }
2424
2425 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2426 Ok(self
2427 .poll_sdk_events()
2428 .await
2429 .into_iter()
2430 .map(|(_, event)| event)
2431 .collect())
2432 }
2433}
2434
2435#[cfg(feature = "adapter-api")]
2436struct HostedRuntimeLease {
2437 connection: HostedHarnessConnection,
2438 _host: std::sync::Arc<HostedHarnessRuntime>,
2439 _registration: LiveRuntimeRegistration,
2440 _server: crate::server::FrontendHttpServer,
2441}
2442
2443#[async_trait::async_trait]
2444#[cfg(feature = "adapter-api")]
2445impl RuntimeConnection for HostedRuntimeLease {
2446 fn handle(&self) -> &crate::RuntimeHandle {
2447 self.connection.handle()
2448 }
2449
2450 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2451 self.connection.send_input(input).await
2452 }
2453
2454 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2455 self.connection.next_event().await
2456 }
2457
2458 async fn interrupt(&mut self) -> crate::Result<()> {
2459 self.connection.interrupt().await
2460 }
2461
2462 async fn steer(&mut self, text: String) -> crate::Result<()> {
2465 self.connection.steer(text).await
2466 }
2467
2468 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2469 self.connection.respond(request_id, response).await
2470 }
2471
2472 async fn close(&mut self) -> crate::Result<()> {
2473 self.connection.close().await
2474 }
2475}
2476
2477struct InventoryWork {
2480 params: HarnessInventoryParams,
2481 supported: Vec<crate::HarnessSupportDescriptor>,
2482 global_counts: Option<BTreeMap<String, usize>>,
2483 workspace_counts: Option<BTreeMap<String, usize>>,
2484}
2485
2486async fn run_session_mutation(
2493 verb: crate::SessionVerb,
2494 mutation: &crate::SessionMutation,
2495) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2496 let door =
2503 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2504 if let crate::SessionDoor::Http = door {
2505 return crate::sessions_control::mutate(verb, mutation)
2506 .await
2507 .map_err(session_control_error);
2508 }
2509 let mutation = mutation.clone();
2510 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2511 .await
2512 .map_err(|error| {
2513 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2514 })?
2515 .map_err(session_control_error)
2516}
2517
2518async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2521 let InventoryWork {
2522 params,
2523 supported,
2524 global_counts,
2525 workspace_counts,
2526 } = work;
2527 let probes = supported.into_iter().map(|descriptor| {
2528 let global = global_counts
2529 .as_ref()
2530 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2531 let workspace = workspace_counts
2532 .as_ref()
2533 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2534 probe_harness(descriptor, ¶ms, global, workspace)
2535 });
2536 let harnesses = futures::future::join_all(probes).await;
2537 serde_json::to_value(HarnessInventoryReport {
2538 probe: params.probe,
2539 workspace: params.workspace,
2540 harnesses,
2541 })
2542 .map_err(|error| ServiceError::Operation(error.to_string()))
2543}
2544
2545async fn probe_harness(
2546 descriptor: crate::HarnessSupportDescriptor,
2547 params: &HarnessInventoryParams,
2548 global: Option<usize>,
2549 workspace: Option<usize>,
2550) -> LocalHarness {
2551 let launch = descriptor.runtime.default_launch.as_ref();
2552 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2557 .then(crate::orchestrator::daemon_entry)
2558 .and_then(Result::ok);
2559 let executable = match &orchestrator_entry {
2560 Some(entry) => Some(entry.clone()),
2561 None => launch.and_then(|launch| find_executable(&launch.program)),
2562 };
2563 let installed = executable.is_some();
2564 let version = if params.skip_versions || orchestrator_entry.is_some() {
2565 None
2568 } else {
2569 match executable.as_deref() {
2570 Some(path) => executable_version(path).await,
2571 None => None,
2572 }
2573 };
2574 let configured = auth_evidence(descriptor.id.as_str());
2575 let mut auth = if configured {
2576 HarnessAuthState::Configured
2577 } else if matches!(
2578 descriptor.id.as_str(),
2579 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2580 ) {
2581 HarnessAuthState::Required
2586 } else {
2587 HarnessAuthState::Unknown
2588 };
2589 let mut runtime = if installed {
2590 HarnessRuntimeState::Degraded
2591 } else {
2592 HarnessRuntimeState::Unavailable
2593 };
2594 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2595 let mut reason = (!installed).then(|| {
2596 if is_orchestrator {
2597 format!(
2598 "{} is supported but its daemon entry `{}` was not found",
2599 descriptor.display_name,
2600 crate::orchestrator::DAEMON_ENTRY
2601 )
2602 } else {
2603 format!(
2604 "{} is supported but `{}` was not found on PATH",
2605 descriptor.display_name,
2606 launch
2607 .map(|launch| launch.program.as_str())
2608 .unwrap_or("executable")
2609 )
2610 }
2611 });
2612 let mut repair = (!installed).then(|| {
2613 if is_orchestrator {
2614 format!(
2615 "Install the `supercode-orchestrator` package so `{}` resolves.",
2616 crate::orchestrator::DAEMON_ENTRY
2617 )
2618 } else {
2619 format!(
2620 "Install {} and ensure `{}` is on PATH.",
2621 descriptor.display_name,
2622 launch
2623 .map(|launch| launch.program.as_str())
2624 .unwrap_or("its executable")
2625 )
2626 }
2627 });
2628
2629 if installed && params.probe == HarnessProbeLevel::Handshake {
2630 let backend_params = RuntimeBackendParams {
2631 harness: descriptor.id.clone(),
2632 protocol: None,
2633 launch: None,
2634 base_url: None,
2635 policy: RuntimePolicy::Default,
2636 };
2637 match runtime_backend(&backend_params) {
2638 Ok(backend) => {
2639 let cwd = params
2640 .workspace
2641 .clone()
2642 .or_else(|| std::env::current_dir().ok())
2643 .unwrap_or_else(|| PathBuf::from("."));
2644 let isolated = descriptor
2645 .runtime
2646 .default_launch
2647 .clone()
2648 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2649 let Some(isolated) = isolated else {
2650 reason = Some(
2651 "No-prompt runtime handshake could not create its isolated harness home."
2652 .into(),
2653 );
2654 repair = Some(
2655 "Check temporary-directory permissions, then run the handshake probe again."
2656 .into(),
2657 );
2658 let running = probe_running_instance(descriptor.id.as_str());
2659 return LocalHarness {
2660 gateway: gateway_health(
2661 descriptor.id.as_str(),
2662 installed,
2663 running.as_ref(),
2664 version.as_deref(),
2665 ),
2666 id: descriptor.id,
2667 display_name: descriptor.display_name,
2668 supported: true,
2669 installed,
2670 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2671 version,
2672 auth,
2673 runtime,
2674 protocol: descriptor.runtime.protocol,
2675 capabilities: descriptor.runtime.capabilities.clone(),
2676 effective_capabilities: descriptor.runtime.capabilities,
2677 sessions: HarnessSessionCounts { global, workspace },
2678 running,
2679 reason,
2680 repair,
2681 };
2682 };
2683 match tokio::time::timeout(
2684 Duration::from_secs(30),
2685 backend.start(RuntimeStartRequest {
2686 cwd,
2687 launch: Some(isolated.launch.clone()),
2688 mcp_servers: Vec::new(),
2689 approval_policy: None,
2690 }),
2691 )
2692 .await
2693 {
2694 Ok(Ok(mut connection)) => {
2695 match stabilize_handshake(connection.as_mut()).await {
2696 Ok(()) => {
2697 auth = HarnessAuthState::Ready;
2698 runtime = HarnessRuntimeState::Ready;
2699 reason = Some(
2700 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2701 .into(),
2702 );
2703 repair = None;
2704 }
2705 Err(message) => {
2706 auth = if looks_like_auth_error(&message) {
2707 HarnessAuthState::Required
2708 } else if configured {
2709 HarnessAuthState::Configured
2710 } else {
2711 HarnessAuthState::Unknown
2712 };
2713 reason = Some(format!(
2714 "No-prompt runtime handshake became unhealthy during startup: {message}"
2715 ));
2716 repair = Some(if auth == HarnessAuthState::Required {
2717 format!(
2718 "Run `{}` interactively once and complete sign-in, then probe again.",
2719 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2720 )
2721 } else {
2722 "Run the harness directly to inspect its startup failure, then probe again."
2723 .into()
2724 });
2725 }
2726 }
2727 let _ =
2728 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2729 }
2730 Ok(Err(error)) => {
2731 let message = truncate_text(&error.to_string(), 500);
2732 auth = if looks_like_auth_error(&message) {
2733 HarnessAuthState::Required
2734 } else if configured {
2735 HarnessAuthState::Configured
2736 } else {
2737 HarnessAuthState::Unknown
2738 };
2739 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2740 repair = Some(if auth == HarnessAuthState::Required {
2741 format!(
2742 "Run `{}` interactively once and complete sign-in, then probe again.",
2743 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2744 )
2745 } else {
2746 "Check the harness installation and run the handshake probe again."
2747 .into()
2748 });
2749 }
2750 Err(_) => {
2751 reason =
2752 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2753 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2754 }
2755 }
2756 let _ = isolated.cleanup();
2765 tokio::time::sleep(Duration::from_millis(250)).await;
2766 if let Err(error) = isolated.cleanup() {
2767 auth = if configured {
2768 HarnessAuthState::Configured
2769 } else {
2770 HarnessAuthState::Unknown
2771 };
2772 runtime = HarnessRuntimeState::Degraded;
2773 reason = Some(format!(
2774 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2775 ));
2776 repair = Some(
2777 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2778 .into(),
2779 );
2780 }
2781 }
2782 Err(error) => {
2783 reason = Some(error_message(error));
2784 }
2785 }
2786 } else if installed && configured {
2787 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2788 } else if installed && auth == HarnessAuthState::Required {
2789 reason = Some("Executable found, but no native authentication evidence is present.".into());
2790 repair = Some(format!(
2791 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2792 descriptor.id.as_str()
2793 ));
2794 } else if installed {
2795 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2796 repair = Some(format!(
2797 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2798 launch
2799 .map(|launch| launch.program.as_str())
2800 .unwrap_or("the harness")
2801 ));
2802 }
2803
2804 let effective_capabilities = if installed {
2805 descriptor.runtime.capabilities.clone()
2806 } else {
2807 unavailable_capabilities()
2808 };
2809 let running = probe_running_instance(descriptor.id.as_str());
2810 LocalHarness {
2811 gateway: gateway_health(
2812 descriptor.id.as_str(),
2813 installed,
2814 running.as_ref(),
2815 version.as_deref(),
2816 ),
2817 id: descriptor.id,
2818 display_name: descriptor.display_name,
2819 supported: true,
2820 installed,
2821 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2822 version,
2823 auth,
2824 runtime,
2825 protocol: descriptor.runtime.protocol,
2826 capabilities: descriptor.runtime.capabilities,
2827 effective_capabilities,
2828 sessions: HarnessSessionCounts { global, workspace },
2829 running,
2830 reason,
2831 repair,
2832 }
2833}
2834
2835async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2836 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2837 loop {
2838 let now = tokio::time::Instant::now();
2839 if now >= deadline {
2840 return Ok(());
2841 }
2842 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2843 Err(_) => return Ok(()),
2844 Ok(Ok(Some(event))) => {
2845 if let Some(message) = handshake_event_failure(&event) {
2846 return Err(truncate_text(&message, 500));
2847 }
2848 }
2849 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2850 Ok(Err(error)) => return Err(error.to_string()),
2851 }
2852 }
2853}
2854
2855fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2856 let detail = event
2857 .payload
2858 .get("message")
2859 .or_else(|| event.payload.get("line"))
2860 .and_then(Value::as_str)
2861 .unwrap_or(event.kind.as_str());
2862 match event.kind.as_str() {
2863 "transport_closed" => Some("runtime transport closed during startup".into()),
2864 "transport_error" => Some(format!("runtime transport error: {detail}")),
2865 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2866 _ => None,
2871 }
2872}
2873
2874fn indexed_claude_window(
2875 locator: &SessionLocator,
2876 options: &SessionLoadOptions,
2877) -> std::result::Result<Option<Value>, ServiceError> {
2878 use supercode_interchange::session::ClaudeReadIndex;
2879 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2882 || options.include_subagents != Some(false)
2883 {
2884 return Ok(None);
2885 }
2886 let crate::StorageLocator::File { path } = &locator.storage else {
2887 return Ok(None);
2888 };
2889 if !ClaudeReadIndex::supports(path)
2890 .map_err(|error| ServiceError::Operation(error.to_string()))?
2891 {
2892 return Ok(None);
2893 }
2894 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2895 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2896 let total = index.len();
2897 let (offset, end) = projected_message_window(total, options);
2898 let session = index
2899 .read_messages(offset..end)
2900 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2901 let summary = index
2902 .read_summary()
2903 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2904 let selected_options = SessionLoadOptions {
2905 message_offset: None,
2906 message_limit: None,
2907 message_tail: None,
2908 ..options.clone()
2909 };
2910 let mut selected = projected_session_json(&session, &selected_options);
2911 selected["raw_record_count"] = json!(index.raw_record_count());
2912 Ok(Some(json!({
2913 "session": selected,
2914 "summary": projected_session_summary(&summary, options),
2915 "window": {
2916 "has_more": offset > 0 || end < total, "has_newer": end < total,
2917 "has_older": offset > 0, "newer_items": index.item_count(end..total),
2918 "offset": offset, "older_items": index.item_count(0..offset),
2919 "returned": end - offset, "total_messages": total,
2920 }
2921 })))
2922}
2923
2924fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
2925 let total_messages = session.messages.len();
2926 let (offset, end) = projected_message_window(total_messages, options);
2927 json!({
2928 "session": projected_session_json(session, options),
2929 "summary": projected_session_summary(session, options),
2930 "window": {
2931 "has_more": offset > 0 || end < total_messages,
2932 "has_newer": end < total_messages,
2933 "has_older": offset > 0,
2934 "newer_items": normalized_item_count(&session.messages[end..]),
2935 "offset": offset,
2936 "older_items": normalized_item_count(&session.messages[..offset]),
2937 "returned": end.saturating_sub(offset),
2938 "total_messages": total_messages,
2939 }
2940 })
2941}
2942
2943fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
2944 messages
2945 .iter()
2946 .map(|message| {
2947 let conversation = usize::from(
2948 matches!(message.role, Role::Assistant | Role::User)
2949 && message_has_content(message),
2950 );
2951 let tool_result =
2952 usize::from(message.role == Role::Tool && message_has_content(message));
2953 conversation + tool_result + message.tool_calls().len()
2954 })
2955 .sum()
2956}
2957
2958fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
2959 let mut conversational = session.messages.iter().filter(|message| {
2960 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
2961 });
2962 let first_message = conversational.clone().next();
2963 let last_message = conversational.next_back();
2964 let mut assistant = session
2965 .messages
2966 .iter()
2967 .filter(|message| message.role == Role::Assistant && message_has_content(message));
2968 let first_assistant_message = assistant.clone().next();
2969 let last_assistant_message = assistant.next_back();
2970 let end_of_turn = session
2971 .messages
2972 .iter()
2973 .rev()
2974 .find(|message| message.role != Role::System)
2975 .is_some_and(|message| {
2976 message.role == Role::Assistant
2977 && message_has_content(message)
2978 && message.tool_calls().is_empty()
2979 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
2981 });
2982 let project = |message: Option<&crate::ChatMessage>| {
2983 message.map(|message| project_inline_media(message_json(message), options))
2984 };
2985 json!({
2986 "end_of_turn": end_of_turn,
2987 "first_assistant_message": project(first_assistant_message),
2988 "first_message": project(first_message),
2989 "last_assistant_message": project(last_assistant_message),
2990 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
2991 "last_message": project(last_message),
2992 })
2993}
2994
2995fn message_has_content(message: &crate::ChatMessage) -> bool {
2996 message
2997 .content
2998 .as_deref()
2999 .is_some_and(|content| !content.trim().is_empty())
3000 || message
3001 .content_parts
3002 .as_ref()
3003 .is_some_and(|parts| !parts.is_empty())
3004}
3005
3006fn message_text(message: &crate::ChatMessage) -> String {
3007 if let Some(content) = &message.content {
3008 return content.clone();
3009 }
3010 message
3011 .content_parts
3012 .as_ref()
3013 .into_iter()
3014 .flatten()
3015 .filter_map(|part| part.get("text").and_then(Value::as_str))
3016 .collect::<Vec<_>>()
3017 .join("\n")
3018}
3019
3020fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3021 let (offset, end) = projected_message_window(session.messages.len(), options);
3022 let messages = session.messages[offset..end]
3023 .iter()
3024 .map(|message| project_inline_media(message_json(message), options))
3025 .collect::<Vec<_>>();
3026 let subagents = if options.include_subagents.unwrap_or(true) {
3027 let subagent_options = SessionLoadOptions {
3032 message_limit: None,
3033 message_offset: None,
3034 message_tail: None,
3035 ..options.clone()
3036 };
3037 session
3038 .subagents
3039 .iter()
3040 .map(|subagent| projected_session_json(subagent, &subagent_options))
3041 .collect::<Vec<_>>()
3042 } else {
3043 Vec::new()
3044 };
3045 json!({
3046 "source": match session.meta.source {
3047 SessionSource::ClaudeCode => "claude_code",
3048 SessionSource::Codex => "codex",
3049 SessionSource::Gemini => "gemini",
3050 SessionSource::Goose => "goose",
3051 SessionSource::Grok => "grok",
3052 SessionSource::Native => "native",
3053 SessionSource::OpenClaw => "openclaw",
3054 SessionSource::Hermes => "hermes",
3055 SessionSource::OpenCode => "opencode",
3056 SessionSource::Pi => "pi",
3057 },
3058 "session_id": session.meta.session_id,
3059 "ended_at": session.meta.ended_at,
3060 "end_reason": session.meta.end_reason,
3061 "model": session.meta.model,
3062 "cwd": session.meta.cwd,
3063 "system_prompt": session.meta.system_prompt,
3064 "agent_id": session.meta.agent_id,
3065 "parent_tool_use_id": session.meta.parent_tool_use_id,
3066 "lineage": session.meta.lineage,
3067 "messages": messages,
3068 "subagents": subagents,
3069 "raw_record_count": session.raw.len(),
3070 "parse_error_lines": session.parse_error_lines,
3071 })
3072}
3073
3074fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3075 if let Some(tail) = options.message_tail {
3076 return (total.saturating_sub(tail), total);
3077 }
3078 let offset = options.message_offset.unwrap_or(0).min(total);
3079 let end = options
3080 .message_limit
3081 .map(|limit| offset.saturating_add(limit).min(total))
3082 .unwrap_or(total);
3083 (offset, end)
3084}
3085
3086fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3087 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3088 return message;
3089 };
3090 for part in parts {
3091 let Some(url) = part
3092 .get("image_url")
3093 .and_then(|image| image.get("url"))
3094 .and_then(Value::as_str)
3095 else {
3096 continue;
3097 };
3098 let Some(rest) = url.strip_prefix("data:") else {
3099 continue;
3100 };
3101 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3102 continue;
3103 };
3104 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3105 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3106 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3107 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3108 || options
3109 .max_inline_media_bytes
3110 .is_some_and(|limit| decoded_bytes > limit);
3111 if should_elide {
3112 *part = json!({
3113 "type": "media_reference",
3114 "media_type": media_type,
3115 "encoding": "base64",
3116 "encoded_bytes": encoded.len(),
3117 "decoded_bytes": decoded_bytes,
3118 "omitted": true,
3119 });
3120 }
3121 }
3122 message
3123}
3124
3125#[derive(Deserialize)]
3126struct LocatorParams {
3127 locator: SessionLocator,
3128 #[serde(default)]
3141 fidelity: Option<Fidelity>,
3142 #[serde(default)]
3145 view: Option<SessionReadView>,
3146}
3147
3148#[derive(Deserialize)]
3149struct SessionReadView {
3150 #[serde(default)]
3153 tail_messages: Option<usize>,
3154 #[serde(default)]
3157 include_subagents: bool,
3158 #[serde(default)]
3160 display_history: bool,
3161 #[serde(default)]
3164 max_message_chars: Option<usize>,
3165}
3166
3167impl LocatorParams {
3168 fn read_fidelity(&self) -> Fidelity {
3169 self.fidelity.unwrap_or(Fidelity::Semantic)
3170 }
3171
3172 fn include_subagents(&self) -> bool {
3173 self.view
3174 .as_ref()
3175 .map(|view| view.include_subagents)
3176 .unwrap_or(true)
3177 }
3178
3179 fn tail_messages(&self) -> Option<usize> {
3180 self.view
3181 .as_ref()
3182 .and_then(|view| view.tail_messages)
3183 .map(|limit| limit.clamp(1, 5_000))
3184 }
3185
3186 fn display_history(&self) -> bool {
3187 self.view.as_ref().is_some_and(|view| view.display_history)
3188 }
3189
3190 fn max_message_chars(&self) -> Option<usize> {
3191 self.view
3192 .as_ref()
3193 .and_then(|view| view.max_message_chars)
3194 .map(|limit| limit.clamp(256, 64_000))
3195 }
3196
3197 fn bound_session(&self, session: &mut Session) {
3198 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3199 }
3200}
3201
3202#[derive(Debug, Clone, Copy, Default, Deserialize)]
3203#[serde(rename_all = "snake_case")]
3204enum InlineMediaMode {
3205 #[default]
3206 Full,
3207 Metadata,
3208}
3209
3210#[derive(Debug, Clone, Default, Deserialize)]
3211#[serde(default)]
3212struct SessionLoadOptions {
3213 include_subagents: Option<bool>,
3214 inline_media: InlineMediaMode,
3215 max_inline_media_bytes: Option<usize>,
3216 message_limit: Option<usize>,
3217 message_offset: Option<usize>,
3218 message_tail: Option<usize>,
3219}
3220
3221impl SessionLoadOptions {
3222 fn validate(&self) -> std::result::Result<(), ServiceError> {
3223 if self.message_tail.is_some()
3224 && (self.message_limit.is_some() || self.message_offset.is_some())
3225 {
3226 return Err(ServiceError::InvalidParams(
3227 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3228 .into(),
3229 ));
3230 }
3231 Ok(())
3232 }
3233}
3234
3235#[derive(Deserialize)]
3236struct LoadSessionParams {
3237 #[serde(flatten)]
3238 read: LocatorParams,
3239 #[serde(default)]
3240 options: Option<SessionLoadOptions>,
3241}
3242
3243#[derive(Deserialize)]
3244struct UnfollowParams {
3245 subscription: String,
3246}
3247
3248#[derive(Debug, Deserialize)]
3249#[serde(deny_unknown_fields)]
3250struct IndexResizeParams {
3251 subscription: String,
3252 limit: usize,
3253}
3254
3255#[derive(Deserialize)]
3256struct ActivitySubscribeParams {
3257 locators: Vec<SessionLocator>,
3258 #[serde(default)]
3259 homes: crate::HarnessHomes,
3260}
3261
3262#[derive(Deserialize)]
3263struct MessageSessionParams {
3264 locator: SessionLocator,
3265 text: String,
3266 #[serde(default)]
3270 from_name: Option<String>,
3271 #[serde(default)]
3273 in_reply_to: Option<String>,
3274 #[serde(default)]
3276 notify_when_idle: bool,
3277 #[serde(default)]
3282 as_user: bool,
3283 #[serde(default)]
3286 homes: crate::HarnessHomes,
3287}
3288
3289#[derive(Deserialize)]
3290#[serde(deny_unknown_fields)]
3291struct InboxParams {
3292 #[serde(default)]
3294 from_name: Option<String>,
3295 #[serde(default)]
3297 address: Option<String>,
3298 #[serde(default)]
3300 all: bool,
3301}
3302
3303#[derive(Deserialize)]
3304#[serde(deny_unknown_fields)]
3305struct HarnessSettingsParams {
3306 harness: String,
3307}
3308
3309#[derive(Deserialize)]
3310#[serde(deny_unknown_fields)]
3311struct ConfigureHarnessParams {
3312 harness: String,
3313 #[serde(default)]
3314 changes: Vec<crate::HarnessSettingChange>,
3315 #[serde(default)]
3316 expected_revision: Option<String>,
3317}
3318
3319fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3320 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3321 Ok(report) => (
3322 serde_json::to_value(report).unwrap_or(Value::Null),
3323 Value::Null,
3324 ),
3325 Err(error) => (
3326 Value::Null,
3327 Value::String(format!(
3328 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3329 )),
3330 ),
3331 }
3332}
3333
3334async fn message_live_session(params: &MessageSessionParams) -> Value {
3348 use crate::mail_route::{Delivered, NoDoor, Refused};
3349 let (inbound_controls, inbound_controls_error) =
3350 claude_inbound_controls_or_error(¶ms.homes);
3351 let refused = |reason: &str, message: String| {
3352 json!({
3353 "delivered_to_bus": false,
3354 "refusal": {"reason": reason, "message": message},
3355 "inbound_controls": inbound_controls,
3356 "inbound_controls_error": inbound_controls_error,
3357 })
3358 };
3359 if params.text.trim().is_empty() {
3360 return refused(
3361 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3362 "refusing to deliver an empty message".into(),
3363 );
3364 }
3365 let sender = match operator_address(params.from_name.as_deref()) {
3366 Ok(sender) => sender,
3367 Err(message) => return refused("invalid_sender", message),
3368 };
3369 let receiver = match crate::mailbox::MailAddress::new(
3370 crate::mailbox::local_machine_name(),
3371 params.locator.harness.as_str(),
3372 ¶ms.locator.session_id,
3373 ) {
3374 Ok(receiver) => receiver,
3375 Err(error) => return refused("delivery_failed", error.to_string()),
3376 };
3377 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3379 && crate::runtime_mail::controlled_runtime(
3380 HarnessId::CLAUDE_CODE,
3381 ¶ms.locator.session_id,
3382 )
3383 .is_none()
3384 {
3385 if let Err(refusal) =
3386 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3387 {
3388 return refused(refusal.reason.as_str(), refusal.message);
3389 }
3390 }
3391 if params.as_user {
3392 return message_as_user(
3393 params,
3394 sender,
3395 receiver,
3396 inbound_controls,
3397 inbound_controls_error,
3398 )
3399 .await;
3400 }
3401 let door = match crate::mail_route::door_for(¶ms.homes, &receiver) {
3402 Ok(door) => door,
3403 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3404 return refused(
3405 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3406 format!(
3407 "no running `{}` session `{}` is reachable; its transcript is persisted only",
3408 params.locator.harness.as_str(),
3409 params.locator.session_id
3410 ),
3411 )
3412 }
3413 };
3414 let mut envelope = match crate::mailbox::Envelope::new(
3415 sender.clone(),
3416 format!("{}@{}", sender.session_id, sender.machine),
3417 crate::mailbox::MailKind::Peer,
3418 crate::mailbox::ReplyVia::Command,
3419 params.text.clone(),
3420 ) {
3421 Ok(envelope) => envelope,
3422 Err(error) => return refused("delivery_failed", error.to_string()),
3423 };
3424 envelope.in_reply_to = params.in_reply_to.clone();
3425 let how = match crate::mail_route::deliver(
3426 &envelope,
3427 &receiver,
3428 &door,
3429 true,
3430 params.notify_when_idle,
3431 )
3432 .await
3433 {
3434 Err(detail) => {
3435 return refused(
3436 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3437 detail,
3438 )
3439 }
3440 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3441 Ok(Err(Refused::TooLong(bytes))) => {
3442 return refused(
3443 "too_long",
3444 format!(
3445 "the message is {bytes} bytes; the limit is {}",
3446 crate::mail_route::MAX_RELAYED_BYTES
3447 ),
3448 )
3449 }
3450 Ok(Ok(delivered)) => match delivered {
3451 Delivered::Steered => "steered",
3452 Delivered::Started => "started",
3453 Delivered::Native { busy: true } => "next_tool_call",
3454 Delivered::Native { busy: false } => "started",
3455 Delivered::Hooked => "hook",
3456 Delivered::Queued => "queued",
3457 Delivered::Stored => "stored",
3458 Delivered::Operator => "filed",
3459 },
3460 };
3461 json!({
3462 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3463 "message_id": envelope.id,
3464 "reply_to": sender.to_string(),
3465 "target": {
3466 "harness": params.locator.harness.as_str(),
3467 "session_id": params.locator.session_id,
3468 "name": match &door {
3469 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
3470 _ => None,
3471 },
3472 },
3473 "delivery": {"door": door.name(), "how": how},
3474 "inbound_controls": inbound_controls,
3475 "inbound_controls_error": inbound_controls_error,
3476 })
3477}
3478
3479async fn message_as_user(
3483 params: &MessageSessionParams,
3484 sender: crate::mailbox::MailAddress,
3485 receiver: crate::mailbox::MailAddress,
3486 inbound_controls: Value,
3487 inbound_controls_error: Value,
3488) -> Value {
3489 let envelope = match crate::mailbox::Envelope::new(
3490 sender.clone(),
3491 format!("{}@{}", sender.session_id, sender.machine),
3492 crate::mailbox::MailKind::User,
3493 crate::mailbox::ReplyVia::None,
3494 params.text.clone(),
3495 ) {
3496 Ok(envelope) => envelope,
3497 Err(error) => {
3498 return json!({
3499 "delivered_to_bus": false,
3500 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
3501 "inbound_controls": inbound_controls,
3502 "inbound_controls_error": inbound_controls_error,
3503 })
3504 }
3505 };
3506 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
3507 Ok(turn) => json!({
3508 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
3509 "message_id": envelope.id,
3510 "target": {
3511 "harness": params.locator.harness.as_str(),
3512 "session_id": params.locator.session_id,
3513 },
3514 "delivery": {
3515 "door": match turn {
3516 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
3517 _ => "pane",
3518 },
3519 "how": turn.as_str(),
3520 },
3521 "inbound_controls": inbound_controls,
3522 "inbound_controls_error": inbound_controls_error,
3523 }),
3524 Err(message) => json!({
3525 "delivered_to_bus": false,
3526 "refusal": {"reason": "no_user_door", "message": message},
3527 "inbound_controls": inbound_controls,
3528 "inbound_controls_error": inbound_controls_error,
3529 }),
3530 }
3531}
3532
3533#[derive(Deserialize)]
3534#[serde(deny_unknown_fields)]
3535struct ActivityUnderParams {
3536 pids: Vec<u32>,
3538 #[serde(default)]
3539 homes: crate::HarnessHomes,
3540}
3541
3542async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
3546 let params = decode::<ActivityUnderParams>(params)?;
3547 if params.pids.len() > 1024 {
3548 return Err(ServiceError::InvalidParams(
3549 "sessions.activity_under accepts at most 1024 pids".into(),
3550 ));
3551 }
3552 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
3553 .await
3554 .map_err(ServiceError::Sdk)?;
3555 Ok(json!({
3556 "activities": found
3557 .into_iter()
3558 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
3559 .collect::<Vec<_>>(),
3560 }))
3561}
3562
3563fn operator_address(
3566 name: Option<&str>,
3567) -> std::result::Result<crate::mailbox::MailAddress, String> {
3568 let name: String = name
3569 .unwrap_or("supercode")
3570 .trim()
3571 .chars()
3572 .map(|character| {
3573 if character.is_whitespace() || character == '@' {
3574 '-'
3575 } else {
3576 character
3577 }
3578 })
3579 .collect();
3580 if name.is_empty() {
3581 return Err("from_name must not be empty".into());
3582 }
3583 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
3584 .map_err(|error| error.to_string())
3585}
3586
3587fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
3591 let address = match (¶ms.address, ¶ms.from_name) {
3592 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
3593 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
3594 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
3595 (Some(_), Some(_)) => {
3596 return Err(ServiceError::InvalidParams(
3597 "sessions.inbox takes from_name or address, not both".into(),
3598 ))
3599 }
3600 };
3601 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
3602 let mailbox =
3603 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
3604 let claimed = mailbox.claim_unread().map_err(operation)?;
3605 let mut messages: Vec<Value> = Vec::new();
3606 if params.all {
3607 for stored in mailbox.list().map_err(operation)? {
3608 if stored.state == crate::mailbox::MailState::Read {
3609 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3610 }
3611 }
3612 }
3613 for stored in &claimed {
3614 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3615 }
3616 for stored in &claimed {
3617 mailbox.acknowledge(stored).map_err(operation)?;
3618 }
3619 Ok(json!({"address": address.to_string(), "messages": messages}))
3620}
3621
3622#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3627struct FollowedSource {
3628 harness: String,
3629 session_id: String,
3630 reported: Option<String>,
3631}
3632
3633#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3634struct ActivitySubscription {
3635 locators: Vec<SessionLocator>,
3636 homes: crate::HarnessHomes,
3637 reported: BTreeMap<(String, String), crate::SessionActivity>,
3638}
3639
3640fn live_descriptor_value(
3647 session: &SessionDescriptor,
3648 doors: &crate::mail_route::LiveSessions,
3649) -> std::result::Result<Value, ServiceError> {
3650 let mut value = serde_json::to_value(session)
3651 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3652 if value.get("title").is_none_or(Value::is_null) {
3654 if let Some(live) = doors.all().iter().find(|live| {
3655 live.address.harness == session.locator.harness.as_str()
3656 && live.address.session_id == session.locator.session_id
3657 }) {
3658 let name = live.name.split('@').next().unwrap_or(&live.name);
3659 if !name.is_empty() {
3660 value["title"] = json!(name);
3661 }
3662 }
3663 }
3664 if let Some(workspace) = &session.cwd {
3665 let source = LiveRuntimeSource {
3666 harness: session.locator.harness.as_str().to_string(),
3667 session_id: session.locator.session_id.clone(),
3668 workspace: workspace.clone(),
3669 };
3670 if let Some(endpoint) = discover_live_runtime(&source)
3671 .map_err(|error| ServiceError::Operation(error.to_string()))?
3672 {
3673 value["live_endpoint"] = json!(endpoint.as_str());
3674 }
3675 }
3676 if let Some(door) = doors.door(
3680 session.locator.harness.as_str(),
3681 &session.locator.session_id,
3682 ) {
3683 value["delivery"] = json!(door);
3684 }
3685 let waiting = doors.all().iter().any(|live| {
3687 live.address.harness == session.locator.harness.as_str()
3688 && live.address.session_id == session.locator.session_id
3689 && live.status == "waiting"
3690 });
3691 if waiting && session.locator.harness.as_str() == "claude-code" {
3692 if let Some(request) = crate::mail_route::pending_request(session.locator.storage.path()) {
3693 value["pending_request"] = request;
3694 }
3695 }
3696 Ok(value)
3697}
3698
3699fn live_index_changes(
3700 changes: Vec<crate::session_index::SessionIndexChange>,
3701 homes: &HarnessHomes,
3702) -> std::result::Result<Vec<Value>, ServiceError> {
3703 use crate::session_index::SessionIndexChange;
3704 let doors = crate::mail_route::LiveSessions::read(homes);
3705 changes
3706 .into_iter()
3707 .map(|change| match change {
3708 SessionIndexChange::Added { descriptor } => Ok(json!({
3709 "kind": "added",
3710 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3711 })),
3712 SessionIndexChange::Updated { descriptor } => Ok(json!({
3713 "kind": "updated",
3714 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3715 })),
3716 SessionIndexChange::Removed { key } => Ok(json!({
3717 "kind": "removed",
3718 "key": key,
3719 })),
3720 })
3721 .collect()
3722}
3723
3724fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
3725 use crate::{SessionPresence, SessionTurnState};
3726 match (activity.presence, activity.turn) {
3727 (SessionPresence::Persisted, _) => None,
3728 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
3729 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
3730 (SessionPresence::Running, SessionTurnState::Unknown)
3734 if activity.evidence.native_state.is_none() =>
3735 {
3736 None
3737 }
3738 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
3739 }
3740}
3741
3742#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
3743#[serde(rename_all = "kebab-case")]
3744enum TransferFormat {
3745 ClaudeCode,
3746 Codex,
3747 #[serde(rename = "opencode", alias = "open-code")]
3748 OpenCode,
3749 Pi,
3750 Grok,
3751 Gemini,
3752 Goose,
3753 Hermes,
3757}
3758
3759impl TransferFormat {
3760 fn id(self) -> &'static str {
3761 match self {
3762 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
3763 Self::Codex => HarnessId::CODEX,
3764 Self::OpenCode => HarnessId::OPENCODE,
3765 Self::Pi => HarnessId::PI,
3766 Self::Grok => HarnessId::GROK,
3767 Self::Gemini => HarnessId::GEMINI,
3768 Self::Goose => HarnessId::GOOSE,
3769 Self::Hermes => HarnessId::HERMES,
3770 }
3771 }
3772}
3773
3774impl From<TransferFormat> for SessionFormat {
3775 fn from(value: TransferFormat) -> Self {
3776 match value {
3777 TransferFormat::ClaudeCode => Self::ClaudeCode,
3778 TransferFormat::Codex => Self::Codex,
3779 TransferFormat::OpenCode => Self::OpenCode,
3780 TransferFormat::Pi => Self::Pi,
3781 TransferFormat::Grok => Self::Grok,
3782 TransferFormat::Gemini => Self::Gemini,
3783 TransferFormat::Goose => Self::Goose,
3784 TransferFormat::Hermes => Self::Codex,
3786 }
3787 }
3788}
3789
3790#[derive(Deserialize)]
3791struct ImportSessionParams {
3792 source_harness: TransferFormat,
3793 content: String,
3794}
3795
3796#[derive(Deserialize)]
3797struct ExportSessionParams {
3798 locator: SessionLocator,
3799 target_harness: TransferFormat,
3800}
3801
3802#[derive(Deserialize)]
3803struct ReduceSessionParams {
3804 locator: SessionLocator,
3805 target_harness: TransferFormat,
3806 #[serde(default = "default_keep_last")]
3807 keep_last: usize,
3808}
3809
3810fn default_keep_last() -> usize {
3811 6
3812}
3813
3814#[derive(Deserialize)]
3815struct BranchSessionParams {
3816 locator: SessionLocator,
3817 #[serde(default)]
3818 target_harness: Option<TransferFormat>,
3819}
3820
3821#[derive(Deserialize)]
3822struct HandoffSessionParams {
3823 locator: SessionLocator,
3824 target_harness: TransferFormat,
3825 #[serde(default)]
3826 cwd: Option<PathBuf>,
3827}
3828
3829#[derive(Deserialize)]
3830struct MaterializeSessionParams {
3831 artifact: crate::native_materialize::MaterializeArtifact,
3832 cwd: PathBuf,
3833 #[serde(default)]
3835 homes: HarnessHomes,
3836}
3837
3838#[derive(Debug, Clone, Copy, Default, Deserialize)]
3841#[serde(rename_all = "snake_case")]
3842enum ResumePolicy {
3843 Default,
3844 #[default]
3845 Yolo,
3846}
3847
3848#[derive(Deserialize)]
3849struct ResumeInstructionsParams {
3850 locator: SessionLocator,
3851 #[serde(default)]
3852 cwd: Option<PathBuf>,
3853 #[serde(default)]
3854 policy: ResumePolicy,
3855}
3856
3857#[derive(Deserialize)]
3859struct WorkflowLoadParams {
3860 from: crate::workflow_doors::WorkflowHarness,
3861 home: PathBuf,
3862}
3863
3864#[derive(Deserialize)]
3867struct OrchestrationLoadParams {
3868 root: PathBuf,
3869 #[serde(default)]
3870 flavor: crate::orchestration_doors::HomeFlavor,
3871}
3872
3873#[derive(Deserialize)]
3876struct OrchestrationSaveParams {
3877 root: PathBuf,
3878 orchestration: crate::orchestration::Orchestration,
3879 #[serde(default)]
3880 vault: BTreeMap<String, String>,
3881}
3882
3883#[derive(Deserialize)]
3885struct OrchestrationCompileParams {
3886 from: crate::orchestration_doors::OrchestrationHarness,
3887 home: PathBuf,
3888}
3889
3890#[derive(Deserialize)]
3894struct OrchestrationDecompileParams {
3895 to: crate::orchestration_doors::OrchestrationHarness,
3896 orchestration: crate::orchestration::Orchestration,
3897 source: PathBuf,
3898 #[serde(default)]
3899 source_flavor: crate::orchestration_doors::SourceFlavor,
3900 dest: PathBuf,
3901 #[serde(default)]
3902 vault: BTreeMap<String, String>,
3903}
3904
3905#[derive(Deserialize)]
3908struct OrchestrationImportParams {
3909 from: crate::orchestration_doors::OrchestrationHarness,
3910 home: PathBuf,
3911 into: PathBuf,
3912}
3913
3914#[derive(Deserialize)]
3917struct OrchestrationExportParams {
3918 to: crate::orchestration_doors::OrchestrationHarness,
3919 root: PathBuf,
3920 dest: PathBuf,
3921}
3922
3923#[derive(Deserialize)]
3925struct JobsGetParams {
3926 harness: String,
3927 id: String,
3928 #[serde(default)]
3929 homes: crate::HarnessHomes,
3930}
3931
3932fn mutate_job(
3940 verb: crate::jobs_control::JobVerb,
3941 params: Value,
3942) -> std::result::Result<Value, ServiceError> {
3943 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
3944 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
3945 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
3946 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3947}
3948
3949fn mutate_skill(
3956 verb: crate::skills_control::SkillVerb,
3957 params: Value,
3958) -> std::result::Result<Value, ServiceError> {
3959 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
3960 if !crate::skills_control::supports_skill_control(&mutation.harness) {
3961 return Err(ServiceError::UnsupportedAction(format!(
3962 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
3963 mutation.harness,
3964 verb.as_str(),
3965 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
3966 )));
3967 }
3968 let outcome =
3969 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
3970 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3971}
3972
3973fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
3975 match error {
3976 crate::skills_control::SkillControlError::Unsupported(message) => {
3977 ServiceError::UnsupportedAction(message)
3978 }
3979 crate::skills_control::SkillControlError::Invalid(message) => {
3980 ServiceError::InvalidParams(message)
3981 }
3982 crate::skills_control::SkillControlError::Failed(message) => {
3983 ServiceError::Operation(message)
3984 }
3985 }
3986}
3987
3988fn mutate_profile(
3996 verb: crate::profiles_control::ProfileVerb,
3997 params: Value,
3998) -> std::result::Result<Value, ServiceError> {
3999 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4000 let outcome =
4001 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4002 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4003}
4004
4005fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4007 match error {
4008 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4009 ServiceError::UnsupportedAction(message)
4010 }
4011 crate::profiles_control::ProfileControlError::Invalid(message) => {
4012 ServiceError::InvalidParams(message)
4013 }
4014 crate::profiles_control::ProfileControlError::Failed(message) => {
4015 ServiceError::Operation(message)
4016 }
4017 }
4018}
4019
4020fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4024 match error {
4025 crate::jobs_control::JobControlError::Unsupported(message) => {
4026 ServiceError::UnsupportedAction(message)
4027 }
4028 crate::jobs_control::JobControlError::Invalid(message) => {
4029 ServiceError::InvalidParams(message)
4030 }
4031 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4032 }
4033}
4034
4035fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4040 match error {
4041 crate::SessionControlError::Unsupported(message) => {
4042 ServiceError::UnsupportedAction(message)
4043 }
4044 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4045 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4046 }
4047}
4048
4049fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4054 if crate::jobs::supports_jobs(harness) {
4055 return Ok(());
4056 }
4057 Err(ServiceError::UnsupportedAction(format!(
4058 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4059 crate::jobs::JOB_HARNESSES.join(", ")
4060 )))
4061}
4062
4063#[derive(Deserialize)]
4065struct RunsGetParams {
4066 harness: String,
4067 id: String,
4068 #[serde(default)]
4069 homes: crate::HarnessHomes,
4070}
4071
4072fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4077 if crate::runs::supports_runs(harness) {
4078 return Ok(());
4079 }
4080 Err(ServiceError::UnsupportedAction(format!(
4081 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4082 crate::runs::RUN_HARNESSES.join(", ")
4083 )))
4084}
4085
4086#[derive(Serialize)]
4087struct SessionArtifact {
4088 source_harness: HarnessId,
4089 target_harness: &'static str,
4090 session_id: Option<String>,
4091 content: String,
4092 suggested_filename: String,
4093 files: Vec<SessionArtifactFile>,
4094 fidelity: Fidelity,
4095 residue: Vec<String>,
4096}
4097
4098#[derive(Serialize)]
4099struct SessionArtifactFile {
4100 path: String,
4101 content: String,
4102 role: ArtifactFileRole,
4103}
4104
4105#[derive(Serialize)]
4106#[serde(rename_all = "snake_case")]
4107enum ArtifactFileRole {
4108 Primary,
4109 Subagent,
4110 Bundle,
4111 SourceRecovery,
4112}
4113
4114#[derive(Serialize)]
4115struct StructuredLaunch {
4116 cwd: PathBuf,
4117 program: String,
4118 arguments: Vec<String>,
4119 env: BTreeMap<String, String>,
4120}
4121
4122struct HandoffInstructions {
4123 launch: StructuredLaunch,
4124 materialize: Option<StructuredLaunch>,
4125 requires_materialization: bool,
4126 note: String,
4127}
4128
4129#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4130#[serde(rename_all = "snake_case")]
4131enum HarnessProbeLevel {
4132 #[default]
4133 Passive,
4134 Handshake,
4135}
4136
4137#[derive(Default, Deserialize)]
4138#[serde(default)]
4139struct HarnessInventoryParams {
4140 harness: Option<HarnessId>,
4141 harnesses: Vec<HarnessId>,
4142 workspace: Option<PathBuf>,
4143 probe: HarnessProbeLevel,
4144 include_sessions: bool,
4145 skip_versions: bool,
4147}
4148
4149#[derive(Deserialize)]
4150struct HarnessAuthenticationParams {
4151 harness: HarnessId,
4152}
4153
4154#[derive(Deserialize)]
4155struct BeginHarnessAuthenticationParams {
4156 harness: HarnessId,
4157 #[serde(default = "local_browser_authentication_environment")]
4158 environment: crate::HarnessAuthenticationEnvironment,
4159 #[serde(default)]
4160 method: Option<crate::HarnessAuthenticationMethodId>,
4161 #[serde(default)]
4162 cwd: Option<PathBuf>,
4163}
4164
4165fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4166 crate::HarnessAuthenticationEnvironment::LocalBrowser
4167}
4168
4169#[derive(Serialize)]
4170struct HarnessInventoryReport {
4171 probe: HarnessProbeLevel,
4172 workspace: Option<PathBuf>,
4173 harnesses: Vec<LocalHarness>,
4174}
4175
4176#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4177#[serde(rename_all = "snake_case")]
4178enum HarnessAuthState {
4179 Ready,
4180 Configured,
4181 Required,
4182 Unknown,
4183}
4184
4185#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4186#[serde(rename_all = "snake_case")]
4187enum HarnessRuntimeState {
4188 Ready,
4189 Degraded,
4190 Unavailable,
4191}
4192
4193#[derive(Serialize)]
4194struct HarnessSessionCounts {
4195 global: Option<usize>,
4196 workspace: Option<usize>,
4197}
4198
4199#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4210#[serde(rename_all = "snake_case")]
4211pub enum GatewayState {
4212 Up,
4213 Down,
4214 Unknown,
4215}
4216
4217#[derive(Debug, Clone, Serialize)]
4219pub struct GatewayHealth {
4220 pub state: GatewayState,
4221 #[serde(skip_serializing_if = "Option::is_none")]
4226 pub endpoint: Option<String>,
4227 #[serde(skip_serializing_if = "Option::is_none")]
4228 pub version: Option<String>,
4229 pub evidence: String,
4231 pub checked_at_ms: u64,
4232}
4233
4234fn openclaw_gateway_endpoint(home: &Path) -> String {
4238 let config_path = home.join(".openclaw/openclaw.json");
4239 let gateway = std::fs::read_to_string(&config_path)
4240 .ok()
4241 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4242 .and_then(|config| config.get("gateway").cloned());
4243 if let Some(url) = gateway
4244 .as_ref()
4245 .and_then(|gateway| gateway.get("url"))
4246 .and_then(serde_json::Value::as_str)
4247 {
4248 return url.to_string();
4249 }
4250 let port = gateway
4251 .as_ref()
4252 .and_then(|gateway| gateway.get("port"))
4253 .and_then(serde_json::Value::as_u64)
4254 .unwrap_or(18789);
4255 format!("ws://127.0.0.1:{port}")
4256}
4257
4258fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4265 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4266 let output = std::process::Command::new(&program)
4267 .args(["gateway", "status"])
4268 .stdin(std::process::Stdio::null())
4269 .output()
4270 .ok()?;
4271 let text = format!(
4272 "{}{}",
4273 String::from_utf8_lossy(&output.stdout),
4274 String::from_utf8_lossy(&output.stderr)
4275 );
4276 let verdict = text.lines().find_map(|line| {
4277 let l = line.trim();
4278 if l.contains("supervised by launchd (PID")
4279 || l.contains("supervised by systemd (PID")
4280 || l.contains("Gateway is running")
4281 || l.contains("process is running")
4282 {
4283 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4284 } else if l.contains("not running") || l.contains("not installed") {
4285 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4286 } else {
4287 None
4288 }
4289 });
4290 verdict
4291}
4292
4293fn gateway_health(
4294 id: &str,
4295 installed: bool,
4296 running: Option<&RunningInstance>,
4297 version: Option<&str>,
4298) -> GatewayHealth {
4299 let checked_at_ms = now_epoch_ms();
4300 let home = supercode_interchange::user_home()
4301 .map(std::path::PathBuf::into_os_string)
4302 .map(PathBuf::from);
4303 match id {
4304 HarnessId::HERMES | HarnessId::OPENCLAW => {
4305 let endpoint = (id == HarnessId::OPENCLAW)
4306 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4307 .flatten();
4308 let (state, evidence) = match running {
4309 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4310 None if !installed => (
4311 GatewayState::Unknown,
4312 format!("`{id}` is not installed; no gateway to probe"),
4313 ),
4314 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4315 Some((state, evidence)) => (state, evidence),
4318 None => (
4319 GatewayState::Down,
4320 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4321 ),
4322 },
4323 None => (
4324 GatewayState::Down,
4325 format!(
4326 "no TCP listener at {}",
4327 endpoint.as_deref().unwrap_or("the gateway endpoint")
4328 ),
4329 ),
4330 };
4331 GatewayHealth {
4332 state,
4333 endpoint,
4334 version: version.map(str::to_string),
4335 evidence,
4336 checked_at_ms,
4337 }
4338 }
4339 HarnessId::ORCHESTRATOR => {
4346 let root = crate::HarnessHomes::default().orchestrator;
4347 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4348 Some(lease) if lease.is_live() => (
4349 GatewayState::Up,
4350 format!(
4351 "`{}` names pid {} (started {}), which is live",
4352 crate::orchestrator::lock_path(&root).display(),
4353 lease.pid,
4354 lease.started_at
4355 ),
4356 ),
4357 Some(lease) => (
4358 GatewayState::Down,
4359 format!(
4360 "stale lease `{}`: pid {} is gone",
4361 crate::orchestrator::lock_path(&root).display(),
4362 lease.pid
4363 ),
4364 ),
4365 None => (
4366 GatewayState::Down,
4367 format!(
4368 "no lease at `{}`; `supercode orchestrator start` writes one",
4369 crate::orchestrator::lock_path(&root).display()
4370 ),
4371 ),
4372 };
4373 GatewayHealth {
4374 state,
4375 endpoint: None,
4376 version: version.map(str::to_string),
4377 evidence,
4378 checked_at_ms,
4379 }
4380 }
4381 _ => GatewayHealth {
4382 state: GatewayState::Unknown,
4383 endpoint: None,
4384 version: version.map(str::to_string),
4385 evidence: format!("`{id}` runs per session, not as a gateway"),
4386 checked_at_ms,
4387 },
4388 }
4389}
4390
4391#[derive(Debug, Clone, Serialize)]
4392struct RunningInstance {
4393 method: RunningInstanceMethod,
4395 evidence: String,
4397 checked_at_ms: u64,
4399}
4400
4401#[derive(Debug, Clone, Copy, Serialize)]
4402#[serde(rename_all = "snake_case")]
4403enum RunningInstanceMethod {
4404 GatewayConnect,
4407 StoreWalActivity,
4410}
4411
4412fn now_epoch_ms() -> u64 {
4413 std::time::SystemTime::now()
4414 .duration_since(std::time::UNIX_EPOCH)
4415 .map(|elapsed| elapsed.as_millis() as u64)
4416 .unwrap_or(0)
4417}
4418
4419fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4423 let config_path = home.join(".openclaw/openclaw.json");
4424 let text = std::fs::read_to_string(&config_path).ok();
4425 let gateway = text
4426 .as_deref()
4427 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4428 .and_then(|config| config.get("gateway").cloned());
4429 let address = gateway
4430 .as_ref()
4431 .and_then(|gateway| gateway.get("url"))
4432 .and_then(serde_json::Value::as_str)
4433 .and_then(|url| {
4434 url.split("://").nth(1).map(|rest| {
4435 rest.trim_end_matches('/')
4436 .split('/')
4437 .next()
4438 .unwrap_or(rest)
4439 .to_string()
4440 })
4441 })
4442 .unwrap_or_else(|| {
4443 let port = gateway
4444 .as_ref()
4445 .and_then(|gateway| gateway.get("port"))
4446 .and_then(serde_json::Value::as_u64)
4447 .unwrap_or(18789);
4448 format!("127.0.0.1:{port}")
4449 });
4450 let reachable = std::net::TcpStream::connect_timeout(
4451 &address.parse().ok()?,
4452 std::time::Duration::from_millis(400),
4453 )
4454 .is_ok();
4455 reachable.then(|| RunningInstance {
4456 method: RunningInstanceMethod::GatewayConnect,
4457 evidence: format!(
4458 "gateway endpoint {address} accepted a TCP connect (from {})",
4459 config_path.display()
4460 ),
4461 checked_at_ms: now_epoch_ms(),
4462 })
4463}
4464
4465fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4470 let wal = home.join(".hermes/state.db-wal");
4471 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4472 let age_ms = std::time::SystemTime::now()
4473 .duration_since(modified)
4474 .map(|age| age.as_millis() as u64)
4475 .unwrap_or(u64::MAX);
4476 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4477 method: RunningInstanceMethod::StoreWalActivity,
4478 evidence: format!(
4479 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4480 wal.display()
4481 ),
4482 checked_at_ms: now_epoch_ms(),
4483 })
4484}
4485
4486fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4488 let home = supercode_interchange::user_home()
4489 .map(std::path::PathBuf::into_os_string)
4490 .map(PathBuf::from)?;
4491 match id {
4492 HarnessId::OPENCLAW => probe_openclaw_running(&home),
4493 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4494 _ => None,
4495 }
4496}
4497
4498#[derive(Serialize)]
4499struct LocalHarness {
4500 id: HarnessId,
4501 display_name: String,
4502 supported: bool,
4503 installed: bool,
4504 executable: Option<String>,
4505 version: Option<String>,
4506 auth: HarnessAuthState,
4507 runtime: HarnessRuntimeState,
4508 protocol: String,
4509 capabilities: crate::RuntimeCapabilities,
4510 effective_capabilities: crate::RuntimeCapabilities,
4511 sessions: HarnessSessionCounts,
4512 #[serde(skip_serializing_if = "Option::is_none")]
4515 running: Option<RunningInstance>,
4516 gateway: GatewayHealth,
4518 reason: Option<String>,
4519 repair: Option<String>,
4520}
4521
4522#[derive(Clone, Deserialize)]
4523struct RuntimeBackendParams {
4524 harness: HarnessId,
4525 #[serde(default)]
4526 protocol: Option<String>,
4527 #[serde(default)]
4528 launch: Option<RuntimeLaunch>,
4529 #[serde(default)]
4530 base_url: Option<String>,
4531 #[serde(default)]
4532 policy: RuntimePolicy,
4533}
4534
4535#[derive(Debug, Clone, Copy, Default, Deserialize)]
4538#[serde(rename_all = "snake_case")]
4539enum RuntimePolicy {
4540 Default,
4541 #[default]
4542 Yolo,
4543}
4544
4545#[derive(Deserialize)]
4546struct RuntimeStartParams {
4547 #[serde(flatten)]
4548 backend: RuntimeBackendParams,
4549 cwd: PathBuf,
4550 #[serde(default)]
4553 mcp_servers: Vec<crate::McpServerLaunch>,
4554 #[serde(default)]
4556 approval_policy: Option<String>,
4557}
4558
4559#[derive(Deserialize)]
4560struct RuntimeAttachParams {
4561 #[serde(flatten)]
4562 backend: RuntimeBackendParams,
4563 runtime_id: String,
4564 #[serde(default)]
4565 cwd: Option<PathBuf>,
4566 #[serde(default)]
4569 mcp_servers: Vec<crate::McpServerLaunch>,
4570 #[serde(default)]
4572 approval_policy: Option<String>,
4573}
4574
4575#[derive(Deserialize)]
4576struct RuntimeConnectionParams {
4577 connection: String,
4578}
4579
4580#[derive(Deserialize)]
4581struct RuntimeInputParams {
4582 connection: String,
4583 text: String,
4584 #[serde(default)]
4585 image_urls: Vec<String>,
4586}
4587
4588const MAX_RUNTIME_IMAGES: usize = 4;
4589const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
4590const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
4591
4592fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
4593 if image_urls.len() > MAX_RUNTIME_IMAGES {
4594 return Err(ServiceError::InvalidParams(format!(
4595 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
4596 )));
4597 }
4598 let mut total = 0usize;
4599 for url in &image_urls {
4600 if !(url.starts_with("data:image/")
4601 || url.starts_with("https://")
4602 || url.starts_with("http://"))
4603 {
4604 return Err(ServiceError::InvalidParams(
4605 "runtime images must be image data URLs or HTTP(S) URLs".into(),
4606 ));
4607 }
4608 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
4609 return Err(ServiceError::InvalidParams(format!(
4610 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
4611 )));
4612 }
4613 total = total.saturating_add(url.len());
4614 }
4615 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
4616 return Err(ServiceError::InvalidParams(format!(
4617 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
4618 )));
4619 }
4620 Ok(image_urls)
4621}
4622
4623#[derive(Deserialize)]
4624struct RuntimeRespondParams {
4625 connection: String,
4626 request_id: Value,
4627 response: Value,
4628}
4629
4630fn default_reduction_store_root() -> PathBuf {
4631 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
4632 return PathBuf::from(root).join("sessions");
4633 }
4634 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
4635 return PathBuf::from(home).join(".supercode").join("sessions");
4636 }
4637 PathBuf::from(".supercode").join("sessions")
4638}
4639
4640fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
4641 let mut output = String::new();
4642 for message in messages {
4643 output.push_str(
4644 &serde_json::to_string(message)
4645 .map_err(|error| ServiceError::Operation(error.to_string()))?,
4646 );
4647 output.push('\n');
4648 }
4649 Ok(output)
4650}
4651
4652fn parse_messages_jsonl(
4653 content: &str,
4654) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
4655 content
4656 .lines()
4657 .enumerate()
4658 .filter(|(_, line)| !line.trim().is_empty())
4659 .map(|(index, line)| {
4660 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
4661 ServiceError::Operation(format!(
4662 "reduced transcript line {} is invalid: {error}",
4663 index + 1
4664 ))
4665 })
4666 })
4667 .collect()
4668}
4669
4670fn reduced_bootstrap_prompt(
4671 source: &SessionLocator,
4672 target: TransferFormat,
4673 view_jsonl: &str,
4674 sidecar_path: &Path,
4675 reduction_log_path: &Path,
4676) -> String {
4677 format!(
4678 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
4679 \n\
4680 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\
4681 \n\
4682 <supercode-reduced-session source-session=\"{source_id}\">\n\
4683 {view_jsonl}\
4684 </supercode-reduced-session>\n\
4685 \n\
4686 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
4687 source_harness = source.harness.as_str(),
4688 target_harness = target.id(),
4689 sidecar = sidecar_path.display(),
4690 log = reduction_log_path.display(),
4691 source_id = source.session_id,
4692 )
4693}
4694
4695fn session_artifact(
4696 locator: &SessionLocator,
4697 session: &Session,
4698 target: TransferFormat,
4699) -> std::result::Result<SessionArtifact, ServiceError> {
4700 session_artifact_with_id(locator, session, target, None)
4701}
4702
4703fn session_artifact_with_id(
4704 locator: &SessionLocator,
4705 session: &Session,
4706 target: TransferFormat,
4707 target_session_id: Option<&str>,
4708) -> std::result::Result<SessionArtifact, ServiceError> {
4709 let format: SessionFormat = target.into();
4710 let diagonal = format.source() == session.meta.source;
4711 crate::residue_store::store_segments(session);
4712 let has_appended_turns = session
4713 .imported_message_count
4714 .is_some_and(|imported| imported < session.messages.len());
4715 let mut restoration = None;
4716 let content = if let Some(id) = target_session_id {
4717 if diagonal && format != SessionFormat::OpenCode {
4718 session
4719 .to_jsonl_spliced(format, Some(id))
4720 .map_err(operation)?
4721 } else {
4722 let mut rewritten = session.clone();
4723 rewritten.meta.session_id = Some(id.to_string());
4724 rewritten.to_jsonl(format).map_err(operation)?
4725 }
4726 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
4727 session.raw_verbatim()
4728 } else if diagonal {
4729 session.to_jsonl_spliced(format, None).map_err(operation)?
4730 } else {
4731 match session
4734 .restore_residue(format, crate::residue_store::lookup)
4735 .map_err(operation)?
4736 {
4737 Some((content, report)) => {
4738 restoration = Some(report);
4739 content
4740 }
4741 None => session.to_jsonl(format).map_err(operation)?,
4742 }
4743 };
4744 let stem = sanitize_filename(
4745 target_session_id
4746 .or(session.meta.session_id.as_deref())
4747 .unwrap_or(&locator.session_id),
4748 );
4749 let suggested_filename = if diagonal && target == TransferFormat::Grok {
4750 "chat_history.jsonl".to_string()
4751 } else if target == TransferFormat::Goose {
4752 format!("{stem}.goose.json")
4753 } else {
4754 format!("{stem}.{}.jsonl", target.id())
4755 };
4756 let mut files = vec![SessionArtifactFile {
4757 path: suggested_filename.clone(),
4758 content: content.clone(),
4759 role: ArtifactFileRole::Primary,
4760 }];
4761 if target == TransferFormat::ClaudeCode {
4762 let bundle_stem = Path::new(&suggested_filename)
4763 .file_stem()
4764 .and_then(|stem| stem.to_str())
4765 .unwrap_or(&stem);
4766 let mut child_paths = BTreeSet::new();
4767 for (index, subagent) in session.subagents.iter().enumerate() {
4768 let agent_id = subagent
4769 .meta
4770 .agent_id
4771 .as_deref()
4772 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
4773 .map(sanitize_filename)
4774 .filter(|id| !id.is_empty())
4775 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4776 let child_has_appended_turns = subagent
4777 .imported_message_count
4778 .is_some_and(|imported| imported < subagent.messages.len());
4779 let child_content = if target_session_id.is_none()
4780 && subagent.meta.source == SessionSource::ClaudeCode
4781 && subagent.raw_is_verbatim
4782 && !child_has_appended_turns
4783 {
4784 subagent.raw_verbatim()
4785 } else if subagent.meta.source == SessionSource::ClaudeCode {
4786 subagent
4787 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
4788 .map_err(operation)?
4789 } else {
4790 let mut child = subagent.clone();
4791 if let Some(id) = target_session_id {
4792 child.meta.session_id = Some(id.to_string());
4793 }
4794 child
4795 .to_jsonl(SessionFormat::ClaudeCode)
4796 .map_err(operation)?
4797 };
4798 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
4799 if !child_paths.insert(path.clone()) {
4800 return Err(ServiceError::Operation(format!(
4801 "Claude subagent ids collide at artifact path `{path}`"
4802 )));
4803 }
4804 files.push(SessionArtifactFile {
4805 path,
4806 content: child_content,
4807 role: ArtifactFileRole::Subagent,
4808 });
4809 }
4810 }
4811 if diagonal && target == TransferFormat::Grok {
4812 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
4813 }
4814 if !diagonal || !session.raw_is_verbatim {
4815 files.push(SessionArtifactFile {
4816 path: "recovery/source.supercode.jsonl".into(),
4817 content: session.to_native_jsonl(),
4818 role: ArtifactFileRole::SourceRecovery,
4819 });
4820 for (index, subagent) in session.subagents.iter().enumerate() {
4821 let id = subagent
4822 .meta
4823 .agent_id
4824 .as_deref()
4825 .map(sanitize_filename)
4826 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4827 files.push(SessionArtifactFile {
4828 path: format!("recovery/subagents/{id}.supercode.jsonl"),
4829 content: subagent.to_native_jsonl(),
4830 role: ArtifactFileRole::SourceRecovery,
4831 });
4832 }
4833 }
4834 if !diagonal && session.meta.source == SessionSource::Grok {
4835 append_grok_bundle_files(
4836 locator,
4837 "recovery/grok/",
4838 ArtifactFileRole::SourceRecovery,
4839 &mut files,
4840 )?;
4841 }
4842 let (fidelity, residue) = if diagonal
4843 && target_session_id.is_none()
4844 && session.raw_is_verbatim
4845 && !has_appended_turns
4846 {
4847 (Fidelity::ByteLossless, Vec::new())
4848 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
4849 (
4850 Fidelity::ValueLossless,
4851 vec![if target_session_id.is_some() {
4852 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
4853 } else {
4854 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
4855 }],
4856 )
4857 } else {
4858 match restoration {
4859 Some(report) if report.rendered_messages == 0 => (
4860 Fidelity::ByteLossless,
4861 vec![format!(
4862 "restored verbatim from this conversation's {} source records in the residue store",
4863 target.id()
4864 )],
4865 ),
4866 Some(report) => (
4867 Fidelity::Semantic,
4868 vec![format!(
4869 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
4870 report.restored_messages,
4871 report.restored_messages + report.rendered_messages,
4872 report.rendered_messages,
4873 target.id()
4874 )],
4875 ),
4876 None => (
4877 Fidelity::Semantic,
4878 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
4879 ),
4880 }
4881 };
4882 Ok(SessionArtifact {
4883 source_harness: locator.harness.clone(),
4884 target_harness: target.id(),
4885 session_id: target_session_id
4886 .map(str::to_string)
4887 .or_else(|| session.meta.session_id.clone()),
4888 content,
4889 suggested_filename,
4890 files,
4891 fidelity,
4892 residue,
4893 })
4894}
4895
4896fn append_grok_bundle_files(
4897 locator: &SessionLocator,
4898 prefix: &str,
4899 role: ArtifactFileRole,
4900 files: &mut Vec<SessionArtifactFile>,
4901) -> std::result::Result<(), ServiceError> {
4902 let primary = locator.storage.path();
4903 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
4904 return Err(ServiceError::Operation(format!(
4905 "Grok bundle locator must name chat_history.jsonl, got {}",
4906 primary.display()
4907 )));
4908 }
4909 let parent = primary.parent().ok_or_else(|| {
4910 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
4911 })?;
4912 for name in ["summary.json", "updates.jsonl"] {
4913 let path = parent.join(name);
4914 let metadata = match std::fs::symlink_metadata(&path) {
4915 Ok(metadata) => metadata,
4916 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
4917 Err(error) => return Err(ServiceError::Operation(error.to_string())),
4918 };
4919 if metadata.file_type().is_symlink() || !metadata.is_file() {
4920 return Err(ServiceError::Operation(format!(
4921 "refusing non-regular Grok bundle member {}",
4922 path.display()
4923 )));
4924 }
4925 let content = std::fs::read_to_string(&path).map_err(|error| {
4926 ServiceError::Operation(format!(
4927 "Grok bundle member {} is not representable as UTF-8: {error}",
4928 path.display()
4929 ))
4930 })?;
4931 files.push(SessionArtifactFile {
4932 path: format!("{prefix}{name}"),
4933 content,
4934 role: match role {
4935 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
4936 _ => ArtifactFileRole::SourceRecovery,
4937 },
4938 });
4939 }
4940 Ok(())
4941}
4942
4943fn handoff_artifact(
4944 locator: &SessionLocator,
4945 session: &Session,
4946 target: TransferFormat,
4947) -> std::result::Result<SessionArtifact, ServiceError> {
4948 let target_session_id = target_session_id(target);
4949 session_artifact_with_id(locator, session, target, Some(&target_session_id))
4950}
4951
4952fn target_session_id(target: TransferFormat) -> String {
4953 let uuid = generated_session_id();
4954 match target {
4955 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
4956 TransferFormat::ClaudeCode
4957 | TransferFormat::Codex
4958 | TransferFormat::Pi
4959 | TransferFormat::Grok
4960 | TransferFormat::Gemini
4961 | TransferFormat::Goose
4962 | TransferFormat::Hermes => uuid,
4963 }
4964}
4965
4966fn sanitize_filename(value: &str) -> String {
4967 let value = value
4968 .chars()
4969 .map(|character| {
4970 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
4971 character
4972 } else {
4973 '-'
4974 }
4975 })
4976 .collect::<String>();
4977 let value = value.trim_matches('-');
4978 if value.is_empty() {
4979 "session".into()
4980 } else {
4981 value.chars().take(100).collect()
4982 }
4983}
4984
4985fn handoff_instructions(
4986 target: TransferFormat,
4987 session_id: &str,
4988 cwd: &Path,
4989) -> HandoffInstructions {
4990 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
4991 cwd: cwd.to_path_buf(),
4992 program: program.into(),
4993 arguments,
4994 env: BTreeMap::new(),
4995 };
4996 match target {
4997 TransferFormat::ClaudeCode => HandoffInstructions {
4998 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
4999 materialize: None,
5000 requires_materialization: true,
5001 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(),
5002 },
5003 TransferFormat::Hermes => HandoffInstructions {
5004 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5005 materialize: None,
5006 requires_materialization: true,
5007 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(),
5008 },
5009 TransferFormat::Codex => HandoffInstructions {
5010 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5011 materialize: None,
5012 requires_materialization: true,
5013 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5014 },
5015 TransferFormat::OpenCode => HandoffInstructions {
5016 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5017 materialize: Some(launch(
5018 "opencode",
5019 vec!["import".into(), "{artifact_path}".into()],
5020 )),
5021 requires_materialization: true,
5022 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5023 },
5024 TransferFormat::Pi => HandoffInstructions {
5025 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5026 materialize: None,
5027 requires_materialization: true,
5028 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5029 },
5030 TransferFormat::Grok => HandoffInstructions {
5031 launch: launch(
5032 "grok",
5033 vec!["--resume".into(), "{materialized_session_id}".into()],
5034 ),
5035 materialize: None,
5036 requires_materialization: true,
5037 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(),
5038 },
5039 TransferFormat::Gemini => HandoffInstructions {
5040 launch: launch(
5041 "gemini",
5042 vec!["--session-file".into(), "{artifact_path}".into()],
5043 ),
5044 materialize: None,
5045 requires_materialization: true,
5046 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(),
5047 },
5048 TransferFormat::Goose => HandoffInstructions {
5049 launch: launch(
5050 "goose",
5051 vec![
5052 "session".into(),
5053 "--resume".into(),
5054 "--session-id".into(),
5055 "{imported_session_id}".into(),
5056 ],
5057 ),
5058 materialize: Some(launch(
5059 "goose",
5060 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5061 )),
5062 requires_materialization: true,
5063 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(),
5064 },
5065 }
5066}
5067
5068fn resume_launch(
5069 harness: &str,
5070 session_id: &str,
5071 cwd: &Path,
5072 policy: ResumePolicy,
5073) -> std::result::Result<StructuredLaunch, ServiceError> {
5074 let mut arguments = Vec::new();
5075 let program = match harness {
5076 HarnessId::GROK => {
5077 if matches!(policy, ResumePolicy::Yolo) {
5078 if crate::support::self_sandbox_supported() {
5079 arguments.extend(["--sandbox".into(), "workspace".into()]);
5080 }
5081 arguments.push("--always-approve".into());
5082 }
5083 arguments.extend(["--resume".into(), session_id.into()]);
5084 "grok"
5085 }
5086 HarnessId::CODEX => {
5087 arguments.extend(crate::startup_prompts::startup_arguments(
5088 harness,
5089 Some(cwd),
5090 &[],
5091 matches!(policy, ResumePolicy::Yolo),
5092 ));
5093 arguments.extend(["resume".into(), session_id.into()]);
5094 "codex"
5095 }
5096 HarnessId::CLAUDE_CODE => {
5097 arguments.extend(crate::startup_prompts::startup_arguments(
5098 harness,
5099 Some(cwd),
5100 &[],
5101 matches!(policy, ResumePolicy::Yolo),
5102 ));
5103 arguments.extend(["--resume".into(), session_id.into()]);
5104 "claude"
5105 }
5106 HarnessId::GEMINI => {
5107 arguments.extend(crate::startup_prompts::startup_arguments(
5108 harness,
5109 Some(cwd),
5110 &[],
5111 matches!(policy, ResumePolicy::Yolo),
5112 ));
5113 arguments.extend(["--resume".into(), session_id.into()]);
5114 "gemini"
5115 }
5116 HarnessId::GOOSE => {
5117 arguments.extend([
5118 "session".into(),
5119 "--resume".into(),
5120 "--session-id".into(),
5121 session_id.into(),
5122 ]);
5123 "goose"
5124 }
5125 HarnessId::PI => {
5126 arguments.extend(crate::startup_prompts::startup_arguments(
5127 harness,
5128 Some(cwd),
5129 &[],
5130 matches!(policy, ResumePolicy::Yolo),
5131 ));
5132 arguments.extend(["--session".into(), session_id.into()]);
5133 "pi"
5134 }
5135 HarnessId::OPENCODE => {
5136 arguments.extend(["--session".into(), session_id.into()]);
5137 "opencode"
5138 }
5139 HarnessId::SUPERCODE => {
5140 if matches!(policy, ResumePolicy::Yolo) {
5141 arguments.push("--dangerous".into());
5142 }
5143 arguments.extend(["resume".into(), session_id.into()]);
5144 "supercode"
5145 }
5146 other => {
5147 return Err(ServiceError::InvalidParams(format!(
5148 "no structured resume launch is registered for harness `{other}`"
5149 )))
5150 }
5151 };
5152 Ok(StructuredLaunch {
5153 cwd: cwd.to_path_buf(),
5154 env: if program == "grok" {
5155 crate::support::grok_home_env()
5156 } else {
5157 BTreeMap::new()
5158 },
5159 program: program.into(),
5160 arguments,
5161 })
5162}
5163
5164fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5170 let digest = blake3::hash(address.as_bytes()).to_hex();
5171 let path = std::env::temp_dir().join(format!(
5172 "supercode-openclaw-gateway-token-{}",
5173 &digest.as_str()[..16]
5174 ));
5175 #[cfg(unix)]
5176 {
5177 use std::io::Write;
5178 use std::os::unix::fs::OpenOptionsExt;
5179 let mut file = std::fs::OpenOptions::new()
5180 .write(true)
5181 .create(true)
5182 .truncate(true)
5183 .mode(0o600)
5184 .open(&path)?;
5185 file.write_all(secret.as_bytes())?;
5186 }
5187 #[cfg(not(unix))]
5188 std::fs::write(&path, secret)?;
5189 Ok(path)
5190}
5191
5192fn open_connect_descriptor(
5198 descriptor: &crate::HarnessSupportDescriptor,
5199 home: &Path,
5200) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5201 let Some(connect) = &descriptor.runtime.connect_launch else {
5202 return Err(ServiceError::InvalidParams(format!(
5203 "harness `{}` has no registered connect-mode launch",
5204 descriptor.id.as_str()
5205 )));
5206 };
5207 let resolved = connect
5208 .resolve(home)
5209 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5210 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5211 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5212 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5213 if let Some(token) = resolved.auth {
5214 backend = backend.with_bearer(token);
5215 }
5216 Ok(Box::new(backend))
5217 }
5218 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5219 let mut env = BTreeMap::new();
5230 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5231 if let Some(token) = resolved.auth {
5232 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5233 .map_err(|error| {
5234 ServiceError::UnsupportedAction(format!(
5235 "could not stage the gateway credential for the bridge: {error}"
5236 ))
5237 })?;
5238 arguments.push("--token-file".into());
5239 arguments.push(token_path.to_string_lossy().into_owned());
5240 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5241 }
5242 let program = descriptor
5247 .runtime
5248 .default_launch
5249 .as_ref()
5250 .map(|launch| launch.program.clone())
5251 .unwrap_or_else(|| "openclaw".into());
5252 let launch = RuntimeLaunch {
5253 program,
5254 arguments,
5255 env,
5256 };
5257 Ok(Box::new(
5258 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5259 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5260 ))
5261 }
5262 _ => Err(ServiceError::UnsupportedAction(format!(
5263 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5264 descriptor.id.as_str(),
5265 connect.protocol
5266 ))),
5267 }
5268}
5269
5270fn registry_connect_descriptor(
5273 params: &RuntimeBackendParams,
5274) -> Option<crate::HarnessSupportDescriptor> {
5275 if params.launch.is_some() || params.base_url.is_some() {
5276 return None;
5277 }
5278 harness_support_registry()
5279 .harnesses
5280 .into_iter()
5281 .find(|descriptor| descriptor.id == params.harness)
5282 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5283}
5284
5285fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5286 supercode_interchange::user_home()
5287 .map(std::path::PathBuf::into_os_string)
5288 .map(PathBuf::from)
5289 .ok_or_else(|| {
5290 ServiceError::UnsupportedAction(
5291 "connect-mode launches need HOME to locate the harness config".into(),
5292 )
5293 })
5294}
5295
5296pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5299 "harness.v1.runtimes.start",
5300 "harness.v1.runtimes.resume",
5301 "harness.v1.runtimes.attach",
5302 "harness.v1.runtimes.attach_existing",
5303];
5304
5305pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5310
5311pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5317
5318pub const DETACHED_METHODS: &[&str] = &[
5324 "harness.v1.harnesses.list",
5325 "harness.v1.harnesses.probe",
5326 "harness.v1.sessions.message",
5327 "harness.v1.sessions.new",
5328 "harness.v1.sessions.reset",
5329 "harness.v1.sessions.archive",
5330 "harness.v1.sessions.delete",
5331];
5332
5333pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5339
5340pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5348
5349async fn within_control_deadline<F: std::future::Future>(
5352 method: &str,
5353 call: F,
5354) -> std::result::Result<F::Output, ServiceError> {
5355 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5356 .await
5357 .map_err(|_| {
5358 ServiceError::Operation(format!(
5359 "`{method}` gave up after {}s: the runtime did not answer",
5360 RUNTIME_CONTROL_DEADLINE.as_secs()
5361 ))
5362 })
5363}
5364
5365pub struct RuntimeOpen {
5369 id: Value,
5370 method: String,
5371 params: Value,
5372}
5373
5374impl RuntimeOpen {
5375 pub async fn open(self) -> OpenedRuntime {
5379 let Self { id, method, params } = self;
5380 let outcome = open_runtime(&method, params).await;
5381 OpenedRuntime { id, outcome }
5382 }
5383}
5384
5385pub struct OpenedRuntime {
5388 id: Value,
5389 outcome: std::result::Result<OpenRuntime, ServiceError>,
5390}
5391
5392pub struct DetachedCall {
5397 id: Value,
5398 method: String,
5399 work: std::result::Result<Work, ServiceError>,
5400}
5401
5402impl DetachedCall {
5403 pub async fn run(self) -> DetachedAnswer {
5406 let Self { id, method, work } = self;
5407 match work {
5408 Ok(Work::Runtime(work)) => {
5413 let (result, returned) = work.run().await;
5414 DetachedAnswer {
5415 response: service_response(id, result),
5416 returned,
5417 }
5418 }
5419 Ok(Work::Free(work)) => {
5420 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5421 Ok(result) => result,
5422 Err(_) => Err(ServiceError::Operation(format!(
5423 "`{method}` gave up after {}s: the harness it waits on did not answer",
5424 DETACHED_CALL_DEADLINE.as_secs()
5425 ))),
5426 };
5427 DetachedAnswer {
5428 response: service_response(id, result),
5429 returned: None,
5430 }
5431 }
5432 Err(error) => DetachedAnswer {
5433 response: service_response(id, Err(error)),
5434 returned: None,
5435 },
5436 }
5437 }
5438}
5439
5440pub struct DetachedAnswer {
5444 response: Value,
5445 returned: Option<ReturnedRuntime>,
5446}
5447
5448impl DetachedAnswer {
5449 pub fn into_response(self) -> Value {
5452 self.response
5453 }
5454}
5455
5456pub struct ReturnedRuntime {
5459 connection: String,
5460 runtime: Box<dyn RuntimeConnection>,
5461}
5462
5463enum Work {
5466 Free(DetachedWork),
5467 Runtime(RuntimeWork),
5468}
5469
5470enum DetachedWork {
5473 Inventory(InventoryWork),
5477 Message(MessageSessionParams),
5479 SessionMutation {
5482 verb: crate::SessionVerb,
5483 mutation: crate::SessionMutation,
5484 },
5485}
5486
5487impl DetachedWork {
5488 async fn run(self) -> std::result::Result<Value, ServiceError> {
5489 match self {
5490 Self::Inventory(work) => run_inventory(work).await,
5491 Self::Message(params) => Ok(message_live_session(¶ms).await),
5492 Self::SessionMutation { verb, mutation } => {
5493 let outcome = run_session_mutation(verb, &mutation).await?;
5494 serde_json::to_value(outcome)
5495 .map_err(|error| ServiceError::Operation(error.to_string()))
5496 }
5497 }
5498 }
5499}
5500
5501enum RuntimeWork {
5503 Close {
5505 runtime: Box<dyn RuntimeConnection>,
5506 process_group: Option<u32>,
5507 },
5508 LiveCommand {
5511 connection: String,
5512 runtime: Box<dyn RuntimeConnection>,
5513 verb: crate::SessionVerb,
5514 mutation: crate::SessionMutation,
5515 command: &'static str,
5516 session: String,
5517 },
5518}
5519
5520type RuntimeWorkAnswer = (
5523 std::result::Result<Value, ServiceError>,
5524 Option<ReturnedRuntime>,
5525);
5526
5527impl RuntimeWork {
5528 async fn run(self) -> RuntimeWorkAnswer {
5529 match self {
5530 Self::Close {
5531 runtime,
5532 process_group,
5533 } => (close_runtime(runtime, process_group).await, None),
5534 Self::LiveCommand {
5535 connection,
5536 mut runtime,
5537 verb,
5538 mutation,
5539 command,
5540 session,
5541 } => {
5542 let result =
5543 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5544 (
5545 result,
5546 Some(ReturnedRuntime {
5547 connection,
5548 runtime,
5549 }),
5550 )
5551 }
5552 }
5553 }
5554}
5555
5556async fn close_runtime(
5559 mut runtime: Box<dyn RuntimeConnection>,
5560 process_group: Option<u32>,
5561) -> std::result::Result<Value, ServiceError> {
5562 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
5563 Ok(result) => {
5564 result.map_err(operation)?;
5565 Ok(json!({"closed": true}))
5566 }
5567 Err(deadline) => {
5568 let killed = kill_runtime_process_group(process_group);
5573 drop(runtime);
5574 Ok(json!({
5575 "closed": true,
5576 "killed": killed,
5577 "detail": error_message(deadline),
5578 }))
5579 }
5580 }
5581}
5582
5583fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
5586 mutation
5587 .session
5588 .clone()
5589 .filter(|value| !value.trim().is_empty())
5590 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
5591}
5592
5593async fn type_live_command(
5597 runtime: &mut dyn RuntimeConnection,
5598 verb: crate::SessionVerb,
5599 mutation: &crate::SessionMutation,
5600 command: &str,
5601 session: String,
5602) -> std::result::Result<Value, ServiceError> {
5603 within_control_deadline(
5604 &format!("sessions.{}", verb.as_str()),
5605 runtime.send_input(RuntimeInput {
5606 text: command.to_string(),
5607 image_urls: Vec::new(),
5608 }),
5609 )
5610 .await?
5611 .map_err(operation)?;
5612 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
5613 .map_err(session_control_error)?;
5614 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
5615}
5616
5617enum OpenRuntime {
5620 Hosted {
5623 runtime: Box<dyn RuntimeConnection>,
5624 capabilities: crate::RuntimeCapabilities,
5625 workspace: PathBuf,
5626 fresh: bool,
5628 },
5629 Joined { runtime: Box<dyn RuntimeConnection> },
5632}
5633
5634async fn open_runtime(
5639 method: &str,
5640 params: Value,
5641) -> std::result::Result<OpenRuntime, ServiceError> {
5642 match tokio::time::timeout(
5643 RUNTIME_OPEN_DEADLINE,
5644 open_runtime_unbounded(method, params),
5645 )
5646 .await
5647 {
5648 Ok(result) => result,
5649 Err(_) => Err(ServiceError::Operation(format!(
5650 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
5651 RUNTIME_OPEN_DEADLINE.as_secs()
5652 ))),
5653 }
5654}
5655
5656async fn open_runtime_unbounded(
5657 method: &str,
5658 params: Value,
5659) -> std::result::Result<OpenRuntime, ServiceError> {
5660 match method {
5661 "harness.v1.runtimes.start" => {
5662 let params = decode::<RuntimeStartParams>(params)?;
5663 let backend = runtime_backend(¶ms.backend)?;
5664 let capabilities = backend.capabilities();
5665 let workspace = params.cwd.clone();
5666 let runtime = backend
5667 .start(RuntimeStartRequest {
5668 cwd: params.cwd,
5669 launch: runtime_launch(¶ms.backend),
5670 mcp_servers: params.mcp_servers,
5671 approval_policy: params.approval_policy,
5672 })
5673 .await
5674 .map_err(operation)?;
5675 Ok(OpenRuntime::Hosted {
5676 runtime,
5677 capabilities,
5678 workspace,
5679 fresh: true,
5680 })
5681 }
5682 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
5683 let params = decode::<RuntimeAttachParams>(params)?;
5684 let backend = runtime_backend(¶ms.backend)?;
5685 let capabilities = backend.capabilities();
5686 let workspace = params
5687 .cwd
5688 .clone()
5689 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
5690 let runtime = backend
5691 .attach(RuntimeAttachRequest {
5692 runtime_id: params.runtime_id,
5693 cwd: params.cwd,
5694 launch: runtime_launch(¶ms.backend),
5695 mcp_servers: params.mcp_servers,
5696 approval_policy: params.approval_policy,
5697 })
5698 .await
5699 .map_err(operation)?;
5700 Ok(OpenRuntime::Hosted {
5701 runtime,
5702 capabilities,
5703 workspace,
5704 fresh: false,
5705 })
5706 }
5707 "harness.v1.runtimes.attach_existing" => {
5708 let params = decode::<RuntimeAttachParams>(params)?;
5709 let backend: Box<dyn RuntimeBackend> = match params
5710 .backend
5711 .base_url
5712 .as_deref()
5713 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
5714 {
5715 Some(endpoint) => {
5716 #[cfg(not(feature = "adapter-api"))]
5717 {
5718 let _ = endpoint;
5719 return Err(ServiceError::UnsupportedAction(
5720 "live HTTP attachment adapter is not compiled".into(),
5721 ));
5722 }
5723 #[cfg(feature = "adapter-api")]
5724 {
5725 let workspace = params.cwd.clone().ok_or_else(|| {
5726 ServiceError::InvalidParams(
5727 "Volter Harness live attach requires the project cwd".into(),
5728 )
5729 })?;
5730 let source = LiveRuntimeSource {
5731 harness: params.backend.harness.as_str().to_string(),
5732 session_id: params.runtime_id.clone(),
5733 workspace,
5734 };
5735 let receipt = resolve_live_runtime(&endpoint, &source)
5736 .map_err(|error| ServiceError::Operation(error.to_string()))?;
5737 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
5738 }
5739 }
5740 None => runtime_backend(¶ms.backend)?,
5741 };
5742 let capabilities = backend.capabilities();
5743 if !capabilities.attach_existing_process {
5744 return Err(ServiceError::Operation(format!(
5745 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
5746 backend.harness().as_str()
5747 )));
5748 }
5749 let runtime = backend
5750 .attach_existing(RuntimeAttachRequest {
5751 runtime_id: params.runtime_id,
5752 cwd: params.cwd,
5753 launch: runtime_launch(¶ms.backend),
5754 mcp_servers: params.mcp_servers,
5755 approval_policy: params.approval_policy,
5756 })
5757 .await
5758 .map_err(operation)?;
5759 Ok(OpenRuntime::Joined { runtime })
5760 }
5761 _ => Err(ServiceError::MethodNotFound),
5762 }
5763}
5764
5765fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
5767 match result {
5768 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
5769 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
5770 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
5771 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
5772 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
5773 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
5774 }
5775}
5776
5777fn runtime_backend(
5778 params: &RuntimeBackendParams,
5779) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5780 if let Some(descriptor) = registry_connect_descriptor(params) {
5781 return open_connect_descriptor(&descriptor, &service_home()?);
5782 }
5783 if params.protocol.as_deref() == Some("acp") {
5784 let launch = params
5785 .launch
5786 .clone()
5787 .or_else(|| {
5788 harness_support_registry()
5789 .harnesses
5790 .into_iter()
5791 .find(|harness| harness.id == params.harness)
5792 .filter(|harness| {
5793 harness.runtime.implementation == ImplementationKind::GenericProtocol
5794 && harness.runtime.protocol.starts_with("acp")
5795 })
5796 .and_then(|harness| harness.runtime.default_launch)
5797 })
5798 .ok_or_else(|| {
5799 ServiceError::InvalidParams(
5800 "an ACP runtime requires `launch` unless the harness has a registered default"
5801 .into(),
5802 )
5803 })?;
5804 let resume_session = harness_support_registry()
5805 .harnesses
5806 .into_iter()
5807 .find(|harness| harness.id == params.harness)
5808 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
5809 return Ok(Box::new(
5810 AcpRuntimeBackend::new(params.harness.clone(), launch)
5811 .with_resume_support(resume_session),
5812 ));
5813 }
5814 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
5815 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
5816 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
5817 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
5818 HarnessId::OPENCODE => match ¶ms.base_url {
5819 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
5820 None => Box::new(OpenCodeRuntimeBackend::new()),
5821 },
5822 harness => {
5823 let descriptor = harness_support_registry()
5824 .harnesses
5825 .into_iter()
5826 .find(|descriptor| descriptor.id.as_str() == harness)
5827 .filter(|descriptor| {
5828 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
5829 && descriptor.runtime.protocol.starts_with("acp")
5830 });
5831 let Some(descriptor) = descriptor else {
5832 return Err(ServiceError::InvalidParams(format!(
5833 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
5834 )));
5835 };
5836 let resume = descriptor.runtime.capabilities.resume_session;
5837 Box::new(
5838 AcpRuntimeBackend::new(
5839 descriptor.id,
5840 descriptor
5841 .runtime
5842 .default_launch
5843 .expect("generic ACP registry entry includes its launch"),
5844 )
5845 .with_resume_support(resume),
5846 )
5847 }
5848 };
5849 Ok(backend)
5850}
5851
5852fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
5853 if let Some(launch) = ¶ms.launch {
5854 return Some(launch.clone());
5855 }
5856 if !matches!(params.policy, RuntimePolicy::Yolo) {
5857 return None;
5858 }
5859 let launch = match params.harness.as_str() {
5860 HarnessId::GROK => RuntimeLaunch {
5861 program: "grok".into(),
5862 arguments: {
5863 let mut arguments: Vec<String> = Vec::new();
5864 if crate::support::self_sandbox_supported() {
5865 arguments.extend(["--sandbox".into(), "workspace".into()]);
5866 }
5867 arguments.extend([
5868 "--always-approve".into(),
5869 "agent".into(),
5870 "--no-leader".into(),
5871 "stdio".into(),
5872 ]);
5873 arguments
5874 },
5875 env: crate::support::grok_env(),
5876 },
5877 HarnessId::CODEX => RuntimeLaunch {
5878 program: "codex".into(),
5879 arguments: vec![
5880 "--dangerously-bypass-approvals-and-sandbox".into(),
5881 "--dangerously-bypass-hook-trust".into(),
5882 "app-server".into(),
5883 ],
5884 env: BTreeMap::new(),
5885 },
5886 HarnessId::CLAUDE_CODE => RuntimeLaunch {
5887 program: "claude".into(),
5888 arguments: vec![
5889 "--dangerously-skip-permissions".into(),
5890 "--print".into(),
5891 "--input-format".into(),
5892 "stream-json".into(),
5893 "--output-format".into(),
5894 "stream-json".into(),
5895 "--verbose".into(),
5896 ],
5897 env: BTreeMap::new(),
5898 },
5899 HarnessId::PI => RuntimeLaunch {
5900 program: "pi".into(),
5901 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
5902 env: BTreeMap::new(),
5903 },
5904 HarnessId::OPENCODE => RuntimeLaunch {
5905 program: "opencode".into(),
5906 arguments: vec!["serve".into()],
5907 env: BTreeMap::new(),
5908 },
5909 HarnessId::GEMINI => RuntimeLaunch {
5910 program: "gemini".into(),
5911 arguments: vec!["--acp".into(), "--yolo".into()],
5912 env: BTreeMap::new(),
5913 },
5914 HarnessId::GOOSE => RuntimeLaunch {
5915 program: "goose".into(),
5916 arguments: vec!["acp".into()],
5917 env: BTreeMap::new(),
5918 },
5919 HarnessId::SUPERCODE => RuntimeLaunch {
5920 program: "supercode".into(),
5921 arguments: vec!["acp".into(), "--dangerous".into()],
5922 env: BTreeMap::new(),
5923 },
5924 _ => return None,
5925 };
5926 Some(launch)
5927}
5928
5929struct IsolatedProbeHome {
5935 launch: RuntimeLaunch,
5936 root: PathBuf,
5937}
5938
5939impl IsolatedProbeHome {
5940 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
5941 let root = std::env::temp_dir().join(format!(
5942 "supercode-harness-probe-{harness}-{}",
5943 generated_session_id()
5944 ));
5945 std::fs::create_dir_all(&root)?;
5946 set_private_dir_permissions(&root)?;
5947
5948 if let Some(source_home) = supercode_interchange::user_home()
5949 .map(std::path::PathBuf::into_os_string)
5950 .map(PathBuf::from)
5951 {
5952 for relative in probe_auth_files(harness) {
5953 copy_probe_file(&source_home, &root, relative)?;
5954 }
5955 }
5956 if harness == HarnessId::SUPERCODE {
5961 let config_home = crate::agent::global_instructions_dir();
5962 for file in ["config.toml", "credentials.toml"] {
5963 copy_probe_path(
5964 &config_home.join(file),
5965 &root.join(".config/supercode").join(file),
5966 )?;
5967 }
5968 }
5969 configure_isolated_probe_auth(harness, &root)?;
5970
5971 let root_text = root.to_string_lossy().into_owned();
5972 for (key, value) in [
5973 ("HOME", root_text.clone()),
5974 (
5975 "XDG_CACHE_HOME",
5976 root.join(".cache").to_string_lossy().into_owned(),
5977 ),
5978 (
5979 "XDG_CONFIG_HOME",
5980 root.join(".config").to_string_lossy().into_owned(),
5981 ),
5982 (
5983 "XDG_DATA_HOME",
5984 root.join(".local/share").to_string_lossy().into_owned(),
5985 ),
5986 ] {
5987 launch.env.insert(key.into(), value);
5988 }
5989 let scoped = match harness {
5990 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
5991 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
5992 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
5993 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
5994 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
5995 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
5996 _ => None,
5997 };
5998 if let Some((key, value)) = scoped {
5999 launch
6000 .env
6001 .insert(key.into(), value.to_string_lossy().into_owned());
6002 }
6003 Ok(Self { launch, root })
6004 }
6005
6006 fn cleanup(&self) -> std::io::Result<()> {
6007 match std::fs::remove_dir_all(&self.root) {
6008 Ok(()) => Ok(()),
6009 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6010 Err(error) => Err(error),
6011 }
6012 }
6013}
6014
6015impl Drop for IsolatedProbeHome {
6016 fn drop(&mut self) {
6017 let _ = self.cleanup();
6018 }
6019}
6020
6021fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6022 match harness {
6023 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6024 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6028 HarnessId::CODEX => &[".codex/auth.json"],
6029 HarnessId::GEMINI => &[
6030 ".gemini/google_accounts.json",
6031 ".gemini/oauth_creds.json",
6032 ".gemini/settings.json",
6033 ],
6034 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6035 HarnessId::OPENCODE => &[
6036 ".config/opencode/auth.json",
6037 ".local/share/opencode/auth.json",
6038 ],
6039 HarnessId::PI => &[".pi/agent/auth.json"],
6040 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6045 _ => &[],
6046 }
6047}
6048
6049fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6050 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6051}
6052
6053fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6054 if !source.is_file() {
6055 return Ok(());
6056 }
6057 if let Some(parent) = destination.parent() {
6058 std::fs::create_dir_all(parent)?;
6059 set_private_dir_permissions(parent)?;
6060 }
6061 std::fs::copy(source, destination)?;
6062 set_private_file_permissions(destination)
6063}
6064
6065fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6066 if harness != HarnessId::GEMINI {
6067 return Ok(());
6068 }
6069 let oauth = probe_home.join(".gemini/oauth_creds.json");
6070 if !oauth.is_file() {
6071 return Ok(());
6072 }
6073 let settings_path = probe_home.join(".gemini/settings.json");
6074 let mut settings = std::fs::read_to_string(&settings_path)
6075 .ok()
6076 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6077 .unwrap_or_else(|| json!({}));
6078 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6079 std::fs::write(
6080 &settings_path,
6081 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6082 )?;
6083 set_private_file_permissions(&settings_path)
6084}
6085
6086#[cfg(unix)]
6087fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6088 use std::os::unix::fs::PermissionsExt;
6089 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6090}
6091
6092#[cfg(not(unix))]
6093fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6094 Ok(())
6095}
6096
6097#[cfg(unix)]
6098fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6099 use std::os::unix::fs::PermissionsExt;
6100 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6101}
6102
6103#[cfg(not(unix))]
6104fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6105 Ok(())
6106}
6107
6108fn find_executable(program: &str) -> Option<PathBuf> {
6109 let candidate = PathBuf::from(program);
6110 if candidate.components().count() > 1 {
6111 return candidate.is_file().then_some(candidate);
6112 }
6113 let path = std::env::var_os("PATH")?;
6114 for directory in std::env::split_paths(&path) {
6115 let candidate = directory.join(program);
6116 if candidate.is_file() {
6117 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6118 }
6119 #[cfg(windows)]
6120 {
6121 for extension in ["exe", "cmd", "bat"] {
6122 let candidate = directory.join(format!("{program}.{extension}"));
6123 if candidate.is_file() {
6124 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6125 }
6126 }
6127 }
6128 }
6129 None
6130}
6131
6132async fn executable_version(executable: &Path) -> Option<String> {
6133 let mut command = tokio::process::Command::new(executable);
6134 command
6135 .arg("--version")
6136 .stdin(std::process::Stdio::null())
6137 .stdout(std::process::Stdio::piped())
6138 .stderr(std::process::Stdio::piped())
6139 .kill_on_drop(true);
6140 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6141 .await
6142 .ok()?
6143 .ok()?;
6144 let stdout = String::from_utf8_lossy(&output.stdout);
6145 let stderr = String::from_utf8_lossy(&output.stderr);
6146 stdout
6147 .lines()
6148 .chain(stderr.lines())
6149 .map(str::trim)
6150 .find(|line| !line.is_empty())
6151 .map(|line| truncate_text(line, 200))
6152}
6153
6154pub(crate) fn auth_evidence(harness: &str) -> bool {
6155 let env_names: &[&str] = match harness {
6156 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6157 HarnessId::CODEX => &["OPENAI_API_KEY"],
6158 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6159 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6160 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6161 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6162 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6163 _ => &[],
6164 };
6165 if env_names
6166 .iter()
6167 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6168 {
6169 return true;
6170 }
6171 let Some(home) = supercode_interchange::user_home()
6172 .map(std::path::PathBuf::into_os_string)
6173 .map(PathBuf::from)
6174 else {
6175 return false;
6176 };
6177 let files: Vec<PathBuf> = match harness {
6178 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6179 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6180 HarnessId::OPENCODE => vec![
6181 home.join(".local/share/opencode/auth.json"),
6182 home.join(".config/opencode/auth.json"),
6183 ],
6184 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6185 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6186 HarnessId::GEMINI => vec![
6187 home.join(".gemini/oauth_creds.json"),
6188 home.join(".gemini/google_accounts.json"),
6189 ],
6190 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6191 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6192 _ => Vec::new(),
6193 };
6194 if files.into_iter().any(|path| {
6195 std::fs::metadata(path)
6196 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6197 .unwrap_or(false)
6198 }) {
6199 return true;
6200 }
6201 if harness == HarnessId::CLAUDE_CODE {
6208 return std::fs::read_to_string(home.join(".claude.json"))
6209 .map(|text| text.contains("\"oauthAccount\""))
6210 .unwrap_or(false);
6211 }
6212 false
6213}
6214
6215fn looks_like_auth_error(message: &str) -> bool {
6216 let message = message.to_ascii_lowercase();
6217 [
6218 "auth",
6219 "login",
6220 "sign in",
6221 "sign-in",
6222 "credential",
6223 "unauthorized",
6224 "forbidden",
6225 "token",
6226 ]
6227 .iter()
6228 .any(|needle| message.contains(needle))
6229}
6230
6231fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6232 crate::RuntimeCapabilities {
6233 start_session: false,
6234 resume_session: false,
6235 attach_existing_process: false,
6236 send_input: false,
6237 stream_events: false,
6238 interrupt: false,
6239 steer: false,
6240 respond_to_requests: false,
6241 }
6242}
6243
6244fn truncate_text(text: &str, max_chars: usize) -> String {
6245 let mut chars = text.chars();
6246 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6247 if chars.next().is_some() {
6248 format!("{truncated}…")
6249 } else {
6250 truncated
6251 }
6252}
6253
6254fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6261 match &handle.endpoint {
6262 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6263 crate::RuntimeEndpoint::Http { .. } => None,
6264 }
6265}
6266
6267fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6273 match process_group {
6274 #[cfg(unix)]
6275 Some(pid) => {
6276 crate::lsp::kill_process_group(pid);
6277 true
6278 }
6279 #[cfg(not(unix))]
6280 Some(_) => false,
6281 None => false,
6282 }
6283}
6284
6285fn error_message(error: ServiceError) -> String {
6286 match error {
6287 ServiceError::InvalidParams(message)
6288 | ServiceError::Operation(message)
6289 | ServiceError::UnsupportedAction(message) => message,
6290 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6291 ServiceError::Sdk(error) => error.to_string(),
6292 }
6293}
6294
6295#[derive(Debug)]
6296enum ServiceError {
6297 InvalidParams(String),
6298 MethodNotFound,
6299 UnsupportedAction(String),
6300 Operation(String),
6301 Sdk(SdkError),
6302}
6303
6304fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6305 match error {
6306 ServiceError::InvalidParams(message) => {
6307 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6308 }
6309 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6310 SdkError::unsupported(operation)
6311 }
6312 ServiceError::Operation(message) => {
6313 let code = if message.contains("already in progress") {
6314 SdkErrorCode::Busy
6315 } else if message.contains("not supported by this runtime") {
6316 SdkErrorCode::UnsupportedAction
6317 } else if message.contains("unknown runtime connection") {
6318 SdkErrorCode::NotFound
6319 } else {
6320 SdkErrorCode::Execution
6321 };
6322 SdkError::new(code, operation, message)
6323 }
6324 ServiceError::Sdk(error) => error,
6325 }
6326}
6327
6328fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6329 let error_code = error.code();
6330 let code = match error_code {
6331 SdkErrorCode::Unauthenticated => -32030,
6332 SdkErrorCode::Unauthorized => -32031,
6333 SdkErrorCode::ControllerRequired => -32032,
6334 SdkErrorCode::LeaseExpired => -32033,
6335 SdkErrorCode::InvalidArgument => -32602,
6336 SdkErrorCode::NotFound => -32004,
6337 SdkErrorCode::Busy => -32000,
6338 SdkErrorCode::UnsupportedAction => -32020,
6339 SdkErrorCode::Execution => -32002,
6340 SdkErrorCode::Transport => -32003,
6341 };
6342 json!({
6343 "jsonrpc": "2.0",
6344 "id": id,
6345 "error": {
6346 "code": code,
6347 "name": error_code,
6348 "operation": error.operation(),
6349 "message": error.to_string(),
6350 },
6351 })
6352}
6353
6354fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6355 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6356}
6357
6358fn operation(error: impl Into<crate::Error>) -> ServiceError {
6359 let error = error.into();
6360 match error {
6361 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6362 error => ServiceError::Operation(error.to_string()),
6363 }
6364}
6365
6366#[derive(Debug, Clone, Deserialize, Default)]
6370#[serde(default)]
6371struct MemoryRequest {
6372 harness: Option<String>,
6374 query: Option<String>,
6376 profile: Option<String>,
6378 session: Option<String>,
6380 full: bool,
6382 regex: bool,
6384 cwd: Option<std::path::PathBuf>,
6386 homes: crate::HarnessHomes,
6388}
6389
6390fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6393 let request = decode::<MemoryRequest>(params)?;
6394 let harness = request
6395 .harness
6396 .clone()
6397 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6398 let to_service = |error: crate::memory::MemoryError| match error {
6399 crate::memory::MemoryError::UnsupportedHarness { .. }
6400 | crate::memory::MemoryError::SessionNotScoped { .. } => {
6401 ServiceError::UnsupportedAction(error.to_string())
6402 }
6403 other => ServiceError::InvalidParams(other.to_string()),
6404 };
6405 match method {
6406 "harness.v1.memory.show" => {
6407 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6408 harness,
6409 profile: request.profile,
6410 session: request.session,
6411 full: request.full,
6412 cwd: request.cwd,
6413 homes: request.homes,
6414 })
6415 .map_err(to_service)?;
6416 Ok(json!({
6417 "schema": crate::memory::MEMORY_SCHEMA,
6418 "documents": documents,
6419 }))
6420 }
6421 "harness.v1.memory.search" => {
6422 let query = request
6423 .query
6424 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6425 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6426 harness,
6427 query,
6428 profile: request.profile,
6429 regex: request.regex,
6430 cwd: request.cwd,
6431 homes: request.homes,
6432 })
6433 .map_err(to_service)?;
6434 Ok(json!({
6435 "schema": crate::memory::MEMORY_SCHEMA,
6436 "matches": matches,
6437 }))
6438 }
6439 _ => Err(ServiceError::MethodNotFound),
6440 }
6441}
6442
6443#[derive(Debug, Clone, Deserialize)]
6447#[serde(default)]
6448struct ProfilesQuery {
6449 harness: Option<String>,
6451 name: Option<String>,
6453 homes: crate::HarnessHomes,
6455}
6456
6457impl Default for ProfilesQuery {
6458 fn default() -> Self {
6459 Self {
6460 harness: None,
6461 name: None,
6462 homes: crate::HarnessHomes::default(),
6463 }
6464 }
6465}
6466
6467fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6470 let query = decode::<ProfilesQuery>(params)?;
6471 let to_service = |error: crate::profiles::ProfileError| match error {
6472 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6473 ServiceError::UnsupportedAction(error.to_string())
6474 }
6475 crate::profiles::ProfileError::NotFound { .. } => {
6476 ServiceError::InvalidParams(error.to_string())
6477 }
6478 };
6479 match method {
6480 "harness.v1.profiles.list" => {
6481 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6482 .map_err(to_service)?;
6483 Ok(json!({
6484 "schema": crate::profiles::PROFILES_SCHEMA,
6485 "profiles": profiles,
6486 }))
6487 }
6488 "harness.v1.profiles.get" => {
6489 let harness = query
6490 .harness
6491 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6492 let name = query
6493 .name
6494 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6495 let profile =
6496 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6497 Ok(json!({
6498 "schema": crate::profiles::PROFILES_SCHEMA,
6499 "profile": profile,
6500 }))
6501 }
6502 _ => Err(ServiceError::MethodNotFound),
6503 }
6504}
6505
6506#[derive(Debug, Clone, Deserialize)]
6510#[serde(default)]
6511struct ChannelsQuery {
6512 harness: Option<String>,
6514 name: Option<String>,
6516 homes: crate::HarnessHomes,
6518}
6519
6520impl Default for ChannelsQuery {
6521 fn default() -> Self {
6522 Self {
6523 harness: None,
6524 name: None,
6525 homes: crate::HarnessHomes::default(),
6526 }
6527 }
6528}
6529
6530#[derive(Debug, Clone, Deserialize)]
6534#[serde(default)]
6535struct RoutesQuery {
6536 harness: Option<String>,
6537 profile: Option<String>,
6539 homes: crate::HarnessHomes,
6540}
6541
6542impl Default for RoutesQuery {
6543 fn default() -> Self {
6544 Self {
6545 harness: None,
6546 profile: None,
6547 homes: crate::HarnessHomes::default(),
6548 }
6549 }
6550}
6551
6552#[derive(Debug, Clone, Deserialize)]
6553#[serde(default)]
6554struct TriggersQuery {
6555 harness: Option<String>,
6556 homes: crate::HarnessHomes,
6557}
6558
6559impl Default for TriggersQuery {
6560 fn default() -> Self {
6561 Self {
6562 harness: None,
6563 homes: crate::HarnessHomes::default(),
6564 }
6565 }
6566}
6567
6568fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
6569 let query = decode::<TriggersQuery>(params)?;
6570 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
6571 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6572 Ok(json!({
6573 "schema": crate::triggers::TRIGGERS_SCHEMA,
6574 "triggers": triggers,
6575 }))
6576}
6577
6578fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
6579 let query = decode::<RoutesQuery>(params)?;
6580 let routes = crate::routes::list_routes(
6581 &query.homes,
6582 query.harness.as_deref(),
6583 query.profile.as_deref(),
6584 )
6585 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6586 Ok(json!({
6587 "schema": crate::routes::ROUTES_SCHEMA,
6588 "routes": routes,
6589 }))
6590}
6591
6592fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6593 let query = decode::<ChannelsQuery>(params)?;
6594 let to_service = |error: crate::channels::ChannelError| match error {
6595 crate::channels::ChannelError::UnsupportedHarness { .. } => {
6596 ServiceError::UnsupportedAction(error.to_string())
6597 }
6598 crate::channels::ChannelError::NotFound { .. } => {
6599 ServiceError::InvalidParams(error.to_string())
6600 }
6601 };
6602 match method {
6603 "harness.v1.channels.list" => {
6604 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
6605 .map_err(to_service)?;
6606 Ok(json!({
6607 "schema": crate::channels::CHANNELS_SCHEMA,
6608 "channels": channels,
6609 }))
6610 }
6611 "harness.v1.channels.status" => {
6612 let harness = query
6613 .harness
6614 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6615 let name = query
6616 .name
6617 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6618 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
6619 .map_err(to_service)?;
6620 Ok(json!({
6621 "schema": crate::channels::CHANNELS_SCHEMA,
6622 "channel": channel,
6623 }))
6624 }
6625 _ => Err(ServiceError::MethodNotFound),
6626 }
6627}
6628
6629fn rpc_error(id: Value, code: i64, message: &str) -> Value {
6630 json!({
6631 "jsonrpc": "2.0",
6632 "id": id,
6633 "error": {"code": code, "message": message},
6634 })
6635}