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