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(activities) = self.activity_monitor.resolve(&locators, &homes).await else {
594 continue;
597 };
598 let Some(subscription) = self.activity_subscriptions.get_mut(&subscription_id) else {
599 continue;
600 };
601 let mut changed = Vec::new();
602 for activity in activities {
603 let key = activity.key();
604 if subscription
605 .reported
606 .get(&key)
607 .is_some_and(|previous| previous.same_state(&activity))
608 {
609 continue;
610 }
611 subscription.reported.insert(key, activity.clone());
612 changed.push(activity);
613 }
614 if !changed.is_empty() {
615 notifications.push(json!({
616 "jsonrpc": "2.0",
617 "method": SESSION_ACTIVITY_EVENT_METHOD,
618 "params": {
619 "subscription": subscription_id,
620 "activities": changed,
621 },
622 }));
623 }
624 }
625 notifications
626 }
627
628 #[cfg(feature = "adapter-api")]
632 pub fn poll_session_indexes(&mut self) -> Vec<Value> {
633 let mut notifications = Vec::new();
634 for (subscription, index) in &mut self.index_subscriptions {
635 let homes = index.homes().clone();
636 match index.poll() {
637 Ok(Some(delta)) => match live_index_changes(delta.changes, &homes) {
638 Ok(changes) => notifications.push(json!({
639 "jsonrpc": "2.0",
640 "method": SESSION_INDEX_EVENT_METHOD,
641 "params": {
642 "subscription": subscription,
643 "revision": delta.revision,
644 "changes": changes,
645 },
646 })),
647 Err(error) => notifications.push(json!({
648 "jsonrpc": "2.0",
649 "method": SESSION_INDEX_EVENT_METHOD,
650 "params": {
651 "subscription": subscription,
652 "error": {"recoverable": true, "message": error_message(error)},
653 },
654 })),
655 },
656 Ok(None) => {}
657 Err(error) => notifications.push(json!({
658 "jsonrpc": "2.0",
659 "method": SESSION_INDEX_EVENT_METHOD,
660 "params": {
661 "subscription": subscription,
662 "error": {"recoverable": true, "message": error},
663 },
664 })),
665 }
666 }
667 notifications
668 }
669
670 #[cfg(feature = "adapter-api")]
671 async fn subscribe_session_activity(
672 &mut self,
673 params: Value,
674 ) -> std::result::Result<Value, ServiceError> {
675 let params = decode::<ActivitySubscribeParams>(params)?;
676 if params.locators.is_empty() {
677 return Err(ServiceError::InvalidParams(
678 "sessions.activity.subscribe requires at least one locator".into(),
679 ));
680 }
681 if params.locators.len() > 2_048 {
682 return Err(ServiceError::InvalidParams(
683 "sessions.activity.subscribe accepts at most 2048 locators".into(),
684 ));
685 }
686 let initial = self
687 .activity_monitor
688 .resolve(¶ms.locators, ¶ms.homes)
689 .await
690 .map_err(ServiceError::Sdk)?;
691 let subscription = format!("activity-sub-{}", self.next_subscription);
692 self.next_subscription += 1;
693 let reported = initial
694 .iter()
695 .cloned()
696 .map(|activity| (activity.key(), activity))
697 .collect();
698 self.activity_subscriptions.insert(
699 subscription.clone(),
700 ActivitySubscription {
701 locators: params.locators,
702 homes: params.homes,
703 reported,
704 },
705 );
706 Ok(json!({"subscription": subscription, "initial": initial}))
707 }
708
709 #[cfg(feature = "adapter-api")]
711 pub async fn poll_runtimes(&mut self) -> Vec<Value> {
712 self.poll_sdk_events()
713 .await
714 .into_iter()
715 .map(|(connection, runtime_event)| {
716 json!({
717 "jsonrpc": "2.0",
718 "method": RUNTIME_EVENT_METHOD,
719 "params": {
720 "connection": connection,
721 "session_id": runtime_event.session_id,
722 "sequence": runtime_event.event.sequence,
723 "event": {
724 "kind": runtime_event.event.kind,
725 "payload": runtime_event.event.payload,
726 },
727 },
728 })
729 })
730 .collect()
731 }
732
733 async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
734 let mut events = Vec::new();
735 let mut closed = Vec::new();
736 let now_ms = crate::approvals::now_ms();
737 for (connection, runtime) in &mut self.runtimes {
738 let session_id = runtime.handle().runtime_id.clone();
739 let harness = runtime.handle().harness.clone();
740 for _ in 0..256 {
745 match tokio::time::timeout(Duration::ZERO, runtime.next_event()).await {
746 Ok(Ok(Some(event))) => {
747 let terminal = event.kind == "transport_closed";
748 self.approvals
752 .observe(connection, &harness, &session_id, &event, now_ms);
753 let next_sequence = self
754 .runtime_sequences
755 .entry(session_id.clone())
756 .or_insert(0);
757 let sequence = event.sequence.unwrap_or_else(|| {
758 *next_sequence = next_sequence.saturating_add(1);
759 *next_sequence
760 });
761 *next_sequence = (*next_sequence).max(sequence);
762 events.push((
763 connection.clone(),
764 SdkRuntimeEvent {
765 session_id: session_id.clone(),
766 event: SdkEvent {
767 sequence,
768 kind: event.kind,
769 payload: event.payload,
770 },
771 },
772 ));
773 if terminal {
774 closed.push(connection.clone());
775 break;
776 }
777 }
778 Ok(Ok(None)) => {
779 let sequence = self
780 .runtime_sequences
781 .entry(session_id.clone())
782 .or_insert(0);
783 *sequence = sequence.saturating_add(1);
784 events.push((
785 connection.clone(),
786 SdkRuntimeEvent {
787 session_id,
788 event: SdkEvent {
789 sequence: *sequence,
790 kind: "transport_closed".into(),
791 payload: json!({"message": "Harness runtime transport closed."}),
792 },
793 },
794 ));
795 closed.push(connection.clone());
796 break;
797 }
798 Err(_) => break,
799 Ok(Err(error)) => {
800 let sequence = self
801 .runtime_sequences
802 .entry(session_id.clone())
803 .or_insert(0);
804 *sequence = sequence.saturating_add(1);
805 events.push((
806 connection.clone(),
807 SdkRuntimeEvent {
808 session_id,
809 event: SdkEvent {
810 sequence: *sequence,
811 kind: "transport_error".into(),
812 payload: json!({"message": error.to_string(), "terminal": true}),
813 },
814 },
815 ));
816 closed.push(connection.clone());
817 break;
818 }
819 }
820 }
821 }
822 for connection in closed {
823 if let Some(runtime) = self.runtimes.remove(&connection) {
824 self.runtime_sequences.remove(&runtime.handle().runtime_id);
825 }
826 self.terminal_launches.remove(&connection);
827 self.approvals.forget(&connection);
830 }
831 events
832 }
833
834 fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
835 match method {
836 "harness.v1.capabilities" => Ok(json!({
837 "version": HARNESS_SERVICE_VERSION,
838 "sdk": self.capabilities(),
839 "methods": HARNESS_SERVICE_METHODS,
840 "notifications": [
841 SESSION_EVENT_METHOD,
842 SESSION_ACTIVITY_EVENT_METHOD,
843 SESSION_INDEX_EVENT_METHOD,
844 RUNTIME_EVENT_METHOD
845 ],
846 "harnesses": harness_support_registry()
847 .harnesses
848 .into_iter()
849 .map(|harness| harness.id)
850 .collect::<Vec<_>>(),
851 })),
852 "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
853 .map_err(|error| ServiceError::Operation(error.to_string())),
854 "harness.v1.profiles.list" | "harness.v1.profiles.get" => profiles_call(method, params),
855 "harness.v1.profiles.create" => {
861 mutate_profile(crate::profiles_control::ProfileVerb::Create, params)
862 }
863 "harness.v1.profiles.delete" => {
864 mutate_profile(crate::profiles_control::ProfileVerb::Delete, params)
865 }
866 "harness.v1.channels.list" | "harness.v1.channels.status" => {
867 channels_call(method, params)
868 }
869 "harness.v1.routes.list" => routes_call(params),
872 "harness.v1.triggers.list" => triggers_call(params),
874 "harness.v1.workflow.load" => {
884 let params = decode::<WorkflowLoadParams>(params)?;
885 let read =
886 crate::workflow_doors::load(params.from, ¶ms.home).map_err(operation)?;
887 serde_json::to_value(read)
888 .map_err(|error| ServiceError::Operation(error.to_string()))
889 }
890 "harness.v1.orchestration.load" => {
891 let params = decode::<OrchestrationLoadParams>(params)?;
892 let read = crate::orchestration_doors::load(¶ms.root, params.flavor)
893 .map_err(operation)?;
894 serde_json::to_value(read)
895 .map_err(|error| ServiceError::Operation(error.to_string()))
896 }
897 "harness.v1.orchestration.save" => {
898 let params = decode::<OrchestrationSaveParams>(params)?;
899 let saved = crate::orchestration_doors::save(
900 ¶ms.root,
901 params.orchestration,
902 params.vault,
903 )
904 .map_err(operation)?;
905 serde_json::to_value(saved)
906 .map_err(|error| ServiceError::Operation(error.to_string()))
907 }
908 "harness.v1.orchestration.compile" => {
909 let params = decode::<OrchestrationCompileParams>(params)?;
910 let read = crate::orchestration_doors::compile(params.from, ¶ms.home)
911 .map_err(operation)?;
912 serde_json::to_value(read)
913 .map_err(|error| ServiceError::Operation(error.to_string()))
914 }
915 "harness.v1.orchestration.decompile" => {
916 let params = decode::<OrchestrationDecompileParams>(params)?;
917 let report = crate::orchestration_doors::decompile(
918 params.to,
919 params.orchestration,
920 ¶ms.source,
921 params.source_flavor,
922 ¶ms.dest,
923 params.vault,
924 )
925 .map_err(operation)?;
926 serde_json::to_value(report)
927 .map_err(|error| ServiceError::Operation(error.to_string()))
928 }
929 "harness.v1.orchestration.import" => {
934 let params = decode::<OrchestrationImportParams>(params)?;
935 let imported =
936 crate::orchestration_doors::import(params.from, ¶ms.home, ¶ms.into)
937 .map_err(operation)?;
938 serde_json::to_value(imported)
939 .map_err(|error| ServiceError::Operation(error.to_string()))
940 }
941 "harness.v1.orchestration.export" => {
942 let params = decode::<OrchestrationExportParams>(params)?;
943 let report =
944 crate::orchestration_doors::export(params.to, ¶ms.root, ¶ms.dest)
945 .map_err(operation)?;
946 serde_json::to_value(report)
947 .map_err(|error| ServiceError::Operation(error.to_string()))
948 }
949 "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
955 "harness.v1.skills.list" => {
960 let query = decode::<crate::skills::SkillsQuery>(params)?;
961 if let Some(harness) = query.harness.as_deref() {
962 if !crate::skills::SKILL_HARNESSES.contains(&harness) {
963 return Err(ServiceError::UnsupportedAction(format!(
964 "`{harness}` has no skills root Volter Harness reads"
965 )));
966 }
967 }
968 serde_json::to_value(crate::skills::list_skills(&query))
969 .map_err(|error| ServiceError::Operation(error.to_string()))
970 }
971 "harness.v1.skills.install" => {
979 mutate_skill(crate::skills_control::SkillVerb::Install, params)
980 }
981 "harness.v1.skills.remove" => {
982 mutate_skill(crate::skills_control::SkillVerb::Remove, params)
983 }
984 "harness.v1.approvals.list" => {
992 let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
993 if let Some(harness) = query.harness.as_deref() {
994 if !crate::approvals::lists_approvals(harness) {
995 return Err(ServiceError::UnsupportedAction(format!(
996 "`{harness}` has no runtime door that carries an approval request"
997 )));
998 }
999 }
1000 serde_json::to_value(self.approvals(&query))
1001 .map_err(|error| ServiceError::Operation(error.to_string()))
1002 }
1003 "harness.v1.sessions.discover" => {
1004 let query = decode::<DiscoveryQuery>(params)?;
1005 let mut page = discover_session_page(&query).map_err(operation)?;
1006 let doors = crate::mail_route::LiveSessions::read(&query.homes);
1011 if let Some(text) = query
1015 .query
1016 .as_deref()
1017 .map(str::trim)
1018 .filter(|text| !text.is_empty())
1019 {
1020 let needle = text.to_lowercase();
1021 for live in doors.all() {
1022 let name = live
1023 .name
1024 .split('@')
1025 .next()
1026 .unwrap_or(&live.name)
1027 .to_lowercase();
1028 if !name.contains(&needle)
1029 || page.sessions.iter().any(|session| {
1030 session.locator.session_id == live.address.session_id
1031 })
1032 {
1033 continue;
1034 }
1035 if query
1036 .limit
1037 .is_some_and(|limit| page.sessions.len() >= limit)
1038 {
1039 break;
1040 }
1041 let mut by_id = query.clone();
1043 by_id.query = Some(live.address.session_id.clone());
1044 by_id.cursor = None;
1045 if let Some(cwd) = &live.cwd {
1046 by_id.workspace = Some(cwd.clone());
1047 by_id.workspace_subtree = false;
1048 by_id.workspace_family = None;
1049 }
1050 if let Ok(found) = discover_session_page(&by_id) {
1051 page.sessions
1052 .extend(found.sessions.into_iter().filter(|session| {
1053 session.locator.session_id == live.address.session_id
1054 }));
1055 }
1056 }
1057 }
1058 let activities = crate::session_activity::resolve_stock_session_activities(
1059 &page
1060 .sessions
1061 .iter()
1062 .map(|session| session.locator.clone())
1063 .collect::<Vec<_>>(),
1064 &query.homes,
1065 )
1066 .into_iter()
1067 .map(|activity| (activity.key(), activity))
1068 .collect::<BTreeMap<_, _>>();
1069 let sessions = page
1070 .sessions
1071 .into_iter()
1072 .map(|session| {
1073 let mut value = live_descriptor_value(&session, &doors)?;
1074 let activity_key = (
1075 session.locator.harness.as_str().to_string(),
1076 session.locator.session_id.clone(),
1077 );
1078 if let Some(activity) = activities.get(&activity_key) {
1079 let mut activity = activity.clone();
1082 let waiting = doors.all().iter().any(|live| {
1083 live.address.harness == activity_key.0
1084 && live.address.session_id == activity_key.1
1085 && live.status == "waiting"
1086 });
1087 if waiting && activity.turn != crate::SessionTurnState::NeedsInput {
1088 activity.turn = crate::SessionTurnState::NeedsInput;
1089 activity.evidence.source = "pane_prompt".into();
1090 }
1091 let activity = &activity;
1092 value["activity"] = serde_json::to_value(activity)
1093 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1094 if let Some(status) = legacy_live_status(activity) {
1095 value["live_status"] = json!(status);
1096 }
1097 }
1098 Ok(value)
1099 })
1100 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1101 let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1102 if query.search_previews {
1105 result["receipt"] = serde_json::to_value(page.receipt)
1106 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1107 }
1108 Ok(result)
1109 }
1110 "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1111 "harness.v1.sessions.usage" => {
1112 let params = decode::<LocatorParams>(params)?;
1113 let session = load_session(¶ms.locator).map_err(operation)?;
1114 Ok(
1115 json!({"harness": params.locator.harness, "session_id": session.meta.session_id,
1116 "usage": recorded_session_usage(&session)}),
1117 )
1118 }
1119 "harness.v1.sessions.recorded_config" => {
1120 let params = decode::<LocatorParams>(params)?;
1121 let session = load_session(¶ms.locator).map_err(operation)?;
1122 let mut config = json!({"harness": params.locator.harness, "session_id": session.meta.session_id,
1123 "cwd": session.meta.cwd, "model": session.meta.model,
1124 "approval_policy": null, "sandbox_policy": null});
1125 for header in &session.meta.codex_headers {
1127 let payload = header.get("payload").unwrap_or(header);
1128 for key in ["cwd", "model", "approval_policy", "sandbox_policy"] {
1129 if let Some(value) = payload.get(key) {
1130 config[key] = value.clone();
1131 }
1132 }
1133 }
1134 Ok(config)
1135 }
1136 "harness.v1.sessions.load" => {
1137 let params = decode::<LoadSessionParams>(params)?;
1138 if let Some(options) = ¶ms.options {
1139 options.validate()?;
1140 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
1141 return Ok(result);
1142 }
1143 return load_session(¶ms.read.locator)
1144 .map(|session| projected_session_result(&session, options))
1145 .map_err(operation);
1146 }
1147 let mut session = if params.read.display_history() {
1148 self.catalog
1149 .load_display_view(
1150 ¶ms.read.locator,
1151 params.read.read_fidelity(),
1152 params.read.tail_messages().unwrap_or(500),
1153 )
1154 .map_err(crate::Error::from)
1155 } else if params.read.include_subagents() {
1156 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
1157 } else {
1158 self.catalog
1159 .load_parent_with_fidelity(
1160 ¶ms.read.locator,
1161 params.read.read_fidelity(),
1162 )
1163 .map_err(crate::Error::from)
1164 }
1165 .map_err(operation)?;
1166 params.read.bound_session(&mut session);
1167 let mut value = normalized_session_json(&session);
1170 value["resolved_profile"] = recorded_session_profile(&session);
1171 value["resolved_profile_history"] =
1172 Value::Array(recorded_session_profile_history(&session));
1173 Ok(json!({"session": value}))
1174 }
1175 "harness.v1.sessions.follow" => {
1176 let params = decode::<LocatorParams>(params)?;
1177 let mut follower = self
1178 .catalog
1179 .follow_read_view(
1180 ¶ms.locator,
1181 params.read_fidelity(),
1182 params.include_subagents(),
1183 params.tail_messages(),
1184 params.max_message_chars(),
1185 params.display_history(),
1186 )
1187 .map_err(operation)?;
1188 let initial = follower
1189 .poll()
1190 .map_err(operation)?
1191 .map(|event| event.to_json());
1192 let subscription = format!("sub-{}", self.next_subscription);
1193 self.next_subscription += 1;
1194 self.followers.insert(subscription.clone(), follower);
1195 self.followed_sources.insert(
1196 subscription.clone(),
1197 FollowedSource {
1198 harness: params.locator.harness.as_str().to_string(),
1199 session_id: params.locator.session_id.clone(),
1200 reported: None,
1201 },
1202 );
1203 Ok(json!({"subscription": subscription, "initial": initial}))
1204 }
1205 "harness.v1.sessions.unfollow" => {
1206 let params = decode::<UnfollowParams>(params)?;
1207 self.followed_sources.remove(¶ms.subscription);
1208 Ok(json!({
1209 "removed": self.followers.remove(¶ms.subscription).is_some()
1210 }))
1211 }
1212 "harness.v1.sessions.activity.unsubscribe" => {
1213 let params = decode::<UnfollowParams>(params)?;
1214 Ok(json!({
1215 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1216 }))
1217 }
1218 "harness.v1.sessions.index.subscribe" => {
1219 let query = decode::<DiscoveryQuery>(params)?;
1220 crate::session_index::validate_query(&query)
1221 .map_err(ServiceError::InvalidParams)?;
1222 let homes = query.homes.clone();
1223 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1224 query,
1225 Arc::clone(&self.index_notifier),
1226 )
1227 .map_err(ServiceError::Operation)?;
1228 let doors = crate::mail_route::LiveSessions::read(&homes);
1229 let initial = initial
1230 .iter()
1231 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1232 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1233 let subscription = format!("index-sub-{}", self.next_subscription);
1234 self.next_subscription += 1;
1235 self.index_subscriptions.insert(subscription.clone(), index);
1236 Ok(json!({
1237 "subscription": subscription,
1238 "revision": 1,
1239 "initial": initial,
1240 }))
1241 }
1242 "harness.v1.sessions.index.resize" => {
1243 let params = decode::<IndexResizeParams>(params)?;
1244 crate::session_index::validate_limit(params.limit)
1245 .map_err(ServiceError::InvalidParams)?;
1246 let index = self
1247 .index_subscriptions
1248 .get_mut(¶ms.subscription)
1249 .ok_or_else(|| {
1250 ServiceError::InvalidParams("unknown session index subscription".into())
1251 })?;
1252 let prepared = index
1253 .prepare_resize(params.limit)
1254 .map_err(ServiceError::Operation)?;
1255 let doors = crate::mail_route::LiveSessions::read(index.homes());
1256 let initial = prepared
1257 .page
1258 .sessions
1259 .iter()
1260 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1261 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1262 let response = json!({
1263 "subscription": params.subscription,
1264 "revision": prepared.revision,
1265 "initial": initial,
1266 "receipt": prepared.page.receipt,
1267 });
1268 index.commit_resize(prepared);
1269 Ok(response)
1270 }
1271 "harness.v1.sessions.index.unsubscribe" => {
1272 let params = decode::<UnfollowParams>(params)?;
1273 Ok(json!({
1274 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1275 }))
1276 }
1277 "harness.v1.sessions.import" => {
1278 let params = decode::<ImportSessionParams>(params)?;
1279 let session = Session::load_str(¶ms.content, params.source_harness.into())
1280 .map_err(operation)?;
1281 Ok(json!({"session": normalized_session_json(&session)}))
1282 }
1283 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1284 let params = decode::<ExportSessionParams>(params)?;
1285 let session = load_session(¶ms.locator).map_err(operation)?;
1286 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1287 if method == "harness.v1.sessions.export"
1288 && params.target_harness == TransferFormat::Hermes
1289 {
1290 let imported = crate::hermes_import::import_into_hermes(&session, None)
1292 .map_err(operation)?;
1293 return Ok(json!({"artifact": artifact, "imported": imported}));
1294 }
1295 Ok(json!({"artifact": artifact}))
1296 }
1297 "harness.v1.sessions.reduce" => {
1298 let params = decode::<ReduceSessionParams>(params)?;
1299 self.reduce_session(params)
1300 }
1301 "harness.v1.sessions.branch" => {
1302 let params = decode::<BranchSessionParams>(params)?;
1303 let session = load_session(¶ms.locator).map_err(operation)?;
1304 let storage = params.locator.storage.path().display().to_string();
1305 let bootstrap_prompt = format!(
1306 "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.",
1307 params.locator.harness.as_str(), params.locator.session_id, storage
1308 );
1309 let artifact = params
1310 .target_harness
1311 .map(|target| session_artifact(¶ms.locator, &session, target))
1312 .transpose()?;
1313 Ok(json!({
1314 "parent": params.locator,
1315 "session": normalized_session_json(&session),
1316 "bootstrap_prompt": bootstrap_prompt,
1317 "artifact": artifact,
1318 }))
1319 }
1320 "harness.v1.sessions.handoff" => {
1321 let params = decode::<HandoffSessionParams>(params)?;
1322 let session = load_session(¶ms.locator).map_err(operation)?;
1323 let cwd = params
1324 .cwd
1325 .or_else(|| session.meta.cwd.clone())
1326 .unwrap_or_else(|| PathBuf::from("."));
1327 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1328 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1329 ServiceError::Operation(
1330 "handoff artifact omitted target session identity".into(),
1331 )
1332 })?;
1333 let instructions =
1334 handoff_instructions(params.target_harness, target_session_id, &cwd);
1335 Ok(json!({
1336 "artifact": artifact,
1337 "launch": instructions.launch,
1338 "materialize": instructions.materialize,
1339 "requires_materialization": instructions.requires_materialization,
1340 "note": instructions.note,
1341 }))
1342 }
1343 "harness.v1.sessions.materialize" => {
1344 let params = decode::<MaterializeSessionParams>(params)?;
1345 for file in ¶ms.artifact.files {
1349 if file.role == "source_recovery"
1350 && file.path == "recovery/source.supercode.jsonl"
1351 {
1352 if let Ok(source) = Session::from_native_str(&file.content) {
1353 crate::residue_store::store_segments(&source);
1354 }
1355 }
1356 }
1357 let locator = crate::native_materialize::materialize_native_artifact(
1358 params.artifact,
1359 ¶ms.cwd,
1360 ¶ms.homes,
1361 )
1362 .map_err(ServiceError::Operation)?;
1363 Ok(json!({"locator": locator}))
1364 }
1365 "harness.v1.jobs.list" => {
1369 let query = decode::<crate::jobs::JobsQuery>(params)?;
1370 if let Some(harness) = query.harness.as_deref() {
1371 refuse_harness_without_jobs(harness, "jobs.list")?;
1372 }
1373 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1374 serde_json::to_value(listing)
1375 .map_err(|error| ServiceError::Operation(error.to_string()))
1376 }
1377 "harness.v1.jobs.get" => {
1378 let params = decode::<JobsGetParams>(params)?;
1379 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1380 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1381 .map_err(operation)?
1382 {
1383 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1384 None => Err(ServiceError::Operation(format!(
1385 "`{}` has no scheduled job `{}`",
1386 params.harness, params.id
1387 ))),
1388 }
1389 }
1390 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1396 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1397 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1398 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1399 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1400 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1401 "harness.v1.jobs.notepad"
1402 | "harness.v1.jobs.notepad_set"
1403 | "harness.v1.jobs.notepad_delete" => {
1404 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1405 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1406 let answer = match method {
1407 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1408 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1409 _ => crate::jobs_notepad::read(&request),
1410 }
1411 .map_err(job_control_error)?;
1412 serde_json::to_value(answer)
1413 .map_err(|error| ServiceError::Operation(error.to_string()))
1414 }
1415 "harness.v1.runs.list" => {
1419 let query = decode::<crate::runs::RunsQuery>(params)?;
1420 if let Some(harness) = query.harness.as_deref() {
1421 refuse_harness_without_runs(harness, "runs.list")?;
1422 }
1423 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1424 serde_json::to_value(listing)
1425 .map_err(|error| ServiceError::Operation(error.to_string()))
1426 }
1427 "harness.v1.runs.get" => {
1428 let params = decode::<RunsGetParams>(params)?;
1429 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1430 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1431 .map_err(operation)?
1432 {
1433 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1434 None => Err(ServiceError::Operation(format!(
1435 "`{}` has no run `{}`",
1436 params.harness, params.id
1437 ))),
1438 }
1439 }
1440 "harness.v1.sessions.resume_instructions" => {
1441 let params = decode::<ResumeInstructionsParams>(params)?;
1442 let session = load_session(¶ms.locator).map_err(operation)?;
1443 let cwd = params
1444 .cwd
1445 .or(session.meta.cwd)
1446 .unwrap_or_else(|| PathBuf::from("."));
1447 let launch = resume_launch(
1448 params.locator.harness.as_str(),
1449 ¶ms.locator.session_id,
1450 &cwd,
1451 params.policy,
1452 )?;
1453 Ok(json!({"launch": launch}))
1454 }
1455 _ => Err(ServiceError::MethodNotFound),
1456 }
1457 }
1458
1459 fn reduce_session(
1460 &self,
1461 params: ReduceSessionParams,
1462 ) -> std::result::Result<Value, ServiceError> {
1463 let session = load_session(¶ms.locator).map_err(operation)?;
1464 if session.messages.is_empty() {
1465 return Err(ServiceError::InvalidParams(
1466 "cannot reduce an empty session".into(),
1467 ));
1468 }
1469 let keep_last = params.keep_last.clamp(1, 128);
1470 let policy = reduce::ReductionPolicy {
1471 clear_turns_older_than: Some(keep_last),
1472 ..Default::default()
1473 };
1474 let (view, log) =
1475 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1476 if log.reductions.is_empty() {
1477 return Err(ServiceError::UnsupportedAction(format!(
1478 "session `{}` is already too small for a meaningful reversible reduction",
1479 params.locator.session_id
1480 )));
1481 }
1482 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1483 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1484 if reduced_tokens >= source_tokens {
1485 return Err(ServiceError::UnsupportedAction(format!(
1486 "session `{}` has no token-reducing reversible projection",
1487 params.locator.session_id
1488 )));
1489 }
1490
1491 let store_root = self
1492 .reduction_store_root
1493 .clone()
1494 .unwrap_or_else(default_reduction_store_root);
1495 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1496 let rescue_id = format!("rescue-{}", generated_session_id());
1497 let imported = session
1498 .imported_message_count
1499 .unwrap_or(session.messages.len())
1500 .min(session.messages.len());
1501 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1502 let view_jsonl = messages_jsonl(&view)?;
1503 let title = format!(
1504 "Reduced {} continuation from {}",
1505 params.target_harness.id(),
1506 params.locator.session_id
1507 );
1508
1509 store
1514 .save_sidecar(&rescue_id, &sidecar_jsonl)
1515 .map_err(operation)?;
1516 store
1517 .save_reduction_log(&rescue_id, &log)
1518 .map_err(operation)?;
1519 store
1520 .save(&rescue_id, &title, &view_jsonl)
1521 .map_err(operation)?;
1522
1523 let source_bytes = serde_json::to_vec(&session.messages)
1524 .map_err(|error| ServiceError::Operation(error.to_string()))?
1525 .len() as u64;
1526 let reduced_bytes = serde_json::to_vec(&view)
1527 .map_err(|error| ServiceError::Operation(error.to_string()))?
1528 .len() as u64;
1529 store
1530 .set_reduction_stats(
1531 &rescue_id,
1532 &title,
1533 source_bytes,
1534 reduced_bytes,
1535 log.reductions.len() as u32,
1536 )
1537 .map_err(operation)?;
1538
1539 let reloaded_sidecar = store
1543 .load_sidecar(&rescue_id)
1544 .map_err(operation)?
1545 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1546 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1547 let reloaded_log = store
1548 .load_reduction_log(&rescue_id)
1549 .map_err(operation)?
1550 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1551 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1552 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1553 let (restamped_view, restamped_log) =
1560 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1561 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1562 return Err(ServiceError::Operation(
1563 "persisted reduction view does not match its durable log and sidecar".into(),
1564 ));
1565 }
1566 if restamped_log != reloaded_log {
1567 return Err(ServiceError::Operation(
1568 "reapplying the durable reduction log changed its identity".into(),
1569 ));
1570 }
1571 let inverted =
1572 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1573 if inverted != session.messages {
1574 return Err(ServiceError::Operation(
1575 "reduction inversion did not restore the source messages byte-exactly".into(),
1576 ));
1577 }
1578
1579 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1580 let sidecar_path = store.sidecar_path(&rescue_id);
1581 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1582 let bootstrap_prompt = reduced_bootstrap_prompt(
1583 ¶ms.locator,
1584 params.target_harness,
1585 &view_jsonl,
1586 &sidecar_path,
1587 &reduction_log_path,
1588 );
1589 let mut reduced_session = session.clone();
1590 reduced_session.meta.session_id = Some(rescue_id.clone());
1591 reduced_session.messages = view;
1592
1593 Ok(json!({
1594 "session": normalized_session_json(&reduced_session),
1595 "bootstrap_prompt": bootstrap_prompt,
1596 "receipt": {
1597 "id": rescue_id,
1598 "sidecar_id": rescue_id,
1599 "source_harness": params.locator.harness,
1600 "target_harness": params.target_harness.id(),
1601 "source_tokens": source_tokens,
1602 "reduced_tokens": reduced_tokens,
1603 "ratio": ratio,
1604 "source_bytes": source_bytes,
1605 "reduced_bytes": reduced_bytes,
1606 "reductions": reloaded_log.reductions.len(),
1607 "sidecar_path": sidecar_path,
1608 "reduction_log_path": reduction_log_path,
1609 "verified": true,
1610 "reversible": true,
1611 }
1612 }))
1613 }
1614
1615 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1633 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1634 return None;
1635 }
1636 let method = request.get("method").and_then(Value::as_str)?;
1637 if !RUNTIME_OPEN_METHODS.contains(&method) {
1638 return None;
1639 }
1640 Some(RuntimeOpen {
1641 id: request.get("id").cloned().unwrap_or(Value::Null),
1642 method: method.to_string(),
1643 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1644 })
1645 }
1646
1647 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1669 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1670 return None;
1671 }
1672 let method = request.get("method").and_then(Value::as_str)?;
1673 if !DETACHED_METHODS.contains(&method) {
1674 return None;
1675 }
1676 let id = request.get("id").cloned().unwrap_or(Value::Null);
1677 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1678 let work = match method {
1679 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1680 .inventory_work(method, params)
1681 .map(DetachedWork::Inventory),
1682 "harness.v1.sessions.message" => {
1683 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1684 }
1685 _ => {
1686 let verb = match method {
1687 "harness.v1.sessions.new" => crate::SessionVerb::New,
1688 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1689 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1690 _ => crate::SessionVerb::Delete,
1691 };
1692 match decode::<crate::SessionMutation>(params) {
1693 Ok(mutation) => {
1694 match crate::sessions_control::door(&mutation.harness, verb) {
1695 Ok(crate::SessionDoor::Live(_)) => return None,
1698 Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1699 Err(error) => Err(session_control_error(error)),
1700 }
1701 }
1702 Err(error) => Err(error),
1703 }
1704 }
1705 };
1706 Some(DetachedCall {
1707 id,
1708 method: method.to_string(),
1709 work: work.map(Work::Free),
1710 })
1711 }
1712
1713 pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1727 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1728 return None;
1729 }
1730 let method = request.get("method").and_then(Value::as_str)?;
1731 let id = request.get("id").cloned().unwrap_or(Value::Null);
1732 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1733 let work = match method {
1734 "harness.v1.runtimes.close" => decode::<RuntimeCloseParams>(params)
1735 .and_then(|params| self.surrender_runtime_of(¶ms))
1736 .map(|(runtime, process_group)| {
1737 Work::Runtime(RuntimeWork::Close {
1738 runtime,
1739 process_group,
1740 })
1741 }),
1742 "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1743 let verb = if method == "harness.v1.sessions.new" {
1744 crate::SessionVerb::New
1745 } else {
1746 crate::SessionVerb::Reset
1747 };
1748 let mutation = decode::<crate::SessionMutation>(params).ok()?;
1749 let Ok(crate::SessionDoor::Live(command)) =
1753 crate::sessions_control::door(&mutation.harness, verb)
1754 else {
1755 return None;
1756 };
1757 let connection = mutation
1758 .connection
1759 .clone()
1760 .filter(|value| !value.trim().is_empty())?;
1761 self.lend_runtime(&connection).map(|runtime| {
1762 let session = live_session_name(runtime.as_ref(), &mutation);
1763 Work::Runtime(RuntimeWork::LiveCommand {
1764 connection,
1765 runtime,
1766 verb,
1767 mutation,
1768 command,
1769 session,
1770 })
1771 })
1772 }
1773 _ => return None,
1774 };
1775 Some(DetachedCall {
1776 id,
1777 method: method.to_string(),
1778 work,
1779 })
1780 }
1781
1782 pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1787 let DetachedAnswer { response, returned } = answer;
1788 if let Some(ReturnedRuntime {
1789 connection,
1790 runtime,
1791 }) = returned
1792 {
1793 self.runtimes_in_flight.remove(&connection);
1794 self.runtimes.insert(connection, runtime);
1795 }
1796 response
1797 }
1798
1799 pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1803 let OpenedRuntime { id, outcome } = opened;
1804 let result = match outcome {
1805 Ok(open) => self.register_open_runtime(open).await,
1806 Err(error) => Err(error),
1807 };
1808 service_response(id, result)
1809 }
1810
1811 async fn register_open_runtime(
1813 &mut self,
1814 open: OpenRuntime,
1815 ) -> std::result::Result<Value, ServiceError> {
1816 match open {
1817 OpenRuntime::Hosted {
1818 runtime,
1819 capabilities,
1820 workspace,
1821 fresh,
1822 } => {
1823 self.insert_hosted_runtime(runtime, capabilities, workspace, fresh)
1824 .await
1825 }
1826 OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1827 }
1828 }
1829
1830 async fn runtime_call(
1831 &mut self,
1832 method: &str,
1833 params: Value,
1834 ) -> std::result::Result<Value, ServiceError> {
1835 match method {
1836 "harness.v1.runtimes.capabilities" => {
1837 let params = decode::<RuntimeBackendParams>(params)?;
1838 #[cfg(feature = "adapter-api")]
1842 if params
1843 .base_url
1844 .as_deref()
1845 .is_some_and(|value| LiveRuntimeEndpoint::parse(value).is_ok())
1846 {
1847 return Ok(json!({
1848 "harness": ¶ms.harness,
1849 "capabilities": SupercodeHttpRuntimeBackend::live_capabilities(),
1850 }));
1851 }
1852 let backend = runtime_backend(¶ms)?;
1853 Ok(json!({
1854 "harness": backend.harness(),
1855 "capabilities": backend.capabilities(),
1856 }))
1857 }
1858 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1859 self.register_open_runtime(open_runtime(method, params).await?)
1860 .await
1861 }
1862 "harness.v1.runtimes.send_input" => {
1863 let params = decode::<RuntimeInputParams>(params)?;
1864 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1865 let runtime = self.runtime_mut(¶ms.connection)?;
1866 let turn_id = within_control_deadline(
1867 method,
1868 runtime.send_input(RuntimeInput {
1869 text: params.text,
1870 image_urls,
1871 }),
1872 )
1873 .await?
1874 .map_err(operation)?;
1875 Ok(json!({"turn_id": turn_id}))
1876 }
1877 "harness.v1.runtimes.interrupt" => {
1878 let params = decode::<RuntimeConnectionParams>(params)?;
1879 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1880 .await?
1881 .map_err(operation)?;
1882 Ok(json!({}))
1883 }
1884 "harness.v1.runtimes.steer" => {
1885 let params = decode::<RuntimeInputParams>(params)?;
1886 if !params.image_urls.is_empty() {
1887 return Err(ServiceError::InvalidParams(
1888 "runtime steering accepts text only".into(),
1889 ));
1890 }
1891 let text = params.text.trim();
1892 if text.is_empty() || text.chars().count() > 50_000 {
1893 return Err(ServiceError::InvalidParams(
1894 "runtime steering requires 1 to 50,000 text characters".into(),
1895 ));
1896 }
1897 within_control_deadline(
1898 method,
1899 self.runtime_mut(¶ms.connection)?
1900 .steer(text.to_string()),
1901 )
1902 .await?
1903 .map_err(operation)?;
1904 Ok(json!({}))
1905 }
1906 "harness.v1.runtimes.respond" => {
1907 let params = decode::<RuntimeRespondParams>(params)?;
1908 let request_id = params.request_id.clone();
1909 within_control_deadline(
1910 method,
1911 self.runtime_mut(¶ms.connection)?
1912 .respond(params.request_id, params.response),
1913 )
1914 .await?
1915 .map_err(operation)?;
1916 self.approvals.answered(¶ms.connection, &request_id);
1918 Ok(json!({}))
1919 }
1920 "harness.v1.runtimes.acquire_control" => {
1921 let params = decode::<RuntimeConnectionParams>(params)?;
1922 let snapshot = within_control_deadline(
1923 method,
1924 self.runtime_mut(¶ms.connection)?.acquire_control(),
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.heartbeat" => {
1932 let params = decode::<RuntimeConnectionParams>(params)?;
1933 let snapshot = within_control_deadline(
1934 method,
1935 self.runtime_mut(¶ms.connection)?.heartbeat(),
1936 )
1937 .await?
1938 .map_err(operation)?;
1939 serde_json::to_value(snapshot)
1940 .map_err(|error| ServiceError::Operation(error.to_string()))
1941 }
1942 "harness.v1.runtimes.detach" => {
1943 let params = decode::<RuntimeConnectionParams>(params)?;
1944 let snapshot =
1945 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1946 .await?
1947 .map_err(operation)?;
1948 serde_json::to_value(snapshot)
1949 .map_err(|error| ServiceError::Operation(error.to_string()))
1950 }
1951 "harness.v1.runtimes.terminal_instructions" => {
1952 let params = decode::<RuntimeConnectionParams>(params)?;
1953 let launch = self
1954 .terminal_launches
1955 .get(¶ms.connection)
1956 .ok_or_else(|| {
1957 ServiceError::Operation(
1958 "this runtime is not hosted for terminal attachment".into(),
1959 )
1960 })?;
1961 Ok(json!({"launch":launch}))
1962 }
1963 "harness.v1.runtimes.close" => {
1964 let params = decode::<RuntimeCloseParams>(params)?;
1965 let (runtime, process_group) = self.surrender_runtime_of(¶ms)?;
1966 close_runtime(runtime, process_group).await
1967 }
1968 _ => Err(ServiceError::MethodNotFound),
1969 }
1970 }
1971
1972 #[cfg(feature = "adapter-api")]
1974 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1975 let params = decode::<MessageSessionParams>(params)?;
1976 Ok(message_live_session(¶ms).await)
1977 }
1978
1979 #[cfg(feature = "adapter-api")]
1980 fn harness_settings_call(
1981 &self,
1982 method: &str,
1983 params: Value,
1984 ) -> std::result::Result<Value, ServiceError> {
1985 let homes = crate::HarnessHomes::default();
1986 match method {
1987 "harness.v1.harnesses.settings" => {
1988 let params = decode::<HarnessSettingsParams>(params)?;
1989 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1990 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1991 serde_json::to_value(report)
1992 .map_err(|error| ServiceError::Operation(error.to_string()))
1993 }
1994 "harness.v1.harnesses.configure" => {
1995 let params = decode::<ConfigureHarnessParams>(params)?;
1996 let report = crate::configure_harness_interop_settings(
1997 &homes,
1998 ¶ms.harness,
1999 ¶ms.changes,
2000 params.expected_revision.as_deref(),
2001 )
2002 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2003 serde_json::to_value(report)
2004 .map_err(|error| ServiceError::Operation(error.to_string()))
2005 }
2006 _ => Err(ServiceError::MethodNotFound),
2007 }
2008 }
2009
2010 fn insert_runtime(
2011 &mut self,
2012 runtime: Box<dyn RuntimeConnection>,
2013 ) -> std::result::Result<Value, ServiceError> {
2014 let connection = format!("runtime-{}", self.next_runtime);
2015 self.next_runtime += 1;
2016 let handle = runtime.handle().clone();
2017 self.runtime_sequences
2018 .entry(handle.runtime_id.clone())
2019 .or_insert(0);
2020 self.runtimes.insert(connection.clone(), runtime);
2021 Ok(json!({"connection": connection, "handle": handle}))
2022 }
2023
2024 #[cfg(feature = "adapter-api")]
2025 async fn insert_hosted_runtime(
2026 &mut self,
2027 runtime: Box<dyn RuntimeConnection>,
2028 capabilities: crate::RuntimeCapabilities,
2029 workspace: PathBuf,
2030 fresh: bool,
2031 ) -> std::result::Result<Value, ServiceError> {
2032 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
2033 let token: std::sync::Arc<str> = crate::server::generate_token().into();
2034 let server = crate::server::run_frontend_http(
2035 host.clone(),
2036 host.frontend_sender(),
2037 "127.0.0.1:0",
2038 token.clone(),
2039 connection.handle().runtime_id.clone(),
2040 )
2041 .await
2042 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2043 let source = LiveRuntimeSource {
2044 harness: connection.handle().harness.as_str().to_string(),
2045 session_id: connection.handle().runtime_id.clone(),
2046 workspace: workspace.clone(),
2047 };
2048 let registration = register_live_runtime_with_metadata(
2049 connection.handle().runtime_id.clone(),
2050 source.clone(),
2051 format!("http://{}", server.address()),
2052 token.to_string(),
2053 crate::LiveRuntimeMetadata {
2054 child_pid: match &connection.handle().endpoint {
2055 crate::RuntimeEndpoint::LocalProcess { pid, command, .. }
2056 if connection.handle().harness.as_str() != "claude-code"
2057 || command.windows(2).any(|args| {
2058 args[0] == "--name"
2059 && args[1].starts_with(crate::claude_relay::RELAY_NAME_PREFIX)
2060 }) =>
2061 {
2062 *pid
2063 }
2064 _ => None,
2065 },
2066 endpoint_capabilities: vec!["http".into(), "acp".into()],
2067 ..Default::default()
2068 },
2069 )
2070 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2071 let endpoint = registration.endpoint().to_string();
2072 let launch = StructuredLaunch {
2073 cwd: workspace,
2074 program: std::env::current_exe()
2078 .ok()
2079 .map(|path| path.to_string_lossy().into_owned())
2080 .unwrap_or_else(|| "supercode".into()),
2081 arguments: vec![
2082 "open".into(),
2083 endpoint,
2084 "--harness".into(),
2085 source.harness,
2086 "--session".into(),
2087 source.session_id,
2088 ],
2089 env: BTreeMap::new(),
2090 };
2091 host.retain_registration(registration);
2092 let lease = HostedRuntimeLease {
2093 connection,
2094 _host: host,
2095 _server: server,
2096 };
2097 let opened = self.insert_runtime(Box::new(lease))?;
2098 let connection_id = opened["connection"]
2099 .as_str()
2100 .expect("insert_runtime returns a connection id")
2101 .to_string();
2102 self.terminal_launches.insert(connection_id, launch);
2103 Ok(opened)
2104 }
2105
2106 #[cfg(not(feature = "adapter-api"))]
2107 async fn insert_hosted_runtime(
2108 &mut self,
2109 runtime: Box<dyn RuntimeConnection>,
2110 _capabilities: crate::RuntimeCapabilities,
2111 _workspace: PathBuf,
2112 _fresh: bool,
2113 ) -> std::result::Result<Value, ServiceError> {
2114 self.insert_runtime(runtime)
2115 }
2116
2117 fn runtime_mut(
2118 &mut self,
2119 connection: &str,
2120 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2121 if self.runtimes_in_flight.contains(connection) {
2122 return Err(self.lent_out(connection));
2123 }
2124 self.runtimes.get_mut(connection).ok_or_else(|| {
2125 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2126 })
2127 }
2128
2129 fn lent_out(&self, connection: &str) -> ServiceError {
2133 ServiceError::Operation(format!(
2134 "runtime connection `{connection}`: a harness turn is already in progress"
2135 ))
2136 }
2137
2138 fn lend_runtime(
2141 &mut self,
2142 connection: &str,
2143 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2144 if self.runtimes_in_flight.contains(connection) {
2145 return Err(self.lent_out(connection));
2146 }
2147 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2148 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2149 })?;
2150 self.runtimes_in_flight.insert(connection.to_string());
2151 Ok(runtime)
2152 }
2153
2154 fn surrender_runtime_of(
2168 &mut self,
2169 params: &RuntimeCloseParams,
2170 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2171 if params.connection.is_empty() {
2172 let expected = params.runtime_id.as_deref().ok_or_else(|| {
2173 ServiceError::InvalidParams("close requires connection or runtime_id".into())
2174 })?;
2175 let connection = self
2176 .runtimes
2177 .iter()
2178 .find(|(_, runtime)| runtime.handle().runtime_id == expected)
2179 .map(|(connection, _)| connection.clone())
2180 .ok_or_else(|| {
2181 ServiceError::InvalidParams(format!("unknown runtime `{expected}`"))
2182 })?;
2183 return self.surrender_runtime(&connection);
2184 }
2185 if let Some(expected) = params.runtime_id.as_deref() {
2186 let held = self
2187 .runtimes
2188 .get(¶ms.connection)
2189 .map(|runtime| runtime.handle().runtime_id.clone());
2190 if held.as_deref() != Some(expected) {
2191 return Err(ServiceError::InvalidParams(format!(
2192 "runtime connection `{}` does not hold runtime `{expected}`",
2193 params.connection
2194 )));
2195 }
2196 }
2197 self.surrender_runtime(¶ms.connection)
2198 }
2199
2200 fn surrender_runtime(
2201 &mut self,
2202 connection: &str,
2203 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2204 if self.runtimes_in_flight.contains(connection) {
2205 return Err(self.lent_out(connection));
2206 }
2207 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2208 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2209 })?;
2210 let process_group = runtime_process_group(runtime.handle());
2211 let runtime_id = runtime.handle().runtime_id.clone();
2212 self.terminal_launches.remove(connection);
2213 self.runtime_sequences.remove(&runtime_id);
2214 self.approvals.forget(connection);
2215 Ok((runtime, process_group))
2216 }
2217
2218 pub fn kill_all_runtime_groups(&self) -> usize {
2229 self.runtimes
2230 .values()
2231 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2232 .count()
2233 }
2234
2235 async fn mutate_session(
2246 &mut self,
2247 verb: crate::SessionVerb,
2248 params: Value,
2249 ) -> std::result::Result<Value, ServiceError> {
2250 let mutation = decode::<crate::SessionMutation>(params)?;
2251 let door = crate::sessions_control::door(&mutation.harness, verb)
2252 .map_err(session_control_error)?;
2253 let outcome = match door {
2254 #[cfg(not(feature = "adapter-api"))]
2258 crate::SessionDoor::Live(command) => {
2259 return Err(ServiceError::Operation(format!(
2260 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2261 session, which needs this build's `adapter-api` feature",
2262 mutation.harness,
2263 verb.as_str()
2264 )));
2265 }
2266 #[cfg(feature = "adapter-api")]
2267 crate::SessionDoor::Live(command) => {
2268 let connection = mutation
2269 .connection
2270 .clone()
2271 .filter(|value| !value.trim().is_empty())
2272 .ok_or_else(|| {
2273 ServiceError::InvalidParams(format!(
2274 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2275 driven session: pass the `connection` of an open runtime \
2276 (`harness.v1.runtimes.start`)",
2277 mutation.harness,
2278 verb.as_str()
2279 ))
2280 })?;
2281 let runtime = self.runtime_mut(&connection)?;
2282 let session = live_session_name(runtime.as_ref(), &mutation);
2283 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2289 .await;
2290 }
2291 _ => run_session_mutation(verb, &mutation).await?,
2292 };
2293 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2294 }
2295
2296 async fn inventory_call(
2300 &self,
2301 method: &str,
2302 params: Value,
2303 ) -> std::result::Result<Value, ServiceError> {
2304 run_inventory(self.inventory_work(method, params)?).await
2305 }
2306
2307 fn inventory_work(
2313 &self,
2314 method: &str,
2315 params: Value,
2316 ) -> std::result::Result<InventoryWork, ServiceError> {
2317 let mut params = decode::<HarnessInventoryParams>(params)?;
2318 if method == "harness.v1.harnesses.probe" {
2319 let harness = params.harness.take().ok_or_else(|| {
2320 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2321 })?;
2322 params.harnesses = vec![harness];
2323 }
2324 let selected = params
2325 .harnesses
2326 .iter()
2327 .map(HarnessId::as_str)
2328 .collect::<std::collections::BTreeSet<_>>();
2329 let supported = harness_support_registry()
2330 .harnesses
2331 .into_iter()
2332 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2333 .collect::<Vec<_>>();
2334 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2335 let known = supported
2336 .iter()
2337 .map(|harness| harness.id.as_str())
2338 .collect::<std::collections::BTreeSet<_>>();
2339 let missing = params
2340 .harnesses
2341 .iter()
2342 .filter(|id| !known.contains(id.as_str()))
2343 .map(HarnessId::as_str)
2344 .collect::<Vec<_>>();
2345 return Err(ServiceError::InvalidParams(format!(
2346 "unknown harness(es): {}",
2347 missing.join(", ")
2348 )));
2349 }
2350 let global_counts = params
2351 .include_sessions
2352 .then(|| self.session_counts(None, ¶ms.harnesses));
2353 let workspace_counts = params
2354 .include_sessions
2355 .then(|| {
2356 params
2357 .workspace
2358 .as_deref()
2359 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2360 })
2361 .flatten();
2362 Ok(InventoryWork {
2363 params,
2364 supported,
2365 global_counts,
2366 workspace_counts,
2367 })
2368 }
2369
2370 #[cfg(feature = "adapter-api")]
2371 async fn harness_authentication_call(
2372 &self,
2373 method: &str,
2374 params: Value,
2375 ) -> std::result::Result<Value, ServiceError> {
2376 match method {
2377 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2378 let params = decode::<HarnessAuthenticationParams>(params)?;
2379 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2380 .map_err(|error| ServiceError::Operation(error.to_string()))
2381 }
2382 "harness.v1.harnesses.auth.begin" => {
2383 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2384 let cwd = params
2385 .cwd
2386 .or_else(|| std::env::current_dir().ok())
2387 .unwrap_or_else(|| PathBuf::from("."));
2388 let plan = crate::harness_authentication_plan(
2389 ¶ms.harness,
2390 params.environment,
2391 params.method,
2392 &cwd,
2393 )
2394 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2395 serde_json::to_value(plan)
2396 .map_err(|error| ServiceError::Operation(error.to_string()))
2397 }
2398 _ => Err(ServiceError::MethodNotFound),
2399 }
2400 }
2401
2402 fn session_counts(
2403 &self,
2404 workspace: Option<&Path>,
2405 harnesses: &[HarnessId],
2406 ) -> BTreeMap<String, usize> {
2407 let mut counts = BTreeMap::new();
2408 for session in self
2409 .catalog
2410 .discover(&DiscoveryQuery {
2411 workspace: workspace.map(Path::to_path_buf),
2412 harnesses: harnesses.to_vec(),
2413 ..DiscoveryQuery::default()
2414 })
2415 .unwrap_or_default()
2416 {
2417 *counts
2418 .entry(session.locator.harness.as_str().to_string())
2419 .or_insert(0) += 1;
2420 }
2421 counts
2422 }
2423}
2424
2425#[async_trait::async_trait]
2426impl SdkService for HarnessSessionService {
2427 fn capabilities(&self) -> SdkCapabilities {
2428 SdkCapabilities::default()
2429 }
2430
2431 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2432 if request.operation == SdkOperation::Events {
2433 let events = self
2434 .poll_sdk_events()
2435 .await
2436 .into_iter()
2437 .map(|(_, event)| event)
2438 .collect::<Vec<_>>();
2439 return serde_json::to_value(events).map_err(|error| {
2440 SdkError::new(
2441 SdkErrorCode::Execution,
2442 request.operation,
2443 error.to_string(),
2444 )
2445 });
2446 }
2447 if self.runtimes.is_empty()
2448 && matches!(
2449 request.operation,
2450 SdkOperation::Input
2451 | SdkOperation::Interrupt
2452 | SdkOperation::Steer
2453 | SdkOperation::Respond
2454 | SdkOperation::Close
2455 )
2456 {
2457 return Err(SdkError::unsupported(request.operation));
2458 }
2459 let method = request
2460 .operation
2461 .method()
2462 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2463 let result = match request.operation {
2464 SdkOperation::Discover
2465 | SdkOperation::Load
2466 | SdkOperation::Export
2467 | SdkOperation::ProfilesList
2468 | SdkOperation::ProfilesGet
2469 | SdkOperation::ProfilesCreate
2470 | SdkOperation::ProfilesDelete
2471 | SdkOperation::SkillsList
2472 | SdkOperation::SkillsInstall
2473 | SdkOperation::SkillsRemove
2474 | SdkOperation::ChannelsList
2475 | SdkOperation::RoutesList
2476 | SdkOperation::TriggersList
2477 | SdkOperation::ChannelsStatus
2478 | SdkOperation::MemoryShow
2479 | SdkOperation::MemorySearch
2480 | SdkOperation::JobsList
2481 | SdkOperation::JobsGet
2482 | SdkOperation::JobsCreate
2483 | SdkOperation::JobsUpdate
2484 | SdkOperation::JobsPause
2485 | SdkOperation::JobsResume
2486 | SdkOperation::JobsRun
2487 | SdkOperation::JobsDelete
2488 | SdkOperation::JobsNotepad
2489 | SdkOperation::JobsNotepadSet
2490 | SdkOperation::JobsNotepadDelete
2491 | SdkOperation::RunsList
2492 | SdkOperation::RunsGet
2493 | SdkOperation::ApprovalsList
2494 | SdkOperation::OrchestrationLoad
2495 | SdkOperation::OrchestrationSave
2496 | SdkOperation::OrchestrationCompile
2497 | SdkOperation::OrchestrationDecompile
2498 | SdkOperation::OrchestrationImport
2499 | SdkOperation::OrchestrationExport
2500 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2501 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2504 SdkOperation::Start
2505 | SdkOperation::Resume
2506 | SdkOperation::Input
2507 | SdkOperation::Interrupt
2508 | SdkOperation::Steer
2509 | SdkOperation::Respond
2510 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2511 SdkOperation::SessionsNew => {
2516 self.mutate_session(crate::SessionVerb::New, request.params)
2517 .await
2518 }
2519 SdkOperation::SessionsReset => {
2520 self.mutate_session(crate::SessionVerb::Reset, request.params)
2521 .await
2522 }
2523 SdkOperation::SessionsArchive => {
2524 self.mutate_session(crate::SessionVerb::Archive, request.params)
2525 .await
2526 }
2527 SdkOperation::SessionsDelete => {
2528 self.mutate_session(crate::SessionVerb::Delete, request.params)
2529 .await
2530 }
2531 SdkOperation::Events => unreachable!("handled before method dispatch"),
2532 };
2533 result.map_err(|error| sdk_error(request.operation, error))
2534 }
2535
2536 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2537 Ok(self
2538 .poll_sdk_events()
2539 .await
2540 .into_iter()
2541 .map(|(_, event)| event)
2542 .collect())
2543 }
2544}
2545
2546#[cfg(feature = "adapter-api")]
2547struct HostedRuntimeLease {
2548 connection: HostedHarnessConnection,
2549 _host: std::sync::Arc<HostedHarnessRuntime>,
2550 _server: crate::server::FrontendHttpServer,
2551}
2552
2553#[async_trait::async_trait]
2554#[cfg(feature = "adapter-api")]
2555impl RuntimeConnection for HostedRuntimeLease {
2556 fn handle(&self) -> &crate::RuntimeHandle {
2557 self.connection.handle()
2558 }
2559
2560 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2561 self.connection.send_input(input).await
2562 }
2563
2564 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2565 self.connection.next_event().await
2566 }
2567
2568 async fn interrupt(&mut self) -> crate::Result<()> {
2569 self.connection.interrupt().await
2570 }
2571
2572 async fn steer(&mut self, text: String) -> crate::Result<()> {
2575 self.connection.steer(text).await
2576 }
2577
2578 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2579 self.connection.respond(request_id, response).await
2580 }
2581
2582 async fn close(&mut self) -> crate::Result<()> {
2583 self.connection.close().await
2584 }
2585}
2586
2587struct InventoryWork {
2590 params: HarnessInventoryParams,
2591 supported: Vec<crate::HarnessSupportDescriptor>,
2592 global_counts: Option<BTreeMap<String, usize>>,
2593 workspace_counts: Option<BTreeMap<String, usize>>,
2594}
2595
2596async fn run_session_mutation(
2603 verb: crate::SessionVerb,
2604 mutation: &crate::SessionMutation,
2605) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2606 let door =
2613 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2614 if let crate::SessionDoor::Http = door {
2615 return crate::sessions_control::mutate(verb, mutation)
2616 .await
2617 .map_err(session_control_error);
2618 }
2619 let mutation = mutation.clone();
2620 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2621 .await
2622 .map_err(|error| {
2623 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2624 })?
2625 .map_err(session_control_error)
2626}
2627
2628async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2631 let InventoryWork {
2632 params,
2633 supported,
2634 global_counts,
2635 workspace_counts,
2636 } = work;
2637 let probes = supported.into_iter().map(|descriptor| {
2638 let global = global_counts
2639 .as_ref()
2640 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2641 let workspace = workspace_counts
2642 .as_ref()
2643 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2644 probe_harness(descriptor, ¶ms, global, workspace)
2645 });
2646 let harnesses = futures::future::join_all(probes).await;
2647 serde_json::to_value(HarnessInventoryReport {
2648 probe: params.probe,
2649 workspace: params.workspace,
2650 harnesses,
2651 })
2652 .map_err(|error| ServiceError::Operation(error.to_string()))
2653}
2654
2655async fn probe_harness(
2656 descriptor: crate::HarnessSupportDescriptor,
2657 params: &HarnessInventoryParams,
2658 global: Option<usize>,
2659 workspace: Option<usize>,
2660) -> LocalHarness {
2661 let launch = descriptor.runtime.default_launch.as_ref();
2662 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2667 .then(crate::orchestrator::daemon_entry)
2668 .and_then(Result::ok);
2669 let executable = match &orchestrator_entry {
2670 Some(entry) => Some(entry.clone()),
2671 None => launch.and_then(|launch| find_executable(&launch.program)),
2672 };
2673 let installed = executable.is_some();
2674 let version = if params.skip_versions || orchestrator_entry.is_some() {
2675 None
2678 } else {
2679 match executable.as_deref() {
2680 Some(path) => executable_version(path).await,
2681 None => None,
2682 }
2683 };
2684 let configured = auth_evidence(descriptor.id.as_str());
2685 let mut auth = if configured {
2686 HarnessAuthState::Configured
2687 } else if matches!(
2688 descriptor.id.as_str(),
2689 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2690 ) {
2691 HarnessAuthState::Required
2696 } else {
2697 HarnessAuthState::Unknown
2698 };
2699 let mut runtime = if installed {
2700 HarnessRuntimeState::Degraded
2701 } else {
2702 HarnessRuntimeState::Unavailable
2703 };
2704 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2705 let mut reason = (!installed).then(|| {
2706 if is_orchestrator {
2707 format!(
2708 "{} is supported but its daemon entry `{}` was not found",
2709 descriptor.display_name,
2710 crate::orchestrator::DAEMON_ENTRY
2711 )
2712 } else {
2713 format!(
2714 "{} is supported but `{}` was not found on PATH",
2715 descriptor.display_name,
2716 launch
2717 .map(|launch| launch.program.as_str())
2718 .unwrap_or("executable")
2719 )
2720 }
2721 });
2722 let mut repair = (!installed).then(|| {
2723 if is_orchestrator {
2724 format!(
2725 "Install the `supercode-orchestrator` package so `{}` resolves.",
2726 crate::orchestrator::DAEMON_ENTRY
2727 )
2728 } else {
2729 format!(
2730 "Install {} and ensure `{}` is on PATH.",
2731 descriptor.display_name,
2732 launch
2733 .map(|launch| launch.program.as_str())
2734 .unwrap_or("its executable")
2735 )
2736 }
2737 });
2738
2739 if installed && params.probe == HarnessProbeLevel::Handshake {
2740 let backend_params = RuntimeBackendParams {
2741 harness: descriptor.id.clone(),
2742 protocol: None,
2743 launch: None,
2744 base_url: None,
2745 policy: RuntimePolicy::Default,
2746 };
2747 match runtime_backend(&backend_params) {
2748 Ok(backend) => {
2749 let cwd = params
2750 .workspace
2751 .clone()
2752 .or_else(|| std::env::current_dir().ok())
2753 .unwrap_or_else(|| PathBuf::from("."));
2754 let isolated = descriptor
2755 .runtime
2756 .default_launch
2757 .clone()
2758 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2759 let Some(isolated) = isolated else {
2760 reason = Some(
2761 "No-prompt runtime handshake could not create its isolated harness home."
2762 .into(),
2763 );
2764 repair = Some(
2765 "Check temporary-directory permissions, then run the handshake probe again."
2766 .into(),
2767 );
2768 let running = probe_running_instance(descriptor.id.as_str());
2769 return LocalHarness {
2770 gateway: gateway_health(
2771 descriptor.id.as_str(),
2772 installed,
2773 running.as_ref(),
2774 version.as_deref(),
2775 ),
2776 id: descriptor.id,
2777 display_name: descriptor.display_name,
2778 supported: true,
2779 installed,
2780 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2781 version,
2782 auth,
2783 runtime,
2784 protocol: descriptor.runtime.protocol,
2785 capabilities: descriptor.runtime.capabilities.clone(),
2786 effective_capabilities: descriptor.runtime.capabilities,
2787 sessions: HarnessSessionCounts { global, workspace },
2788 running,
2789 reason,
2790 repair,
2791 };
2792 };
2793 match tokio::time::timeout(
2794 Duration::from_secs(30),
2795 backend.start(RuntimeStartRequest {
2796 cwd,
2797 launch: Some(isolated.launch.clone()),
2798 mcp_servers: Vec::new(),
2799 approval_policy: None,
2800 }),
2801 )
2802 .await
2803 {
2804 Ok(Ok(mut connection)) => {
2805 match stabilize_handshake(connection.as_mut()).await {
2806 Ok(()) => {
2807 auth = HarnessAuthState::Ready;
2808 runtime = HarnessRuntimeState::Ready;
2809 reason = Some(
2810 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2811 .into(),
2812 );
2813 repair = None;
2814 }
2815 Err(message) => {
2816 auth = if looks_like_auth_error(&message) {
2817 HarnessAuthState::Required
2818 } else if configured {
2819 HarnessAuthState::Configured
2820 } else {
2821 HarnessAuthState::Unknown
2822 };
2823 reason = Some(format!(
2824 "No-prompt runtime handshake became unhealthy during startup: {message}"
2825 ));
2826 repair = Some(if auth == HarnessAuthState::Required {
2827 format!(
2828 "Run `{}` interactively once and complete sign-in, then probe again.",
2829 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2830 )
2831 } else {
2832 "Run the harness directly to inspect its startup failure, then probe again."
2833 .into()
2834 });
2835 }
2836 }
2837 let _ =
2838 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2839 }
2840 Ok(Err(error)) => {
2841 let message = truncate_text(&error.to_string(), 500);
2842 auth = if looks_like_auth_error(&message) {
2843 HarnessAuthState::Required
2844 } else if configured {
2845 HarnessAuthState::Configured
2846 } else {
2847 HarnessAuthState::Unknown
2848 };
2849 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2850 repair = Some(if auth == HarnessAuthState::Required {
2851 format!(
2852 "Run `{}` interactively once and complete sign-in, then probe again.",
2853 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2854 )
2855 } else {
2856 "Check the harness installation and run the handshake probe again."
2857 .into()
2858 });
2859 }
2860 Err(_) => {
2861 reason =
2862 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2863 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2864 }
2865 }
2866 let _ = isolated.cleanup();
2875 tokio::time::sleep(Duration::from_millis(250)).await;
2876 if let Err(error) = isolated.cleanup() {
2877 auth = if configured {
2878 HarnessAuthState::Configured
2879 } else {
2880 HarnessAuthState::Unknown
2881 };
2882 runtime = HarnessRuntimeState::Degraded;
2883 reason = Some(format!(
2884 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2885 ));
2886 repair = Some(
2887 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2888 .into(),
2889 );
2890 }
2891 }
2892 Err(error) => {
2893 reason = Some(error_message(error));
2894 }
2895 }
2896 } else if installed && configured {
2897 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2898 } else if installed && auth == HarnessAuthState::Required {
2899 reason = Some("Executable found, but no native authentication evidence is present.".into());
2900 repair = Some(format!(
2901 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2902 descriptor.id.as_str()
2903 ));
2904 } else if installed {
2905 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2906 repair = Some(format!(
2907 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2908 launch
2909 .map(|launch| launch.program.as_str())
2910 .unwrap_or("the harness")
2911 ));
2912 }
2913
2914 let effective_capabilities = if installed {
2915 descriptor.runtime.capabilities.clone()
2916 } else {
2917 unavailable_capabilities()
2918 };
2919 let running = probe_running_instance(descriptor.id.as_str());
2920 LocalHarness {
2921 gateway: gateway_health(
2922 descriptor.id.as_str(),
2923 installed,
2924 running.as_ref(),
2925 version.as_deref(),
2926 ),
2927 id: descriptor.id,
2928 display_name: descriptor.display_name,
2929 supported: true,
2930 installed,
2931 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2932 version,
2933 auth,
2934 runtime,
2935 protocol: descriptor.runtime.protocol,
2936 capabilities: descriptor.runtime.capabilities,
2937 effective_capabilities,
2938 sessions: HarnessSessionCounts { global, workspace },
2939 running,
2940 reason,
2941 repair,
2942 }
2943}
2944
2945async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2946 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2947 loop {
2948 let now = tokio::time::Instant::now();
2949 if now >= deadline {
2950 return Ok(());
2951 }
2952 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2953 Err(_) => return Ok(()),
2954 Ok(Ok(Some(event))) => {
2955 if let Some(message) = handshake_event_failure(&event) {
2956 return Err(truncate_text(&message, 500));
2957 }
2958 }
2959 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2960 Ok(Err(error)) => return Err(error.to_string()),
2961 }
2962 }
2963}
2964
2965fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2966 let detail = event
2967 .payload
2968 .get("message")
2969 .or_else(|| event.payload.get("line"))
2970 .and_then(Value::as_str)
2971 .unwrap_or(event.kind.as_str());
2972 match event.kind.as_str() {
2973 "transport_closed" => Some("runtime transport closed during startup".into()),
2974 "transport_error" => Some(format!("runtime transport error: {detail}")),
2975 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2976 _ => None,
2981 }
2982}
2983
2984fn indexed_claude_window(
2985 locator: &SessionLocator,
2986 options: &SessionLoadOptions,
2987) -> std::result::Result<Option<Value>, ServiceError> {
2988 use supercode_interchange::session::ClaudeReadIndex;
2989 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2992 || options.include_subagents != Some(false)
2993 {
2994 return Ok(None);
2995 }
2996 let crate::StorageLocator::File { path } = &locator.storage else {
2997 return Ok(None);
2998 };
2999 if !ClaudeReadIndex::supports(path)
3000 .map_err(|error| ServiceError::Operation(error.to_string()))?
3001 {
3002 return Ok(None);
3003 }
3004 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
3005 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3006 let total = index.len();
3007 let (offset, end) = projected_message_window(total, options);
3008 let session = index
3009 .read_messages(offset..end)
3010 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3011 let summary = index
3012 .read_summary()
3013 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3014 let selected_options = SessionLoadOptions {
3015 message_offset: None,
3016 message_limit: None,
3017 message_tail: None,
3018 ..options.clone()
3019 };
3020 let mut selected = projected_session_json(&session, &selected_options);
3021 selected["raw_record_count"] = json!(index.raw_record_count());
3022 Ok(Some(json!({
3023 "session": selected,
3024 "summary": projected_session_summary(&summary, options),
3025 "window": {
3026 "has_more": offset > 0 || end < total, "has_newer": end < total,
3027 "has_older": offset > 0, "newer_items": index.item_count(end..total),
3028 "offset": offset, "older_items": index.item_count(0..offset),
3029 "returned": end - offset, "total_messages": total,
3030 }
3031 })))
3032}
3033
3034fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
3035 let total_messages = session.messages.len();
3036 let (offset, end) = projected_message_window(total_messages, options);
3037 json!({
3038 "session": projected_session_json(session, options),
3039 "summary": projected_session_summary(session, options),
3040 "window": {
3041 "has_more": offset > 0 || end < total_messages,
3042 "has_newer": end < total_messages,
3043 "has_older": offset > 0,
3044 "newer_items": normalized_item_count(&session.messages[end..]),
3045 "offset": offset,
3046 "older_items": normalized_item_count(&session.messages[..offset]),
3047 "returned": end.saturating_sub(offset),
3048 "total_messages": total_messages,
3049 }
3050 })
3051}
3052
3053fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
3054 messages
3055 .iter()
3056 .map(|message| {
3057 let conversation = usize::from(
3058 matches!(message.role, Role::Assistant | Role::User)
3059 && message_has_content(message),
3060 );
3061 let tool_result =
3062 usize::from(message.role == Role::Tool && message_has_content(message));
3063 conversation + tool_result + message.tool_calls().len()
3064 })
3065 .sum()
3066}
3067
3068fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
3069 let mut conversational = session.messages.iter().filter(|message| {
3070 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
3071 });
3072 let first_message = conversational.clone().next();
3073 let last_message = conversational.next_back();
3074 let mut assistant = session
3075 .messages
3076 .iter()
3077 .filter(|message| message.role == Role::Assistant && message_has_content(message));
3078 let first_assistant_message = assistant.clone().next();
3079 let last_assistant_message = assistant.next_back();
3080 let end_of_turn = session
3081 .messages
3082 .iter()
3083 .rev()
3084 .find(|message| message.role != Role::System)
3085 .is_some_and(|message| {
3086 message.role == Role::Assistant
3087 && message_has_content(message)
3088 && message.tool_calls().is_empty()
3089 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
3091 });
3092 let project = |message: Option<&crate::ChatMessage>| {
3093 message.map(|message| project_inline_media(message_json(message), options))
3094 };
3095 json!({
3096 "end_of_turn": end_of_turn,
3097 "first_assistant_message": project(first_assistant_message),
3098 "first_message": project(first_message),
3099 "last_assistant_message": project(last_assistant_message),
3100 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3101 "last_message": project(last_message),
3102 })
3103}
3104
3105fn message_has_content(message: &crate::ChatMessage) -> bool {
3106 message
3107 .content
3108 .as_deref()
3109 .is_some_and(|content| !content.trim().is_empty())
3110 || message
3111 .content_parts
3112 .as_ref()
3113 .is_some_and(|parts| !parts.is_empty())
3114}
3115
3116fn message_text(message: &crate::ChatMessage) -> String {
3117 if let Some(content) = &message.content {
3118 return content.clone();
3119 }
3120 message
3121 .content_parts
3122 .as_ref()
3123 .into_iter()
3124 .flatten()
3125 .filter_map(|part| part.get("text").and_then(Value::as_str))
3126 .collect::<Vec<_>>()
3127 .join("\n")
3128}
3129
3130fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3131 let (offset, end) = projected_message_window(session.messages.len(), options);
3132 let messages = session.messages[offset..end]
3133 .iter()
3134 .map(|message| project_inline_media(message_json(message), options))
3135 .collect::<Vec<_>>();
3136 let subagents = if options.include_subagents.unwrap_or(true) {
3137 let subagent_options = SessionLoadOptions {
3142 message_limit: None,
3143 message_offset: None,
3144 message_tail: None,
3145 ..options.clone()
3146 };
3147 session
3148 .subagents
3149 .iter()
3150 .map(|subagent| projected_session_json(subagent, &subagent_options))
3151 .collect::<Vec<_>>()
3152 } else {
3153 Vec::new()
3154 };
3155 json!({
3156 "source": match session.meta.source {
3157 SessionSource::ClaudeCode => "claude_code",
3158 SessionSource::Codex => "codex",
3159 SessionSource::Gemini => "gemini",
3160 SessionSource::Goose => "goose",
3161 SessionSource::Grok => "grok",
3162 SessionSource::Native => "native",
3163 SessionSource::OpenClaw => "openclaw",
3164 SessionSource::Hermes => "hermes",
3165 SessionSource::OpenCode => "opencode",
3166 SessionSource::Pi => "pi",
3167 },
3168 "session_id": session.meta.session_id,
3169 "ended_at": session.meta.ended_at,
3170 "end_reason": session.meta.end_reason,
3171 "model": session.meta.model,
3172 "resolved_profile": recorded_session_profile(session),
3173 "resolved_profile_history": recorded_session_profile_history(session),
3174 "cwd": session.meta.cwd,
3175 "system_prompt": session.meta.system_prompt,
3176 "agent_id": session.meta.agent_id,
3177 "parent_tool_use_id": session.meta.parent_tool_use_id,
3178 "lineage": session.meta.lineage,
3179 "messages": messages,
3180 "subagents": subagents,
3181 "raw_record_count": session.raw.len(),
3182 "parse_error_lines": session.parse_error_lines,
3183 })
3184}
3185
3186fn recorded_session_usage(session: &Session) -> Vec<Value> {
3191 let mut days: std::collections::BTreeMap<(String, String), [u64; 4]> =
3192 std::collections::BTreeMap::new();
3193 let mut add =
3194 |at: Option<&str>, model: &str, input: u64, output: u64, read: u64, write: u64| {
3195 let day = at
3196 .and_then(|t| t.get(..10))
3197 .unwrap_or("unknown")
3198 .to_string();
3199 let entry = days.entry((day, model.to_string())).or_insert([0; 4]);
3200 entry[0] += input;
3201 entry[1] += output;
3202 entry[2] += read;
3203 entry[3] += write;
3204 };
3205 let records = session
3206 .raw
3207 .iter()
3208 .filter_map(|line| serde_json::from_str::<Value>(line).ok());
3209 let n = |v: Option<&Value>| v.and_then(Value::as_u64).unwrap_or(0);
3210 match session.meta.source {
3211 SessionSource::ClaudeCode => {
3212 let mut responses: std::collections::BTreeMap<String, Value> =
3213 std::collections::BTreeMap::new();
3214 for record in records {
3215 let Some(message) = record.get("message") else {
3216 continue;
3217 };
3218 let (Some(id), Some(_)) = (
3219 message.get("id").and_then(Value::as_str),
3220 message.get("usage"),
3221 ) else {
3222 continue;
3223 };
3224 if message.get("model").and_then(Value::as_str) == Some("<synthetic>") {
3225 continue;
3226 }
3227 responses.insert(id.to_string(), record.clone());
3228 }
3229 for record in responses.values() {
3230 let message = &record["message"];
3231 let usage = &message["usage"];
3232 add(
3233 record.get("timestamp").and_then(Value::as_str),
3234 message
3235 .get("model")
3236 .and_then(Value::as_str)
3237 .unwrap_or("unknown"),
3238 n(usage.get("input_tokens")),
3239 n(usage.get("output_tokens")),
3240 n(usage.get("cache_read_input_tokens")),
3241 n(usage.get("cache_creation_input_tokens")),
3242 );
3243 }
3244 }
3245 SessionSource::Codex => {
3246 let mut model = session
3247 .meta
3248 .model
3249 .clone()
3250 .unwrap_or_else(|| "unknown".into());
3251 let mut previous = [0u64; 3];
3252 for record in records {
3253 let payload = record.get("payload").unwrap_or(&Value::Null);
3254 if record.get("type").and_then(Value::as_str) == Some("turn_context") {
3255 if let Some(m) = payload.get("model").and_then(Value::as_str) {
3256 model = m.to_string();
3257 }
3258 continue;
3259 }
3260 if payload.get("type").and_then(Value::as_str) != Some("token_count") {
3261 continue;
3262 }
3263 let Some(total) = payload.get("info").and_then(|i| i.get("total_token_usage"))
3264 else {
3265 continue;
3266 };
3267 let now = [
3268 n(total.get("input_tokens")),
3269 n(total.get("output_tokens")),
3270 n(total.get("cached_input_tokens")),
3271 ];
3272 if now == previous {
3273 continue;
3274 }
3275 add(
3276 record.get("timestamp").and_then(Value::as_str),
3277 &model,
3278 now[0].saturating_sub(previous[0]),
3279 now[1].saturating_sub(previous[1]),
3280 now[2].saturating_sub(previous[2]),
3281 0,
3282 );
3283 previous = now;
3284 }
3285 }
3286 _ => {}
3287 }
3288 days.into_iter()
3289 .map(|((day, model), [input, output, read, write])| {
3290 let mut usage = json!({"at": format!("{day}T00:00:00Z"), "model": model,
3291 "input_tokens": input, "output_tokens": output,
3292 "cache_read_tokens": read, "cache_write_tokens": write});
3293 if let Some(price) = crate::pricing::built_in(&model) {
3296 usage["cost"] = json!({"amount": price.cost_usd(input + read + write, output),
3297 "currency": "USD", "source": "rate_table"});
3298 }
3299 usage
3300 })
3301 .collect()
3302}
3303
3304fn recorded_session_profile(session: &Session) -> Value {
3307 let harness = match session.meta.source {
3308 SessionSource::ClaudeCode => "claude-code",
3309 SessionSource::Codex => "codex",
3310 SessionSource::Gemini => "gemini",
3311 SessionSource::Goose => "goose",
3312 SessionSource::Grok => "grok",
3313 SessionSource::Native => "native",
3314 SessionSource::OpenClaw => "openclaw",
3315 SessionSource::Hermes => "hermes",
3316 SessionSource::OpenCode => "opencode",
3317 SessionSource::Pi => "pi",
3318 };
3319 let mut model = session
3320 .messages
3321 .iter()
3322 .rev()
3323 .find_map(|message| message.metadata.get("model").cloned())
3324 .or_else(|| session.meta.model.clone());
3325 let mut cwd = session
3326 .meta
3327 .cwd
3328 .as_ref()
3329 .map(|path| path.to_string_lossy().into_owned());
3330 let mut version: Option<String> = None;
3331 let mut approval = Value::Null;
3332 let mut sandbox = Value::Null;
3333 for header in &session.meta.codex_headers {
3334 let payload = header.get("payload").unwrap_or(header);
3335 if let Some(value) = payload.get("cli_version").and_then(Value::as_str) {
3336 version = Some(value.to_owned());
3337 }
3338 if let Some(value) = payload.get("model").and_then(Value::as_str) {
3339 model = Some(value.to_owned());
3340 }
3341 if let Some(value) = payload.get("cwd").and_then(Value::as_str) {
3342 cwd = Some(value.to_owned());
3343 }
3344 if let Some(value) = payload.get("approval_policy") {
3345 approval = value.clone();
3346 }
3347 if let Some(value) = payload.get("sandbox_policy") {
3348 sandbox = value.clone();
3349 }
3350 }
3351 if session.meta.source == SessionSource::ClaudeCode {
3352 version = session
3353 .raw
3354 .iter()
3355 .rev()
3356 .filter_map(|line| serde_json::from_str::<Value>(line).ok())
3357 .find_map(|record| {
3358 record
3359 .get("version")
3360 .and_then(Value::as_str)
3361 .map(str::to_owned)
3362 });
3363 }
3364 json!({"worker":{"harness":harness,"model":model,"cwd":cwd,"binary":{"version":version}},
3365 "permissions":{"native":{"approval":approval,"sandbox":sandbox}},"source":"native-session"})
3366}
3367
3368fn recorded_session_profile_history(session: &Session) -> Vec<Value> {
3371 let mut profile = json!({"worker":{"harness":if session.meta.source == SessionSource::Codex {"codex"} else {"claude-code"},
3372 "model":null,"cwd":null,"binary":{"version":null}},
3373 "permissions":{"native":{"approval":null,"sandbox":null}},"source":"native-session"});
3374 let records: Box<dyn Iterator<Item = (usize, Value)> + '_> =
3375 if session.meta.source == SessionSource::Codex {
3376 Box::new(session.meta.codex_headers.iter().cloned().enumerate())
3377 } else if session.meta.source == SessionSource::ClaudeCode {
3378 Box::new(session.raw.iter().enumerate().filter_map(|(index, line)| {
3379 serde_json::from_str::<Value>(line)
3380 .ok()
3381 .map(|record| (index, record))
3382 }))
3383 } else {
3384 return Vec::new();
3385 };
3386 let mut history = Vec::new();
3387 for (index, record) in records {
3388 let before = profile.clone();
3389 let payload = record.get("payload").unwrap_or(&record);
3390 for (native, target) in [("cwd", "cwd"), ("model", "model")] {
3391 if let Some(value) = payload.get(native).and_then(Value::as_str) {
3392 profile["worker"][target] = json!(value);
3393 }
3394 }
3395 let version = payload
3396 .get("cli_version")
3397 .or_else(|| record.get("version"))
3398 .and_then(Value::as_str);
3399 if let Some(value) = version {
3400 profile["worker"]["binary"]["version"] = json!(value);
3401 }
3402 if let Some(value) = record
3403 .get("message")
3404 .and_then(|m| m.get("model"))
3405 .and_then(Value::as_str)
3406 {
3407 profile["worker"]["model"] = json!(value);
3408 }
3409 for (native, target) in [
3410 ("approval_policy", "approval"),
3411 ("sandbox_policy", "sandbox"),
3412 ] {
3413 if let Some(value) = payload.get(native) {
3414 profile["permissions"]["native"][target] = value.clone();
3415 }
3416 }
3417 if profile != before {
3418 history.push(json!({"native_key":format!("profile-record:{index}"),"at":record.get("timestamp").and_then(Value::as_str),"profile":profile}));
3419 }
3420 }
3421 history
3422}
3423
3424fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3425 if let Some(tail) = options.message_tail {
3426 return (total.saturating_sub(tail), total);
3427 }
3428 let offset = options.message_offset.unwrap_or(0).min(total);
3429 let end = options
3430 .message_limit
3431 .map(|limit| offset.saturating_add(limit).min(total))
3432 .unwrap_or(total);
3433 (offset, end)
3434}
3435
3436fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3437 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3438 return message;
3439 };
3440 for part in parts {
3441 let Some(url) = part
3442 .get("image_url")
3443 .and_then(|image| image.get("url"))
3444 .and_then(Value::as_str)
3445 else {
3446 continue;
3447 };
3448 let Some(rest) = url.strip_prefix("data:") else {
3449 continue;
3450 };
3451 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3452 continue;
3453 };
3454 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3455 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3456 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3457 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3458 || options
3459 .max_inline_media_bytes
3460 .is_some_and(|limit| decoded_bytes > limit);
3461 if should_elide {
3462 *part = json!({
3463 "type": "media_reference",
3464 "media_type": media_type,
3465 "encoding": "base64",
3466 "encoded_bytes": encoded.len(),
3467 "decoded_bytes": decoded_bytes,
3468 "omitted": true,
3469 });
3470 }
3471 }
3472 message
3473}
3474
3475#[derive(Deserialize)]
3476struct LocatorParams {
3477 locator: SessionLocator,
3478 #[serde(default)]
3491 fidelity: Option<Fidelity>,
3492 #[serde(default)]
3495 view: Option<SessionReadView>,
3496}
3497
3498#[derive(Deserialize)]
3499struct SessionReadView {
3500 #[serde(default)]
3503 tail_messages: Option<usize>,
3504 #[serde(default)]
3507 include_subagents: bool,
3508 #[serde(default)]
3510 display_history: bool,
3511 #[serde(default)]
3514 max_message_chars: Option<usize>,
3515}
3516
3517impl LocatorParams {
3518 fn read_fidelity(&self) -> Fidelity {
3519 self.fidelity.unwrap_or(Fidelity::Semantic)
3520 }
3521
3522 fn include_subagents(&self) -> bool {
3523 self.view
3524 .as_ref()
3525 .map(|view| view.include_subagents)
3526 .unwrap_or(true)
3527 }
3528
3529 fn tail_messages(&self) -> Option<usize> {
3530 self.view
3531 .as_ref()
3532 .and_then(|view| view.tail_messages)
3533 .map(|limit| limit.clamp(1, 5_000))
3534 }
3535
3536 fn display_history(&self) -> bool {
3537 self.view.as_ref().is_some_and(|view| view.display_history)
3538 }
3539
3540 fn max_message_chars(&self) -> Option<usize> {
3541 self.view
3542 .as_ref()
3543 .and_then(|view| view.max_message_chars)
3544 .map(|limit| limit.clamp(256, 64_000))
3545 }
3546
3547 fn bound_session(&self, session: &mut Session) {
3548 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3549 }
3550}
3551
3552#[derive(Debug, Clone, Copy, Default, Deserialize)]
3553#[serde(rename_all = "snake_case")]
3554enum InlineMediaMode {
3555 #[default]
3556 Full,
3557 Metadata,
3558}
3559
3560#[derive(Debug, Clone, Default, Deserialize)]
3561#[serde(default)]
3562struct SessionLoadOptions {
3563 include_subagents: Option<bool>,
3564 inline_media: InlineMediaMode,
3565 max_inline_media_bytes: Option<usize>,
3566 message_limit: Option<usize>,
3567 message_offset: Option<usize>,
3568 message_tail: Option<usize>,
3569}
3570
3571impl SessionLoadOptions {
3572 fn validate(&self) -> std::result::Result<(), ServiceError> {
3573 if self.message_tail.is_some()
3574 && (self.message_limit.is_some() || self.message_offset.is_some())
3575 {
3576 return Err(ServiceError::InvalidParams(
3577 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3578 .into(),
3579 ));
3580 }
3581 Ok(())
3582 }
3583}
3584
3585#[derive(Deserialize)]
3586struct LoadSessionParams {
3587 #[serde(flatten)]
3588 read: LocatorParams,
3589 #[serde(default)]
3590 options: Option<SessionLoadOptions>,
3591}
3592
3593#[derive(Deserialize)]
3594struct UnfollowParams {
3595 subscription: String,
3596}
3597
3598#[derive(Debug, Deserialize)]
3599#[serde(deny_unknown_fields)]
3600struct IndexResizeParams {
3601 subscription: String,
3602 limit: usize,
3603}
3604
3605#[derive(Deserialize)]
3606struct ActivitySubscribeParams {
3607 locators: Vec<SessionLocator>,
3608 #[serde(default)]
3609 homes: crate::HarnessHomes,
3610}
3611
3612#[derive(Deserialize)]
3613struct MessageSessionParams {
3614 locator: SessionLocator,
3615 text: String,
3616 #[serde(default)]
3617 subject: Option<String>,
3618 #[serde(default)]
3620 idempotency_key: Option<String>,
3621 #[serde(default)]
3623 channel: bool,
3624 #[serde(default)]
3628 from_name: Option<String>,
3629 #[serde(default)]
3631 voice_for: Option<crate::mailbox::MailAddress>,
3632 #[serde(default)]
3634 sender_name: Option<String>,
3635 #[serde(default)]
3637 in_reply_to: Option<String>,
3638 #[serde(default)]
3640 notify_when_idle: bool,
3641 #[serde(default)]
3646 as_user: bool,
3647 #[serde(default)]
3650 homes: crate::HarnessHomes,
3651}
3652
3653#[derive(Deserialize)]
3654#[serde(deny_unknown_fields)]
3655struct InboxParams {
3656 #[serde(default)]
3658 from_name: Option<String>,
3659 #[serde(default)]
3661 address: Option<String>,
3662 #[serde(default)]
3664 all: bool,
3665}
3666
3667#[derive(Deserialize)]
3668#[serde(deny_unknown_fields)]
3669struct HarnessSettingsParams {
3670 harness: String,
3671}
3672
3673#[derive(Deserialize)]
3674#[serde(deny_unknown_fields)]
3675struct ConfigureHarnessParams {
3676 harness: String,
3677 #[serde(default)]
3678 changes: Vec<crate::HarnessSettingChange>,
3679 #[serde(default)]
3680 expected_revision: Option<String>,
3681}
3682
3683fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3684 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3685 Ok(report) => (
3686 serde_json::to_value(report).unwrap_or(Value::Null),
3687 Value::Null,
3688 ),
3689 Err(error) => (
3690 Value::Null,
3691 Value::String(format!(
3692 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3693 )),
3694 ),
3695 }
3696}
3697
3698async fn message_live_session(params: &MessageSessionParams) -> Value {
3712 use crate::mail_route::{Delivered, NoDoor, Refused};
3713 let (inbound_controls, inbound_controls_error) =
3714 claude_inbound_controls_or_error(¶ms.homes);
3715 let refused = |reason: &str, message: String| {
3716 json!({
3717 "delivered_to_bus": false,
3718 "refusal": {"reason": reason, "message": message},
3719 "inbound_controls": inbound_controls,
3720 "inbound_controls_error": inbound_controls_error,
3721 })
3722 };
3723 if params.text.trim().is_empty() {
3724 return refused(
3725 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3726 "refusing to deliver an empty message".into(),
3727 );
3728 }
3729 let sender = match operator_address(params.from_name.as_deref()) {
3730 Ok(sender) => sender,
3731 Err(message) => return refused("invalid_sender", message),
3732 };
3733 let receiver = match crate::mailbox::MailAddress::new(
3734 crate::mailbox::local_machine_name(),
3735 params.locator.harness.as_str(),
3736 ¶ms.locator.session_id,
3737 ) {
3738 Ok(receiver) => receiver,
3739 Err(error) => return refused("delivery_failed", error.to_string()),
3740 };
3741 if params
3742 .voice_for
3743 .as_ref()
3744 .is_some_and(|voice_for| voice_for != &receiver)
3745 {
3746 return refused(
3747 "invalid_sender",
3748 "a voice front must represent the receiving session".into(),
3749 );
3750 }
3751 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3753 && crate::runtime_mail::controlled_runtime(
3754 HarnessId::CLAUDE_CODE,
3755 ¶ms.locator.session_id,
3756 )
3757 .is_none()
3758 {
3759 if let Err(refusal) =
3760 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3761 {
3762 return refused(refusal.reason.as_str(), refusal.message);
3763 }
3764 }
3765 if params.as_user {
3766 return message_as_user(
3767 params,
3768 sender,
3769 receiver,
3770 inbound_controls,
3771 inbound_controls_error,
3772 )
3773 .await;
3774 }
3775 let door = match crate::mail_route::door_for(¶ms.homes, &receiver) {
3776 Ok(door) => door,
3777 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3778 return refused(
3779 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3780 format!(
3781 "no running `{}` session `{}` is reachable; its transcript is persisted only",
3782 params.locator.harness.as_str(),
3783 params.locator.session_id
3784 ),
3785 )
3786 }
3787 };
3788 let mut envelope = match crate::mailbox::Envelope::new(
3789 sender.clone(),
3790 params
3791 .sender_name
3792 .clone()
3793 .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3794 if params.channel {
3795 crate::mailbox::MailKind::Channel
3796 } else {
3797 crate::mailbox::MailKind::Peer
3798 },
3799 crate::mailbox::ReplyVia::Command,
3800 params.text.clone(),
3801 ) {
3802 Ok(envelope) => envelope,
3803 Err(error) => return refused("delivery_failed", error.to_string()),
3804 };
3805 if let Some(key) = ¶ms.idempotency_key {
3806 if key.is_empty() || key.len() > 256 {
3807 return refused("invalid_key", "idempotency_key needs 1–256 bytes".into());
3808 }
3809 let identity = serde_json::to_vec(&(sender.to_string(), receiver.to_string(), key))
3810 .expect("string tuple serializes");
3811 envelope.id = format!("m-{}", &blake3::hash(&identity).to_hex()[..24]);
3812 let mailbox = match crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &receiver) {
3813 Ok(mailbox) => mailbox,
3814 Err(error) => return refused("delivery_failed", error.to_string()),
3815 };
3816 match mailbox.find(&envelope.id) {
3817 Ok(Some(previous)) => {
3818 let saved = previous.envelope;
3819 let subject = params.subject.as_deref().map(|value| {
3820 value
3821 .lines()
3822 .next()
3823 .unwrap_or("")
3824 .chars()
3825 .take(200)
3826 .collect::<String>()
3827 });
3828 if saved.body != envelope.body
3829 || saved.subject != subject
3830 || saved.kind != envelope.kind
3831 || saved.from_name != envelope.from_name
3832 || saved.in_reply_to != params.in_reply_to
3833 || saved.voice_for != params.voice_for
3834 {
3835 return refused(
3836 "idempotency_conflict",
3837 "this key already names different mail".into(),
3838 );
3839 }
3840 }
3841 Ok(None) => {}
3842 Err(error) => return refused("delivery_failed", error.to_string()),
3843 }
3844 }
3845 envelope.in_reply_to = params.in_reply_to.clone();
3846 envelope.voice_for = params.voice_for.clone();
3847 envelope.subject = params.subject.as_deref().map(|value| {
3848 value
3849 .lines()
3850 .next()
3851 .unwrap_or("")
3852 .chars()
3853 .take(200)
3854 .collect()
3855 });
3856 let how = match crate::mail_route::deliver(
3857 &envelope,
3858 &receiver,
3859 &door,
3860 true,
3861 params.notify_when_idle,
3862 )
3863 .await
3864 {
3865 Err(detail) => {
3866 return refused(
3867 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3868 detail,
3869 )
3870 }
3871 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3872 Ok(Err(Refused::TooLong(bytes))) => {
3873 return refused(
3874 "too_long",
3875 format!(
3876 "the message is {bytes} bytes; the limit is {}",
3877 crate::mail_route::MAX_RELAYED_BYTES
3878 ),
3879 )
3880 }
3881 Ok(Ok(delivered)) => match delivered {
3882 Delivered::Steered => "steered",
3883 Delivered::Started => "started",
3884 Delivered::Native { busy: true } => "next_tool_call",
3885 Delivered::Native { busy: false } => "started",
3886 Delivered::Hooked => "hook",
3887 Delivered::HookWoken => "woken",
3888 Delivered::Queued => "queued",
3889 Delivered::Stored => "stored",
3890 Delivered::Operator => "filed",
3891 },
3892 };
3893 json!({
3894 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3895 "message_id": envelope.id,
3896 "reply_to": sender.to_string(),
3897 "target": {
3898 "harness": params.locator.harness.as_str(),
3899 "session_id": params.locator.session_id,
3900 "name": match &door {
3901 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
3902 _ => None,
3903 },
3904 },
3905 "delivery": {"door": door.name(), "how": how},
3906 "inbound_controls": inbound_controls,
3907 "inbound_controls_error": inbound_controls_error,
3908 })
3909}
3910
3911async fn message_as_user(
3915 params: &MessageSessionParams,
3916 sender: crate::mailbox::MailAddress,
3917 receiver: crate::mailbox::MailAddress,
3918 inbound_controls: Value,
3919 inbound_controls_error: Value,
3920) -> Value {
3921 let envelope = match crate::mailbox::Envelope::new(
3922 sender.clone(),
3923 format!("{}@{}", sender.session_id, sender.machine),
3924 crate::mailbox::MailKind::User,
3925 crate::mailbox::ReplyVia::None,
3926 params.text.clone(),
3927 ) {
3928 Ok(envelope) => envelope,
3929 Err(error) => {
3930 return json!({
3931 "delivered_to_bus": false,
3932 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
3933 "inbound_controls": inbound_controls,
3934 "inbound_controls_error": inbound_controls_error,
3935 })
3936 }
3937 };
3938 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
3939 Ok(turn) => json!({
3940 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
3941 "message_id": envelope.id,
3942 "target": {
3943 "harness": params.locator.harness.as_str(),
3944 "session_id": params.locator.session_id,
3945 },
3946 "delivery": {
3947 "door": match turn {
3948 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
3949 _ => "pane",
3950 },
3951 "how": turn.as_str(),
3952 },
3953 "inbound_controls": inbound_controls,
3954 "inbound_controls_error": inbound_controls_error,
3955 }),
3956 Err(message) => json!({
3957 "delivered_to_bus": false,
3958 "refusal": {"reason": "no_user_door", "message": message},
3959 "inbound_controls": inbound_controls,
3960 "inbound_controls_error": inbound_controls_error,
3961 }),
3962 }
3963}
3964
3965#[derive(Deserialize)]
3966#[serde(deny_unknown_fields)]
3967struct ActivityUnderParams {
3968 pids: Vec<u32>,
3970 #[serde(default)]
3971 homes: crate::HarnessHomes,
3972}
3973
3974async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
3978 let params = decode::<ActivityUnderParams>(params)?;
3979 if params.pids.len() > 1024 {
3980 return Err(ServiceError::InvalidParams(
3981 "sessions.activity_under accepts at most 1024 pids".into(),
3982 ));
3983 }
3984 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
3985 .await
3986 .map_err(ServiceError::Sdk)?;
3987 Ok(json!({
3988 "activities": found
3989 .into_iter()
3990 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
3991 .collect::<Vec<_>>(),
3992 }))
3993}
3994
3995fn operator_address(
3998 name: Option<&str>,
3999) -> std::result::Result<crate::mailbox::MailAddress, String> {
4000 let name: String = name
4001 .unwrap_or("supercode")
4002 .trim()
4003 .chars()
4004 .map(|character| {
4005 if character.is_whitespace() || character == '@' {
4006 '-'
4007 } else {
4008 character
4009 }
4010 })
4011 .collect();
4012 if name.is_empty() {
4013 return Err("from_name must not be empty".into());
4014 }
4015 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
4016 .map_err(|error| error.to_string())
4017}
4018
4019fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
4023 let address = match (¶ms.address, ¶ms.from_name) {
4024 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
4025 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
4026 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
4027 (Some(_), Some(_)) => {
4028 return Err(ServiceError::InvalidParams(
4029 "sessions.inbox takes from_name or address, not both".into(),
4030 ))
4031 }
4032 };
4033 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
4034 let mailbox =
4035 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
4036 let claimed = mailbox.claim_unread().map_err(operation)?;
4037 let mut messages: Vec<Value> = Vec::new();
4038 if params.all {
4039 for stored in mailbox.list().map_err(operation)? {
4040 if stored.state == crate::mailbox::MailState::Read {
4041 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4042 }
4043 }
4044 }
4045 for stored in &claimed {
4046 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4047 }
4048 for stored in &claimed {
4049 mailbox.acknowledge(stored).map_err(operation)?;
4050 }
4051 Ok(json!({"address": address.to_string(), "messages": messages}))
4052}
4053
4054#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4059struct FollowedSource {
4060 harness: String,
4061 session_id: String,
4062 reported: Option<String>,
4063}
4064
4065#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4066struct ActivitySubscription {
4067 locators: Vec<SessionLocator>,
4068 homes: crate::HarnessHomes,
4069 reported: BTreeMap<(String, String), crate::SessionActivity>,
4070}
4071
4072fn live_descriptor_value(
4079 session: &SessionDescriptor,
4080 doors: &crate::mail_route::LiveSessions,
4081) -> std::result::Result<Value, ServiceError> {
4082 let mut value = serde_json::to_value(session)
4083 .map_err(|error| ServiceError::Operation(error.to_string()))?;
4084 if value.get("title").is_none_or(Value::is_null) {
4086 if let Some(live) = doors.all().iter().find(|live| {
4087 live.address.harness == session.locator.harness.as_str()
4088 && live.address.session_id == session.locator.session_id
4089 }) {
4090 let name = live.name.split('@').next().unwrap_or(&live.name);
4091 if !name.is_empty() {
4092 value["title"] = json!(name);
4093 }
4094 }
4095 }
4096 if let Some(workspace) = &session.cwd {
4097 let source = LiveRuntimeSource {
4098 harness: session.locator.harness.as_str().to_string(),
4099 session_id: session.locator.session_id.clone(),
4100 workspace: workspace.clone(),
4101 };
4102 if let Some(endpoint) = discover_live_runtime(&source)
4103 .map_err(|error| ServiceError::Operation(error.to_string()))?
4104 {
4105 value["live_endpoint"] = json!(endpoint.as_str());
4106 }
4107 }
4108 if let Some(door) = doors.door(
4112 session.locator.harness.as_str(),
4113 &session.locator.session_id,
4114 ) {
4115 value["delivery"] = json!(door);
4116 }
4117 if let Some(live) = doors.all().iter().find(|live| {
4120 live.address.harness == session.locator.harness.as_str()
4121 && live.address.session_id == session.locator.session_id
4122 && live.status == "waiting"
4123 }) {
4124 if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
4125 value["pending_request"] = request;
4126 }
4127 }
4128 Ok(value)
4129}
4130
4131fn live_index_changes(
4132 changes: Vec<crate::session_index::SessionIndexChange>,
4133 homes: &HarnessHomes,
4134) -> std::result::Result<Vec<Value>, ServiceError> {
4135 use crate::session_index::SessionIndexChange;
4136 let doors = crate::mail_route::LiveSessions::read(homes);
4137 changes
4138 .into_iter()
4139 .map(|change| match change {
4140 SessionIndexChange::Added { descriptor } => Ok(json!({
4141 "kind": "added",
4142 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4143 })),
4144 SessionIndexChange::Updated { descriptor } => Ok(json!({
4145 "kind": "updated",
4146 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4147 })),
4148 SessionIndexChange::Removed { key } => Ok(json!({
4149 "kind": "removed",
4150 "key": key,
4151 })),
4152 })
4153 .collect()
4154}
4155
4156fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
4157 use crate::{SessionPresence, SessionTurnState};
4158 match (activity.presence, activity.turn) {
4159 (SessionPresence::Persisted, _) => None,
4160 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
4161 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
4162 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
4164 (SessionPresence::Running, SessionTurnState::Unknown)
4168 if activity.evidence.native_state.is_none() =>
4169 {
4170 None
4171 }
4172 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
4173 }
4174}
4175
4176#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
4177#[serde(rename_all = "kebab-case")]
4178enum TransferFormat {
4179 ClaudeCode,
4180 Codex,
4181 #[serde(rename = "opencode", alias = "open-code")]
4182 OpenCode,
4183 Pi,
4184 Grok,
4185 Gemini,
4186 Goose,
4187 Hermes,
4191}
4192
4193impl TransferFormat {
4194 fn id(self) -> &'static str {
4195 match self {
4196 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
4197 Self::Codex => HarnessId::CODEX,
4198 Self::OpenCode => HarnessId::OPENCODE,
4199 Self::Pi => HarnessId::PI,
4200 Self::Grok => HarnessId::GROK,
4201 Self::Gemini => HarnessId::GEMINI,
4202 Self::Goose => HarnessId::GOOSE,
4203 Self::Hermes => HarnessId::HERMES,
4204 }
4205 }
4206}
4207
4208impl From<TransferFormat> for SessionFormat {
4209 fn from(value: TransferFormat) -> Self {
4210 match value {
4211 TransferFormat::ClaudeCode => Self::ClaudeCode,
4212 TransferFormat::Codex => Self::Codex,
4213 TransferFormat::OpenCode => Self::OpenCode,
4214 TransferFormat::Pi => Self::Pi,
4215 TransferFormat::Grok => Self::Grok,
4216 TransferFormat::Gemini => Self::Gemini,
4217 TransferFormat::Goose => Self::Goose,
4218 TransferFormat::Hermes => Self::Codex,
4220 }
4221 }
4222}
4223
4224#[derive(Deserialize)]
4225struct ImportSessionParams {
4226 source_harness: TransferFormat,
4227 content: String,
4228}
4229
4230#[derive(Deserialize)]
4231struct ExportSessionParams {
4232 locator: SessionLocator,
4233 target_harness: TransferFormat,
4234}
4235
4236#[derive(Deserialize)]
4237struct ReduceSessionParams {
4238 locator: SessionLocator,
4239 target_harness: TransferFormat,
4240 #[serde(default = "default_keep_last")]
4241 keep_last: usize,
4242}
4243
4244fn default_keep_last() -> usize {
4245 6
4246}
4247
4248#[derive(Deserialize)]
4249struct BranchSessionParams {
4250 locator: SessionLocator,
4251 #[serde(default)]
4252 target_harness: Option<TransferFormat>,
4253}
4254
4255#[derive(Deserialize)]
4256struct HandoffSessionParams {
4257 locator: SessionLocator,
4258 target_harness: TransferFormat,
4259 #[serde(default)]
4260 cwd: Option<PathBuf>,
4261}
4262
4263#[derive(Deserialize)]
4264struct MaterializeSessionParams {
4265 artifact: crate::native_materialize::MaterializeArtifact,
4266 cwd: PathBuf,
4267 #[serde(default)]
4269 homes: HarnessHomes,
4270}
4271
4272#[derive(Debug, Clone, Copy, Default, Deserialize)]
4275#[serde(rename_all = "snake_case")]
4276enum ResumePolicy {
4277 Default,
4278 #[default]
4279 Yolo,
4280}
4281
4282#[derive(Deserialize)]
4283struct ResumeInstructionsParams {
4284 locator: SessionLocator,
4285 #[serde(default)]
4286 cwd: Option<PathBuf>,
4287 #[serde(default)]
4288 policy: ResumePolicy,
4289}
4290
4291#[derive(Deserialize)]
4293struct WorkflowLoadParams {
4294 from: crate::workflow_doors::WorkflowHarness,
4295 home: PathBuf,
4296}
4297
4298#[derive(Deserialize)]
4301struct OrchestrationLoadParams {
4302 root: PathBuf,
4303 #[serde(default)]
4304 flavor: crate::orchestration_doors::HomeFlavor,
4305}
4306
4307#[derive(Deserialize)]
4310struct OrchestrationSaveParams {
4311 root: PathBuf,
4312 orchestration: crate::orchestration::Orchestration,
4313 #[serde(default)]
4314 vault: BTreeMap<String, String>,
4315}
4316
4317#[derive(Deserialize)]
4319struct OrchestrationCompileParams {
4320 from: crate::orchestration_doors::OrchestrationHarness,
4321 home: PathBuf,
4322}
4323
4324#[derive(Deserialize)]
4328struct OrchestrationDecompileParams {
4329 to: crate::orchestration_doors::OrchestrationHarness,
4330 orchestration: crate::orchestration::Orchestration,
4331 source: PathBuf,
4332 #[serde(default)]
4333 source_flavor: crate::orchestration_doors::SourceFlavor,
4334 dest: PathBuf,
4335 #[serde(default)]
4336 vault: BTreeMap<String, String>,
4337}
4338
4339#[derive(Deserialize)]
4342struct OrchestrationImportParams {
4343 from: crate::orchestration_doors::OrchestrationHarness,
4344 home: PathBuf,
4345 into: PathBuf,
4346}
4347
4348#[derive(Deserialize)]
4351struct OrchestrationExportParams {
4352 to: crate::orchestration_doors::OrchestrationHarness,
4353 root: PathBuf,
4354 dest: PathBuf,
4355}
4356
4357#[derive(Deserialize)]
4359struct JobsGetParams {
4360 harness: String,
4361 id: String,
4362 #[serde(default)]
4363 homes: crate::HarnessHomes,
4364}
4365
4366fn mutate_job(
4374 verb: crate::jobs_control::JobVerb,
4375 params: Value,
4376) -> std::result::Result<Value, ServiceError> {
4377 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4378 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4379 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4380 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4381}
4382
4383fn mutate_skill(
4390 verb: crate::skills_control::SkillVerb,
4391 params: Value,
4392) -> std::result::Result<Value, ServiceError> {
4393 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4394 if !crate::skills_control::supports_skill_control(&mutation.harness) {
4395 return Err(ServiceError::UnsupportedAction(format!(
4396 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4397 mutation.harness,
4398 verb.as_str(),
4399 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4400 )));
4401 }
4402 let outcome =
4403 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4404 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4405}
4406
4407fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4409 match error {
4410 crate::skills_control::SkillControlError::Unsupported(message) => {
4411 ServiceError::UnsupportedAction(message)
4412 }
4413 crate::skills_control::SkillControlError::Invalid(message) => {
4414 ServiceError::InvalidParams(message)
4415 }
4416 crate::skills_control::SkillControlError::Failed(message) => {
4417 ServiceError::Operation(message)
4418 }
4419 }
4420}
4421
4422fn mutate_profile(
4430 verb: crate::profiles_control::ProfileVerb,
4431 params: Value,
4432) -> std::result::Result<Value, ServiceError> {
4433 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4434 let outcome =
4435 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4436 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4437}
4438
4439fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4441 match error {
4442 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4443 ServiceError::UnsupportedAction(message)
4444 }
4445 crate::profiles_control::ProfileControlError::Invalid(message) => {
4446 ServiceError::InvalidParams(message)
4447 }
4448 crate::profiles_control::ProfileControlError::Failed(message) => {
4449 ServiceError::Operation(message)
4450 }
4451 }
4452}
4453
4454fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4458 match error {
4459 crate::jobs_control::JobControlError::Unsupported(message) => {
4460 ServiceError::UnsupportedAction(message)
4461 }
4462 crate::jobs_control::JobControlError::Invalid(message) => {
4463 ServiceError::InvalidParams(message)
4464 }
4465 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4466 }
4467}
4468
4469fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4474 match error {
4475 crate::SessionControlError::Unsupported(message) => {
4476 ServiceError::UnsupportedAction(message)
4477 }
4478 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4479 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4480 }
4481}
4482
4483fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4488 if crate::jobs::supports_jobs(harness) {
4489 return Ok(());
4490 }
4491 Err(ServiceError::UnsupportedAction(format!(
4492 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4493 crate::jobs::JOB_HARNESSES.join(", ")
4494 )))
4495}
4496
4497#[derive(Deserialize)]
4499struct RunsGetParams {
4500 harness: String,
4501 id: String,
4502 #[serde(default)]
4503 homes: crate::HarnessHomes,
4504}
4505
4506fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4511 if crate::runs::supports_runs(harness) {
4512 return Ok(());
4513 }
4514 Err(ServiceError::UnsupportedAction(format!(
4515 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4516 crate::runs::RUN_HARNESSES.join(", ")
4517 )))
4518}
4519
4520#[derive(Serialize)]
4521struct SessionArtifact {
4522 source_harness: HarnessId,
4523 target_harness: &'static str,
4524 session_id: Option<String>,
4525 content: String,
4526 suggested_filename: String,
4527 files: Vec<SessionArtifactFile>,
4528 fidelity: Fidelity,
4529 residue: Vec<String>,
4530}
4531
4532#[derive(Serialize)]
4533struct SessionArtifactFile {
4534 path: String,
4535 content: String,
4536 role: ArtifactFileRole,
4537}
4538
4539#[derive(Serialize)]
4540#[serde(rename_all = "snake_case")]
4541enum ArtifactFileRole {
4542 Primary,
4543 Subagent,
4544 Bundle,
4545 SourceRecovery,
4546}
4547
4548#[derive(Serialize)]
4549struct StructuredLaunch {
4550 cwd: PathBuf,
4551 program: String,
4552 arguments: Vec<String>,
4553 env: BTreeMap<String, String>,
4554}
4555
4556struct HandoffInstructions {
4557 launch: StructuredLaunch,
4558 materialize: Option<StructuredLaunch>,
4559 requires_materialization: bool,
4560 note: String,
4561}
4562
4563#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4564#[serde(rename_all = "snake_case")]
4565enum HarnessProbeLevel {
4566 #[default]
4567 Passive,
4568 Handshake,
4569}
4570
4571#[derive(Default, Deserialize)]
4572#[serde(default)]
4573struct HarnessInventoryParams {
4574 harness: Option<HarnessId>,
4575 harnesses: Vec<HarnessId>,
4576 workspace: Option<PathBuf>,
4577 probe: HarnessProbeLevel,
4578 include_sessions: bool,
4579 skip_versions: bool,
4581}
4582
4583#[derive(Deserialize)]
4584struct HarnessAuthenticationParams {
4585 harness: HarnessId,
4586}
4587
4588#[derive(Deserialize)]
4589struct BeginHarnessAuthenticationParams {
4590 harness: HarnessId,
4591 #[serde(default = "local_browser_authentication_environment")]
4592 environment: crate::HarnessAuthenticationEnvironment,
4593 #[serde(default)]
4594 method: Option<crate::HarnessAuthenticationMethodId>,
4595 #[serde(default)]
4596 cwd: Option<PathBuf>,
4597}
4598
4599fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4600 crate::HarnessAuthenticationEnvironment::LocalBrowser
4601}
4602
4603#[derive(Serialize)]
4604struct HarnessInventoryReport {
4605 probe: HarnessProbeLevel,
4606 workspace: Option<PathBuf>,
4607 harnesses: Vec<LocalHarness>,
4608}
4609
4610#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4611#[serde(rename_all = "snake_case")]
4612enum HarnessAuthState {
4613 Ready,
4614 Configured,
4615 Required,
4616 Unknown,
4617}
4618
4619#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4620#[serde(rename_all = "snake_case")]
4621enum HarnessRuntimeState {
4622 Ready,
4623 Degraded,
4624 Unavailable,
4625}
4626
4627#[derive(Serialize)]
4628struct HarnessSessionCounts {
4629 global: Option<usize>,
4630 workspace: Option<usize>,
4631}
4632
4633#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4644#[serde(rename_all = "snake_case")]
4645pub enum GatewayState {
4646 Up,
4647 Down,
4648 Unknown,
4649}
4650
4651#[derive(Debug, Clone, Serialize)]
4653pub struct GatewayHealth {
4654 pub state: GatewayState,
4655 #[serde(skip_serializing_if = "Option::is_none")]
4660 pub endpoint: Option<String>,
4661 #[serde(skip_serializing_if = "Option::is_none")]
4662 pub version: Option<String>,
4663 pub evidence: String,
4665 pub checked_at_ms: u64,
4666}
4667
4668fn openclaw_gateway_endpoint(home: &Path) -> String {
4672 let config_path = home.join(".openclaw/openclaw.json");
4673 let gateway = std::fs::read_to_string(&config_path)
4674 .ok()
4675 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4676 .and_then(|config| config.get("gateway").cloned());
4677 if let Some(url) = gateway
4678 .as_ref()
4679 .and_then(|gateway| gateway.get("url"))
4680 .and_then(serde_json::Value::as_str)
4681 {
4682 return url.to_string();
4683 }
4684 let port = gateway
4685 .as_ref()
4686 .and_then(|gateway| gateway.get("port"))
4687 .and_then(serde_json::Value::as_u64)
4688 .unwrap_or(18789);
4689 format!("ws://127.0.0.1:{port}")
4690}
4691
4692fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4699 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4700 let output = std::process::Command::new(&program)
4701 .args(["gateway", "status"])
4702 .stdin(std::process::Stdio::null())
4703 .output()
4704 .ok()?;
4705 let text = format!(
4706 "{}{}",
4707 String::from_utf8_lossy(&output.stdout),
4708 String::from_utf8_lossy(&output.stderr)
4709 );
4710 let verdict = text.lines().find_map(|line| {
4711 let l = line.trim();
4712 if l.contains("supervised by launchd (PID")
4713 || l.contains("supervised by systemd (PID")
4714 || l.contains("Gateway is running")
4715 || l.contains("process is running")
4716 {
4717 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4718 } else if l.contains("not running") || l.contains("not installed") {
4719 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4720 } else {
4721 None
4722 }
4723 });
4724 verdict
4725}
4726
4727fn gateway_health(
4728 id: &str,
4729 installed: bool,
4730 running: Option<&RunningInstance>,
4731 version: Option<&str>,
4732) -> GatewayHealth {
4733 let checked_at_ms = now_epoch_ms();
4734 let home = supercode_interchange::user_home()
4735 .map(std::path::PathBuf::into_os_string)
4736 .map(PathBuf::from);
4737 match id {
4738 HarnessId::HERMES | HarnessId::OPENCLAW => {
4739 let endpoint = (id == HarnessId::OPENCLAW)
4740 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4741 .flatten();
4742 let (state, evidence) = match running {
4743 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4744 None if !installed => (
4745 GatewayState::Unknown,
4746 format!("`{id}` is not installed; no gateway to probe"),
4747 ),
4748 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4749 Some((state, evidence)) => (state, evidence),
4752 None => (
4753 GatewayState::Down,
4754 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4755 ),
4756 },
4757 None => (
4758 GatewayState::Down,
4759 format!(
4760 "no TCP listener at {}",
4761 endpoint.as_deref().unwrap_or("the gateway endpoint")
4762 ),
4763 ),
4764 };
4765 GatewayHealth {
4766 state,
4767 endpoint,
4768 version: version.map(str::to_string),
4769 evidence,
4770 checked_at_ms,
4771 }
4772 }
4773 HarnessId::ORCHESTRATOR => {
4780 let root = crate::HarnessHomes::default().orchestrator;
4781 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4782 Some(lease) if lease.is_live() => (
4783 GatewayState::Up,
4784 format!(
4785 "`{}` names pid {} (started {}), which is live",
4786 crate::orchestrator::lock_path(&root).display(),
4787 lease.pid,
4788 lease.started_at
4789 ),
4790 ),
4791 Some(lease) => (
4792 GatewayState::Down,
4793 format!(
4794 "stale lease `{}`: pid {} is gone",
4795 crate::orchestrator::lock_path(&root).display(),
4796 lease.pid
4797 ),
4798 ),
4799 None => (
4800 GatewayState::Down,
4801 format!(
4802 "no lease at `{}`; `supercode orchestrator start` writes one",
4803 crate::orchestrator::lock_path(&root).display()
4804 ),
4805 ),
4806 };
4807 GatewayHealth {
4808 state,
4809 endpoint: None,
4810 version: version.map(str::to_string),
4811 evidence,
4812 checked_at_ms,
4813 }
4814 }
4815 _ => GatewayHealth {
4816 state: GatewayState::Unknown,
4817 endpoint: None,
4818 version: version.map(str::to_string),
4819 evidence: format!("`{id}` runs per session, not as a gateway"),
4820 checked_at_ms,
4821 },
4822 }
4823}
4824
4825#[derive(Debug, Clone, Serialize)]
4826struct RunningInstance {
4827 method: RunningInstanceMethod,
4829 evidence: String,
4831 checked_at_ms: u64,
4833}
4834
4835#[derive(Debug, Clone, Copy, Serialize)]
4836#[serde(rename_all = "snake_case")]
4837enum RunningInstanceMethod {
4838 GatewayConnect,
4841 StoreWalActivity,
4844}
4845
4846fn now_epoch_ms() -> u64 {
4847 std::time::SystemTime::now()
4848 .duration_since(std::time::UNIX_EPOCH)
4849 .map(|elapsed| elapsed.as_millis() as u64)
4850 .unwrap_or(0)
4851}
4852
4853fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4857 let config_path = home.join(".openclaw/openclaw.json");
4858 let text = std::fs::read_to_string(&config_path).ok();
4859 let gateway = text
4860 .as_deref()
4861 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4862 .and_then(|config| config.get("gateway").cloned());
4863 let address = gateway
4864 .as_ref()
4865 .and_then(|gateway| gateway.get("url"))
4866 .and_then(serde_json::Value::as_str)
4867 .and_then(|url| {
4868 url.split("://").nth(1).map(|rest| {
4869 rest.trim_end_matches('/')
4870 .split('/')
4871 .next()
4872 .unwrap_or(rest)
4873 .to_string()
4874 })
4875 })
4876 .unwrap_or_else(|| {
4877 let port = gateway
4878 .as_ref()
4879 .and_then(|gateway| gateway.get("port"))
4880 .and_then(serde_json::Value::as_u64)
4881 .unwrap_or(18789);
4882 format!("127.0.0.1:{port}")
4883 });
4884 let reachable = std::net::TcpStream::connect_timeout(
4885 &address.parse().ok()?,
4886 std::time::Duration::from_millis(400),
4887 )
4888 .is_ok();
4889 reachable.then(|| RunningInstance {
4890 method: RunningInstanceMethod::GatewayConnect,
4891 evidence: format!(
4892 "gateway endpoint {address} accepted a TCP connect (from {})",
4893 config_path.display()
4894 ),
4895 checked_at_ms: now_epoch_ms(),
4896 })
4897}
4898
4899fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4904 let wal = home.join(".hermes/state.db-wal");
4905 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4906 let age_ms = std::time::SystemTime::now()
4907 .duration_since(modified)
4908 .map(|age| age.as_millis() as u64)
4909 .unwrap_or(u64::MAX);
4910 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4911 method: RunningInstanceMethod::StoreWalActivity,
4912 evidence: format!(
4913 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4914 wal.display()
4915 ),
4916 checked_at_ms: now_epoch_ms(),
4917 })
4918}
4919
4920fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4922 let home = supercode_interchange::user_home()
4923 .map(std::path::PathBuf::into_os_string)
4924 .map(PathBuf::from)?;
4925 match id {
4926 HarnessId::OPENCLAW => probe_openclaw_running(&home),
4927 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4928 _ => None,
4929 }
4930}
4931
4932#[derive(Serialize)]
4933struct LocalHarness {
4934 id: HarnessId,
4935 display_name: String,
4936 supported: bool,
4937 installed: bool,
4938 executable: Option<String>,
4939 version: Option<String>,
4940 auth: HarnessAuthState,
4941 runtime: HarnessRuntimeState,
4942 protocol: String,
4943 capabilities: crate::RuntimeCapabilities,
4944 effective_capabilities: crate::RuntimeCapabilities,
4945 sessions: HarnessSessionCounts,
4946 #[serde(skip_serializing_if = "Option::is_none")]
4949 running: Option<RunningInstance>,
4950 gateway: GatewayHealth,
4952 reason: Option<String>,
4953 repair: Option<String>,
4954}
4955
4956#[derive(Clone, Deserialize)]
4957struct RuntimeBackendParams {
4958 harness: HarnessId,
4959 #[serde(default)]
4960 protocol: Option<String>,
4961 #[serde(default)]
4962 launch: Option<RuntimeLaunch>,
4963 #[serde(default)]
4964 base_url: Option<String>,
4965 #[serde(default)]
4966 policy: RuntimePolicy,
4967}
4968
4969#[derive(Debug, Clone, Copy, Default, Deserialize)]
4972#[serde(rename_all = "snake_case")]
4973enum RuntimePolicy {
4974 Default,
4975 #[default]
4976 Yolo,
4977}
4978
4979#[derive(Deserialize)]
4980struct RuntimeStartParams {
4981 #[serde(flatten)]
4982 backend: RuntimeBackendParams,
4983 cwd: PathBuf,
4984 #[serde(default)]
4987 mcp_servers: Vec<crate::McpServerLaunch>,
4988 #[serde(default)]
4990 approval_policy: Option<String>,
4991}
4992
4993#[derive(Deserialize)]
4994struct RuntimeAttachParams {
4995 #[serde(flatten)]
4996 backend: RuntimeBackendParams,
4997 runtime_id: String,
4998 #[serde(default)]
4999 cwd: Option<PathBuf>,
5000 #[serde(default)]
5003 mcp_servers: Vec<crate::McpServerLaunch>,
5004 #[serde(default)]
5006 approval_policy: Option<String>,
5007}
5008
5009#[derive(Deserialize)]
5010struct RuntimeConnectionParams {
5011 connection: String,
5012}
5013
5014#[derive(Deserialize)]
5017struct RuntimeCloseParams {
5018 #[serde(default)]
5019 connection: String,
5020 #[serde(default)]
5021 runtime_id: Option<String>,
5022}
5023
5024#[derive(Deserialize)]
5025struct RuntimeInputParams {
5026 connection: String,
5027 text: String,
5028 #[serde(default)]
5029 image_urls: Vec<String>,
5030}
5031
5032const MAX_RUNTIME_IMAGES: usize = 4;
5033const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
5034const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
5035
5036fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
5037 if image_urls.len() > MAX_RUNTIME_IMAGES {
5038 return Err(ServiceError::InvalidParams(format!(
5039 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
5040 )));
5041 }
5042 let mut total = 0usize;
5043 for url in &image_urls {
5044 if !(url.starts_with("data:image/")
5045 || url.starts_with("https://")
5046 || url.starts_with("http://"))
5047 {
5048 return Err(ServiceError::InvalidParams(
5049 "runtime images must be image data URLs or HTTP(S) URLs".into(),
5050 ));
5051 }
5052 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
5053 return Err(ServiceError::InvalidParams(format!(
5054 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
5055 )));
5056 }
5057 total = total.saturating_add(url.len());
5058 }
5059 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
5060 return Err(ServiceError::InvalidParams(format!(
5061 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
5062 )));
5063 }
5064 Ok(image_urls)
5065}
5066
5067#[derive(Deserialize)]
5068struct RuntimeRespondParams {
5069 connection: String,
5070 request_id: Value,
5071 response: Value,
5072}
5073
5074fn default_reduction_store_root() -> PathBuf {
5075 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
5076 return PathBuf::from(root).join("sessions");
5077 }
5078 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
5079 return PathBuf::from(home).join(".supercode").join("sessions");
5080 }
5081 PathBuf::from(".supercode").join("sessions")
5082}
5083
5084fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
5085 let mut output = String::new();
5086 for message in messages {
5087 output.push_str(
5088 &serde_json::to_string(message)
5089 .map_err(|error| ServiceError::Operation(error.to_string()))?,
5090 );
5091 output.push('\n');
5092 }
5093 Ok(output)
5094}
5095
5096fn parse_messages_jsonl(
5097 content: &str,
5098) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
5099 content
5100 .lines()
5101 .enumerate()
5102 .filter(|(_, line)| !line.trim().is_empty())
5103 .map(|(index, line)| {
5104 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
5105 ServiceError::Operation(format!(
5106 "reduced transcript line {} is invalid: {error}",
5107 index + 1
5108 ))
5109 })
5110 })
5111 .collect()
5112}
5113
5114fn reduced_bootstrap_prompt(
5115 source: &SessionLocator,
5116 target: TransferFormat,
5117 view_jsonl: &str,
5118 sidecar_path: &Path,
5119 reduction_log_path: &Path,
5120) -> String {
5121 format!(
5122 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
5123 \n\
5124 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\
5125 \n\
5126 <supercode-reduced-session source-session=\"{source_id}\">\n\
5127 {view_jsonl}\
5128 </supercode-reduced-session>\n\
5129 \n\
5130 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
5131 source_harness = source.harness.as_str(),
5132 target_harness = target.id(),
5133 sidecar = sidecar_path.display(),
5134 log = reduction_log_path.display(),
5135 source_id = source.session_id,
5136 )
5137}
5138
5139fn session_artifact(
5140 locator: &SessionLocator,
5141 session: &Session,
5142 target: TransferFormat,
5143) -> std::result::Result<SessionArtifact, ServiceError> {
5144 session_artifact_with_id(locator, session, target, None)
5145}
5146
5147fn session_artifact_with_id(
5148 locator: &SessionLocator,
5149 session: &Session,
5150 target: TransferFormat,
5151 target_session_id: Option<&str>,
5152) -> std::result::Result<SessionArtifact, ServiceError> {
5153 let format: SessionFormat = target.into();
5154 let diagonal = format.source() == session.meta.source;
5155 crate::residue_store::store_segments(session);
5156 let has_appended_turns = session
5157 .imported_message_count
5158 .is_some_and(|imported| imported < session.messages.len());
5159 let mut restoration = None;
5160 let content = if let Some(id) = target_session_id {
5161 if diagonal && format != SessionFormat::OpenCode {
5162 session
5163 .to_jsonl_spliced(format, Some(id))
5164 .map_err(operation)?
5165 } else {
5166 let mut rewritten = session.clone();
5167 rewritten.meta.session_id = Some(id.to_string());
5168 rewritten.to_jsonl(format).map_err(operation)?
5169 }
5170 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
5171 session.raw_verbatim()
5172 } else if diagonal {
5173 session.to_jsonl_spliced(format, None).map_err(operation)?
5174 } else {
5175 match session
5178 .restore_residue(format, crate::residue_store::lookup)
5179 .map_err(operation)?
5180 {
5181 Some((content, report)) => {
5182 restoration = Some(report);
5183 content
5184 }
5185 None => session.to_jsonl(format).map_err(operation)?,
5186 }
5187 };
5188 let stem = sanitize_filename(
5189 target_session_id
5190 .or(session.meta.session_id.as_deref())
5191 .unwrap_or(&locator.session_id),
5192 );
5193 let suggested_filename = if diagonal && target == TransferFormat::Grok {
5194 "chat_history.jsonl".to_string()
5195 } else if target == TransferFormat::Goose {
5196 format!("{stem}.goose.json")
5197 } else {
5198 format!("{stem}.{}.jsonl", target.id())
5199 };
5200 let mut files = vec![SessionArtifactFile {
5201 path: suggested_filename.clone(),
5202 content: content.clone(),
5203 role: ArtifactFileRole::Primary,
5204 }];
5205 if target == TransferFormat::ClaudeCode {
5206 let bundle_stem = Path::new(&suggested_filename)
5207 .file_stem()
5208 .and_then(|stem| stem.to_str())
5209 .unwrap_or(&stem);
5210 let mut child_paths = BTreeSet::new();
5211 for (index, subagent) in session.subagents.iter().enumerate() {
5212 let agent_id = subagent
5213 .meta
5214 .agent_id
5215 .as_deref()
5216 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
5217 .map(sanitize_filename)
5218 .filter(|id| !id.is_empty())
5219 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5220 let child_has_appended_turns = subagent
5221 .imported_message_count
5222 .is_some_and(|imported| imported < subagent.messages.len());
5223 let child_content = if target_session_id.is_none()
5224 && subagent.meta.source == SessionSource::ClaudeCode
5225 && subagent.raw_is_verbatim
5226 && !child_has_appended_turns
5227 {
5228 subagent.raw_verbatim()
5229 } else if subagent.meta.source == SessionSource::ClaudeCode {
5230 subagent
5231 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
5232 .map_err(operation)?
5233 } else {
5234 let mut child = subagent.clone();
5235 if let Some(id) = target_session_id {
5236 child.meta.session_id = Some(id.to_string());
5237 }
5238 child
5239 .to_jsonl(SessionFormat::ClaudeCode)
5240 .map_err(operation)?
5241 };
5242 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
5243 if !child_paths.insert(path.clone()) {
5244 return Err(ServiceError::Operation(format!(
5245 "Claude subagent ids collide at artifact path `{path}`"
5246 )));
5247 }
5248 files.push(SessionArtifactFile {
5249 path,
5250 content: child_content,
5251 role: ArtifactFileRole::Subagent,
5252 });
5253 }
5254 }
5255 if diagonal && target == TransferFormat::Grok {
5256 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
5257 }
5258 if !diagonal || !session.raw_is_verbatim {
5259 files.push(SessionArtifactFile {
5260 path: "recovery/source.supercode.jsonl".into(),
5261 content: session.to_native_jsonl(),
5262 role: ArtifactFileRole::SourceRecovery,
5263 });
5264 for (index, subagent) in session.subagents.iter().enumerate() {
5265 let id = subagent
5266 .meta
5267 .agent_id
5268 .as_deref()
5269 .map(sanitize_filename)
5270 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5271 files.push(SessionArtifactFile {
5272 path: format!("recovery/subagents/{id}.supercode.jsonl"),
5273 content: subagent.to_native_jsonl(),
5274 role: ArtifactFileRole::SourceRecovery,
5275 });
5276 }
5277 }
5278 if !diagonal && session.meta.source == SessionSource::Grok {
5279 append_grok_bundle_files(
5280 locator,
5281 "recovery/grok/",
5282 ArtifactFileRole::SourceRecovery,
5283 &mut files,
5284 )?;
5285 }
5286 let (fidelity, residue) = if diagonal
5287 && target_session_id.is_none()
5288 && session.raw_is_verbatim
5289 && !has_appended_turns
5290 {
5291 (Fidelity::ByteLossless, Vec::new())
5292 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
5293 (
5294 Fidelity::ValueLossless,
5295 vec![if target_session_id.is_some() {
5296 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
5297 } else {
5298 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
5299 }],
5300 )
5301 } else {
5302 match restoration {
5303 Some(report) if report.rendered_messages == 0 => (
5304 Fidelity::ByteLossless,
5305 vec![format!(
5306 "restored verbatim from this conversation's {} source records in the residue store",
5307 target.id()
5308 )],
5309 ),
5310 Some(report) => (
5311 Fidelity::Semantic,
5312 vec![format!(
5313 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
5314 report.restored_messages,
5315 report.restored_messages + report.rendered_messages,
5316 report.rendered_messages,
5317 target.id()
5318 )],
5319 ),
5320 None => (
5321 Fidelity::Semantic,
5322 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
5323 ),
5324 }
5325 };
5326 Ok(SessionArtifact {
5327 source_harness: locator.harness.clone(),
5328 target_harness: target.id(),
5329 session_id: target_session_id
5330 .map(str::to_string)
5331 .or_else(|| session.meta.session_id.clone()),
5332 content,
5333 suggested_filename,
5334 files,
5335 fidelity,
5336 residue,
5337 })
5338}
5339
5340fn append_grok_bundle_files(
5341 locator: &SessionLocator,
5342 prefix: &str,
5343 role: ArtifactFileRole,
5344 files: &mut Vec<SessionArtifactFile>,
5345) -> std::result::Result<(), ServiceError> {
5346 let primary = locator.storage.path();
5347 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5348 return Err(ServiceError::Operation(format!(
5349 "Grok bundle locator must name chat_history.jsonl, got {}",
5350 primary.display()
5351 )));
5352 }
5353 let parent = primary.parent().ok_or_else(|| {
5354 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5355 })?;
5356 for name in ["summary.json", "updates.jsonl"] {
5357 let path = parent.join(name);
5358 let metadata = match std::fs::symlink_metadata(&path) {
5359 Ok(metadata) => metadata,
5360 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5361 Err(error) => return Err(ServiceError::Operation(error.to_string())),
5362 };
5363 if metadata.file_type().is_symlink() || !metadata.is_file() {
5364 return Err(ServiceError::Operation(format!(
5365 "refusing non-regular Grok bundle member {}",
5366 path.display()
5367 )));
5368 }
5369 let content = std::fs::read_to_string(&path).map_err(|error| {
5370 ServiceError::Operation(format!(
5371 "Grok bundle member {} is not representable as UTF-8: {error}",
5372 path.display()
5373 ))
5374 })?;
5375 files.push(SessionArtifactFile {
5376 path: format!("{prefix}{name}"),
5377 content,
5378 role: match role {
5379 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5380 _ => ArtifactFileRole::SourceRecovery,
5381 },
5382 });
5383 }
5384 Ok(())
5385}
5386
5387fn handoff_artifact(
5388 locator: &SessionLocator,
5389 session: &Session,
5390 target: TransferFormat,
5391) -> std::result::Result<SessionArtifact, ServiceError> {
5392 let target_session_id = target_session_id(target);
5393 session_artifact_with_id(locator, session, target, Some(&target_session_id))
5394}
5395
5396fn target_session_id(target: TransferFormat) -> String {
5397 let uuid = generated_session_id();
5398 match target {
5399 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5400 TransferFormat::ClaudeCode
5401 | TransferFormat::Codex
5402 | TransferFormat::Pi
5403 | TransferFormat::Grok
5404 | TransferFormat::Gemini
5405 | TransferFormat::Goose
5406 | TransferFormat::Hermes => uuid,
5407 }
5408}
5409
5410fn sanitize_filename(value: &str) -> String {
5411 let value = value
5412 .chars()
5413 .map(|character| {
5414 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5415 character
5416 } else {
5417 '-'
5418 }
5419 })
5420 .collect::<String>();
5421 let value = value.trim_matches('-');
5422 if value.is_empty() {
5423 "session".into()
5424 } else {
5425 value.chars().take(100).collect()
5426 }
5427}
5428
5429fn handoff_instructions(
5430 target: TransferFormat,
5431 session_id: &str,
5432 cwd: &Path,
5433) -> HandoffInstructions {
5434 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5435 cwd: cwd.to_path_buf(),
5436 program: program.into(),
5437 arguments,
5438 env: BTreeMap::new(),
5439 };
5440 match target {
5441 TransferFormat::ClaudeCode => HandoffInstructions {
5442 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5443 materialize: None,
5444 requires_materialization: true,
5445 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(),
5446 },
5447 TransferFormat::Hermes => HandoffInstructions {
5448 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5449 materialize: None,
5450 requires_materialization: true,
5451 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(),
5452 },
5453 TransferFormat::Codex => HandoffInstructions {
5454 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5455 materialize: None,
5456 requires_materialization: true,
5457 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5458 },
5459 TransferFormat::OpenCode => HandoffInstructions {
5460 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5461 materialize: Some(launch(
5462 "opencode",
5463 vec!["import".into(), "{artifact_path}".into()],
5464 )),
5465 requires_materialization: true,
5466 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5467 },
5468 TransferFormat::Pi => HandoffInstructions {
5469 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5470 materialize: None,
5471 requires_materialization: true,
5472 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5473 },
5474 TransferFormat::Grok => HandoffInstructions {
5475 launch: launch(
5476 "grok",
5477 vec!["--resume".into(), "{materialized_session_id}".into()],
5478 ),
5479 materialize: None,
5480 requires_materialization: true,
5481 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(),
5482 },
5483 TransferFormat::Gemini => HandoffInstructions {
5484 launch: launch(
5485 "gemini",
5486 vec!["--session-file".into(), "{artifact_path}".into()],
5487 ),
5488 materialize: None,
5489 requires_materialization: true,
5490 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(),
5491 },
5492 TransferFormat::Goose => HandoffInstructions {
5493 launch: launch(
5494 "goose",
5495 vec![
5496 "session".into(),
5497 "--resume".into(),
5498 "--session-id".into(),
5499 "{imported_session_id}".into(),
5500 ],
5501 ),
5502 materialize: Some(launch(
5503 "goose",
5504 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5505 )),
5506 requires_materialization: true,
5507 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(),
5508 },
5509 }
5510}
5511
5512fn resume_launch(
5513 harness: &str,
5514 session_id: &str,
5515 cwd: &Path,
5516 policy: ResumePolicy,
5517) -> std::result::Result<StructuredLaunch, ServiceError> {
5518 let mut arguments = Vec::new();
5519 let program = match harness {
5520 HarnessId::GROK => {
5521 if matches!(policy, ResumePolicy::Yolo) {
5522 if crate::support::self_sandbox_supported() {
5523 arguments.extend(["--sandbox".into(), "workspace".into()]);
5524 }
5525 arguments.push("--always-approve".into());
5526 }
5527 arguments.extend(["--resume".into(), session_id.into()]);
5528 "grok"
5529 }
5530 HarnessId::CODEX => {
5531 arguments.extend(crate::startup_prompts::startup_arguments(
5532 harness,
5533 Some(cwd),
5534 &[],
5535 matches!(policy, ResumePolicy::Yolo),
5536 ));
5537 arguments.extend(["resume".into(), session_id.into()]);
5538 "codex"
5539 }
5540 HarnessId::CLAUDE_CODE => {
5541 arguments.extend(crate::startup_prompts::startup_arguments(
5542 harness,
5543 Some(cwd),
5544 &[],
5545 matches!(policy, ResumePolicy::Yolo),
5546 ));
5547 arguments.extend(["--resume".into(), session_id.into()]);
5548 "claude"
5549 }
5550 HarnessId::GEMINI => {
5551 arguments.extend(crate::startup_prompts::startup_arguments(
5552 harness,
5553 Some(cwd),
5554 &[],
5555 matches!(policy, ResumePolicy::Yolo),
5556 ));
5557 arguments.extend(["--resume".into(), session_id.into()]);
5558 "gemini"
5559 }
5560 HarnessId::GOOSE => {
5561 arguments.extend([
5562 "session".into(),
5563 "--resume".into(),
5564 "--session-id".into(),
5565 session_id.into(),
5566 ]);
5567 "goose"
5568 }
5569 HarnessId::PI => {
5570 arguments.extend(crate::startup_prompts::startup_arguments(
5571 harness,
5572 Some(cwd),
5573 &[],
5574 matches!(policy, ResumePolicy::Yolo),
5575 ));
5576 arguments.extend(["--session".into(), session_id.into()]);
5577 "pi"
5578 }
5579 HarnessId::OPENCODE => {
5580 arguments.extend(["--session".into(), session_id.into()]);
5581 "opencode"
5582 }
5583 HarnessId::SUPERCODE => {
5584 if matches!(policy, ResumePolicy::Yolo) {
5585 arguments.push("--dangerous".into());
5586 }
5587 arguments.extend(["resume".into(), session_id.into()]);
5588 "supercode"
5589 }
5590 other => {
5591 return Err(ServiceError::InvalidParams(format!(
5592 "no structured resume launch is registered for harness `{other}`"
5593 )))
5594 }
5595 };
5596 Ok(StructuredLaunch {
5597 cwd: cwd.to_path_buf(),
5598 env: if program == "grok" {
5599 crate::support::grok_home_env()
5600 } else {
5601 BTreeMap::new()
5602 },
5603 program: program.into(),
5604 arguments,
5605 })
5606}
5607
5608fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5614 let digest = blake3::hash(address.as_bytes()).to_hex();
5615 let path = std::env::temp_dir().join(format!(
5616 "supercode-openclaw-gateway-token-{}",
5617 &digest.as_str()[..16]
5618 ));
5619 #[cfg(unix)]
5620 {
5621 use std::io::Write;
5622 use std::os::unix::fs::OpenOptionsExt;
5623 let mut file = std::fs::OpenOptions::new()
5624 .write(true)
5625 .create(true)
5626 .truncate(true)
5627 .mode(0o600)
5628 .open(&path)?;
5629 file.write_all(secret.as_bytes())?;
5630 }
5631 #[cfg(not(unix))]
5632 std::fs::write(&path, secret)?;
5633 Ok(path)
5634}
5635
5636fn open_connect_descriptor(
5642 descriptor: &crate::HarnessSupportDescriptor,
5643 home: &Path,
5644) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5645 let Some(connect) = &descriptor.runtime.connect_launch else {
5646 return Err(ServiceError::InvalidParams(format!(
5647 "harness `{}` has no registered connect-mode launch",
5648 descriptor.id.as_str()
5649 )));
5650 };
5651 let resolved = connect
5652 .resolve(home)
5653 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5654 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5655 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5656 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5657 if let Some(token) = resolved.auth {
5658 backend = backend.with_bearer(token);
5659 }
5660 Ok(Box::new(backend))
5661 }
5662 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5663 let mut env = BTreeMap::new();
5674 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5675 if let Some(token) = resolved.auth {
5676 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5677 .map_err(|error| {
5678 ServiceError::UnsupportedAction(format!(
5679 "could not stage the gateway credential for the bridge: {error}"
5680 ))
5681 })?;
5682 arguments.push("--token-file".into());
5683 arguments.push(token_path.to_string_lossy().into_owned());
5684 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5685 }
5686 let program = descriptor
5691 .runtime
5692 .default_launch
5693 .as_ref()
5694 .map(|launch| launch.program.clone())
5695 .unwrap_or_else(|| "openclaw".into());
5696 let launch = RuntimeLaunch {
5697 program,
5698 arguments,
5699 env,
5700 };
5701 Ok(Box::new(
5702 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5703 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5704 ))
5705 }
5706 _ => Err(ServiceError::UnsupportedAction(format!(
5707 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5708 descriptor.id.as_str(),
5709 connect.protocol
5710 ))),
5711 }
5712}
5713
5714fn registry_connect_descriptor(
5717 params: &RuntimeBackendParams,
5718) -> Option<crate::HarnessSupportDescriptor> {
5719 if params.launch.is_some() || params.base_url.is_some() {
5720 return None;
5721 }
5722 harness_support_registry()
5723 .harnesses
5724 .into_iter()
5725 .find(|descriptor| descriptor.id == params.harness)
5726 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5727}
5728
5729fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5730 supercode_interchange::user_home()
5731 .map(std::path::PathBuf::into_os_string)
5732 .map(PathBuf::from)
5733 .ok_or_else(|| {
5734 ServiceError::UnsupportedAction(
5735 "connect-mode launches need HOME to locate the harness config".into(),
5736 )
5737 })
5738}
5739
5740pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5743 "harness.v1.runtimes.start",
5744 "harness.v1.runtimes.resume",
5745 "harness.v1.runtimes.attach",
5746 "harness.v1.runtimes.attach_existing",
5747];
5748
5749pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5754
5755pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5761
5762pub const DETACHED_METHODS: &[&str] = &[
5768 "harness.v1.harnesses.list",
5769 "harness.v1.harnesses.probe",
5770 "harness.v1.sessions.message",
5771 "harness.v1.sessions.new",
5772 "harness.v1.sessions.reset",
5773 "harness.v1.sessions.archive",
5774 "harness.v1.sessions.delete",
5775];
5776
5777pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5783
5784pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5792
5793async fn within_control_deadline<F: std::future::Future>(
5796 method: &str,
5797 call: F,
5798) -> std::result::Result<F::Output, ServiceError> {
5799 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5800 .await
5801 .map_err(|_| {
5802 ServiceError::Operation(format!(
5803 "`{method}` gave up after {}s: the runtime did not answer",
5804 RUNTIME_CONTROL_DEADLINE.as_secs()
5805 ))
5806 })
5807}
5808
5809pub struct RuntimeOpen {
5813 id: Value,
5814 method: String,
5815 params: Value,
5816}
5817
5818impl RuntimeOpen {
5819 pub async fn open(self) -> OpenedRuntime {
5823 let Self { id, method, params } = self;
5824 let outcome = open_runtime(&method, params).await;
5825 OpenedRuntime { id, outcome }
5826 }
5827}
5828
5829pub struct OpenedRuntime {
5832 id: Value,
5833 outcome: std::result::Result<OpenRuntime, ServiceError>,
5834}
5835
5836pub struct DetachedCall {
5841 id: Value,
5842 method: String,
5843 work: std::result::Result<Work, ServiceError>,
5844}
5845
5846impl DetachedCall {
5847 pub async fn run(self) -> DetachedAnswer {
5850 let Self { id, method, work } = self;
5851 match work {
5852 Ok(Work::Runtime(work)) => {
5857 let (result, returned) = work.run().await;
5858 DetachedAnswer {
5859 response: service_response(id, result),
5860 returned,
5861 }
5862 }
5863 Ok(Work::Free(work)) => {
5864 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5865 Ok(result) => result,
5866 Err(_) => Err(ServiceError::Operation(format!(
5867 "`{method}` gave up after {}s: the harness it waits on did not answer",
5868 DETACHED_CALL_DEADLINE.as_secs()
5869 ))),
5870 };
5871 DetachedAnswer {
5872 response: service_response(id, result),
5873 returned: None,
5874 }
5875 }
5876 Err(error) => DetachedAnswer {
5877 response: service_response(id, Err(error)),
5878 returned: None,
5879 },
5880 }
5881 }
5882}
5883
5884pub struct DetachedAnswer {
5888 response: Value,
5889 returned: Option<ReturnedRuntime>,
5890}
5891
5892impl DetachedAnswer {
5893 pub fn into_response(self) -> Value {
5896 self.response
5897 }
5898}
5899
5900pub struct ReturnedRuntime {
5903 connection: String,
5904 runtime: Box<dyn RuntimeConnection>,
5905}
5906
5907enum Work {
5910 Free(DetachedWork),
5911 Runtime(RuntimeWork),
5912}
5913
5914enum DetachedWork {
5917 Inventory(InventoryWork),
5921 Message(MessageSessionParams),
5923 SessionMutation {
5926 verb: crate::SessionVerb,
5927 mutation: crate::SessionMutation,
5928 },
5929}
5930
5931impl DetachedWork {
5932 async fn run(self) -> std::result::Result<Value, ServiceError> {
5933 match self {
5934 Self::Inventory(work) => run_inventory(work).await,
5935 Self::Message(params) => Ok(message_live_session(¶ms).await),
5936 Self::SessionMutation { verb, mutation } => {
5937 let outcome = run_session_mutation(verb, &mutation).await?;
5938 serde_json::to_value(outcome)
5939 .map_err(|error| ServiceError::Operation(error.to_string()))
5940 }
5941 }
5942 }
5943}
5944
5945enum RuntimeWork {
5947 Close {
5949 runtime: Box<dyn RuntimeConnection>,
5950 process_group: Option<u32>,
5951 },
5952 LiveCommand {
5955 connection: String,
5956 runtime: Box<dyn RuntimeConnection>,
5957 verb: crate::SessionVerb,
5958 mutation: crate::SessionMutation,
5959 command: &'static str,
5960 session: String,
5961 },
5962}
5963
5964type RuntimeWorkAnswer = (
5967 std::result::Result<Value, ServiceError>,
5968 Option<ReturnedRuntime>,
5969);
5970
5971impl RuntimeWork {
5972 async fn run(self) -> RuntimeWorkAnswer {
5973 match self {
5974 Self::Close {
5975 runtime,
5976 process_group,
5977 } => (close_runtime(runtime, process_group).await, None),
5978 Self::LiveCommand {
5979 connection,
5980 mut runtime,
5981 verb,
5982 mutation,
5983 command,
5984 session,
5985 } => {
5986 let result =
5987 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5988 (
5989 result,
5990 Some(ReturnedRuntime {
5991 connection,
5992 runtime,
5993 }),
5994 )
5995 }
5996 }
5997 }
5998}
5999
6000async fn close_runtime(
6003 mut runtime: Box<dyn RuntimeConnection>,
6004 process_group: Option<u32>,
6005) -> std::result::Result<Value, ServiceError> {
6006 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
6007 Ok(result) => {
6008 result.map_err(operation)?;
6009 Ok(json!({"closed": true}))
6010 }
6011 Err(deadline) => {
6012 let killed = kill_runtime_process_group(process_group);
6017 drop(runtime);
6018 Ok(json!({
6019 "closed": true,
6020 "killed": killed,
6021 "detail": error_message(deadline),
6022 }))
6023 }
6024 }
6025}
6026
6027fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
6030 mutation
6031 .session
6032 .clone()
6033 .filter(|value| !value.trim().is_empty())
6034 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
6035}
6036
6037async fn type_live_command(
6041 runtime: &mut dyn RuntimeConnection,
6042 verb: crate::SessionVerb,
6043 mutation: &crate::SessionMutation,
6044 command: &str,
6045 session: String,
6046) -> std::result::Result<Value, ServiceError> {
6047 within_control_deadline(
6048 &format!("sessions.{}", verb.as_str()),
6049 runtime.send_input(RuntimeInput {
6050 text: command.to_string(),
6051 image_urls: Vec::new(),
6052 }),
6053 )
6054 .await?
6055 .map_err(operation)?;
6056 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
6057 .map_err(session_control_error)?;
6058 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
6059}
6060
6061enum OpenRuntime {
6064 Hosted {
6067 runtime: Box<dyn RuntimeConnection>,
6068 capabilities: crate::RuntimeCapabilities,
6069 workspace: PathBuf,
6070 fresh: bool,
6072 },
6073 Joined { runtime: Box<dyn RuntimeConnection> },
6076}
6077
6078async fn open_runtime(
6083 method: &str,
6084 params: Value,
6085) -> std::result::Result<OpenRuntime, ServiceError> {
6086 match tokio::time::timeout(
6087 RUNTIME_OPEN_DEADLINE,
6088 open_runtime_unbounded(method, params),
6089 )
6090 .await
6091 {
6092 Ok(result) => result,
6093 Err(_) => Err(ServiceError::Operation(format!(
6094 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
6095 RUNTIME_OPEN_DEADLINE.as_secs()
6096 ))),
6097 }
6098}
6099
6100async fn open_runtime_unbounded(
6101 method: &str,
6102 params: Value,
6103) -> std::result::Result<OpenRuntime, ServiceError> {
6104 match method {
6105 "harness.v1.runtimes.start" => {
6106 let params = decode::<RuntimeStartParams>(params)?;
6107 let backend = runtime_backend(¶ms.backend)?;
6108 let capabilities = backend.capabilities();
6109 let workspace = params.cwd.clone();
6110 let runtime = backend
6111 .start(RuntimeStartRequest {
6112 cwd: params.cwd,
6113 launch: runtime_launch(¶ms.backend),
6114 mcp_servers: params.mcp_servers,
6115 approval_policy: params.approval_policy,
6116 })
6117 .await
6118 .map_err(operation)?;
6119 Ok(OpenRuntime::Hosted {
6120 runtime,
6121 capabilities,
6122 workspace,
6123 fresh: true,
6124 })
6125 }
6126 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
6127 let params = decode::<RuntimeAttachParams>(params)?;
6128 let backend = runtime_backend(¶ms.backend)?;
6129 let capabilities = backend.capabilities();
6130 let workspace = params
6131 .cwd
6132 .clone()
6133 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
6134 let runtime = backend
6135 .attach(RuntimeAttachRequest {
6136 runtime_id: params.runtime_id,
6137 cwd: params.cwd,
6138 launch: runtime_launch(¶ms.backend),
6139 mcp_servers: params.mcp_servers,
6140 approval_policy: params.approval_policy,
6141 })
6142 .await
6143 .map_err(operation)?;
6144 Ok(OpenRuntime::Hosted {
6145 runtime,
6146 capabilities,
6147 workspace,
6148 fresh: false,
6149 })
6150 }
6151 "harness.v1.runtimes.attach_existing" => {
6152 let params = decode::<RuntimeAttachParams>(params)?;
6153 let backend: Box<dyn RuntimeBackend> = match params
6154 .backend
6155 .base_url
6156 .as_deref()
6157 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
6158 {
6159 Some(endpoint) => {
6160 #[cfg(not(feature = "adapter-api"))]
6161 {
6162 let _ = endpoint;
6163 return Err(ServiceError::UnsupportedAction(
6164 "live HTTP attachment adapter is not compiled".into(),
6165 ));
6166 }
6167 #[cfg(feature = "adapter-api")]
6168 {
6169 let workspace = params.cwd.clone().ok_or_else(|| {
6170 ServiceError::InvalidParams(
6171 "Volter Harness live attach requires the project cwd".into(),
6172 )
6173 })?;
6174 let source = LiveRuntimeSource {
6175 harness: params.backend.harness.as_str().to_string(),
6176 session_id: params.runtime_id.clone(),
6177 workspace,
6178 };
6179 let receipt = resolve_live_runtime(&endpoint, &source)
6180 .map_err(|error| ServiceError::Operation(error.to_string()))?;
6181 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
6182 }
6183 }
6184 None => runtime_backend(¶ms.backend)?,
6185 };
6186 let capabilities = backend.capabilities();
6187 if !capabilities.attach_existing_process {
6188 return Err(ServiceError::Operation(format!(
6189 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
6190 backend.harness().as_str()
6191 )));
6192 }
6193 let runtime = backend
6194 .attach_existing(RuntimeAttachRequest {
6195 runtime_id: params.runtime_id,
6196 cwd: params.cwd,
6197 launch: runtime_launch(¶ms.backend),
6198 mcp_servers: params.mcp_servers,
6199 approval_policy: params.approval_policy,
6200 })
6201 .await
6202 .map_err(operation)?;
6203 Ok(OpenRuntime::Joined { runtime })
6204 }
6205 _ => Err(ServiceError::MethodNotFound),
6206 }
6207}
6208
6209fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
6211 match result {
6212 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
6213 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
6214 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
6215 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
6216 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
6217 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
6218 }
6219}
6220
6221fn runtime_backend(
6222 params: &RuntimeBackendParams,
6223) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
6224 if let Some(descriptor) = registry_connect_descriptor(params) {
6225 return open_connect_descriptor(&descriptor, &service_home()?);
6226 }
6227 if params.protocol.as_deref() == Some("acp") {
6228 let launch = params
6229 .launch
6230 .clone()
6231 .or_else(|| {
6232 harness_support_registry()
6233 .harnesses
6234 .into_iter()
6235 .find(|harness| harness.id == params.harness)
6236 .filter(|harness| {
6237 harness.runtime.implementation == ImplementationKind::GenericProtocol
6238 && harness.runtime.protocol.starts_with("acp")
6239 })
6240 .and_then(|harness| harness.runtime.default_launch)
6241 })
6242 .ok_or_else(|| {
6243 ServiceError::InvalidParams(
6244 "an ACP runtime requires `launch` unless the harness has a registered default"
6245 .into(),
6246 )
6247 })?;
6248 let resume_session = harness_support_registry()
6249 .harnesses
6250 .into_iter()
6251 .find(|harness| harness.id == params.harness)
6252 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
6253 return Ok(Box::new(
6254 AcpRuntimeBackend::new(params.harness.clone(), launch)
6255 .with_resume_support(resume_session),
6256 ));
6257 }
6258 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
6259 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
6260 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
6261 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
6262 HarnessId::OPENCODE => match ¶ms.base_url {
6263 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
6264 None => Box::new(OpenCodeRuntimeBackend::new()),
6265 },
6266 harness => {
6267 let descriptor = harness_support_registry()
6268 .harnesses
6269 .into_iter()
6270 .find(|descriptor| descriptor.id.as_str() == harness)
6271 .filter(|descriptor| {
6272 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
6273 && descriptor.runtime.protocol.starts_with("acp")
6274 });
6275 let Some(descriptor) = descriptor else {
6276 return Err(ServiceError::InvalidParams(format!(
6277 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
6278 )));
6279 };
6280 let resume = descriptor.runtime.capabilities.resume_session;
6281 Box::new(
6282 AcpRuntimeBackend::new(
6283 descriptor.id,
6284 descriptor
6285 .runtime
6286 .default_launch
6287 .expect("generic ACP registry entry includes its launch"),
6288 )
6289 .with_resume_support(resume),
6290 )
6291 }
6292 };
6293 Ok(backend)
6294}
6295
6296fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
6297 if let Some(launch) = ¶ms.launch {
6298 return Some(launch.clone());
6299 }
6300 if !matches!(params.policy, RuntimePolicy::Yolo) {
6301 return None;
6302 }
6303 let launch = match params.harness.as_str() {
6304 HarnessId::GROK => RuntimeLaunch {
6305 program: "grok".into(),
6306 arguments: {
6307 let mut arguments: Vec<String> = Vec::new();
6308 if crate::support::self_sandbox_supported() {
6309 arguments.extend(["--sandbox".into(), "workspace".into()]);
6310 }
6311 arguments.extend([
6312 "--always-approve".into(),
6313 "agent".into(),
6314 "--no-leader".into(),
6315 "stdio".into(),
6316 ]);
6317 arguments
6318 },
6319 env: crate::support::grok_env(),
6320 },
6321 HarnessId::CODEX => RuntimeLaunch {
6322 program: "codex".into(),
6323 arguments: vec![
6324 "--dangerously-bypass-approvals-and-sandbox".into(),
6325 "--dangerously-bypass-hook-trust".into(),
6326 "app-server".into(),
6327 ],
6328 env: BTreeMap::new(),
6329 },
6330 HarnessId::CLAUDE_CODE => RuntimeLaunch {
6331 program: "claude".into(),
6332 arguments: vec![
6333 "--dangerously-skip-permissions".into(),
6334 "--print".into(),
6335 "--input-format".into(),
6336 "stream-json".into(),
6337 "--output-format".into(),
6338 "stream-json".into(),
6339 "--verbose".into(),
6340 ],
6341 env: BTreeMap::new(),
6342 },
6343 HarnessId::PI => RuntimeLaunch {
6344 program: "pi".into(),
6345 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
6346 env: BTreeMap::new(),
6347 },
6348 HarnessId::OPENCODE => RuntimeLaunch {
6349 program: "opencode".into(),
6350 arguments: vec!["serve".into()],
6351 env: BTreeMap::new(),
6352 },
6353 HarnessId::GEMINI => RuntimeLaunch {
6354 program: "gemini".into(),
6355 arguments: vec!["--acp".into(), "--yolo".into()],
6356 env: BTreeMap::new(),
6357 },
6358 HarnessId::GOOSE => RuntimeLaunch {
6359 program: "goose".into(),
6360 arguments: vec!["acp".into()],
6361 env: BTreeMap::new(),
6362 },
6363 HarnessId::SUPERCODE => RuntimeLaunch {
6364 program: "supercode".into(),
6365 arguments: vec!["acp".into(), "--dangerous".into()],
6366 env: BTreeMap::new(),
6367 },
6368 _ => return None,
6369 };
6370 Some(launch)
6371}
6372
6373struct IsolatedProbeHome {
6379 launch: RuntimeLaunch,
6380 root: PathBuf,
6381}
6382
6383impl IsolatedProbeHome {
6384 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6385 let root = std::env::temp_dir().join(format!(
6386 "supercode-harness-probe-{harness}-{}",
6387 generated_session_id()
6388 ));
6389 std::fs::create_dir_all(&root)?;
6390 set_private_dir_permissions(&root)?;
6391
6392 if let Some(source_home) = supercode_interchange::user_home()
6393 .map(std::path::PathBuf::into_os_string)
6394 .map(PathBuf::from)
6395 {
6396 for relative in probe_auth_files(harness) {
6397 copy_probe_file(&source_home, &root, relative)?;
6398 }
6399 }
6400 if harness == HarnessId::SUPERCODE {
6405 let config_home = crate::agent::global_instructions_dir();
6406 for file in ["config.toml", "credentials.toml"] {
6407 copy_probe_path(
6408 &config_home.join(file),
6409 &root.join(".config/supercode").join(file),
6410 )?;
6411 }
6412 }
6413 configure_isolated_probe_auth(harness, &root)?;
6414
6415 let root_text = root.to_string_lossy().into_owned();
6416 for (key, value) in [
6417 ("HOME", root_text.clone()),
6418 (
6419 "XDG_CACHE_HOME",
6420 root.join(".cache").to_string_lossy().into_owned(),
6421 ),
6422 (
6423 "XDG_CONFIG_HOME",
6424 root.join(".config").to_string_lossy().into_owned(),
6425 ),
6426 (
6427 "XDG_DATA_HOME",
6428 root.join(".local/share").to_string_lossy().into_owned(),
6429 ),
6430 ] {
6431 launch.env.insert(key.into(), value);
6432 }
6433 let scoped = match harness {
6434 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6435 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6436 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6437 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6438 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6439 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6440 _ => None,
6441 };
6442 if let Some((key, value)) = scoped {
6443 launch
6444 .env
6445 .insert(key.into(), value.to_string_lossy().into_owned());
6446 }
6447 Ok(Self { launch, root })
6448 }
6449
6450 fn cleanup(&self) -> std::io::Result<()> {
6451 match std::fs::remove_dir_all(&self.root) {
6452 Ok(()) => Ok(()),
6453 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6454 Err(error) => Err(error),
6455 }
6456 }
6457}
6458
6459impl Drop for IsolatedProbeHome {
6460 fn drop(&mut self) {
6461 let _ = self.cleanup();
6462 }
6463}
6464
6465fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6466 match harness {
6467 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6468 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6472 HarnessId::CODEX => &[".codex/auth.json"],
6473 HarnessId::GEMINI => &[
6474 ".gemini/google_accounts.json",
6475 ".gemini/oauth_creds.json",
6476 ".gemini/settings.json",
6477 ],
6478 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6479 HarnessId::OPENCODE => &[
6480 ".config/opencode/auth.json",
6481 ".local/share/opencode/auth.json",
6482 ],
6483 HarnessId::PI => &[".pi/agent/auth.json"],
6484 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6489 _ => &[],
6490 }
6491}
6492
6493fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6494 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6495}
6496
6497fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6498 if !source.is_file() {
6499 return Ok(());
6500 }
6501 if let Some(parent) = destination.parent() {
6502 std::fs::create_dir_all(parent)?;
6503 set_private_dir_permissions(parent)?;
6504 }
6505 std::fs::copy(source, destination)?;
6506 set_private_file_permissions(destination)
6507}
6508
6509fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6510 if harness != HarnessId::GEMINI {
6511 return Ok(());
6512 }
6513 let oauth = probe_home.join(".gemini/oauth_creds.json");
6514 if !oauth.is_file() {
6515 return Ok(());
6516 }
6517 let settings_path = probe_home.join(".gemini/settings.json");
6518 let mut settings = std::fs::read_to_string(&settings_path)
6519 .ok()
6520 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6521 .unwrap_or_else(|| json!({}));
6522 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6523 std::fs::write(
6524 &settings_path,
6525 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6526 )?;
6527 set_private_file_permissions(&settings_path)
6528}
6529
6530#[cfg(unix)]
6531fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6532 use std::os::unix::fs::PermissionsExt;
6533 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6534}
6535
6536#[cfg(not(unix))]
6537fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6538 Ok(())
6539}
6540
6541#[cfg(unix)]
6542fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6543 use std::os::unix::fs::PermissionsExt;
6544 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6545}
6546
6547#[cfg(not(unix))]
6548fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6549 Ok(())
6550}
6551
6552fn find_executable(program: &str) -> Option<PathBuf> {
6553 let candidate = PathBuf::from(program);
6554 if candidate.components().count() > 1 {
6555 return candidate.is_file().then_some(candidate);
6556 }
6557 let path = std::env::var_os("PATH")?;
6558 for directory in std::env::split_paths(&path) {
6559 let candidate = directory.join(program);
6560 if candidate.is_file() {
6561 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6562 }
6563 #[cfg(windows)]
6564 {
6565 for extension in ["exe", "cmd", "bat"] {
6566 let candidate = directory.join(format!("{program}.{extension}"));
6567 if candidate.is_file() {
6568 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6569 }
6570 }
6571 }
6572 }
6573 None
6574}
6575
6576async fn executable_version(executable: &Path) -> Option<String> {
6577 let mut command = tokio::process::Command::new(executable);
6578 command
6579 .arg("--version")
6580 .stdin(std::process::Stdio::null())
6581 .stdout(std::process::Stdio::piped())
6582 .stderr(std::process::Stdio::piped())
6583 .kill_on_drop(true);
6584 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6585 .await
6586 .ok()?
6587 .ok()?;
6588 let stdout = String::from_utf8_lossy(&output.stdout);
6589 let stderr = String::from_utf8_lossy(&output.stderr);
6590 stdout
6591 .lines()
6592 .chain(stderr.lines())
6593 .map(str::trim)
6594 .find(|line| !line.is_empty())
6595 .map(|line| truncate_text(line, 200))
6596}
6597
6598pub(crate) fn auth_evidence(harness: &str) -> bool {
6599 let env_names: &[&str] = match harness {
6600 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6601 HarnessId::CODEX => &["OPENAI_API_KEY"],
6602 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6603 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6604 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6605 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6606 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6607 _ => &[],
6608 };
6609 if env_names
6610 .iter()
6611 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6612 {
6613 return true;
6614 }
6615 let Some(home) = supercode_interchange::user_home()
6616 .map(std::path::PathBuf::into_os_string)
6617 .map(PathBuf::from)
6618 else {
6619 return false;
6620 };
6621 let files: Vec<PathBuf> = match harness {
6622 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6623 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6624 HarnessId::OPENCODE => vec![
6625 home.join(".local/share/opencode/auth.json"),
6626 home.join(".config/opencode/auth.json"),
6627 ],
6628 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6629 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6630 HarnessId::GEMINI => vec![
6631 home.join(".gemini/oauth_creds.json"),
6632 home.join(".gemini/google_accounts.json"),
6633 ],
6634 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6635 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6636 _ => Vec::new(),
6637 };
6638 if files.into_iter().any(|path| {
6639 std::fs::metadata(path)
6640 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6641 .unwrap_or(false)
6642 }) {
6643 return true;
6644 }
6645 if harness == HarnessId::CLAUDE_CODE {
6652 return std::fs::read_to_string(home.join(".claude.json"))
6653 .map(|text| text.contains("\"oauthAccount\""))
6654 .unwrap_or(false);
6655 }
6656 false
6657}
6658
6659fn looks_like_auth_error(message: &str) -> bool {
6660 let message = message.to_ascii_lowercase();
6661 [
6662 "auth",
6663 "login",
6664 "sign in",
6665 "sign-in",
6666 "credential",
6667 "unauthorized",
6668 "forbidden",
6669 "token",
6670 ]
6671 .iter()
6672 .any(|needle| message.contains(needle))
6673}
6674
6675fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6676 crate::RuntimeCapabilities {
6677 start_session: false,
6678 resume_session: false,
6679 attach_existing_process: false,
6680 send_input: false,
6681 stream_events: false,
6682 interrupt: false,
6683 steer: false,
6684 respond_to_requests: false,
6685 }
6686}
6687
6688fn truncate_text(text: &str, max_chars: usize) -> String {
6689 let mut chars = text.chars();
6690 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6691 if chars.next().is_some() {
6692 format!("{truncated}…")
6693 } else {
6694 truncated
6695 }
6696}
6697
6698fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6705 match &handle.endpoint {
6706 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6707 crate::RuntimeEndpoint::Http { .. } => None,
6708 }
6709}
6710
6711fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6717 match process_group {
6718 #[cfg(unix)]
6719 Some(pid) => {
6720 crate::lsp::kill_process_group(pid);
6721 true
6722 }
6723 #[cfg(not(unix))]
6724 Some(_) => false,
6725 None => false,
6726 }
6727}
6728
6729fn error_message(error: ServiceError) -> String {
6730 match error {
6731 ServiceError::InvalidParams(message)
6732 | ServiceError::Operation(message)
6733 | ServiceError::UnsupportedAction(message) => message,
6734 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6735 ServiceError::Sdk(error) => error.to_string(),
6736 }
6737}
6738
6739#[derive(Debug)]
6740enum ServiceError {
6741 InvalidParams(String),
6742 MethodNotFound,
6743 UnsupportedAction(String),
6744 Operation(String),
6745 Sdk(SdkError),
6746}
6747
6748fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6749 match error {
6750 ServiceError::InvalidParams(message) => {
6751 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6752 }
6753 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6754 SdkError::unsupported(operation)
6755 }
6756 ServiceError::Operation(message) => {
6757 let code = if message.contains("already in progress") {
6758 SdkErrorCode::Busy
6759 } else if message.contains("not supported by this runtime") {
6760 SdkErrorCode::UnsupportedAction
6761 } else if message.contains("unknown runtime connection") {
6762 SdkErrorCode::NotFound
6763 } else {
6764 SdkErrorCode::Execution
6765 };
6766 SdkError::new(code, operation, message)
6767 }
6768 ServiceError::Sdk(error) => error,
6769 }
6770}
6771
6772fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6773 let error_code = error.code();
6774 let code = match error_code {
6775 SdkErrorCode::Unauthenticated => -32030,
6776 SdkErrorCode::Unauthorized => -32031,
6777 SdkErrorCode::ControllerRequired => -32032,
6778 SdkErrorCode::LeaseExpired => -32033,
6779 SdkErrorCode::InvalidArgument => -32602,
6780 SdkErrorCode::NotFound => -32004,
6781 SdkErrorCode::Busy => -32000,
6782 SdkErrorCode::UnsupportedAction => -32020,
6783 SdkErrorCode::Execution => -32002,
6784 SdkErrorCode::Transport => -32003,
6785 };
6786 json!({
6787 "jsonrpc": "2.0",
6788 "id": id,
6789 "error": {
6790 "code": code,
6791 "name": error_code,
6792 "operation": error.operation(),
6793 "message": error.to_string(),
6794 },
6795 })
6796}
6797
6798fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6799 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6800}
6801
6802fn operation(error: impl Into<crate::Error>) -> ServiceError {
6803 let error = error.into();
6804 match error {
6805 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6806 error => ServiceError::Operation(error.to_string()),
6807 }
6808}
6809
6810#[derive(Debug, Clone, Deserialize, Default)]
6814#[serde(default)]
6815struct MemoryRequest {
6816 harness: Option<String>,
6818 query: Option<String>,
6820 profile: Option<String>,
6822 session: Option<String>,
6824 full: bool,
6826 regex: bool,
6828 cwd: Option<std::path::PathBuf>,
6830 homes: crate::HarnessHomes,
6832}
6833
6834fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6837 let request = decode::<MemoryRequest>(params)?;
6838 let harness = request
6839 .harness
6840 .clone()
6841 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6842 let to_service = |error: crate::memory::MemoryError| match error {
6843 crate::memory::MemoryError::UnsupportedHarness { .. }
6844 | crate::memory::MemoryError::SessionNotScoped { .. } => {
6845 ServiceError::UnsupportedAction(error.to_string())
6846 }
6847 other => ServiceError::InvalidParams(other.to_string()),
6848 };
6849 match method {
6850 "harness.v1.memory.show" => {
6851 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6852 harness,
6853 profile: request.profile,
6854 session: request.session,
6855 full: request.full,
6856 cwd: request.cwd,
6857 homes: request.homes,
6858 })
6859 .map_err(to_service)?;
6860 Ok(json!({
6861 "schema": crate::memory::MEMORY_SCHEMA,
6862 "documents": documents,
6863 }))
6864 }
6865 "harness.v1.memory.search" => {
6866 let query = request
6867 .query
6868 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6869 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6870 harness,
6871 query,
6872 profile: request.profile,
6873 regex: request.regex,
6874 cwd: request.cwd,
6875 homes: request.homes,
6876 })
6877 .map_err(to_service)?;
6878 Ok(json!({
6879 "schema": crate::memory::MEMORY_SCHEMA,
6880 "matches": matches,
6881 }))
6882 }
6883 _ => Err(ServiceError::MethodNotFound),
6884 }
6885}
6886
6887#[derive(Debug, Clone, Deserialize)]
6891#[serde(default)]
6892struct ProfilesQuery {
6893 harness: Option<String>,
6895 name: Option<String>,
6897 homes: crate::HarnessHomes,
6899}
6900
6901impl Default for ProfilesQuery {
6902 fn default() -> Self {
6903 Self {
6904 harness: None,
6905 name: None,
6906 homes: crate::HarnessHomes::default(),
6907 }
6908 }
6909}
6910
6911fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6914 let query = decode::<ProfilesQuery>(params)?;
6915 let to_service = |error: crate::profiles::ProfileError| match error {
6916 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6917 ServiceError::UnsupportedAction(error.to_string())
6918 }
6919 crate::profiles::ProfileError::NotFound { .. } => {
6920 ServiceError::InvalidParams(error.to_string())
6921 }
6922 };
6923 match method {
6924 "harness.v1.profiles.list" => {
6925 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6926 .map_err(to_service)?;
6927 Ok(json!({
6928 "schema": crate::profiles::PROFILES_SCHEMA,
6929 "profiles": profiles,
6930 }))
6931 }
6932 "harness.v1.profiles.get" => {
6933 let harness = query
6934 .harness
6935 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6936 let name = query
6937 .name
6938 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6939 let profile =
6940 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6941 Ok(json!({
6942 "schema": crate::profiles::PROFILES_SCHEMA,
6943 "profile": profile,
6944 }))
6945 }
6946 _ => Err(ServiceError::MethodNotFound),
6947 }
6948}
6949
6950#[derive(Debug, Clone, Deserialize)]
6954#[serde(default)]
6955struct ChannelsQuery {
6956 harness: Option<String>,
6958 name: Option<String>,
6960 homes: crate::HarnessHomes,
6962}
6963
6964impl Default for ChannelsQuery {
6965 fn default() -> Self {
6966 Self {
6967 harness: None,
6968 name: None,
6969 homes: crate::HarnessHomes::default(),
6970 }
6971 }
6972}
6973
6974#[derive(Debug, Clone, Deserialize)]
6978#[serde(default)]
6979struct RoutesQuery {
6980 harness: Option<String>,
6981 profile: Option<String>,
6983 homes: crate::HarnessHomes,
6984}
6985
6986impl Default for RoutesQuery {
6987 fn default() -> Self {
6988 Self {
6989 harness: None,
6990 profile: None,
6991 homes: crate::HarnessHomes::default(),
6992 }
6993 }
6994}
6995
6996#[derive(Debug, Clone, Deserialize)]
6997#[serde(default)]
6998struct TriggersQuery {
6999 harness: Option<String>,
7000 homes: crate::HarnessHomes,
7001}
7002
7003impl Default for TriggersQuery {
7004 fn default() -> Self {
7005 Self {
7006 harness: None,
7007 homes: crate::HarnessHomes::default(),
7008 }
7009 }
7010}
7011
7012fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
7013 let query = decode::<TriggersQuery>(params)?;
7014 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
7015 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7016 Ok(json!({
7017 "schema": crate::triggers::TRIGGERS_SCHEMA,
7018 "triggers": triggers,
7019 }))
7020}
7021
7022fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
7023 let query = decode::<RoutesQuery>(params)?;
7024 let routes = crate::routes::list_routes(
7025 &query.homes,
7026 query.harness.as_deref(),
7027 query.profile.as_deref(),
7028 )
7029 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7030 Ok(json!({
7031 "schema": crate::routes::ROUTES_SCHEMA,
7032 "routes": routes,
7033 }))
7034}
7035
7036fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7037 let query = decode::<ChannelsQuery>(params)?;
7038 let to_service = |error: crate::channels::ChannelError| match error {
7039 crate::channels::ChannelError::UnsupportedHarness { .. } => {
7040 ServiceError::UnsupportedAction(error.to_string())
7041 }
7042 crate::channels::ChannelError::NotFound { .. } => {
7043 ServiceError::InvalidParams(error.to_string())
7044 }
7045 };
7046 match method {
7047 "harness.v1.channels.list" => {
7048 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
7049 .map_err(to_service)?;
7050 Ok(json!({
7051 "schema": crate::channels::CHANNELS_SCHEMA,
7052 "channels": channels,
7053 }))
7054 }
7055 "harness.v1.channels.status" => {
7056 let harness = query
7057 .harness
7058 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7059 let name = query
7060 .name
7061 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7062 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
7063 .map_err(to_service)?;
7064 Ok(json!({
7065 "schema": crate::channels::CHANNELS_SCHEMA,
7066 "channel": channel,
7067 }))
7068 }
7069 _ => Err(ServiceError::MethodNotFound),
7070 }
7071}
7072
7073fn rpc_error(id: Value, code: i64, message: &str) -> Value {
7074 json!({
7075 "jsonrpc": "2.0",
7076 "id": id,
7077 "error": {"code": code, "message": message},
7078 })
7079}