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