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