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 recorded_config(¶ms.locator).map_err(operation)
1122 }
1123 "harness.v1.sessions.load" => {
1124 let params = decode::<LoadSessionParams>(params)?;
1125 if let Some(options) = ¶ms.options {
1126 options.validate()?;
1127 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
1128 return Ok(result);
1129 }
1130 return load_session(¶ms.read.locator)
1131 .map(|session| projected_session_result(&session, options))
1132 .map_err(operation);
1133 }
1134 let mut session = if params.read.display_history() {
1135 self.catalog
1136 .load_display_view(
1137 ¶ms.read.locator,
1138 params.read.read_fidelity(),
1139 params.read.tail_messages().unwrap_or(500),
1140 )
1141 .map_err(crate::Error::from)
1142 } else if params.read.include_subagents() {
1143 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
1144 } else {
1145 self.catalog
1146 .load_parent_with_fidelity(
1147 ¶ms.read.locator,
1148 params.read.read_fidelity(),
1149 )
1150 .map_err(crate::Error::from)
1151 }
1152 .map_err(operation)?;
1153 params.read.bound_session(&mut session);
1154 let mut value = normalized_session_json(&session);
1157 value["resolved_profile"] = recorded_session_profile(&session);
1158 value["resolved_profile_history"] =
1159 Value::Array(recorded_session_profile_history(&session));
1160 Ok(json!({"session": value}))
1161 }
1162 "harness.v1.sessions.follow" => {
1163 let params = decode::<LocatorParams>(params)?;
1164 let mut follower = self
1165 .catalog
1166 .follow_read_view(
1167 ¶ms.locator,
1168 params.read_fidelity(),
1169 params.include_subagents(),
1170 params.tail_messages(),
1171 params.max_message_chars(),
1172 params.display_history(),
1173 )
1174 .map_err(operation)?;
1175 let initial = follower
1176 .poll()
1177 .map_err(operation)?
1178 .map(|event| event.to_json());
1179 let subscription = format!("sub-{}", self.next_subscription);
1180 self.next_subscription += 1;
1181 self.followers.insert(subscription.clone(), follower);
1182 self.followed_sources.insert(
1183 subscription.clone(),
1184 FollowedSource {
1185 harness: params.locator.harness.as_str().to_string(),
1186 session_id: params.locator.session_id.clone(),
1187 reported: None,
1188 },
1189 );
1190 Ok(json!({"subscription": subscription, "initial": initial}))
1191 }
1192 "harness.v1.sessions.unfollow" => {
1193 let params = decode::<UnfollowParams>(params)?;
1194 self.followed_sources.remove(¶ms.subscription);
1195 Ok(json!({
1196 "removed": self.followers.remove(¶ms.subscription).is_some()
1197 }))
1198 }
1199 "harness.v1.sessions.activity.unsubscribe" => {
1200 let params = decode::<UnfollowParams>(params)?;
1201 Ok(json!({
1202 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1203 }))
1204 }
1205 "harness.v1.sessions.index.subscribe" => {
1206 let query = decode::<DiscoveryQuery>(params)?;
1207 crate::session_index::validate_query(&query)
1208 .map_err(ServiceError::InvalidParams)?;
1209 let homes = query.homes.clone();
1210 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1211 query,
1212 Arc::clone(&self.index_notifier),
1213 )
1214 .map_err(ServiceError::Operation)?;
1215 let doors = crate::mail_route::LiveSessions::read(&homes);
1216 let initial = initial
1217 .iter()
1218 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1219 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1220 let subscription = format!("index-sub-{}", self.next_subscription);
1221 self.next_subscription += 1;
1222 self.index_subscriptions.insert(subscription.clone(), index);
1223 Ok(json!({
1224 "subscription": subscription,
1225 "revision": 1,
1226 "initial": initial,
1227 }))
1228 }
1229 "harness.v1.sessions.index.resize" => {
1230 let params = decode::<IndexResizeParams>(params)?;
1231 crate::session_index::validate_limit(params.limit)
1232 .map_err(ServiceError::InvalidParams)?;
1233 let index = self
1234 .index_subscriptions
1235 .get_mut(¶ms.subscription)
1236 .ok_or_else(|| {
1237 ServiceError::InvalidParams("unknown session index subscription".into())
1238 })?;
1239 let prepared = index
1240 .prepare_resize(params.limit)
1241 .map_err(ServiceError::Operation)?;
1242 let doors = crate::mail_route::LiveSessions::read(index.homes());
1243 let initial = prepared
1244 .page
1245 .sessions
1246 .iter()
1247 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1248 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1249 let response = json!({
1250 "subscription": params.subscription,
1251 "revision": prepared.revision,
1252 "initial": initial,
1253 "receipt": prepared.page.receipt,
1254 });
1255 index.commit_resize(prepared);
1256 Ok(response)
1257 }
1258 "harness.v1.sessions.index.unsubscribe" => {
1259 let params = decode::<UnfollowParams>(params)?;
1260 Ok(json!({
1261 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1262 }))
1263 }
1264 "harness.v1.sessions.import" => {
1265 let params = decode::<ImportSessionParams>(params)?;
1266 let session = Session::load_str(¶ms.content, params.source_harness.into())
1267 .map_err(operation)?;
1268 Ok(json!({"session": normalized_session_json(&session)}))
1269 }
1270 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1271 let params = decode::<ExportSessionParams>(params)?;
1272 let session = load_session(¶ms.locator).map_err(operation)?;
1273 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1274 if method == "harness.v1.sessions.export"
1275 && params.target_harness == TransferFormat::Hermes
1276 {
1277 let imported = crate::hermes_import::import_into_hermes(&session, None)
1279 .map_err(operation)?;
1280 return Ok(json!({"artifact": artifact, "imported": imported}));
1281 }
1282 Ok(json!({"artifact": artifact}))
1283 }
1284 "harness.v1.sessions.reduce" => {
1285 let params = decode::<ReduceSessionParams>(params)?;
1286 self.reduce_session(params)
1287 }
1288 "harness.v1.sessions.branch" => {
1289 let params = decode::<BranchSessionParams>(params)?;
1290 let session = load_session(¶ms.locator).map_err(operation)?;
1291 let storage = params.locator.storage.path().display().to_string();
1292 let bootstrap_prompt = format!(
1293 "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.",
1294 params.locator.harness.as_str(), params.locator.session_id, storage
1295 );
1296 let artifact = params
1297 .target_harness
1298 .map(|target| session_artifact(¶ms.locator, &session, target))
1299 .transpose()?;
1300 Ok(json!({
1301 "parent": params.locator,
1302 "session": normalized_session_json(&session),
1303 "bootstrap_prompt": bootstrap_prompt,
1304 "artifact": artifact,
1305 }))
1306 }
1307 "harness.v1.sessions.handoff" => {
1308 let params = decode::<HandoffSessionParams>(params)?;
1309 let session = load_session(¶ms.locator).map_err(operation)?;
1310 let cwd = params
1311 .cwd
1312 .or_else(|| session.meta.cwd.clone())
1313 .unwrap_or_else(|| PathBuf::from("."));
1314 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1315 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1316 ServiceError::Operation(
1317 "handoff artifact omitted target session identity".into(),
1318 )
1319 })?;
1320 let instructions =
1321 handoff_instructions(params.target_harness, target_session_id, &cwd);
1322 Ok(json!({
1323 "artifact": artifact,
1324 "launch": instructions.launch,
1325 "materialize": instructions.materialize,
1326 "requires_materialization": instructions.requires_materialization,
1327 "note": instructions.note,
1328 }))
1329 }
1330 "harness.v1.sessions.materialize" => {
1331 let params = decode::<MaterializeSessionParams>(params)?;
1332 for file in ¶ms.artifact.files {
1336 if file.role == "source_recovery"
1337 && file.path == "recovery/source.supercode.jsonl"
1338 {
1339 if let Ok(source) = Session::from_native_str(&file.content) {
1340 crate::residue_store::store_segments(&source);
1341 }
1342 }
1343 }
1344 let locator = crate::native_materialize::materialize_native_artifact(
1345 params.artifact,
1346 ¶ms.cwd,
1347 ¶ms.homes,
1348 )
1349 .map_err(ServiceError::Operation)?;
1350 Ok(json!({"locator": locator}))
1351 }
1352 "harness.v1.jobs.list" => {
1356 let query = decode::<crate::jobs::JobsQuery>(params)?;
1357 if let Some(harness) = query.harness.as_deref() {
1358 refuse_harness_without_jobs(harness, "jobs.list")?;
1359 }
1360 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1361 serde_json::to_value(listing)
1362 .map_err(|error| ServiceError::Operation(error.to_string()))
1363 }
1364 "harness.v1.jobs.get" => {
1365 let params = decode::<JobsGetParams>(params)?;
1366 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1367 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1368 .map_err(operation)?
1369 {
1370 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1371 None => Err(ServiceError::Operation(format!(
1372 "`{}` has no scheduled job `{}`",
1373 params.harness, params.id
1374 ))),
1375 }
1376 }
1377 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1383 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1384 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1385 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1386 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1387 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1388 "harness.v1.jobs.notepad"
1389 | "harness.v1.jobs.notepad_set"
1390 | "harness.v1.jobs.notepad_delete" => {
1391 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1392 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1393 let answer = match method {
1394 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1395 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1396 _ => crate::jobs_notepad::read(&request),
1397 }
1398 .map_err(job_control_error)?;
1399 serde_json::to_value(answer)
1400 .map_err(|error| ServiceError::Operation(error.to_string()))
1401 }
1402 "harness.v1.runs.list" => {
1406 let query = decode::<crate::runs::RunsQuery>(params)?;
1407 if let Some(harness) = query.harness.as_deref() {
1408 refuse_harness_without_runs(harness, "runs.list")?;
1409 }
1410 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1411 serde_json::to_value(listing)
1412 .map_err(|error| ServiceError::Operation(error.to_string()))
1413 }
1414 "harness.v1.runs.get" => {
1415 let params = decode::<RunsGetParams>(params)?;
1416 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1417 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1418 .map_err(operation)?
1419 {
1420 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1421 None => Err(ServiceError::Operation(format!(
1422 "`{}` has no run `{}`",
1423 params.harness, params.id
1424 ))),
1425 }
1426 }
1427 "harness.v1.sessions.resume_instructions" => {
1428 let params = decode::<ResumeInstructionsParams>(params)?;
1429 let session = load_session(¶ms.locator).map_err(operation)?;
1430 let cwd = params
1431 .cwd
1432 .or(session.meta.cwd)
1433 .unwrap_or_else(|| PathBuf::from("."));
1434 let launch = resume_launch(
1435 params.locator.harness.as_str(),
1436 ¶ms.locator.session_id,
1437 &cwd,
1438 params.policy,
1439 )?;
1440 Ok(json!({"launch": launch}))
1441 }
1442 _ => Err(ServiceError::MethodNotFound),
1443 }
1444 }
1445
1446 fn reduce_session(
1447 &self,
1448 params: ReduceSessionParams,
1449 ) -> std::result::Result<Value, ServiceError> {
1450 let session = load_session(¶ms.locator).map_err(operation)?;
1451 if session.messages.is_empty() {
1452 return Err(ServiceError::InvalidParams(
1453 "cannot reduce an empty session".into(),
1454 ));
1455 }
1456 let keep_last = params.keep_last.clamp(1, 128);
1457 let policy = reduce::ReductionPolicy {
1458 clear_turns_older_than: Some(keep_last),
1459 ..Default::default()
1460 };
1461 let (view, log) =
1462 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1463 if log.reductions.is_empty() {
1464 return Err(ServiceError::UnsupportedAction(format!(
1465 "session `{}` is already too small for a meaningful reversible reduction",
1466 params.locator.session_id
1467 )));
1468 }
1469 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1470 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1471 if reduced_tokens >= source_tokens {
1472 return Err(ServiceError::UnsupportedAction(format!(
1473 "session `{}` has no token-reducing reversible projection",
1474 params.locator.session_id
1475 )));
1476 }
1477
1478 let store_root = self
1479 .reduction_store_root
1480 .clone()
1481 .unwrap_or_else(default_reduction_store_root);
1482 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1483 let rescue_id = format!("rescue-{}", generated_session_id());
1484 let imported = session
1485 .imported_message_count
1486 .unwrap_or(session.messages.len())
1487 .min(session.messages.len());
1488 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1489 let view_jsonl = messages_jsonl(&view)?;
1490 let title = format!(
1491 "Reduced {} continuation from {}",
1492 params.target_harness.id(),
1493 params.locator.session_id
1494 );
1495
1496 store
1501 .save_sidecar(&rescue_id, &sidecar_jsonl)
1502 .map_err(operation)?;
1503 store
1504 .save_reduction_log(&rescue_id, &log)
1505 .map_err(operation)?;
1506 store
1507 .save(&rescue_id, &title, &view_jsonl)
1508 .map_err(operation)?;
1509
1510 let source_bytes = serde_json::to_vec(&session.messages)
1511 .map_err(|error| ServiceError::Operation(error.to_string()))?
1512 .len() as u64;
1513 let reduced_bytes = serde_json::to_vec(&view)
1514 .map_err(|error| ServiceError::Operation(error.to_string()))?
1515 .len() as u64;
1516 store
1517 .set_reduction_stats(
1518 &rescue_id,
1519 &title,
1520 source_bytes,
1521 reduced_bytes,
1522 log.reductions.len() as u32,
1523 )
1524 .map_err(operation)?;
1525
1526 let reloaded_sidecar = store
1530 .load_sidecar(&rescue_id)
1531 .map_err(operation)?
1532 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1533 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1534 let reloaded_log = store
1535 .load_reduction_log(&rescue_id)
1536 .map_err(operation)?
1537 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1538 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1539 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1540 let (restamped_view, restamped_log) =
1547 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1548 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1549 return Err(ServiceError::Operation(
1550 "persisted reduction view does not match its durable log and sidecar".into(),
1551 ));
1552 }
1553 if restamped_log != reloaded_log {
1554 return Err(ServiceError::Operation(
1555 "reapplying the durable reduction log changed its identity".into(),
1556 ));
1557 }
1558 let inverted =
1559 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1560 if inverted != session.messages {
1561 return Err(ServiceError::Operation(
1562 "reduction inversion did not restore the source messages byte-exactly".into(),
1563 ));
1564 }
1565
1566 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1567 let sidecar_path = store.sidecar_path(&rescue_id);
1568 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1569 let bootstrap_prompt = reduced_bootstrap_prompt(
1570 ¶ms.locator,
1571 params.target_harness,
1572 &view_jsonl,
1573 &sidecar_path,
1574 &reduction_log_path,
1575 );
1576 let mut reduced_session = session.clone();
1577 reduced_session.meta.session_id = Some(rescue_id.clone());
1578 reduced_session.messages = view;
1579
1580 Ok(json!({
1581 "session": normalized_session_json(&reduced_session),
1582 "bootstrap_prompt": bootstrap_prompt,
1583 "receipt": {
1584 "id": rescue_id,
1585 "sidecar_id": rescue_id,
1586 "source_harness": params.locator.harness,
1587 "target_harness": params.target_harness.id(),
1588 "source_tokens": source_tokens,
1589 "reduced_tokens": reduced_tokens,
1590 "ratio": ratio,
1591 "source_bytes": source_bytes,
1592 "reduced_bytes": reduced_bytes,
1593 "reductions": reloaded_log.reductions.len(),
1594 "sidecar_path": sidecar_path,
1595 "reduction_log_path": reduction_log_path,
1596 "verified": true,
1597 "reversible": true,
1598 }
1599 }))
1600 }
1601
1602 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1620 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1621 return None;
1622 }
1623 let method = request.get("method").and_then(Value::as_str)?;
1624 if !RUNTIME_OPEN_METHODS.contains(&method) {
1625 return None;
1626 }
1627 Some(RuntimeOpen {
1628 id: request.get("id").cloned().unwrap_or(Value::Null),
1629 method: method.to_string(),
1630 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1631 })
1632 }
1633
1634 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1656 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1657 return None;
1658 }
1659 let method = request.get("method").and_then(Value::as_str)?;
1660 if !DETACHED_METHODS.contains(&method) {
1661 return None;
1662 }
1663 let id = request.get("id").cloned().unwrap_or(Value::Null);
1664 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1665 let work = match method {
1666 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1667 .inventory_work(method, params)
1668 .map(DetachedWork::Inventory),
1669 "harness.v1.sessions.message" => {
1670 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1671 }
1672 _ => {
1673 let verb = match method {
1674 "harness.v1.sessions.new" => crate::SessionVerb::New,
1675 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1676 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1677 _ => crate::SessionVerb::Delete,
1678 };
1679 match decode::<crate::SessionMutation>(params) {
1680 Ok(mutation) => {
1681 match crate::sessions_control::door(&mutation.harness, verb) {
1682 Ok(crate::SessionDoor::Live(_)) => return None,
1685 Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1686 Err(error) => Err(session_control_error(error)),
1687 }
1688 }
1689 Err(error) => Err(error),
1690 }
1691 }
1692 };
1693 Some(DetachedCall {
1694 id,
1695 method: method.to_string(),
1696 work: work.map(Work::Free),
1697 })
1698 }
1699
1700 pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1714 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1715 return None;
1716 }
1717 let method = request.get("method").and_then(Value::as_str)?;
1718 let id = request.get("id").cloned().unwrap_or(Value::Null);
1719 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1720 let work = match method {
1721 "harness.v1.runtimes.close" => decode::<RuntimeCloseParams>(params)
1722 .and_then(|params| self.surrender_runtime_of(¶ms))
1723 .map(|(runtime, process_group)| {
1724 Work::Runtime(RuntimeWork::Close {
1725 runtime,
1726 process_group,
1727 })
1728 }),
1729 "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1730 let verb = if method == "harness.v1.sessions.new" {
1731 crate::SessionVerb::New
1732 } else {
1733 crate::SessionVerb::Reset
1734 };
1735 let mutation = decode::<crate::SessionMutation>(params).ok()?;
1736 let Ok(crate::SessionDoor::Live(command)) =
1740 crate::sessions_control::door(&mutation.harness, verb)
1741 else {
1742 return None;
1743 };
1744 let connection = mutation
1745 .connection
1746 .clone()
1747 .filter(|value| !value.trim().is_empty())?;
1748 self.lend_runtime(&connection).map(|runtime| {
1749 let session = live_session_name(runtime.as_ref(), &mutation);
1750 Work::Runtime(RuntimeWork::LiveCommand {
1751 connection,
1752 runtime,
1753 verb,
1754 mutation,
1755 command,
1756 session,
1757 })
1758 })
1759 }
1760 _ => return None,
1761 };
1762 Some(DetachedCall {
1763 id,
1764 method: method.to_string(),
1765 work,
1766 })
1767 }
1768
1769 pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1774 let DetachedAnswer { response, returned } = answer;
1775 if let Some(ReturnedRuntime {
1776 connection,
1777 runtime,
1778 }) = returned
1779 {
1780 self.runtimes_in_flight.remove(&connection);
1781 self.runtimes.insert(connection, runtime);
1782 }
1783 response
1784 }
1785
1786 pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1790 let OpenedRuntime { id, outcome } = opened;
1791 let result = match outcome {
1792 Ok(open) => self.register_open_runtime(open).await,
1793 Err(error) => Err(error),
1794 };
1795 service_response(id, result)
1796 }
1797
1798 async fn register_open_runtime(
1800 &mut self,
1801 open: OpenRuntime,
1802 ) -> std::result::Result<Value, ServiceError> {
1803 match open {
1804 OpenRuntime::Hosted {
1805 runtime,
1806 capabilities,
1807 workspace,
1808 fresh,
1809 } => {
1810 self.insert_hosted_runtime(runtime, capabilities, workspace, fresh)
1811 .await
1812 }
1813 OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1814 }
1815 }
1816
1817 async fn runtime_call(
1818 &mut self,
1819 method: &str,
1820 params: Value,
1821 ) -> std::result::Result<Value, ServiceError> {
1822 match method {
1823 "harness.v1.runtimes.capabilities" => {
1824 let params = decode::<RuntimeBackendParams>(params)?;
1825 #[cfg(feature = "adapter-api")]
1829 if params
1830 .base_url
1831 .as_deref()
1832 .is_some_and(|value| LiveRuntimeEndpoint::parse(value).is_ok())
1833 {
1834 return Ok(json!({
1835 "harness": ¶ms.harness,
1836 "capabilities": SupercodeHttpRuntimeBackend::live_capabilities(),
1837 }));
1838 }
1839 let backend = runtime_backend(¶ms)?;
1840 Ok(json!({
1841 "harness": backend.harness(),
1842 "capabilities": backend.capabilities(),
1843 }))
1844 }
1845 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1846 self.register_open_runtime(open_runtime(method, params).await?)
1847 .await
1848 }
1849 "harness.v1.runtimes.send_input" => {
1850 let params = decode::<RuntimeInputParams>(params)?;
1851 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1852 let runtime = self.runtime_mut(¶ms.connection)?;
1853 let turn_id = within_control_deadline(
1854 method,
1855 runtime.send_input(RuntimeInput {
1856 text: params.text,
1857 image_urls,
1858 }),
1859 )
1860 .await?
1861 .map_err(operation)?;
1862 Ok(json!({"turn_id": turn_id}))
1863 }
1864 "harness.v1.runtimes.interrupt" => {
1865 let params = decode::<RuntimeConnectionParams>(params)?;
1866 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1867 .await?
1868 .map_err(operation)?;
1869 Ok(json!({}))
1870 }
1871 "harness.v1.runtimes.steer" => {
1872 let params = decode::<RuntimeInputParams>(params)?;
1873 if !params.image_urls.is_empty() {
1874 return Err(ServiceError::InvalidParams(
1875 "runtime steering accepts text only".into(),
1876 ));
1877 }
1878 let text = params.text.trim();
1879 if text.is_empty() || text.chars().count() > 50_000 {
1880 return Err(ServiceError::InvalidParams(
1881 "runtime steering requires 1 to 50,000 text characters".into(),
1882 ));
1883 }
1884 within_control_deadline(
1885 method,
1886 self.runtime_mut(¶ms.connection)?
1887 .steer(text.to_string()),
1888 )
1889 .await?
1890 .map_err(operation)?;
1891 Ok(json!({}))
1892 }
1893 "harness.v1.runtimes.respond" => {
1894 let params = decode::<RuntimeRespondParams>(params)?;
1895 let request_id = params.request_id.clone();
1896 within_control_deadline(
1897 method,
1898 self.runtime_mut(¶ms.connection)?
1899 .respond(params.request_id, params.response),
1900 )
1901 .await?
1902 .map_err(operation)?;
1903 self.approvals.answered(¶ms.connection, &request_id);
1905 Ok(json!({}))
1906 }
1907 "harness.v1.runtimes.acquire_control" => {
1908 let params = decode::<RuntimeConnectionParams>(params)?;
1909 let snapshot = within_control_deadline(
1910 method,
1911 self.runtime_mut(¶ms.connection)?.acquire_control(),
1912 )
1913 .await?
1914 .map_err(operation)?;
1915 serde_json::to_value(snapshot)
1916 .map_err(|error| ServiceError::Operation(error.to_string()))
1917 }
1918 "harness.v1.runtimes.heartbeat" => {
1919 let params = decode::<RuntimeConnectionParams>(params)?;
1920 let snapshot = within_control_deadline(
1921 method,
1922 self.runtime_mut(¶ms.connection)?.heartbeat(),
1923 )
1924 .await?
1925 .map_err(operation)?;
1926 serde_json::to_value(snapshot)
1927 .map_err(|error| ServiceError::Operation(error.to_string()))
1928 }
1929 "harness.v1.runtimes.detach" => {
1930 let params = decode::<RuntimeConnectionParams>(params)?;
1931 let snapshot =
1932 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1933 .await?
1934 .map_err(operation)?;
1935 serde_json::to_value(snapshot)
1936 .map_err(|error| ServiceError::Operation(error.to_string()))
1937 }
1938 "harness.v1.runtimes.terminal_instructions" => {
1939 let params = decode::<RuntimeConnectionParams>(params)?;
1940 let launch = self
1941 .terminal_launches
1942 .get(¶ms.connection)
1943 .ok_or_else(|| {
1944 ServiceError::Operation(
1945 "this runtime is not hosted for terminal attachment".into(),
1946 )
1947 })?;
1948 Ok(json!({"launch":launch}))
1949 }
1950 "harness.v1.runtimes.close" => {
1951 let params = decode::<RuntimeCloseParams>(params)?;
1952 let (runtime, process_group) = self.surrender_runtime_of(¶ms)?;
1953 close_runtime(runtime, process_group).await
1954 }
1955 _ => Err(ServiceError::MethodNotFound),
1956 }
1957 }
1958
1959 #[cfg(feature = "adapter-api")]
1961 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1962 let params = decode::<MessageSessionParams>(params)?;
1963 Ok(message_live_session(¶ms).await)
1964 }
1965
1966 #[cfg(feature = "adapter-api")]
1967 fn harness_settings_call(
1968 &self,
1969 method: &str,
1970 params: Value,
1971 ) -> std::result::Result<Value, ServiceError> {
1972 let homes = crate::HarnessHomes::default();
1973 match method {
1974 "harness.v1.harnesses.settings" => {
1975 let params = decode::<HarnessSettingsParams>(params)?;
1976 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1977 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1978 serde_json::to_value(report)
1979 .map_err(|error| ServiceError::Operation(error.to_string()))
1980 }
1981 "harness.v1.harnesses.configure" => {
1982 let params = decode::<ConfigureHarnessParams>(params)?;
1983 let report = crate::configure_harness_interop_settings(
1984 &homes,
1985 ¶ms.harness,
1986 ¶ms.changes,
1987 params.expected_revision.as_deref(),
1988 )
1989 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1990 serde_json::to_value(report)
1991 .map_err(|error| ServiceError::Operation(error.to_string()))
1992 }
1993 _ => Err(ServiceError::MethodNotFound),
1994 }
1995 }
1996
1997 fn insert_runtime(
1998 &mut self,
1999 runtime: Box<dyn RuntimeConnection>,
2000 ) -> std::result::Result<Value, ServiceError> {
2001 let connection = format!("runtime-{}", self.next_runtime);
2002 self.next_runtime += 1;
2003 let handle = runtime.handle().clone();
2004 self.runtime_sequences
2005 .entry(handle.runtime_id.clone())
2006 .or_insert(0);
2007 self.runtimes.insert(connection.clone(), runtime);
2008 Ok(json!({"connection": connection, "handle": handle}))
2009 }
2010
2011 #[cfg(feature = "adapter-api")]
2012 async fn insert_hosted_runtime(
2013 &mut self,
2014 runtime: Box<dyn RuntimeConnection>,
2015 capabilities: crate::RuntimeCapabilities,
2016 workspace: PathBuf,
2017 fresh: bool,
2018 ) -> std::result::Result<Value, ServiceError> {
2019 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
2020 let token: std::sync::Arc<str> = crate::server::generate_token().into();
2021 let server = crate::server::run_frontend_http(
2022 host.clone(),
2023 host.frontend_sender(),
2024 "127.0.0.1:0",
2025 token.clone(),
2026 connection.handle().runtime_id.clone(),
2027 )
2028 .await
2029 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2030 let source = LiveRuntimeSource {
2031 harness: connection.handle().harness.as_str().to_string(),
2032 session_id: connection.handle().runtime_id.clone(),
2033 workspace: workspace.clone(),
2034 };
2035 let registration = register_live_runtime_with_metadata(
2036 connection.handle().runtime_id.clone(),
2037 source.clone(),
2038 format!("http://{}", server.address()),
2039 token.to_string(),
2040 crate::LiveRuntimeMetadata {
2041 child_pid: match &connection.handle().endpoint {
2042 crate::RuntimeEndpoint::LocalProcess { pid, command, .. }
2043 if connection.handle().harness.as_str() != "claude-code"
2044 || command.windows(2).any(|args| {
2045 args[0] == "--name"
2046 && args[1].starts_with(crate::claude_relay::RELAY_NAME_PREFIX)
2047 }) =>
2048 {
2049 *pid
2050 }
2051 _ => None,
2052 },
2053 endpoint_capabilities: vec!["http".into(), "acp".into()],
2054 ..Default::default()
2055 },
2056 )
2057 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2058 let endpoint = registration.endpoint().to_string();
2059 let launch = StructuredLaunch {
2060 cwd: workspace,
2061 program: std::env::current_exe()
2065 .ok()
2066 .map(|path| path.to_string_lossy().into_owned())
2067 .unwrap_or_else(|| "supercode".into()),
2068 arguments: vec![
2069 "open".into(),
2070 endpoint,
2071 "--harness".into(),
2072 source.harness,
2073 "--session".into(),
2074 source.session_id,
2075 ],
2076 env: BTreeMap::new(),
2077 };
2078 host.retain_registration(registration);
2079 let lease = HostedRuntimeLease {
2080 connection,
2081 _host: host,
2082 _server: server,
2083 };
2084 let opened = self.insert_runtime(Box::new(lease))?;
2085 let connection_id = opened["connection"]
2086 .as_str()
2087 .expect("insert_runtime returns a connection id")
2088 .to_string();
2089 self.terminal_launches.insert(connection_id, launch);
2090 Ok(opened)
2091 }
2092
2093 #[cfg(not(feature = "adapter-api"))]
2094 async fn insert_hosted_runtime(
2095 &mut self,
2096 runtime: Box<dyn RuntimeConnection>,
2097 _capabilities: crate::RuntimeCapabilities,
2098 _workspace: PathBuf,
2099 _fresh: bool,
2100 ) -> std::result::Result<Value, ServiceError> {
2101 self.insert_runtime(runtime)
2102 }
2103
2104 fn runtime_mut(
2105 &mut self,
2106 connection: &str,
2107 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2108 if self.runtimes_in_flight.contains(connection) {
2109 return Err(self.lent_out(connection));
2110 }
2111 self.runtimes.get_mut(connection).ok_or_else(|| {
2112 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2113 })
2114 }
2115
2116 fn lent_out(&self, connection: &str) -> ServiceError {
2120 ServiceError::Operation(format!(
2121 "runtime connection `{connection}`: a harness turn is already in progress"
2122 ))
2123 }
2124
2125 fn lend_runtime(
2128 &mut self,
2129 connection: &str,
2130 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2131 if self.runtimes_in_flight.contains(connection) {
2132 return Err(self.lent_out(connection));
2133 }
2134 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2135 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2136 })?;
2137 self.runtimes_in_flight.insert(connection.to_string());
2138 Ok(runtime)
2139 }
2140
2141 fn surrender_runtime_of(
2155 &mut self,
2156 params: &RuntimeCloseParams,
2157 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2158 if params.connection.is_empty() {
2159 let expected = params.runtime_id.as_deref().ok_or_else(|| {
2160 ServiceError::InvalidParams("close requires connection or runtime_id".into())
2161 })?;
2162 let connection = self
2163 .runtimes
2164 .iter()
2165 .find(|(_, runtime)| runtime.handle().runtime_id == expected)
2166 .map(|(connection, _)| connection.clone())
2167 .ok_or_else(|| {
2168 ServiceError::InvalidParams(format!("unknown runtime `{expected}`"))
2169 })?;
2170 return self.surrender_runtime(&connection);
2171 }
2172 if let Some(expected) = params.runtime_id.as_deref() {
2173 let held = self
2174 .runtimes
2175 .get(¶ms.connection)
2176 .map(|runtime| runtime.handle().runtime_id.clone());
2177 if held.as_deref() != Some(expected) {
2178 return Err(ServiceError::InvalidParams(format!(
2179 "runtime connection `{}` does not hold runtime `{expected}`",
2180 params.connection
2181 )));
2182 }
2183 }
2184 self.surrender_runtime(¶ms.connection)
2185 }
2186
2187 fn surrender_runtime(
2188 &mut self,
2189 connection: &str,
2190 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2191 if self.runtimes_in_flight.contains(connection) {
2192 return Err(self.lent_out(connection));
2193 }
2194 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2195 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2196 })?;
2197 let process_group = runtime_process_group(runtime.handle());
2198 let runtime_id = runtime.handle().runtime_id.clone();
2199 self.terminal_launches.remove(connection);
2200 self.runtime_sequences.remove(&runtime_id);
2201 self.approvals.forget(connection);
2202 Ok((runtime, process_group))
2203 }
2204
2205 pub fn kill_all_runtime_groups(&self) -> usize {
2216 self.runtimes
2217 .values()
2218 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2219 .count()
2220 }
2221
2222 async fn mutate_session(
2233 &mut self,
2234 verb: crate::SessionVerb,
2235 params: Value,
2236 ) -> std::result::Result<Value, ServiceError> {
2237 let mutation = decode::<crate::SessionMutation>(params)?;
2238 let door = crate::sessions_control::door(&mutation.harness, verb)
2239 .map_err(session_control_error)?;
2240 let outcome = match door {
2241 #[cfg(not(feature = "adapter-api"))]
2245 crate::SessionDoor::Live(command) => {
2246 return Err(ServiceError::Operation(format!(
2247 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2248 session, which needs this build's `adapter-api` feature",
2249 mutation.harness,
2250 verb.as_str()
2251 )));
2252 }
2253 #[cfg(feature = "adapter-api")]
2254 crate::SessionDoor::Live(command) => {
2255 let connection = mutation
2256 .connection
2257 .clone()
2258 .filter(|value| !value.trim().is_empty())
2259 .ok_or_else(|| {
2260 ServiceError::InvalidParams(format!(
2261 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2262 driven session: pass the `connection` of an open runtime \
2263 (`harness.v1.runtimes.start`)",
2264 mutation.harness,
2265 verb.as_str()
2266 ))
2267 })?;
2268 let runtime = self.runtime_mut(&connection)?;
2269 let session = live_session_name(runtime.as_ref(), &mutation);
2270 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2276 .await;
2277 }
2278 _ => run_session_mutation(verb, &mutation).await?,
2279 };
2280 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2281 }
2282
2283 async fn inventory_call(
2287 &self,
2288 method: &str,
2289 params: Value,
2290 ) -> std::result::Result<Value, ServiceError> {
2291 run_inventory(self.inventory_work(method, params)?).await
2292 }
2293
2294 fn inventory_work(
2300 &self,
2301 method: &str,
2302 params: Value,
2303 ) -> std::result::Result<InventoryWork, ServiceError> {
2304 let mut params = decode::<HarnessInventoryParams>(params)?;
2305 if method == "harness.v1.harnesses.probe" {
2306 let harness = params.harness.take().ok_or_else(|| {
2307 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2308 })?;
2309 params.harnesses = vec![harness];
2310 }
2311 let selected = params
2312 .harnesses
2313 .iter()
2314 .map(HarnessId::as_str)
2315 .collect::<std::collections::BTreeSet<_>>();
2316 let supported = harness_support_registry()
2317 .harnesses
2318 .into_iter()
2319 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2320 .collect::<Vec<_>>();
2321 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2322 let known = supported
2323 .iter()
2324 .map(|harness| harness.id.as_str())
2325 .collect::<std::collections::BTreeSet<_>>();
2326 let missing = params
2327 .harnesses
2328 .iter()
2329 .filter(|id| !known.contains(id.as_str()))
2330 .map(HarnessId::as_str)
2331 .collect::<Vec<_>>();
2332 return Err(ServiceError::InvalidParams(format!(
2333 "unknown harness(es): {}",
2334 missing.join(", ")
2335 )));
2336 }
2337 let global_counts = params
2338 .include_sessions
2339 .then(|| self.session_counts(None, ¶ms.harnesses));
2340 let workspace_counts = params
2341 .include_sessions
2342 .then(|| {
2343 params
2344 .workspace
2345 .as_deref()
2346 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2347 })
2348 .flatten();
2349 Ok(InventoryWork {
2350 params,
2351 supported,
2352 global_counts,
2353 workspace_counts,
2354 })
2355 }
2356
2357 #[cfg(feature = "adapter-api")]
2358 async fn harness_authentication_call(
2359 &self,
2360 method: &str,
2361 params: Value,
2362 ) -> std::result::Result<Value, ServiceError> {
2363 match method {
2364 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2365 let params = decode::<HarnessAuthenticationParams>(params)?;
2366 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2367 .map_err(|error| ServiceError::Operation(error.to_string()))
2368 }
2369 "harness.v1.harnesses.auth.begin" => {
2370 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2371 let cwd = params
2372 .cwd
2373 .or_else(|| std::env::current_dir().ok())
2374 .unwrap_or_else(|| PathBuf::from("."));
2375 let plan = crate::harness_authentication_plan(
2376 ¶ms.harness,
2377 params.environment,
2378 params.method,
2379 &cwd,
2380 )
2381 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2382 serde_json::to_value(plan)
2383 .map_err(|error| ServiceError::Operation(error.to_string()))
2384 }
2385 _ => Err(ServiceError::MethodNotFound),
2386 }
2387 }
2388
2389 fn session_counts(
2390 &self,
2391 workspace: Option<&Path>,
2392 harnesses: &[HarnessId],
2393 ) -> BTreeMap<String, usize> {
2394 let mut counts = BTreeMap::new();
2395 for session in self
2396 .catalog
2397 .discover(&DiscoveryQuery {
2398 workspace: workspace.map(Path::to_path_buf),
2399 harnesses: harnesses.to_vec(),
2400 ..DiscoveryQuery::default()
2401 })
2402 .unwrap_or_default()
2403 {
2404 *counts
2405 .entry(session.locator.harness.as_str().to_string())
2406 .or_insert(0) += 1;
2407 }
2408 counts
2409 }
2410}
2411
2412#[async_trait::async_trait]
2413impl SdkService for HarnessSessionService {
2414 fn capabilities(&self) -> SdkCapabilities {
2415 SdkCapabilities::default()
2416 }
2417
2418 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2419 if request.operation == SdkOperation::Events {
2420 let events = self
2421 .poll_sdk_events()
2422 .await
2423 .into_iter()
2424 .map(|(_, event)| event)
2425 .collect::<Vec<_>>();
2426 return serde_json::to_value(events).map_err(|error| {
2427 SdkError::new(
2428 SdkErrorCode::Execution,
2429 request.operation,
2430 error.to_string(),
2431 )
2432 });
2433 }
2434 if self.runtimes.is_empty()
2435 && matches!(
2436 request.operation,
2437 SdkOperation::Input
2438 | SdkOperation::Interrupt
2439 | SdkOperation::Steer
2440 | SdkOperation::Respond
2441 | SdkOperation::Close
2442 )
2443 {
2444 return Err(SdkError::unsupported(request.operation));
2445 }
2446 let method = request
2447 .operation
2448 .method()
2449 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2450 let result = match request.operation {
2451 SdkOperation::Discover
2452 | SdkOperation::Load
2453 | SdkOperation::Export
2454 | SdkOperation::ProfilesList
2455 | SdkOperation::ProfilesGet
2456 | SdkOperation::ProfilesCreate
2457 | SdkOperation::ProfilesDelete
2458 | SdkOperation::SkillsList
2459 | SdkOperation::SkillsInstall
2460 | SdkOperation::SkillsRemove
2461 | SdkOperation::ChannelsList
2462 | SdkOperation::RoutesList
2463 | SdkOperation::TriggersList
2464 | SdkOperation::ChannelsStatus
2465 | SdkOperation::MemoryShow
2466 | SdkOperation::MemorySearch
2467 | SdkOperation::JobsList
2468 | SdkOperation::JobsGet
2469 | SdkOperation::JobsCreate
2470 | SdkOperation::JobsUpdate
2471 | SdkOperation::JobsPause
2472 | SdkOperation::JobsResume
2473 | SdkOperation::JobsRun
2474 | SdkOperation::JobsDelete
2475 | SdkOperation::JobsNotepad
2476 | SdkOperation::JobsNotepadSet
2477 | SdkOperation::JobsNotepadDelete
2478 | SdkOperation::RunsList
2479 | SdkOperation::RunsGet
2480 | SdkOperation::ApprovalsList
2481 | SdkOperation::OrchestrationLoad
2482 | SdkOperation::OrchestrationSave
2483 | SdkOperation::OrchestrationCompile
2484 | SdkOperation::OrchestrationDecompile
2485 | SdkOperation::OrchestrationImport
2486 | SdkOperation::OrchestrationExport
2487 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2488 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2491 SdkOperation::Start
2492 | SdkOperation::Resume
2493 | SdkOperation::Input
2494 | SdkOperation::Interrupt
2495 | SdkOperation::Steer
2496 | SdkOperation::Respond
2497 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2498 SdkOperation::SessionsNew => {
2503 self.mutate_session(crate::SessionVerb::New, request.params)
2504 .await
2505 }
2506 SdkOperation::SessionsReset => {
2507 self.mutate_session(crate::SessionVerb::Reset, request.params)
2508 .await
2509 }
2510 SdkOperation::SessionsArchive => {
2511 self.mutate_session(crate::SessionVerb::Archive, request.params)
2512 .await
2513 }
2514 SdkOperation::SessionsDelete => {
2515 self.mutate_session(crate::SessionVerb::Delete, request.params)
2516 .await
2517 }
2518 SdkOperation::Events => unreachable!("handled before method dispatch"),
2519 };
2520 result.map_err(|error| sdk_error(request.operation, error))
2521 }
2522
2523 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2524 Ok(self
2525 .poll_sdk_events()
2526 .await
2527 .into_iter()
2528 .map(|(_, event)| event)
2529 .collect())
2530 }
2531}
2532
2533#[cfg(feature = "adapter-api")]
2534struct HostedRuntimeLease {
2535 connection: HostedHarnessConnection,
2536 _host: std::sync::Arc<HostedHarnessRuntime>,
2537 _server: crate::server::FrontendHttpServer,
2538}
2539
2540#[async_trait::async_trait]
2541#[cfg(feature = "adapter-api")]
2542impl RuntimeConnection for HostedRuntimeLease {
2543 fn handle(&self) -> &crate::RuntimeHandle {
2544 self.connection.handle()
2545 }
2546
2547 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2548 self.connection.send_input(input).await
2549 }
2550
2551 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2552 self.connection.next_event().await
2553 }
2554
2555 async fn interrupt(&mut self) -> crate::Result<()> {
2556 self.connection.interrupt().await
2557 }
2558
2559 async fn steer(&mut self, text: String) -> crate::Result<()> {
2562 self.connection.steer(text).await
2563 }
2564
2565 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2566 self.connection.respond(request_id, response).await
2567 }
2568
2569 async fn close(&mut self) -> crate::Result<()> {
2570 self.connection.close().await
2571 }
2572}
2573
2574struct InventoryWork {
2577 params: HarnessInventoryParams,
2578 supported: Vec<crate::HarnessSupportDescriptor>,
2579 global_counts: Option<BTreeMap<String, usize>>,
2580 workspace_counts: Option<BTreeMap<String, usize>>,
2581}
2582
2583async fn run_session_mutation(
2590 verb: crate::SessionVerb,
2591 mutation: &crate::SessionMutation,
2592) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2593 let door =
2600 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2601 if let crate::SessionDoor::Http = door {
2602 return crate::sessions_control::mutate(verb, mutation)
2603 .await
2604 .map_err(session_control_error);
2605 }
2606 let mutation = mutation.clone();
2607 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2608 .await
2609 .map_err(|error| {
2610 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2611 })?
2612 .map_err(session_control_error)
2613}
2614
2615async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2618 let InventoryWork {
2619 params,
2620 supported,
2621 global_counts,
2622 workspace_counts,
2623 } = work;
2624 let probes = supported.into_iter().map(|descriptor| {
2625 let global = global_counts
2626 .as_ref()
2627 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2628 let workspace = workspace_counts
2629 .as_ref()
2630 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2631 probe_harness(descriptor, ¶ms, global, workspace)
2632 });
2633 let harnesses = futures::future::join_all(probes).await;
2634 serde_json::to_value(HarnessInventoryReport {
2635 probe: params.probe,
2636 workspace: params.workspace,
2637 harnesses,
2638 })
2639 .map_err(|error| ServiceError::Operation(error.to_string()))
2640}
2641
2642async fn probe_harness(
2643 descriptor: crate::HarnessSupportDescriptor,
2644 params: &HarnessInventoryParams,
2645 global: Option<usize>,
2646 workspace: Option<usize>,
2647) -> LocalHarness {
2648 let launch = descriptor.runtime.default_launch.as_ref();
2649 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2654 .then(crate::orchestrator::daemon_entry)
2655 .and_then(Result::ok);
2656 let executable = match &orchestrator_entry {
2657 Some(entry) => Some(entry.clone()),
2658 None => launch.and_then(|launch| find_executable(&launch.program)),
2659 };
2660 let installed = executable.is_some();
2661 let version = if params.skip_versions || orchestrator_entry.is_some() {
2662 None
2665 } else {
2666 match executable.as_deref() {
2667 Some(path) => executable_version(path).await,
2668 None => None,
2669 }
2670 };
2671 let configured = auth_evidence(descriptor.id.as_str());
2672 let mut auth = if configured {
2673 HarnessAuthState::Configured
2674 } else if matches!(
2675 descriptor.id.as_str(),
2676 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2677 ) {
2678 HarnessAuthState::Required
2683 } else {
2684 HarnessAuthState::Unknown
2685 };
2686 let mut runtime = if installed {
2687 HarnessRuntimeState::Degraded
2688 } else {
2689 HarnessRuntimeState::Unavailable
2690 };
2691 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2692 let mut reason = (!installed).then(|| {
2693 if is_orchestrator {
2694 format!(
2695 "{} is supported but its daemon entry `{}` was not found",
2696 descriptor.display_name,
2697 crate::orchestrator::DAEMON_ENTRY
2698 )
2699 } else {
2700 format!(
2701 "{} is supported but `{}` was not found on PATH",
2702 descriptor.display_name,
2703 launch
2704 .map(|launch| launch.program.as_str())
2705 .unwrap_or("executable")
2706 )
2707 }
2708 });
2709 let mut repair = (!installed).then(|| {
2710 if is_orchestrator {
2711 format!(
2712 "Install the `supercode-orchestrator` package so `{}` resolves.",
2713 crate::orchestrator::DAEMON_ENTRY
2714 )
2715 } else {
2716 format!(
2717 "Install {} and ensure `{}` is on PATH.",
2718 descriptor.display_name,
2719 launch
2720 .map(|launch| launch.program.as_str())
2721 .unwrap_or("its executable")
2722 )
2723 }
2724 });
2725
2726 if installed && params.probe == HarnessProbeLevel::Handshake {
2727 let backend_params = RuntimeBackendParams {
2728 harness: descriptor.id.clone(),
2729 protocol: None,
2730 launch: None,
2731 base_url: None,
2732 policy: RuntimePolicy::Default,
2733 };
2734 match runtime_backend(&backend_params) {
2735 Ok(backend) => {
2736 let cwd = params
2737 .workspace
2738 .clone()
2739 .or_else(|| std::env::current_dir().ok())
2740 .unwrap_or_else(|| PathBuf::from("."));
2741 let isolated = descriptor
2742 .runtime
2743 .default_launch
2744 .clone()
2745 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2746 let Some(isolated) = isolated else {
2747 reason = Some(
2748 "No-prompt runtime handshake could not create its isolated harness home."
2749 .into(),
2750 );
2751 repair = Some(
2752 "Check temporary-directory permissions, then run the handshake probe again."
2753 .into(),
2754 );
2755 let running = probe_running_instance(descriptor.id.as_str());
2756 return LocalHarness {
2757 gateway: gateway_health(
2758 descriptor.id.as_str(),
2759 installed,
2760 running.as_ref(),
2761 version.as_deref(),
2762 ),
2763 id: descriptor.id,
2764 display_name: descriptor.display_name,
2765 supported: true,
2766 installed,
2767 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2768 version,
2769 auth,
2770 runtime,
2771 protocol: descriptor.runtime.protocol,
2772 capabilities: descriptor.runtime.capabilities.clone(),
2773 effective_capabilities: descriptor.runtime.capabilities,
2774 sessions: HarnessSessionCounts { global, workspace },
2775 running,
2776 reason,
2777 repair,
2778 };
2779 };
2780 match tokio::time::timeout(
2781 Duration::from_secs(30),
2782 backend.start(RuntimeStartRequest {
2783 cwd,
2784 launch: Some(isolated.launch.clone()),
2785 mcp_servers: Vec::new(),
2786 approval_policy: None,
2787 }),
2788 )
2789 .await
2790 {
2791 Ok(Ok(mut connection)) => {
2792 match stabilize_handshake(connection.as_mut()).await {
2793 Ok(()) => {
2794 auth = HarnessAuthState::Ready;
2795 runtime = HarnessRuntimeState::Ready;
2796 reason = Some(
2797 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2798 .into(),
2799 );
2800 repair = None;
2801 }
2802 Err(message) => {
2803 auth = if looks_like_auth_error(&message) {
2804 HarnessAuthState::Required
2805 } else if configured {
2806 HarnessAuthState::Configured
2807 } else {
2808 HarnessAuthState::Unknown
2809 };
2810 reason = Some(format!(
2811 "No-prompt runtime handshake became unhealthy during startup: {message}"
2812 ));
2813 repair = Some(if auth == HarnessAuthState::Required {
2814 format!(
2815 "Run `{}` interactively once and complete sign-in, then probe again.",
2816 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2817 )
2818 } else {
2819 "Run the harness directly to inspect its startup failure, then probe again."
2820 .into()
2821 });
2822 }
2823 }
2824 let _ =
2825 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2826 }
2827 Ok(Err(error)) => {
2828 let message = truncate_text(&error.to_string(), 500);
2829 auth = if looks_like_auth_error(&message) {
2830 HarnessAuthState::Required
2831 } else if configured {
2832 HarnessAuthState::Configured
2833 } else {
2834 HarnessAuthState::Unknown
2835 };
2836 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2837 repair = Some(if auth == HarnessAuthState::Required {
2838 format!(
2839 "Run `{}` interactively once and complete sign-in, then probe again.",
2840 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2841 )
2842 } else {
2843 "Check the harness installation and run the handshake probe again."
2844 .into()
2845 });
2846 }
2847 Err(_) => {
2848 reason =
2849 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2850 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2851 }
2852 }
2853 let _ = isolated.cleanup();
2862 tokio::time::sleep(Duration::from_millis(250)).await;
2863 if let Err(error) = isolated.cleanup() {
2864 auth = if configured {
2865 HarnessAuthState::Configured
2866 } else {
2867 HarnessAuthState::Unknown
2868 };
2869 runtime = HarnessRuntimeState::Degraded;
2870 reason = Some(format!(
2871 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2872 ));
2873 repair = Some(
2874 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2875 .into(),
2876 );
2877 }
2878 }
2879 Err(error) => {
2880 reason = Some(error_message(error));
2881 }
2882 }
2883 } else if installed && configured {
2884 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2885 } else if installed && auth == HarnessAuthState::Required {
2886 reason = Some("Executable found, but no native authentication evidence is present.".into());
2887 repair = Some(format!(
2888 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2889 descriptor.id.as_str()
2890 ));
2891 } else if installed {
2892 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2893 repair = Some(format!(
2894 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2895 launch
2896 .map(|launch| launch.program.as_str())
2897 .unwrap_or("the harness")
2898 ));
2899 }
2900
2901 let effective_capabilities = if installed {
2902 descriptor.runtime.capabilities.clone()
2903 } else {
2904 unavailable_capabilities()
2905 };
2906 let running = probe_running_instance(descriptor.id.as_str());
2907 LocalHarness {
2908 gateway: gateway_health(
2909 descriptor.id.as_str(),
2910 installed,
2911 running.as_ref(),
2912 version.as_deref(),
2913 ),
2914 id: descriptor.id,
2915 display_name: descriptor.display_name,
2916 supported: true,
2917 installed,
2918 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2919 version,
2920 auth,
2921 runtime,
2922 protocol: descriptor.runtime.protocol,
2923 capabilities: descriptor.runtime.capabilities,
2924 effective_capabilities,
2925 sessions: HarnessSessionCounts { global, workspace },
2926 running,
2927 reason,
2928 repair,
2929 }
2930}
2931
2932async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2933 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2934 loop {
2935 let now = tokio::time::Instant::now();
2936 if now >= deadline {
2937 return Ok(());
2938 }
2939 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2940 Err(_) => return Ok(()),
2941 Ok(Ok(Some(event))) => {
2942 if let Some(message) = handshake_event_failure(&event) {
2943 return Err(truncate_text(&message, 500));
2944 }
2945 }
2946 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2947 Ok(Err(error)) => return Err(error.to_string()),
2948 }
2949 }
2950}
2951
2952fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2953 let detail = event
2954 .payload
2955 .get("message")
2956 .or_else(|| event.payload.get("line"))
2957 .and_then(Value::as_str)
2958 .unwrap_or(event.kind.as_str());
2959 match event.kind.as_str() {
2960 "transport_closed" => Some("runtime transport closed during startup".into()),
2961 "transport_error" => Some(format!("runtime transport error: {detail}")),
2962 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2963 _ => None,
2968 }
2969}
2970
2971fn indexed_claude_window(
2972 locator: &SessionLocator,
2973 options: &SessionLoadOptions,
2974) -> std::result::Result<Option<Value>, ServiceError> {
2975 use supercode_interchange::session::ClaudeReadIndex;
2976 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2979 || options.include_subagents != Some(false)
2980 {
2981 return Ok(None);
2982 }
2983 let crate::StorageLocator::File { path } = &locator.storage else {
2984 return Ok(None);
2985 };
2986 if !ClaudeReadIndex::supports(path)
2987 .map_err(|error| ServiceError::Operation(error.to_string()))?
2988 {
2989 return Ok(None);
2990 }
2991 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2992 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2993 let total = index.len();
2994 let (offset, end) = projected_message_window(total, options);
2995 let session = index
2996 .read_messages(offset..end)
2997 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2998 let summary = index
2999 .read_summary()
3000 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3001 let selected_options = SessionLoadOptions {
3002 message_offset: None,
3003 message_limit: None,
3004 message_tail: None,
3005 ..options.clone()
3006 };
3007 let mut selected = projected_session_json(&session, &selected_options);
3008 selected["raw_record_count"] = json!(index.raw_record_count());
3009 Ok(Some(json!({
3010 "session": selected,
3011 "summary": projected_session_summary(&summary, options),
3012 "window": {
3013 "has_more": offset > 0 || end < total, "has_newer": end < total,
3014 "has_older": offset > 0, "newer_items": index.item_count(end..total),
3015 "offset": offset, "older_items": index.item_count(0..offset),
3016 "returned": end - offset, "total_messages": total,
3017 }
3018 })))
3019}
3020
3021fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
3022 let total_messages = session.messages.len();
3023 let (offset, end) = projected_message_window(total_messages, options);
3024 json!({
3025 "session": projected_session_json(session, options),
3026 "summary": projected_session_summary(session, options),
3027 "window": {
3028 "has_more": offset > 0 || end < total_messages,
3029 "has_newer": end < total_messages,
3030 "has_older": offset > 0,
3031 "newer_items": normalized_item_count(&session.messages[end..]),
3032 "offset": offset,
3033 "older_items": normalized_item_count(&session.messages[..offset]),
3034 "returned": end.saturating_sub(offset),
3035 "total_messages": total_messages,
3036 }
3037 })
3038}
3039
3040fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
3041 messages
3042 .iter()
3043 .map(|message| {
3044 let conversation = usize::from(
3045 matches!(message.role, Role::Assistant | Role::User)
3046 && message_has_content(message),
3047 );
3048 let tool_result =
3049 usize::from(message.role == Role::Tool && message_has_content(message));
3050 conversation + tool_result + message.tool_calls().len()
3051 })
3052 .sum()
3053}
3054
3055fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
3056 let mut conversational = session.messages.iter().filter(|message| {
3057 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
3058 });
3059 let first_message = conversational.clone().next();
3060 let last_message = conversational.next_back();
3061 let mut assistant = session
3062 .messages
3063 .iter()
3064 .filter(|message| message.role == Role::Assistant && message_has_content(message));
3065 let first_assistant_message = assistant.clone().next();
3066 let last_assistant_message = assistant.next_back();
3067 let end_of_turn = session
3068 .messages
3069 .iter()
3070 .rev()
3071 .find(|message| message.role != Role::System)
3072 .is_some_and(|message| {
3073 message.role == Role::Assistant
3074 && message_has_content(message)
3075 && message.tool_calls().is_empty()
3076 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
3078 });
3079 let project = |message: Option<&crate::ChatMessage>| {
3080 message.map(|message| project_inline_media(message_json(message), options))
3081 };
3082 json!({
3083 "end_of_turn": end_of_turn,
3084 "first_assistant_message": project(first_assistant_message),
3085 "first_message": project(first_message),
3086 "last_assistant_message": project(last_assistant_message),
3087 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3088 "last_message": project(last_message),
3089 })
3090}
3091
3092fn message_has_content(message: &crate::ChatMessage) -> bool {
3093 message
3094 .content
3095 .as_deref()
3096 .is_some_and(|content| !content.trim().is_empty())
3097 || message
3098 .content_parts
3099 .as_ref()
3100 .is_some_and(|parts| !parts.is_empty())
3101}
3102
3103fn message_text(message: &crate::ChatMessage) -> String {
3104 if let Some(content) = &message.content {
3105 return content.clone();
3106 }
3107 message
3108 .content_parts
3109 .as_ref()
3110 .into_iter()
3111 .flatten()
3112 .filter_map(|part| part.get("text").and_then(Value::as_str))
3113 .collect::<Vec<_>>()
3114 .join("\n")
3115}
3116
3117fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3118 let (offset, end) = projected_message_window(session.messages.len(), options);
3119 let messages = session.messages[offset..end]
3120 .iter()
3121 .map(|message| project_inline_media(message_json(message), options))
3122 .collect::<Vec<_>>();
3123 let subagents = if options.include_subagents.unwrap_or(true) {
3124 let subagent_options = SessionLoadOptions {
3129 message_limit: None,
3130 message_offset: None,
3131 message_tail: None,
3132 ..options.clone()
3133 };
3134 session
3135 .subagents
3136 .iter()
3137 .map(|subagent| projected_session_json(subagent, &subagent_options))
3138 .collect::<Vec<_>>()
3139 } else {
3140 Vec::new()
3141 };
3142 json!({
3143 "source": match session.meta.source {
3144 SessionSource::ClaudeCode => "claude_code",
3145 SessionSource::Codex => "codex",
3146 SessionSource::Gemini => "gemini",
3147 SessionSource::Goose => "goose",
3148 SessionSource::Grok => "grok",
3149 SessionSource::Native => "native",
3150 SessionSource::OpenClaw => "openclaw",
3151 SessionSource::Hermes => "hermes",
3152 SessionSource::OpenCode => "opencode",
3153 SessionSource::Pi => "pi",
3154 },
3155 "session_id": session.meta.session_id,
3156 "ended_at": session.meta.ended_at,
3157 "end_reason": session.meta.end_reason,
3158 "model": session.meta.model,
3159 "resolved_profile": recorded_session_profile(session),
3160 "resolved_profile_history": recorded_session_profile_history(session),
3161 "cwd": session.meta.cwd,
3162 "system_prompt": session.meta.system_prompt,
3163 "agent_id": session.meta.agent_id,
3164 "parent_tool_use_id": session.meta.parent_tool_use_id,
3165 "lineage": session.meta.lineage,
3166 "messages": messages,
3167 "subagents": subagents,
3168 "raw_record_count": session.raw.len(),
3169 "parse_error_lines": session.parse_error_lines,
3170 })
3171}
3172
3173pub fn recorded_config(locator: &SessionLocator) -> crate::Result<Value> {
3181 let session = load_session(locator)?;
3182 let mut config = json!({"harness": locator.harness, "session_id": session.meta.session_id,
3183 "cwd": session.meta.cwd, "model": session.meta.model,
3184 "approval_policy": null, "sandbox_policy": null, "permission_mode": null});
3185 for header in &session.meta.codex_headers {
3186 let payload = header.get("payload").unwrap_or(header);
3187 for key in ["cwd", "model", "approval_policy", "sandbox_policy"] {
3188 if let Some(value) = payload.get(key) {
3189 config[key] = value.clone();
3190 }
3191 }
3192 }
3193 if let Ok(manifest) = crate::claude_runtime_state::ClaudeRuntimeManifest::from_session(&session)
3195 {
3196 config["permission_mode"] = json!(manifest.posture.permission_mode);
3197 }
3198 Ok(config)
3199}
3200
3201fn recorded_session_usage(session: &Session) -> Vec<Value> {
3202 let mut days: std::collections::BTreeMap<(String, String), [u64; 4]> =
3203 std::collections::BTreeMap::new();
3204 let mut add =
3205 |at: Option<&str>, model: &str, input: u64, output: u64, read: u64, write: u64| {
3206 let day = at
3207 .and_then(|t| t.get(..10))
3208 .unwrap_or("unknown")
3209 .to_string();
3210 let entry = days.entry((day, model.to_string())).or_insert([0; 4]);
3211 entry[0] += input;
3212 entry[1] += output;
3213 entry[2] += read;
3214 entry[3] += write;
3215 };
3216 let records = session
3217 .raw
3218 .iter()
3219 .filter_map(|line| serde_json::from_str::<Value>(line).ok());
3220 let n = |v: Option<&Value>| v.and_then(Value::as_u64).unwrap_or(0);
3221 match session.meta.source {
3222 SessionSource::ClaudeCode => {
3223 let mut responses: std::collections::BTreeMap<String, Value> =
3224 std::collections::BTreeMap::new();
3225 for record in records {
3226 let Some(message) = record.get("message") else {
3227 continue;
3228 };
3229 let (Some(id), Some(_)) = (
3230 message.get("id").and_then(Value::as_str),
3231 message.get("usage"),
3232 ) else {
3233 continue;
3234 };
3235 if message.get("model").and_then(Value::as_str) == Some("<synthetic>") {
3236 continue;
3237 }
3238 responses.insert(id.to_string(), record.clone());
3239 }
3240 for record in responses.values() {
3241 let message = &record["message"];
3242 let usage = &message["usage"];
3243 add(
3244 record.get("timestamp").and_then(Value::as_str),
3245 message
3246 .get("model")
3247 .and_then(Value::as_str)
3248 .unwrap_or("unknown"),
3249 n(usage.get("input_tokens")),
3250 n(usage.get("output_tokens")),
3251 n(usage.get("cache_read_input_tokens")),
3252 n(usage.get("cache_creation_input_tokens")),
3253 );
3254 }
3255 }
3256 SessionSource::Codex => {
3257 let mut model = session
3258 .meta
3259 .model
3260 .clone()
3261 .unwrap_or_else(|| "unknown".into());
3262 let mut previous = [0u64; 3];
3263 for record in records {
3264 let payload = record.get("payload").unwrap_or(&Value::Null);
3265 if record.get("type").and_then(Value::as_str) == Some("turn_context") {
3266 if let Some(m) = payload.get("model").and_then(Value::as_str) {
3267 model = m.to_string();
3268 }
3269 continue;
3270 }
3271 if payload.get("type").and_then(Value::as_str) != Some("token_count") {
3272 continue;
3273 }
3274 let Some(total) = payload.get("info").and_then(|i| i.get("total_token_usage"))
3275 else {
3276 continue;
3277 };
3278 let now = [
3279 n(total.get("input_tokens")),
3280 n(total.get("output_tokens")),
3281 n(total.get("cached_input_tokens")),
3282 ];
3283 if now == previous {
3284 continue;
3285 }
3286 add(
3287 record.get("timestamp").and_then(Value::as_str),
3288 &model,
3289 now[0].saturating_sub(previous[0]),
3290 now[1].saturating_sub(previous[1]),
3291 now[2].saturating_sub(previous[2]),
3292 0,
3293 );
3294 previous = now;
3295 }
3296 }
3297 _ => {}
3298 }
3299 days.into_iter()
3300 .map(|((day, model), [input, output, read, write])| {
3301 let mut usage = json!({"at": format!("{day}T00:00:00Z"), "model": model,
3302 "input_tokens": input, "output_tokens": output,
3303 "cache_read_tokens": read, "cache_write_tokens": write});
3304 if let Some(price) = crate::pricing::built_in(&model) {
3307 usage["cost"] = json!({"amount": price.cost_usd(input + read + write, output),
3308 "currency": "USD", "source": "rate_table"});
3309 }
3310 usage
3311 })
3312 .collect()
3313}
3314
3315fn recorded_session_profile(session: &Session) -> Value {
3318 let harness = match session.meta.source {
3319 SessionSource::ClaudeCode => "claude-code",
3320 SessionSource::Codex => "codex",
3321 SessionSource::Gemini => "gemini",
3322 SessionSource::Goose => "goose",
3323 SessionSource::Grok => "grok",
3324 SessionSource::Native => "native",
3325 SessionSource::OpenClaw => "openclaw",
3326 SessionSource::Hermes => "hermes",
3327 SessionSource::OpenCode => "opencode",
3328 SessionSource::Pi => "pi",
3329 };
3330 let mut model = session
3331 .messages
3332 .iter()
3333 .rev()
3334 .find_map(|message| message.metadata.get("model").cloned())
3335 .or_else(|| session.meta.model.clone());
3336 let mut cwd = session
3337 .meta
3338 .cwd
3339 .as_ref()
3340 .map(|path| path.to_string_lossy().into_owned());
3341 let mut version: Option<String> = None;
3342 let mut approval = Value::Null;
3343 let mut sandbox = Value::Null;
3344 for header in &session.meta.codex_headers {
3345 let payload = header.get("payload").unwrap_or(header);
3346 if let Some(value) = payload.get("cli_version").and_then(Value::as_str) {
3347 version = Some(value.to_owned());
3348 }
3349 if let Some(value) = payload.get("model").and_then(Value::as_str) {
3350 model = Some(value.to_owned());
3351 }
3352 if let Some(value) = payload.get("cwd").and_then(Value::as_str) {
3353 cwd = Some(value.to_owned());
3354 }
3355 if let Some(value) = payload.get("approval_policy") {
3356 approval = value.clone();
3357 }
3358 if let Some(value) = payload.get("sandbox_policy") {
3359 sandbox = value.clone();
3360 }
3361 }
3362 if session.meta.source == SessionSource::ClaudeCode {
3363 version = session
3364 .raw
3365 .iter()
3366 .rev()
3367 .filter_map(|line| serde_json::from_str::<Value>(line).ok())
3368 .find_map(|record| {
3369 record
3370 .get("version")
3371 .and_then(Value::as_str)
3372 .map(str::to_owned)
3373 });
3374 }
3375 json!({"worker":{"harness":harness,"model":model,"cwd":cwd,"binary":{"version":version}},
3376 "permissions":{"native":{"approval":approval,"sandbox":sandbox}},"source":"native-session"})
3377}
3378
3379fn recorded_session_profile_history(session: &Session) -> Vec<Value> {
3382 let mut profile = json!({"worker":{"harness":if session.meta.source == SessionSource::Codex {"codex"} else {"claude-code"},
3383 "model":null,"cwd":null,"binary":{"version":null}},
3384 "permissions":{"native":{"approval":null,"sandbox":null}},"source":"native-session"});
3385 let records: Box<dyn Iterator<Item = (usize, Value)> + '_> =
3386 if session.meta.source == SessionSource::Codex {
3387 Box::new(session.meta.codex_headers.iter().cloned().enumerate())
3388 } else if session.meta.source == SessionSource::ClaudeCode {
3389 Box::new(session.raw.iter().enumerate().filter_map(|(index, line)| {
3390 serde_json::from_str::<Value>(line)
3391 .ok()
3392 .map(|record| (index, record))
3393 }))
3394 } else {
3395 return Vec::new();
3396 };
3397 let mut history = Vec::new();
3398 for (index, record) in records {
3399 let before = profile.clone();
3400 let payload = record.get("payload").unwrap_or(&record);
3401 for (native, target) in [("cwd", "cwd"), ("model", "model")] {
3402 if let Some(value) = payload.get(native).and_then(Value::as_str) {
3403 profile["worker"][target] = json!(value);
3404 }
3405 }
3406 let version = payload
3407 .get("cli_version")
3408 .or_else(|| record.get("version"))
3409 .and_then(Value::as_str);
3410 if let Some(value) = version {
3411 profile["worker"]["binary"]["version"] = json!(value);
3412 }
3413 if let Some(value) = record
3414 .get("message")
3415 .and_then(|m| m.get("model"))
3416 .and_then(Value::as_str)
3417 {
3418 profile["worker"]["model"] = json!(value);
3419 }
3420 for (native, target) in [
3421 ("approval_policy", "approval"),
3422 ("sandbox_policy", "sandbox"),
3423 ] {
3424 if let Some(value) = payload.get(native) {
3425 profile["permissions"]["native"][target] = value.clone();
3426 }
3427 }
3428 if profile != before {
3429 history.push(json!({"native_key":format!("profile-record:{index}"),"at":record.get("timestamp").and_then(Value::as_str),"profile":profile}));
3430 }
3431 }
3432 history
3433}
3434
3435fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3436 if let Some(tail) = options.message_tail {
3437 return (total.saturating_sub(tail), total);
3438 }
3439 let offset = options.message_offset.unwrap_or(0).min(total);
3440 let end = options
3441 .message_limit
3442 .map(|limit| offset.saturating_add(limit).min(total))
3443 .unwrap_or(total);
3444 (offset, end)
3445}
3446
3447fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3448 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3449 return message;
3450 };
3451 for part in parts {
3452 let Some(url) = part
3453 .get("image_url")
3454 .and_then(|image| image.get("url"))
3455 .and_then(Value::as_str)
3456 else {
3457 continue;
3458 };
3459 let Some(rest) = url.strip_prefix("data:") else {
3460 continue;
3461 };
3462 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3463 continue;
3464 };
3465 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3466 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3467 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3468 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3469 || options
3470 .max_inline_media_bytes
3471 .is_some_and(|limit| decoded_bytes > limit);
3472 if should_elide {
3473 *part = json!({
3474 "type": "media_reference",
3475 "media_type": media_type,
3476 "encoding": "base64",
3477 "encoded_bytes": encoded.len(),
3478 "decoded_bytes": decoded_bytes,
3479 "omitted": true,
3480 });
3481 }
3482 }
3483 message
3484}
3485
3486#[derive(Deserialize)]
3487struct LocatorParams {
3488 locator: SessionLocator,
3489 #[serde(default)]
3502 fidelity: Option<Fidelity>,
3503 #[serde(default)]
3506 view: Option<SessionReadView>,
3507}
3508
3509#[derive(Deserialize)]
3510struct SessionReadView {
3511 #[serde(default)]
3514 tail_messages: Option<usize>,
3515 #[serde(default)]
3518 include_subagents: bool,
3519 #[serde(default)]
3521 display_history: bool,
3522 #[serde(default)]
3525 max_message_chars: Option<usize>,
3526}
3527
3528impl LocatorParams {
3529 fn read_fidelity(&self) -> Fidelity {
3530 self.fidelity.unwrap_or(Fidelity::Semantic)
3531 }
3532
3533 fn include_subagents(&self) -> bool {
3534 self.view
3535 .as_ref()
3536 .map(|view| view.include_subagents)
3537 .unwrap_or(true)
3538 }
3539
3540 fn tail_messages(&self) -> Option<usize> {
3541 self.view
3542 .as_ref()
3543 .and_then(|view| view.tail_messages)
3544 .map(|limit| limit.clamp(1, 5_000))
3545 }
3546
3547 fn display_history(&self) -> bool {
3548 self.view.as_ref().is_some_and(|view| view.display_history)
3549 }
3550
3551 fn max_message_chars(&self) -> Option<usize> {
3552 self.view
3553 .as_ref()
3554 .and_then(|view| view.max_message_chars)
3555 .map(|limit| limit.clamp(256, 64_000))
3556 }
3557
3558 fn bound_session(&self, session: &mut Session) {
3559 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3560 }
3561}
3562
3563#[derive(Debug, Clone, Copy, Default, Deserialize)]
3564#[serde(rename_all = "snake_case")]
3565enum InlineMediaMode {
3566 #[default]
3567 Full,
3568 Metadata,
3569}
3570
3571#[derive(Debug, Clone, Default, Deserialize)]
3572#[serde(default)]
3573struct SessionLoadOptions {
3574 include_subagents: Option<bool>,
3575 inline_media: InlineMediaMode,
3576 max_inline_media_bytes: Option<usize>,
3577 message_limit: Option<usize>,
3578 message_offset: Option<usize>,
3579 message_tail: Option<usize>,
3580}
3581
3582impl SessionLoadOptions {
3583 fn validate(&self) -> std::result::Result<(), ServiceError> {
3584 if self.message_tail.is_some()
3585 && (self.message_limit.is_some() || self.message_offset.is_some())
3586 {
3587 return Err(ServiceError::InvalidParams(
3588 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3589 .into(),
3590 ));
3591 }
3592 Ok(())
3593 }
3594}
3595
3596#[derive(Deserialize)]
3597struct LoadSessionParams {
3598 #[serde(flatten)]
3599 read: LocatorParams,
3600 #[serde(default)]
3601 options: Option<SessionLoadOptions>,
3602}
3603
3604#[derive(Deserialize)]
3605struct UnfollowParams {
3606 subscription: String,
3607}
3608
3609#[derive(Debug, Deserialize)]
3610#[serde(deny_unknown_fields)]
3611struct IndexResizeParams {
3612 subscription: String,
3613 limit: usize,
3614}
3615
3616#[derive(Deserialize)]
3617struct ActivitySubscribeParams {
3618 locators: Vec<SessionLocator>,
3619 #[serde(default)]
3620 homes: crate::HarnessHomes,
3621}
3622
3623#[derive(Debug, Clone, Deserialize)]
3627struct MessageTarget {
3628 harness: HarnessId,
3629 session_id: String,
3630 #[serde(default)]
3631 #[allow(dead_code)]
3632 storage: Option<supercode_interchange::catalog::StorageLocator>,
3633}
3634
3635#[derive(Deserialize)]
3636struct MessageSessionParams {
3637 locator: MessageTarget,
3638 text: String,
3639 #[serde(default)]
3640 subject: Option<String>,
3641 #[serde(default)]
3643 idempotency_key: Option<String>,
3644 #[serde(default)]
3646 channel: bool,
3647 #[serde(default)]
3651 from_name: Option<String>,
3652 #[serde(default)]
3654 voice_for: Option<crate::mailbox::MailAddress>,
3655 #[serde(default)]
3657 sender_name: Option<String>,
3658 #[serde(default)]
3660 in_reply_to: Option<String>,
3661 #[serde(default)]
3664 marker: Option<String>,
3665 #[serde(default)]
3668 reply_marker: Option<String>,
3669 #[serde(default)]
3671 surface: Option<String>,
3672 #[serde(default)]
3674 notify_when_idle: bool,
3675 #[serde(default)]
3680 as_user: bool,
3681 #[serde(default)]
3684 homes: crate::HarnessHomes,
3685}
3686
3687#[derive(Deserialize)]
3688#[serde(deny_unknown_fields)]
3689struct InboxParams {
3690 #[serde(default)]
3692 from_name: Option<String>,
3693 #[serde(default)]
3695 address: Option<String>,
3696 #[serde(default)]
3698 all: bool,
3699}
3700
3701#[derive(Deserialize)]
3702#[serde(deny_unknown_fields)]
3703struct HarnessSettingsParams {
3704 harness: String,
3705}
3706
3707#[derive(Deserialize)]
3708#[serde(deny_unknown_fields)]
3709struct ConfigureHarnessParams {
3710 harness: String,
3711 #[serde(default)]
3712 changes: Vec<crate::HarnessSettingChange>,
3713 #[serde(default)]
3714 expected_revision: Option<String>,
3715}
3716
3717fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3718 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3719 Ok(report) => (
3720 serde_json::to_value(report).unwrap_or(Value::Null),
3721 Value::Null,
3722 ),
3723 Err(error) => (
3724 Value::Null,
3725 Value::String(format!(
3726 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3727 )),
3728 ),
3729 }
3730}
3731
3732async fn message_live_session(params: &MessageSessionParams) -> Value {
3746 use crate::mail_route::{Delivered, NoDoor, Refused};
3747 let (inbound_controls, inbound_controls_error) =
3748 claude_inbound_controls_or_error(¶ms.homes);
3749 let refused = |reason: &str, message: String| {
3750 json!({
3751 "delivered_to_bus": false,
3752 "refusal": {"reason": reason, "message": message},
3753 "inbound_controls": inbound_controls,
3754 "inbound_controls_error": inbound_controls_error,
3755 })
3756 };
3757 if params.text.trim().is_empty() {
3758 return refused(
3759 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3760 "refusing to deliver an empty message".into(),
3761 );
3762 }
3763 let sender = match operator_address(params.from_name.as_deref()) {
3764 Ok(sender) => sender,
3765 Err(message) => return refused("invalid_sender", message),
3766 };
3767 let receiver = match crate::mailbox::MailAddress::new(
3768 crate::mailbox::local_machine_name(),
3769 params.locator.harness.as_str(),
3770 ¶ms.locator.session_id,
3771 ) {
3772 Ok(receiver) => receiver,
3773 Err(error) => return refused("delivery_failed", error.to_string()),
3774 };
3775 if params
3776 .voice_for
3777 .as_ref()
3778 .is_some_and(|voice_for| voice_for != &receiver)
3779 {
3780 return refused(
3781 "invalid_sender",
3782 "a voice front must represent the receiving session".into(),
3783 );
3784 }
3785 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3788 && !crate::mail_agent::resumable(&receiver)
3789 && crate::runtime_mail::controlled_runtime(
3790 HarnessId::CLAUDE_CODE,
3791 ¶ms.locator.session_id,
3792 )
3793 .is_none()
3794 {
3795 if let Err(refusal) =
3796 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3797 {
3798 return refused(refusal.reason.as_str(), refusal.message);
3799 }
3800 }
3801 if params.as_user {
3802 let agent_address =
3805 (receiver.harness == crate::mail_agent::AGENT_HARNESS).then(|| receiver.clone());
3806 let receiver = if receiver.harness == crate::mail_agent::AGENT_HARNESS {
3807 match crate::mail_agent::load(&receiver.session_id) {
3808 Some(agent) => agent.main_session,
3809 None => {
3810 return refused(
3811 "delivery_failed",
3812 format!(
3813 "no agent named {} is declared on this machine",
3814 receiver.session_id
3815 ),
3816 )
3817 }
3818 }
3819 } else {
3820 receiver
3821 };
3822 return message_as_user(
3823 params,
3824 sender,
3825 receiver,
3826 agent_address,
3827 inbound_controls,
3828 inbound_controls_error,
3829 )
3830 .await;
3831 }
3832 let door = crate::mail_route::door_for(¶ms.homes, &receiver);
3835 let mut envelope = match crate::mailbox::Envelope::new(
3836 sender.clone(),
3837 params
3838 .sender_name
3839 .clone()
3840 .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3841 if params.channel {
3842 crate::mailbox::MailKind::Channel
3843 } else {
3844 crate::mailbox::MailKind::Peer
3845 },
3846 crate::mailbox::ReplyVia::Command,
3847 params.text.clone(),
3848 ) {
3849 Ok(envelope) => envelope,
3850 Err(error) => return refused("delivery_failed", error.to_string()),
3851 };
3852 if let Some(key) = ¶ms.idempotency_key {
3853 if key.is_empty() || key.len() > 256 {
3854 return refused("invalid_key", "idempotency_key needs 1–256 bytes".into());
3855 }
3856 let identity = serde_json::to_vec(&(sender.to_string(), receiver.to_string(), key))
3857 .expect("string tuple serializes");
3858 envelope.id = format!("m-{}", &blake3::hash(&identity).to_hex()[..24]);
3859 let mailbox = match crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &receiver) {
3860 Ok(mailbox) => mailbox,
3861 Err(error) => return refused("delivery_failed", error.to_string()),
3862 };
3863 match mailbox.find(&envelope.id) {
3864 Ok(Some(previous)) => {
3865 let saved = previous.envelope;
3866 let subject = params.subject.as_deref().map(|value| {
3867 value
3868 .lines()
3869 .next()
3870 .unwrap_or("")
3871 .chars()
3872 .take(200)
3873 .collect::<String>()
3874 });
3875 if saved.body != envelope.body
3876 || saved.subject != subject
3877 || saved.kind != envelope.kind
3878 || saved.from_name != envelope.from_name
3879 || saved.in_reply_to != params.in_reply_to
3880 || saved.voice_for != params.voice_for
3881 {
3882 return refused(
3883 "idempotency_conflict",
3884 "this key already names different mail".into(),
3885 );
3886 }
3887 }
3888 Ok(None) => {}
3889 Err(error) => return refused("delivery_failed", error.to_string()),
3890 }
3891 }
3892 envelope.in_reply_to = params.in_reply_to.clone();
3893 envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
3894 envelope.voice_for = params.voice_for.clone();
3895 envelope.subject = params.subject.as_deref().map(|value| {
3896 value
3897 .lines()
3898 .next()
3899 .unwrap_or("")
3900 .chars()
3901 .take(200)
3902 .collect()
3903 });
3904 let channel = crate::mail_agent::Channel {
3905 marker: params.marker.clone(),
3906 reply_marker: params.reply_marker.clone(),
3907 surface: params.surface.clone(),
3908 };
3909 match crate::mail_agent::plan(&mut envelope, &receiver, &channel) {
3910 Err(error) => return refused("delivery_failed", error.to_string()),
3911 Ok(Some(plan)) => {
3912 let caller = crate::mail_route::Caller {
3913 address: sender.clone(),
3914 name: envelope.from_name.clone(),
3915 };
3916 let outcome = crate::mail_send::deliver_planned(
3917 ¶ms.homes,
3918 &caller,
3919 &envelope,
3920 &plan,
3921 false,
3922 params.notify_when_idle,
3923 )
3924 .await;
3925 return json!({
3926 "delivered_to_bus": outcome.code == 0,
3927 "message_id": envelope.id,
3928 "reply_to": sender.to_string(),
3929 "thread": plan.thread.as_ref().map(|thread| thread.id.clone()),
3930 "target": {
3931 "harness": params.locator.harness.as_str(),
3932 "session_id": params.locator.session_id,
3933 },
3934 "delivery": {"door": "agent", "how": outcome.text},
3935 "inbound_controls": inbound_controls,
3936 "inbound_controls_error": inbound_controls_error,
3937 });
3938 }
3939 Ok(None) => {}
3940 }
3941 let door = match door {
3942 Ok(door) => door,
3943 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3944 return refused(
3945 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3946 format!(
3947 "no running `{}` session `{}` is reachable; its transcript is persisted only",
3948 params.locator.harness.as_str(),
3949 params.locator.session_id
3950 ),
3951 )
3952 }
3953 };
3954 let how = match crate::mail_route::deliver(
3955 &envelope,
3956 &receiver,
3957 &door,
3958 true,
3959 params.notify_when_idle,
3960 )
3961 .await
3962 {
3963 Err(detail) => {
3964 return refused(
3965 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3966 detail,
3967 )
3968 }
3969 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3970 Ok(Err(Refused::TooLong(bytes))) => {
3971 return refused(
3972 "too_long",
3973 format!(
3974 "the message is {bytes} bytes; the limit is {}",
3975 crate::mail_route::MAX_RELAYED_BYTES
3976 ),
3977 )
3978 }
3979 Ok(Ok(delivered)) => match delivered {
3980 Delivered::Steered => "steered",
3981 Delivered::Started => "started",
3982 Delivered::Native { busy: true } => "next_tool_call",
3983 Delivered::Native { busy: false } => "started",
3984 Delivered::Hooked => "hook",
3985 Delivered::HookWoken => "woken",
3986 Delivered::Queued => "queued",
3987 Delivered::Stored => "stored",
3988 Delivered::Operator => "filed",
3989 },
3990 };
3991 json!({
3992 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3993 "message_id": envelope.id,
3994 "reply_to": sender.to_string(),
3995 "target": {
3996 "harness": params.locator.harness.as_str(),
3997 "session_id": params.locator.session_id,
3998 "name": match &door {
3999 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
4000 _ => None,
4001 },
4002 },
4003 "delivery": {"door": door.name(), "how": how},
4004 "inbound_controls": inbound_controls,
4005 "inbound_controls_error": inbound_controls_error,
4006 })
4007}
4008
4009async fn message_as_user(
4013 params: &MessageSessionParams,
4014 sender: crate::mailbox::MailAddress,
4015 receiver: crate::mailbox::MailAddress,
4016 agent: Option<crate::mailbox::MailAddress>,
4017 inbound_controls: Value,
4018 inbound_controls_error: Value,
4019) -> Value {
4020 let mut envelope = match crate::mailbox::Envelope::new(
4021 sender.clone(),
4022 format!("{}@{}", sender.session_id, sender.machine),
4023 crate::mailbox::MailKind::User,
4024 crate::mailbox::ReplyVia::None,
4025 params.text.clone(),
4026 ) {
4027 Ok(envelope) => envelope,
4028 Err(error) => {
4029 return json!({
4030 "delivered_to_bus": false,
4031 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4032 "inbound_controls": inbound_controls,
4033 "inbound_controls_error": inbound_controls_error,
4034 })
4035 }
4036 };
4037 let mut receiver = receiver;
4042 if let Some(agent) = &agent {
4043 envelope.in_reply_to = params.in_reply_to.clone();
4044 envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
4045 let channel = crate::mail_agent::Channel::default();
4046 let planned = match crate::mail_agent::plan(&mut envelope, agent, &channel) {
4047 Ok(planned) => planned,
4048 Err(error) => {
4049 return json!({
4050 "delivered_to_bus": false,
4051 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4052 "inbound_controls": inbound_controls,
4053 "inbound_controls_error": inbound_controls_error,
4054 })
4055 }
4056 };
4057 if let Some(planned) = planned {
4058 let typed_to = planned
4059 .recipients
4060 .iter()
4061 .find(|r| r.wake && r.address.harness != "operator")
4062 .map(|r| r.address.clone());
4063 if let Some(typed_to) = typed_to {
4064 receiver = typed_to;
4065 }
4066 if let Some(agent) = &planned.agent {
4067 let _ = crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), agent)
4069 .and_then(|mailbox| mailbox.deliver_read(&envelope));
4070 }
4071 let mut copy = envelope.clone();
4074 copy.kind = crate::mailbox::MailKind::Typed;
4075 for other in planned
4076 .recipients
4077 .iter()
4078 .filter(|r| r.address != receiver && r.address != sender)
4079 {
4080 let _ = crate::mail_agent::file_unread(&other.address, ©);
4081 }
4082 if let Some(account_manager) = crate::mail_agent::owners_account_manager() {
4085 let main = account_manager.main_session.clone();
4086 let already =
4087 main == receiver || planned.recipients.iter().any(|r| r.address == main);
4088 if account_manager.name != agent.session_id && !already {
4089 let _ = crate::mail_agent::file_unread(&main, ©);
4090 }
4091 }
4092 }
4093 }
4094 if agent.is_some()
4097 && crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4098 && crate::mail_agent::resumable(&receiver)
4099 {
4100 if let Err(error) = crate::mail_agent::resume(&receiver) {
4101 return json!({
4102 "delivered_to_bus": false,
4103 "refusal": {"reason": "no_user_door", "message": format!("{receiver} was not resumed: {error}")},
4104 "inbound_controls": inbound_controls,
4105 "inbound_controls_error": inbound_controls_error,
4106 });
4107 }
4108 let started = std::time::Instant::now();
4109 while crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4110 && started.elapsed() < crate::mail_agent::RESUME_WAIT
4111 {
4112 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
4113 }
4114 }
4115 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
4116 Ok(turn) => json!({
4117 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
4118 "message_id": envelope.id,
4119 "target": {
4120 "harness": params.locator.harness.as_str(),
4121 "session_id": params.locator.session_id,
4122 },
4123 "delivery": {
4124 "door": match turn {
4125 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
4126 _ => "pane",
4127 },
4128 "how": turn.as_str(),
4129 },
4130 "inbound_controls": inbound_controls,
4131 "inbound_controls_error": inbound_controls_error,
4132 }),
4133 Err(message) => json!({
4134 "delivered_to_bus": false,
4135 "refusal": {"reason": "no_user_door", "message": message},
4136 "inbound_controls": inbound_controls,
4137 "inbound_controls_error": inbound_controls_error,
4138 }),
4139 }
4140}
4141
4142#[derive(Deserialize)]
4143#[serde(deny_unknown_fields)]
4144struct ActivityUnderParams {
4145 pids: Vec<u32>,
4147 #[serde(default)]
4148 homes: crate::HarnessHomes,
4149}
4150
4151async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
4155 let params = decode::<ActivityUnderParams>(params)?;
4156 if params.pids.len() > 1024 {
4157 return Err(ServiceError::InvalidParams(
4158 "sessions.activity_under accepts at most 1024 pids".into(),
4159 ));
4160 }
4161 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
4162 .await
4163 .map_err(ServiceError::Sdk)?;
4164 Ok(json!({
4165 "activities": found
4166 .into_iter()
4167 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
4168 .collect::<Vec<_>>(),
4169 }))
4170}
4171
4172fn operator_address(
4175 name: Option<&str>,
4176) -> std::result::Result<crate::mailbox::MailAddress, String> {
4177 let name: String = name
4178 .unwrap_or("supercode")
4179 .trim()
4180 .chars()
4181 .map(|character| {
4182 if character.is_whitespace() || character == '@' {
4183 '-'
4184 } else {
4185 character
4186 }
4187 })
4188 .collect();
4189 if name.is_empty() {
4190 return Err("from_name must not be empty".into());
4191 }
4192 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
4193 .map_err(|error| error.to_string())
4194}
4195
4196fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
4200 let address = match (¶ms.address, ¶ms.from_name) {
4201 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
4202 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
4203 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
4204 (Some(_), Some(_)) => {
4205 return Err(ServiceError::InvalidParams(
4206 "sessions.inbox takes from_name or address, not both".into(),
4207 ))
4208 }
4209 };
4210 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
4211 let mailbox =
4212 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
4213 let claimed = mailbox.claim_unread().map_err(operation)?;
4214 let mut messages: Vec<Value> = Vec::new();
4215 if params.all {
4216 for stored in mailbox.list().map_err(operation)? {
4217 if stored.state == crate::mailbox::MailState::Read {
4218 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4219 }
4220 }
4221 }
4222 for stored in &claimed {
4223 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4224 }
4225 for stored in &claimed {
4226 mailbox.acknowledge(stored).map_err(operation)?;
4227 }
4228 Ok(json!({"address": address.to_string(), "messages": messages}))
4229}
4230
4231#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4236struct FollowedSource {
4237 harness: String,
4238 session_id: String,
4239 reported: Option<String>,
4240}
4241
4242#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4243struct ActivitySubscription {
4244 locators: Vec<SessionLocator>,
4245 homes: crate::HarnessHomes,
4246 reported: BTreeMap<(String, String), crate::SessionActivity>,
4247}
4248
4249fn live_descriptor_value(
4256 session: &SessionDescriptor,
4257 doors: &crate::mail_route::LiveSessions,
4258) -> std::result::Result<Value, ServiceError> {
4259 let mut value = serde_json::to_value(session)
4260 .map_err(|error| ServiceError::Operation(error.to_string()))?;
4261 if value.get("title").is_none_or(Value::is_null) {
4263 if let Some(live) = doors.all().iter().find(|live| {
4264 live.address.harness == session.locator.harness.as_str()
4265 && live.address.session_id == session.locator.session_id
4266 }) {
4267 let name = live.name.split('@').next().unwrap_or(&live.name);
4268 if !name.is_empty() {
4269 value["title"] = json!(name);
4270 }
4271 }
4272 }
4273 if let Some(workspace) = &session.cwd {
4274 let source = LiveRuntimeSource {
4275 harness: session.locator.harness.as_str().to_string(),
4276 session_id: session.locator.session_id.clone(),
4277 workspace: workspace.clone(),
4278 };
4279 if let Some(endpoint) = discover_live_runtime(&source)
4280 .map_err(|error| ServiceError::Operation(error.to_string()))?
4281 {
4282 value["live_endpoint"] = json!(endpoint.as_str());
4283 }
4284 }
4285 if let Some(door) = doors.door(
4289 session.locator.harness.as_str(),
4290 &session.locator.session_id,
4291 ) {
4292 value["delivery"] = json!(door);
4293 }
4294 if let Some(live) = doors.all().iter().find(|live| {
4297 live.address.harness == session.locator.harness.as_str()
4298 && live.address.session_id == session.locator.session_id
4299 && live.status == "waiting"
4300 }) {
4301 if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
4302 value["pending_request"] = request;
4303 }
4304 }
4305 Ok(value)
4306}
4307
4308fn live_index_changes(
4309 changes: Vec<crate::session_index::SessionIndexChange>,
4310 homes: &HarnessHomes,
4311) -> std::result::Result<Vec<Value>, ServiceError> {
4312 use crate::session_index::SessionIndexChange;
4313 let doors = crate::mail_route::LiveSessions::read(homes);
4314 changes
4315 .into_iter()
4316 .map(|change| match change {
4317 SessionIndexChange::Added { descriptor } => Ok(json!({
4318 "kind": "added",
4319 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4320 })),
4321 SessionIndexChange::Updated { descriptor } => Ok(json!({
4322 "kind": "updated",
4323 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4324 })),
4325 SessionIndexChange::Removed { key } => Ok(json!({
4326 "kind": "removed",
4327 "key": key,
4328 })),
4329 })
4330 .collect()
4331}
4332
4333fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
4334 use crate::{SessionPresence, SessionTurnState};
4335 match (activity.presence, activity.turn) {
4336 (SessionPresence::Persisted, _) => None,
4337 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
4338 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
4339 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
4341 (SessionPresence::Running, SessionTurnState::Unknown)
4345 if activity.evidence.native_state.is_none() =>
4346 {
4347 None
4348 }
4349 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
4350 }
4351}
4352
4353#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
4354#[serde(rename_all = "kebab-case")]
4355enum TransferFormat {
4356 ClaudeCode,
4357 Codex,
4358 #[serde(rename = "opencode", alias = "open-code")]
4359 OpenCode,
4360 Pi,
4361 Grok,
4362 Gemini,
4363 Goose,
4364 Hermes,
4368}
4369
4370impl TransferFormat {
4371 fn id(self) -> &'static str {
4372 match self {
4373 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
4374 Self::Codex => HarnessId::CODEX,
4375 Self::OpenCode => HarnessId::OPENCODE,
4376 Self::Pi => HarnessId::PI,
4377 Self::Grok => HarnessId::GROK,
4378 Self::Gemini => HarnessId::GEMINI,
4379 Self::Goose => HarnessId::GOOSE,
4380 Self::Hermes => HarnessId::HERMES,
4381 }
4382 }
4383}
4384
4385impl From<TransferFormat> for SessionFormat {
4386 fn from(value: TransferFormat) -> Self {
4387 match value {
4388 TransferFormat::ClaudeCode => Self::ClaudeCode,
4389 TransferFormat::Codex => Self::Codex,
4390 TransferFormat::OpenCode => Self::OpenCode,
4391 TransferFormat::Pi => Self::Pi,
4392 TransferFormat::Grok => Self::Grok,
4393 TransferFormat::Gemini => Self::Gemini,
4394 TransferFormat::Goose => Self::Goose,
4395 TransferFormat::Hermes => Self::Codex,
4397 }
4398 }
4399}
4400
4401#[derive(Deserialize)]
4402struct ImportSessionParams {
4403 source_harness: TransferFormat,
4404 content: String,
4405}
4406
4407#[derive(Deserialize)]
4408struct ExportSessionParams {
4409 locator: SessionLocator,
4410 target_harness: TransferFormat,
4411}
4412
4413#[derive(Deserialize)]
4414struct ReduceSessionParams {
4415 locator: SessionLocator,
4416 target_harness: TransferFormat,
4417 #[serde(default = "default_keep_last")]
4418 keep_last: usize,
4419}
4420
4421fn default_keep_last() -> usize {
4422 6
4423}
4424
4425#[derive(Deserialize)]
4426struct BranchSessionParams {
4427 locator: SessionLocator,
4428 #[serde(default)]
4429 target_harness: Option<TransferFormat>,
4430}
4431
4432#[derive(Deserialize)]
4433struct HandoffSessionParams {
4434 locator: SessionLocator,
4435 target_harness: TransferFormat,
4436 #[serde(default)]
4437 cwd: Option<PathBuf>,
4438}
4439
4440#[derive(Deserialize)]
4441struct MaterializeSessionParams {
4442 artifact: crate::native_materialize::MaterializeArtifact,
4443 cwd: PathBuf,
4444 #[serde(default)]
4446 homes: HarnessHomes,
4447}
4448
4449#[derive(Debug, Clone, Copy, Default, Deserialize)]
4452#[serde(rename_all = "snake_case")]
4453enum ResumePolicy {
4454 Default,
4455 #[default]
4456 Yolo,
4457}
4458
4459#[derive(Deserialize)]
4460struct ResumeInstructionsParams {
4461 locator: SessionLocator,
4462 #[serde(default)]
4463 cwd: Option<PathBuf>,
4464 #[serde(default)]
4465 policy: ResumePolicy,
4466}
4467
4468#[derive(Deserialize)]
4470struct WorkflowLoadParams {
4471 from: crate::workflow_doors::WorkflowHarness,
4472 home: PathBuf,
4473}
4474
4475#[derive(Deserialize)]
4478struct OrchestrationLoadParams {
4479 root: PathBuf,
4480 #[serde(default)]
4481 flavor: crate::orchestration_doors::HomeFlavor,
4482}
4483
4484#[derive(Deserialize)]
4487struct OrchestrationSaveParams {
4488 root: PathBuf,
4489 orchestration: crate::orchestration::Orchestration,
4490 #[serde(default)]
4491 vault: BTreeMap<String, String>,
4492}
4493
4494#[derive(Deserialize)]
4496struct OrchestrationCompileParams {
4497 from: crate::orchestration_doors::OrchestrationHarness,
4498 home: PathBuf,
4499}
4500
4501#[derive(Deserialize)]
4505struct OrchestrationDecompileParams {
4506 to: crate::orchestration_doors::OrchestrationHarness,
4507 orchestration: crate::orchestration::Orchestration,
4508 source: PathBuf,
4509 #[serde(default)]
4510 source_flavor: crate::orchestration_doors::SourceFlavor,
4511 dest: PathBuf,
4512 #[serde(default)]
4513 vault: BTreeMap<String, String>,
4514}
4515
4516#[derive(Deserialize)]
4519struct OrchestrationImportParams {
4520 from: crate::orchestration_doors::OrchestrationHarness,
4521 home: PathBuf,
4522 into: PathBuf,
4523}
4524
4525#[derive(Deserialize)]
4528struct OrchestrationExportParams {
4529 to: crate::orchestration_doors::OrchestrationHarness,
4530 root: PathBuf,
4531 dest: PathBuf,
4532}
4533
4534#[derive(Deserialize)]
4536struct JobsGetParams {
4537 harness: String,
4538 id: String,
4539 #[serde(default)]
4540 homes: crate::HarnessHomes,
4541}
4542
4543fn mutate_job(
4551 verb: crate::jobs_control::JobVerb,
4552 params: Value,
4553) -> std::result::Result<Value, ServiceError> {
4554 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4555 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4556 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4557 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4558}
4559
4560fn mutate_skill(
4567 verb: crate::skills_control::SkillVerb,
4568 params: Value,
4569) -> std::result::Result<Value, ServiceError> {
4570 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4571 if !crate::skills_control::supports_skill_control(&mutation.harness) {
4572 return Err(ServiceError::UnsupportedAction(format!(
4573 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4574 mutation.harness,
4575 verb.as_str(),
4576 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4577 )));
4578 }
4579 let outcome =
4580 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4581 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4582}
4583
4584fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4586 match error {
4587 crate::skills_control::SkillControlError::Unsupported(message) => {
4588 ServiceError::UnsupportedAction(message)
4589 }
4590 crate::skills_control::SkillControlError::Invalid(message) => {
4591 ServiceError::InvalidParams(message)
4592 }
4593 crate::skills_control::SkillControlError::Failed(message) => {
4594 ServiceError::Operation(message)
4595 }
4596 }
4597}
4598
4599fn mutate_profile(
4607 verb: crate::profiles_control::ProfileVerb,
4608 params: Value,
4609) -> std::result::Result<Value, ServiceError> {
4610 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4611 let outcome =
4612 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4613 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4614}
4615
4616fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4618 match error {
4619 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4620 ServiceError::UnsupportedAction(message)
4621 }
4622 crate::profiles_control::ProfileControlError::Invalid(message) => {
4623 ServiceError::InvalidParams(message)
4624 }
4625 crate::profiles_control::ProfileControlError::Failed(message) => {
4626 ServiceError::Operation(message)
4627 }
4628 }
4629}
4630
4631fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4635 match error {
4636 crate::jobs_control::JobControlError::Unsupported(message) => {
4637 ServiceError::UnsupportedAction(message)
4638 }
4639 crate::jobs_control::JobControlError::Invalid(message) => {
4640 ServiceError::InvalidParams(message)
4641 }
4642 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4643 }
4644}
4645
4646fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4651 match error {
4652 crate::SessionControlError::Unsupported(message) => {
4653 ServiceError::UnsupportedAction(message)
4654 }
4655 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4656 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4657 }
4658}
4659
4660fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4665 if crate::jobs::supports_jobs(harness) {
4666 return Ok(());
4667 }
4668 Err(ServiceError::UnsupportedAction(format!(
4669 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4670 crate::jobs::JOB_HARNESSES.join(", ")
4671 )))
4672}
4673
4674#[derive(Deserialize)]
4676struct RunsGetParams {
4677 harness: String,
4678 id: String,
4679 #[serde(default)]
4680 homes: crate::HarnessHomes,
4681}
4682
4683fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4688 if crate::runs::supports_runs(harness) {
4689 return Ok(());
4690 }
4691 Err(ServiceError::UnsupportedAction(format!(
4692 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4693 crate::runs::RUN_HARNESSES.join(", ")
4694 )))
4695}
4696
4697#[derive(Serialize)]
4698struct SessionArtifact {
4699 source_harness: HarnessId,
4700 target_harness: &'static str,
4701 session_id: Option<String>,
4702 content: String,
4703 suggested_filename: String,
4704 files: Vec<SessionArtifactFile>,
4705 fidelity: Fidelity,
4706 residue: Vec<String>,
4707}
4708
4709#[derive(Serialize)]
4710struct SessionArtifactFile {
4711 path: String,
4712 content: String,
4713 role: ArtifactFileRole,
4714}
4715
4716#[derive(Serialize)]
4717#[serde(rename_all = "snake_case")]
4718enum ArtifactFileRole {
4719 Primary,
4720 Subagent,
4721 Bundle,
4722 SourceRecovery,
4723}
4724
4725#[derive(Serialize)]
4726struct StructuredLaunch {
4727 cwd: PathBuf,
4728 program: String,
4729 arguments: Vec<String>,
4730 env: BTreeMap<String, String>,
4731}
4732
4733struct HandoffInstructions {
4734 launch: StructuredLaunch,
4735 materialize: Option<StructuredLaunch>,
4736 requires_materialization: bool,
4737 note: String,
4738}
4739
4740#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4741#[serde(rename_all = "snake_case")]
4742enum HarnessProbeLevel {
4743 #[default]
4744 Passive,
4745 Handshake,
4746}
4747
4748#[derive(Default, Deserialize)]
4749#[serde(default)]
4750struct HarnessInventoryParams {
4751 harness: Option<HarnessId>,
4752 harnesses: Vec<HarnessId>,
4753 workspace: Option<PathBuf>,
4754 probe: HarnessProbeLevel,
4755 include_sessions: bool,
4756 skip_versions: bool,
4758}
4759
4760#[derive(Deserialize)]
4761struct HarnessAuthenticationParams {
4762 harness: HarnessId,
4763}
4764
4765#[derive(Deserialize)]
4766struct BeginHarnessAuthenticationParams {
4767 harness: HarnessId,
4768 #[serde(default = "local_browser_authentication_environment")]
4769 environment: crate::HarnessAuthenticationEnvironment,
4770 #[serde(default)]
4771 method: Option<crate::HarnessAuthenticationMethodId>,
4772 #[serde(default)]
4773 cwd: Option<PathBuf>,
4774}
4775
4776fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4777 crate::HarnessAuthenticationEnvironment::LocalBrowser
4778}
4779
4780#[derive(Serialize)]
4781struct HarnessInventoryReport {
4782 probe: HarnessProbeLevel,
4783 workspace: Option<PathBuf>,
4784 harnesses: Vec<LocalHarness>,
4785}
4786
4787#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4788#[serde(rename_all = "snake_case")]
4789enum HarnessAuthState {
4790 Ready,
4791 Configured,
4792 Required,
4793 Unknown,
4794}
4795
4796#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4797#[serde(rename_all = "snake_case")]
4798enum HarnessRuntimeState {
4799 Ready,
4800 Degraded,
4801 Unavailable,
4802}
4803
4804#[derive(Serialize)]
4805struct HarnessSessionCounts {
4806 global: Option<usize>,
4807 workspace: Option<usize>,
4808}
4809
4810#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4821#[serde(rename_all = "snake_case")]
4822pub enum GatewayState {
4823 Up,
4824 Down,
4825 Unknown,
4826}
4827
4828#[derive(Debug, Clone, Serialize)]
4830pub struct GatewayHealth {
4831 pub state: GatewayState,
4832 #[serde(skip_serializing_if = "Option::is_none")]
4837 pub endpoint: Option<String>,
4838 #[serde(skip_serializing_if = "Option::is_none")]
4839 pub version: Option<String>,
4840 pub evidence: String,
4842 pub checked_at_ms: u64,
4843}
4844
4845fn openclaw_gateway_endpoint(home: &Path) -> String {
4849 let config_path = home.join(".openclaw/openclaw.json");
4850 let gateway = std::fs::read_to_string(&config_path)
4851 .ok()
4852 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4853 .and_then(|config| config.get("gateway").cloned());
4854 if let Some(url) = gateway
4855 .as_ref()
4856 .and_then(|gateway| gateway.get("url"))
4857 .and_then(serde_json::Value::as_str)
4858 {
4859 return url.to_string();
4860 }
4861 let port = gateway
4862 .as_ref()
4863 .and_then(|gateway| gateway.get("port"))
4864 .and_then(serde_json::Value::as_u64)
4865 .unwrap_or(18789);
4866 format!("ws://127.0.0.1:{port}")
4867}
4868
4869fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4876 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4877 let output = std::process::Command::new(&program)
4878 .args(["gateway", "status"])
4879 .stdin(std::process::Stdio::null())
4880 .output()
4881 .ok()?;
4882 let text = format!(
4883 "{}{}",
4884 String::from_utf8_lossy(&output.stdout),
4885 String::from_utf8_lossy(&output.stderr)
4886 );
4887 let verdict = text.lines().find_map(|line| {
4888 let l = line.trim();
4889 if l.contains("supervised by launchd (PID")
4890 || l.contains("supervised by systemd (PID")
4891 || l.contains("Gateway is running")
4892 || l.contains("process is running")
4893 {
4894 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4895 } else if l.contains("not running") || l.contains("not installed") {
4896 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4897 } else {
4898 None
4899 }
4900 });
4901 verdict
4902}
4903
4904fn gateway_health(
4905 id: &str,
4906 installed: bool,
4907 running: Option<&RunningInstance>,
4908 version: Option<&str>,
4909) -> GatewayHealth {
4910 let checked_at_ms = now_epoch_ms();
4911 let home = supercode_interchange::user_home()
4912 .map(std::path::PathBuf::into_os_string)
4913 .map(PathBuf::from);
4914 match id {
4915 HarnessId::HERMES | HarnessId::OPENCLAW => {
4916 let endpoint = (id == HarnessId::OPENCLAW)
4917 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4918 .flatten();
4919 let (state, evidence) = match running {
4920 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4921 None if !installed => (
4922 GatewayState::Unknown,
4923 format!("`{id}` is not installed; no gateway to probe"),
4924 ),
4925 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4926 Some((state, evidence)) => (state, evidence),
4929 None => (
4930 GatewayState::Down,
4931 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4932 ),
4933 },
4934 None => (
4935 GatewayState::Down,
4936 format!(
4937 "no TCP listener at {}",
4938 endpoint.as_deref().unwrap_or("the gateway endpoint")
4939 ),
4940 ),
4941 };
4942 GatewayHealth {
4943 state,
4944 endpoint,
4945 version: version.map(str::to_string),
4946 evidence,
4947 checked_at_ms,
4948 }
4949 }
4950 HarnessId::ORCHESTRATOR => {
4957 let root = crate::HarnessHomes::default().orchestrator;
4958 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4959 Some(lease) if lease.is_live() => (
4960 GatewayState::Up,
4961 format!(
4962 "`{}` names pid {} (started {}), which is live",
4963 crate::orchestrator::lock_path(&root).display(),
4964 lease.pid,
4965 lease.started_at
4966 ),
4967 ),
4968 Some(lease) => (
4969 GatewayState::Down,
4970 format!(
4971 "stale lease `{}`: pid {} is gone",
4972 crate::orchestrator::lock_path(&root).display(),
4973 lease.pid
4974 ),
4975 ),
4976 None => (
4977 GatewayState::Down,
4978 format!(
4979 "no lease at `{}`; `supercode orchestrator start` writes one",
4980 crate::orchestrator::lock_path(&root).display()
4981 ),
4982 ),
4983 };
4984 GatewayHealth {
4985 state,
4986 endpoint: None,
4987 version: version.map(str::to_string),
4988 evidence,
4989 checked_at_ms,
4990 }
4991 }
4992 _ => GatewayHealth {
4993 state: GatewayState::Unknown,
4994 endpoint: None,
4995 version: version.map(str::to_string),
4996 evidence: format!("`{id}` runs per session, not as a gateway"),
4997 checked_at_ms,
4998 },
4999 }
5000}
5001
5002#[derive(Debug, Clone, Serialize)]
5003struct RunningInstance {
5004 method: RunningInstanceMethod,
5006 evidence: String,
5008 checked_at_ms: u64,
5010}
5011
5012#[derive(Debug, Clone, Copy, Serialize)]
5013#[serde(rename_all = "snake_case")]
5014enum RunningInstanceMethod {
5015 GatewayConnect,
5018 StoreWalActivity,
5021}
5022
5023fn now_epoch_ms() -> u64 {
5024 std::time::SystemTime::now()
5025 .duration_since(std::time::UNIX_EPOCH)
5026 .map(|elapsed| elapsed.as_millis() as u64)
5027 .unwrap_or(0)
5028}
5029
5030fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
5034 let config_path = home.join(".openclaw/openclaw.json");
5035 let text = std::fs::read_to_string(&config_path).ok();
5036 let gateway = text
5037 .as_deref()
5038 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
5039 .and_then(|config| config.get("gateway").cloned());
5040 let address = gateway
5041 .as_ref()
5042 .and_then(|gateway| gateway.get("url"))
5043 .and_then(serde_json::Value::as_str)
5044 .and_then(|url| {
5045 url.split("://").nth(1).map(|rest| {
5046 rest.trim_end_matches('/')
5047 .split('/')
5048 .next()
5049 .unwrap_or(rest)
5050 .to_string()
5051 })
5052 })
5053 .unwrap_or_else(|| {
5054 let port = gateway
5055 .as_ref()
5056 .and_then(|gateway| gateway.get("port"))
5057 .and_then(serde_json::Value::as_u64)
5058 .unwrap_or(18789);
5059 format!("127.0.0.1:{port}")
5060 });
5061 let reachable = std::net::TcpStream::connect_timeout(
5062 &address.parse().ok()?,
5063 std::time::Duration::from_millis(400),
5064 )
5065 .is_ok();
5066 reachable.then(|| RunningInstance {
5067 method: RunningInstanceMethod::GatewayConnect,
5068 evidence: format!(
5069 "gateway endpoint {address} accepted a TCP connect (from {})",
5070 config_path.display()
5071 ),
5072 checked_at_ms: now_epoch_ms(),
5073 })
5074}
5075
5076fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
5081 let wal = home.join(".hermes/state.db-wal");
5082 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
5083 let age_ms = std::time::SystemTime::now()
5084 .duration_since(modified)
5085 .map(|age| age.as_millis() as u64)
5086 .unwrap_or(u64::MAX);
5087 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
5088 method: RunningInstanceMethod::StoreWalActivity,
5089 evidence: format!(
5090 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
5091 wal.display()
5092 ),
5093 checked_at_ms: now_epoch_ms(),
5094 })
5095}
5096
5097fn probe_running_instance(id: &str) -> Option<RunningInstance> {
5099 let home = supercode_interchange::user_home()
5100 .map(std::path::PathBuf::into_os_string)
5101 .map(PathBuf::from)?;
5102 match id {
5103 HarnessId::OPENCLAW => probe_openclaw_running(&home),
5104 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
5105 _ => None,
5106 }
5107}
5108
5109#[derive(Serialize)]
5110struct LocalHarness {
5111 id: HarnessId,
5112 display_name: String,
5113 supported: bool,
5114 installed: bool,
5115 executable: Option<String>,
5116 version: Option<String>,
5117 auth: HarnessAuthState,
5118 runtime: HarnessRuntimeState,
5119 protocol: String,
5120 capabilities: crate::RuntimeCapabilities,
5121 effective_capabilities: crate::RuntimeCapabilities,
5122 sessions: HarnessSessionCounts,
5123 #[serde(skip_serializing_if = "Option::is_none")]
5126 running: Option<RunningInstance>,
5127 gateway: GatewayHealth,
5129 reason: Option<String>,
5130 repair: Option<String>,
5131}
5132
5133#[derive(Clone, Deserialize)]
5134struct RuntimeBackendParams {
5135 harness: HarnessId,
5136 #[serde(default)]
5137 protocol: Option<String>,
5138 #[serde(default)]
5139 launch: Option<RuntimeLaunch>,
5140 #[serde(default)]
5141 base_url: Option<String>,
5142 #[serde(default)]
5143 policy: RuntimePolicy,
5144}
5145
5146#[derive(Debug, Clone, Copy, Default, Deserialize)]
5149#[serde(rename_all = "snake_case")]
5150enum RuntimePolicy {
5151 Default,
5152 #[default]
5153 Yolo,
5154}
5155
5156#[derive(Deserialize)]
5157struct RuntimeStartParams {
5158 #[serde(flatten)]
5159 backend: RuntimeBackendParams,
5160 cwd: PathBuf,
5161 #[serde(default)]
5164 mcp_servers: Vec<crate::McpServerLaunch>,
5165 #[serde(default)]
5167 approval_policy: Option<String>,
5168}
5169
5170#[derive(Deserialize)]
5171struct RuntimeAttachParams {
5172 #[serde(flatten)]
5173 backend: RuntimeBackendParams,
5174 runtime_id: String,
5175 #[serde(default)]
5176 cwd: Option<PathBuf>,
5177 #[serde(default)]
5180 mcp_servers: Vec<crate::McpServerLaunch>,
5181 #[serde(default)]
5183 approval_policy: Option<String>,
5184}
5185
5186#[derive(Deserialize)]
5187struct RuntimeConnectionParams {
5188 connection: String,
5189}
5190
5191#[derive(Deserialize)]
5194struct RuntimeCloseParams {
5195 #[serde(default)]
5196 connection: String,
5197 #[serde(default)]
5198 runtime_id: Option<String>,
5199}
5200
5201#[derive(Deserialize)]
5202struct RuntimeInputParams {
5203 connection: String,
5204 text: String,
5205 #[serde(default)]
5206 image_urls: Vec<String>,
5207}
5208
5209const MAX_RUNTIME_IMAGES: usize = 4;
5210const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
5211const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
5212
5213fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
5214 if image_urls.len() > MAX_RUNTIME_IMAGES {
5215 return Err(ServiceError::InvalidParams(format!(
5216 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
5217 )));
5218 }
5219 let mut total = 0usize;
5220 for url in &image_urls {
5221 if !(url.starts_with("data:image/")
5222 || url.starts_with("https://")
5223 || url.starts_with("http://"))
5224 {
5225 return Err(ServiceError::InvalidParams(
5226 "runtime images must be image data URLs or HTTP(S) URLs".into(),
5227 ));
5228 }
5229 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
5230 return Err(ServiceError::InvalidParams(format!(
5231 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
5232 )));
5233 }
5234 total = total.saturating_add(url.len());
5235 }
5236 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
5237 return Err(ServiceError::InvalidParams(format!(
5238 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
5239 )));
5240 }
5241 Ok(image_urls)
5242}
5243
5244#[derive(Deserialize)]
5245struct RuntimeRespondParams {
5246 connection: String,
5247 request_id: Value,
5248 response: Value,
5249}
5250
5251fn default_reduction_store_root() -> PathBuf {
5252 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
5253 return PathBuf::from(root).join("sessions");
5254 }
5255 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
5256 return PathBuf::from(home).join(".supercode").join("sessions");
5257 }
5258 PathBuf::from(".supercode").join("sessions")
5259}
5260
5261fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
5262 let mut output = String::new();
5263 for message in messages {
5264 output.push_str(
5265 &serde_json::to_string(message)
5266 .map_err(|error| ServiceError::Operation(error.to_string()))?,
5267 );
5268 output.push('\n');
5269 }
5270 Ok(output)
5271}
5272
5273fn parse_messages_jsonl(
5274 content: &str,
5275) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
5276 content
5277 .lines()
5278 .enumerate()
5279 .filter(|(_, line)| !line.trim().is_empty())
5280 .map(|(index, line)| {
5281 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
5282 ServiceError::Operation(format!(
5283 "reduced transcript line {} is invalid: {error}",
5284 index + 1
5285 ))
5286 })
5287 })
5288 .collect()
5289}
5290
5291fn reduced_bootstrap_prompt(
5292 source: &SessionLocator,
5293 target: TransferFormat,
5294 view_jsonl: &str,
5295 sidecar_path: &Path,
5296 reduction_log_path: &Path,
5297) -> String {
5298 format!(
5299 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
5300 \n\
5301 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\
5302 \n\
5303 <supercode-reduced-session source-session=\"{source_id}\">\n\
5304 {view_jsonl}\
5305 </supercode-reduced-session>\n\
5306 \n\
5307 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
5308 source_harness = source.harness.as_str(),
5309 target_harness = target.id(),
5310 sidecar = sidecar_path.display(),
5311 log = reduction_log_path.display(),
5312 source_id = source.session_id,
5313 )
5314}
5315
5316fn session_artifact(
5317 locator: &SessionLocator,
5318 session: &Session,
5319 target: TransferFormat,
5320) -> std::result::Result<SessionArtifact, ServiceError> {
5321 session_artifact_with_id(locator, session, target, None)
5322}
5323
5324fn session_artifact_with_id(
5325 locator: &SessionLocator,
5326 session: &Session,
5327 target: TransferFormat,
5328 target_session_id: Option<&str>,
5329) -> std::result::Result<SessionArtifact, ServiceError> {
5330 let format: SessionFormat = target.into();
5331 let diagonal = format.source() == session.meta.source;
5332 crate::residue_store::store_segments(session);
5333 let has_appended_turns = session
5334 .imported_message_count
5335 .is_some_and(|imported| imported < session.messages.len());
5336 let mut restoration = None;
5337 let content = if let Some(id) = target_session_id {
5338 if diagonal && format != SessionFormat::OpenCode {
5339 session
5340 .to_jsonl_spliced(format, Some(id))
5341 .map_err(operation)?
5342 } else {
5343 let mut rewritten = session.clone();
5344 rewritten.meta.session_id = Some(id.to_string());
5345 rewritten.to_jsonl(format).map_err(operation)?
5346 }
5347 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
5348 session.raw_verbatim()
5349 } else if diagonal {
5350 session.to_jsonl_spliced(format, None).map_err(operation)?
5351 } else {
5352 match session
5355 .restore_residue(format, crate::residue_store::lookup)
5356 .map_err(operation)?
5357 {
5358 Some((content, report)) => {
5359 restoration = Some(report);
5360 content
5361 }
5362 None => session.to_jsonl(format).map_err(operation)?,
5363 }
5364 };
5365 let stem = sanitize_filename(
5366 target_session_id
5367 .or(session.meta.session_id.as_deref())
5368 .unwrap_or(&locator.session_id),
5369 );
5370 let suggested_filename = if diagonal && target == TransferFormat::Grok {
5371 "chat_history.jsonl".to_string()
5372 } else if target == TransferFormat::Goose {
5373 format!("{stem}.goose.json")
5374 } else {
5375 format!("{stem}.{}.jsonl", target.id())
5376 };
5377 let mut files = vec![SessionArtifactFile {
5378 path: suggested_filename.clone(),
5379 content: content.clone(),
5380 role: ArtifactFileRole::Primary,
5381 }];
5382 if target == TransferFormat::ClaudeCode {
5383 let bundle_stem = Path::new(&suggested_filename)
5384 .file_stem()
5385 .and_then(|stem| stem.to_str())
5386 .unwrap_or(&stem);
5387 let mut child_paths = BTreeSet::new();
5388 for (index, subagent) in session.subagents.iter().enumerate() {
5389 let agent_id = subagent
5390 .meta
5391 .agent_id
5392 .as_deref()
5393 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
5394 .map(sanitize_filename)
5395 .filter(|id| !id.is_empty())
5396 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5397 let child_has_appended_turns = subagent
5398 .imported_message_count
5399 .is_some_and(|imported| imported < subagent.messages.len());
5400 let child_content = if target_session_id.is_none()
5401 && subagent.meta.source == SessionSource::ClaudeCode
5402 && subagent.raw_is_verbatim
5403 && !child_has_appended_turns
5404 {
5405 subagent.raw_verbatim()
5406 } else if subagent.meta.source == SessionSource::ClaudeCode {
5407 subagent
5408 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
5409 .map_err(operation)?
5410 } else {
5411 let mut child = subagent.clone();
5412 if let Some(id) = target_session_id {
5413 child.meta.session_id = Some(id.to_string());
5414 }
5415 child
5416 .to_jsonl(SessionFormat::ClaudeCode)
5417 .map_err(operation)?
5418 };
5419 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
5420 if !child_paths.insert(path.clone()) {
5421 return Err(ServiceError::Operation(format!(
5422 "Claude subagent ids collide at artifact path `{path}`"
5423 )));
5424 }
5425 files.push(SessionArtifactFile {
5426 path,
5427 content: child_content,
5428 role: ArtifactFileRole::Subagent,
5429 });
5430 }
5431 }
5432 if diagonal && target == TransferFormat::Grok {
5433 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
5434 }
5435 if !diagonal || !session.raw_is_verbatim {
5436 files.push(SessionArtifactFile {
5437 path: "recovery/source.supercode.jsonl".into(),
5438 content: session.to_native_jsonl(),
5439 role: ArtifactFileRole::SourceRecovery,
5440 });
5441 for (index, subagent) in session.subagents.iter().enumerate() {
5442 let id = subagent
5443 .meta
5444 .agent_id
5445 .as_deref()
5446 .map(sanitize_filename)
5447 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5448 files.push(SessionArtifactFile {
5449 path: format!("recovery/subagents/{id}.supercode.jsonl"),
5450 content: subagent.to_native_jsonl(),
5451 role: ArtifactFileRole::SourceRecovery,
5452 });
5453 }
5454 }
5455 if !diagonal && session.meta.source == SessionSource::Grok {
5456 append_grok_bundle_files(
5457 locator,
5458 "recovery/grok/",
5459 ArtifactFileRole::SourceRecovery,
5460 &mut files,
5461 )?;
5462 }
5463 let (fidelity, residue) = if diagonal
5464 && target_session_id.is_none()
5465 && session.raw_is_verbatim
5466 && !has_appended_turns
5467 {
5468 (Fidelity::ByteLossless, Vec::new())
5469 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
5470 (
5471 Fidelity::ValueLossless,
5472 vec![if target_session_id.is_some() {
5473 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
5474 } else {
5475 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
5476 }],
5477 )
5478 } else {
5479 match restoration {
5480 Some(report) if report.rendered_messages == 0 => (
5481 Fidelity::ByteLossless,
5482 vec![format!(
5483 "restored verbatim from this conversation's {} source records in the residue store",
5484 target.id()
5485 )],
5486 ),
5487 Some(report) => (
5488 Fidelity::Semantic,
5489 vec![format!(
5490 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
5491 report.restored_messages,
5492 report.restored_messages + report.rendered_messages,
5493 report.rendered_messages,
5494 target.id()
5495 )],
5496 ),
5497 None => (
5498 Fidelity::Semantic,
5499 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
5500 ),
5501 }
5502 };
5503 Ok(SessionArtifact {
5504 source_harness: locator.harness.clone(),
5505 target_harness: target.id(),
5506 session_id: target_session_id
5507 .map(str::to_string)
5508 .or_else(|| session.meta.session_id.clone()),
5509 content,
5510 suggested_filename,
5511 files,
5512 fidelity,
5513 residue,
5514 })
5515}
5516
5517fn append_grok_bundle_files(
5518 locator: &SessionLocator,
5519 prefix: &str,
5520 role: ArtifactFileRole,
5521 files: &mut Vec<SessionArtifactFile>,
5522) -> std::result::Result<(), ServiceError> {
5523 let primary = locator.storage.path();
5524 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5525 return Err(ServiceError::Operation(format!(
5526 "Grok bundle locator must name chat_history.jsonl, got {}",
5527 primary.display()
5528 )));
5529 }
5530 let parent = primary.parent().ok_or_else(|| {
5531 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5532 })?;
5533 for name in ["summary.json", "updates.jsonl"] {
5534 let path = parent.join(name);
5535 let metadata = match std::fs::symlink_metadata(&path) {
5536 Ok(metadata) => metadata,
5537 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5538 Err(error) => return Err(ServiceError::Operation(error.to_string())),
5539 };
5540 if metadata.file_type().is_symlink() || !metadata.is_file() {
5541 return Err(ServiceError::Operation(format!(
5542 "refusing non-regular Grok bundle member {}",
5543 path.display()
5544 )));
5545 }
5546 let content = std::fs::read_to_string(&path).map_err(|error| {
5547 ServiceError::Operation(format!(
5548 "Grok bundle member {} is not representable as UTF-8: {error}",
5549 path.display()
5550 ))
5551 })?;
5552 files.push(SessionArtifactFile {
5553 path: format!("{prefix}{name}"),
5554 content,
5555 role: match role {
5556 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5557 _ => ArtifactFileRole::SourceRecovery,
5558 },
5559 });
5560 }
5561 Ok(())
5562}
5563
5564fn handoff_artifact(
5565 locator: &SessionLocator,
5566 session: &Session,
5567 target: TransferFormat,
5568) -> std::result::Result<SessionArtifact, ServiceError> {
5569 let target_session_id = target_session_id(target);
5570 session_artifact_with_id(locator, session, target, Some(&target_session_id))
5571}
5572
5573fn target_session_id(target: TransferFormat) -> String {
5574 let uuid = generated_session_id();
5575 match target {
5576 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5577 TransferFormat::ClaudeCode
5578 | TransferFormat::Codex
5579 | TransferFormat::Pi
5580 | TransferFormat::Grok
5581 | TransferFormat::Gemini
5582 | TransferFormat::Goose
5583 | TransferFormat::Hermes => uuid,
5584 }
5585}
5586
5587fn sanitize_filename(value: &str) -> String {
5588 let value = value
5589 .chars()
5590 .map(|character| {
5591 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5592 character
5593 } else {
5594 '-'
5595 }
5596 })
5597 .collect::<String>();
5598 let value = value.trim_matches('-');
5599 if value.is_empty() {
5600 "session".into()
5601 } else {
5602 value.chars().take(100).collect()
5603 }
5604}
5605
5606fn handoff_instructions(
5607 target: TransferFormat,
5608 session_id: &str,
5609 cwd: &Path,
5610) -> HandoffInstructions {
5611 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5612 cwd: cwd.to_path_buf(),
5613 program: program.into(),
5614 arguments,
5615 env: BTreeMap::new(),
5616 };
5617 match target {
5618 TransferFormat::ClaudeCode => HandoffInstructions {
5619 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5620 materialize: None,
5621 requires_materialization: true,
5622 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(),
5623 },
5624 TransferFormat::Hermes => HandoffInstructions {
5625 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5626 materialize: None,
5627 requires_materialization: true,
5628 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(),
5629 },
5630 TransferFormat::Codex => HandoffInstructions {
5631 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5632 materialize: None,
5633 requires_materialization: true,
5634 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5635 },
5636 TransferFormat::OpenCode => HandoffInstructions {
5637 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5638 materialize: Some(launch(
5639 "opencode",
5640 vec!["import".into(), "{artifact_path}".into()],
5641 )),
5642 requires_materialization: true,
5643 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5644 },
5645 TransferFormat::Pi => HandoffInstructions {
5646 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5647 materialize: None,
5648 requires_materialization: true,
5649 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5650 },
5651 TransferFormat::Grok => HandoffInstructions {
5652 launch: launch(
5653 "grok",
5654 vec!["--resume".into(), "{materialized_session_id}".into()],
5655 ),
5656 materialize: None,
5657 requires_materialization: true,
5658 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(),
5659 },
5660 TransferFormat::Gemini => HandoffInstructions {
5661 launch: launch(
5662 "gemini",
5663 vec!["--session-file".into(), "{artifact_path}".into()],
5664 ),
5665 materialize: None,
5666 requires_materialization: true,
5667 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(),
5668 },
5669 TransferFormat::Goose => HandoffInstructions {
5670 launch: launch(
5671 "goose",
5672 vec![
5673 "session".into(),
5674 "--resume".into(),
5675 "--session-id".into(),
5676 "{imported_session_id}".into(),
5677 ],
5678 ),
5679 materialize: Some(launch(
5680 "goose",
5681 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5682 )),
5683 requires_materialization: true,
5684 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(),
5685 },
5686 }
5687}
5688
5689fn resume_launch(
5690 harness: &str,
5691 session_id: &str,
5692 cwd: &Path,
5693 policy: ResumePolicy,
5694) -> std::result::Result<StructuredLaunch, ServiceError> {
5695 let mut arguments = Vec::new();
5696 let program = match harness {
5697 HarnessId::GROK => {
5698 if matches!(policy, ResumePolicy::Yolo) {
5699 if crate::support::self_sandbox_supported() {
5700 arguments.extend(["--sandbox".into(), "workspace".into()]);
5701 }
5702 arguments.push("--always-approve".into());
5703 }
5704 arguments.extend(["--resume".into(), session_id.into()]);
5705 "grok"
5706 }
5707 HarnessId::CODEX => {
5708 arguments.extend(crate::startup_prompts::startup_arguments(
5709 harness,
5710 Some(cwd),
5711 &[],
5712 matches!(policy, ResumePolicy::Yolo),
5713 ));
5714 arguments.extend(["resume".into(), session_id.into()]);
5715 "codex"
5716 }
5717 HarnessId::CLAUDE_CODE => {
5718 arguments.extend(crate::startup_prompts::startup_arguments(
5719 harness,
5720 Some(cwd),
5721 &[],
5722 matches!(policy, ResumePolicy::Yolo),
5723 ));
5724 arguments.extend(["--resume".into(), session_id.into()]);
5725 "claude"
5726 }
5727 HarnessId::GEMINI => {
5728 arguments.extend(crate::startup_prompts::startup_arguments(
5729 harness,
5730 Some(cwd),
5731 &[],
5732 matches!(policy, ResumePolicy::Yolo),
5733 ));
5734 arguments.extend(["--resume".into(), session_id.into()]);
5735 "gemini"
5736 }
5737 HarnessId::GOOSE => {
5738 arguments.extend([
5739 "session".into(),
5740 "--resume".into(),
5741 "--session-id".into(),
5742 session_id.into(),
5743 ]);
5744 "goose"
5745 }
5746 HarnessId::PI => {
5747 arguments.extend(crate::startup_prompts::startup_arguments(
5748 harness,
5749 Some(cwd),
5750 &[],
5751 matches!(policy, ResumePolicy::Yolo),
5752 ));
5753 arguments.extend(["--session".into(), session_id.into()]);
5754 "pi"
5755 }
5756 HarnessId::OPENCODE => {
5757 arguments.extend(["--session".into(), session_id.into()]);
5758 "opencode"
5759 }
5760 HarnessId::SUPERCODE => {
5761 if matches!(policy, ResumePolicy::Yolo) {
5762 arguments.push("--dangerous".into());
5763 }
5764 arguments.extend(["resume".into(), session_id.into()]);
5765 "supercode"
5766 }
5767 other => {
5768 return Err(ServiceError::InvalidParams(format!(
5769 "no structured resume launch is registered for harness `{other}`"
5770 )))
5771 }
5772 };
5773 Ok(StructuredLaunch {
5774 cwd: cwd.to_path_buf(),
5775 env: if program == "grok" {
5776 crate::support::grok_home_env()
5777 } else {
5778 BTreeMap::new()
5779 },
5780 program: program.into(),
5781 arguments,
5782 })
5783}
5784
5785fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5791 let digest = blake3::hash(address.as_bytes()).to_hex();
5792 let path = std::env::temp_dir().join(format!(
5793 "supercode-openclaw-gateway-token-{}",
5794 &digest.as_str()[..16]
5795 ));
5796 #[cfg(unix)]
5797 {
5798 use std::io::Write;
5799 use std::os::unix::fs::OpenOptionsExt;
5800 let mut file = std::fs::OpenOptions::new()
5801 .write(true)
5802 .create(true)
5803 .truncate(true)
5804 .mode(0o600)
5805 .open(&path)?;
5806 file.write_all(secret.as_bytes())?;
5807 }
5808 #[cfg(not(unix))]
5809 std::fs::write(&path, secret)?;
5810 Ok(path)
5811}
5812
5813fn open_connect_descriptor(
5819 descriptor: &crate::HarnessSupportDescriptor,
5820 home: &Path,
5821) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5822 let Some(connect) = &descriptor.runtime.connect_launch else {
5823 return Err(ServiceError::InvalidParams(format!(
5824 "harness `{}` has no registered connect-mode launch",
5825 descriptor.id.as_str()
5826 )));
5827 };
5828 let resolved = connect
5829 .resolve(home)
5830 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5831 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5832 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5833 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5834 if let Some(token) = resolved.auth {
5835 backend = backend.with_bearer(token);
5836 }
5837 Ok(Box::new(backend))
5838 }
5839 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5840 let mut env = BTreeMap::new();
5851 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5852 if let Some(token) = resolved.auth {
5853 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5854 .map_err(|error| {
5855 ServiceError::UnsupportedAction(format!(
5856 "could not stage the gateway credential for the bridge: {error}"
5857 ))
5858 })?;
5859 arguments.push("--token-file".into());
5860 arguments.push(token_path.to_string_lossy().into_owned());
5861 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5862 }
5863 let program = descriptor
5868 .runtime
5869 .default_launch
5870 .as_ref()
5871 .map(|launch| launch.program.clone())
5872 .unwrap_or_else(|| "openclaw".into());
5873 let launch = RuntimeLaunch {
5874 program,
5875 arguments,
5876 env,
5877 };
5878 Ok(Box::new(
5879 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5880 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5881 ))
5882 }
5883 _ => Err(ServiceError::UnsupportedAction(format!(
5884 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5885 descriptor.id.as_str(),
5886 connect.protocol
5887 ))),
5888 }
5889}
5890
5891fn registry_connect_descriptor(
5894 params: &RuntimeBackendParams,
5895) -> Option<crate::HarnessSupportDescriptor> {
5896 if params.launch.is_some() || params.base_url.is_some() {
5897 return None;
5898 }
5899 harness_support_registry()
5900 .harnesses
5901 .into_iter()
5902 .find(|descriptor| descriptor.id == params.harness)
5903 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5904}
5905
5906fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5907 supercode_interchange::user_home()
5908 .map(std::path::PathBuf::into_os_string)
5909 .map(PathBuf::from)
5910 .ok_or_else(|| {
5911 ServiceError::UnsupportedAction(
5912 "connect-mode launches need HOME to locate the harness config".into(),
5913 )
5914 })
5915}
5916
5917pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5920 "harness.v1.runtimes.start",
5921 "harness.v1.runtimes.resume",
5922 "harness.v1.runtimes.attach",
5923 "harness.v1.runtimes.attach_existing",
5924];
5925
5926pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5931
5932pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5938
5939pub const DETACHED_METHODS: &[&str] = &[
5945 "harness.v1.harnesses.list",
5946 "harness.v1.harnesses.probe",
5947 "harness.v1.sessions.message",
5948 "harness.v1.sessions.new",
5949 "harness.v1.sessions.reset",
5950 "harness.v1.sessions.archive",
5951 "harness.v1.sessions.delete",
5952];
5953
5954pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5960
5961pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5969
5970async fn within_control_deadline<F: std::future::Future>(
5973 method: &str,
5974 call: F,
5975) -> std::result::Result<F::Output, ServiceError> {
5976 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5977 .await
5978 .map_err(|_| {
5979 ServiceError::Operation(format!(
5980 "`{method}` gave up after {}s: the runtime did not answer",
5981 RUNTIME_CONTROL_DEADLINE.as_secs()
5982 ))
5983 })
5984}
5985
5986pub struct RuntimeOpen {
5990 id: Value,
5991 method: String,
5992 params: Value,
5993}
5994
5995impl RuntimeOpen {
5996 pub async fn open(self) -> OpenedRuntime {
6000 let Self { id, method, params } = self;
6001 let outcome = open_runtime(&method, params).await;
6002 OpenedRuntime { id, outcome }
6003 }
6004}
6005
6006pub struct OpenedRuntime {
6009 id: Value,
6010 outcome: std::result::Result<OpenRuntime, ServiceError>,
6011}
6012
6013pub struct DetachedCall {
6018 id: Value,
6019 method: String,
6020 work: std::result::Result<Work, ServiceError>,
6021}
6022
6023impl DetachedCall {
6024 pub async fn run(self) -> DetachedAnswer {
6027 let Self { id, method, work } = self;
6028 match work {
6029 Ok(Work::Runtime(work)) => {
6034 let (result, returned) = work.run().await;
6035 DetachedAnswer {
6036 response: service_response(id, result),
6037 returned,
6038 }
6039 }
6040 Ok(Work::Free(work)) => {
6041 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
6042 Ok(result) => result,
6043 Err(_) => Err(ServiceError::Operation(format!(
6044 "`{method}` gave up after {}s: the harness it waits on did not answer",
6045 DETACHED_CALL_DEADLINE.as_secs()
6046 ))),
6047 };
6048 DetachedAnswer {
6049 response: service_response(id, result),
6050 returned: None,
6051 }
6052 }
6053 Err(error) => DetachedAnswer {
6054 response: service_response(id, Err(error)),
6055 returned: None,
6056 },
6057 }
6058 }
6059}
6060
6061pub struct DetachedAnswer {
6065 response: Value,
6066 returned: Option<ReturnedRuntime>,
6067}
6068
6069impl DetachedAnswer {
6070 pub fn into_response(self) -> Value {
6073 self.response
6074 }
6075}
6076
6077pub struct ReturnedRuntime {
6080 connection: String,
6081 runtime: Box<dyn RuntimeConnection>,
6082}
6083
6084enum Work {
6087 Free(DetachedWork),
6088 Runtime(RuntimeWork),
6089}
6090
6091enum DetachedWork {
6094 Inventory(InventoryWork),
6098 Message(MessageSessionParams),
6100 SessionMutation {
6103 verb: crate::SessionVerb,
6104 mutation: crate::SessionMutation,
6105 },
6106}
6107
6108impl DetachedWork {
6109 async fn run(self) -> std::result::Result<Value, ServiceError> {
6110 match self {
6111 Self::Inventory(work) => run_inventory(work).await,
6112 Self::Message(params) => Ok(message_live_session(¶ms).await),
6113 Self::SessionMutation { verb, mutation } => {
6114 let outcome = run_session_mutation(verb, &mutation).await?;
6115 serde_json::to_value(outcome)
6116 .map_err(|error| ServiceError::Operation(error.to_string()))
6117 }
6118 }
6119 }
6120}
6121
6122enum RuntimeWork {
6124 Close {
6126 runtime: Box<dyn RuntimeConnection>,
6127 process_group: Option<u32>,
6128 },
6129 LiveCommand {
6132 connection: String,
6133 runtime: Box<dyn RuntimeConnection>,
6134 verb: crate::SessionVerb,
6135 mutation: crate::SessionMutation,
6136 command: &'static str,
6137 session: String,
6138 },
6139}
6140
6141type RuntimeWorkAnswer = (
6144 std::result::Result<Value, ServiceError>,
6145 Option<ReturnedRuntime>,
6146);
6147
6148impl RuntimeWork {
6149 async fn run(self) -> RuntimeWorkAnswer {
6150 match self {
6151 Self::Close {
6152 runtime,
6153 process_group,
6154 } => (close_runtime(runtime, process_group).await, None),
6155 Self::LiveCommand {
6156 connection,
6157 mut runtime,
6158 verb,
6159 mutation,
6160 command,
6161 session,
6162 } => {
6163 let result =
6164 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
6165 (
6166 result,
6167 Some(ReturnedRuntime {
6168 connection,
6169 runtime,
6170 }),
6171 )
6172 }
6173 }
6174 }
6175}
6176
6177async fn close_runtime(
6180 mut runtime: Box<dyn RuntimeConnection>,
6181 process_group: Option<u32>,
6182) -> std::result::Result<Value, ServiceError> {
6183 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
6184 Ok(result) => {
6185 result.map_err(operation)?;
6186 Ok(json!({"closed": true}))
6187 }
6188 Err(deadline) => {
6189 let killed = kill_runtime_process_group(process_group);
6194 drop(runtime);
6195 Ok(json!({
6196 "closed": true,
6197 "killed": killed,
6198 "detail": error_message(deadline),
6199 }))
6200 }
6201 }
6202}
6203
6204fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
6207 mutation
6208 .session
6209 .clone()
6210 .filter(|value| !value.trim().is_empty())
6211 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
6212}
6213
6214async fn type_live_command(
6218 runtime: &mut dyn RuntimeConnection,
6219 verb: crate::SessionVerb,
6220 mutation: &crate::SessionMutation,
6221 command: &str,
6222 session: String,
6223) -> std::result::Result<Value, ServiceError> {
6224 within_control_deadline(
6225 &format!("sessions.{}", verb.as_str()),
6226 runtime.send_input(RuntimeInput {
6227 text: command.to_string(),
6228 image_urls: Vec::new(),
6229 }),
6230 )
6231 .await?
6232 .map_err(operation)?;
6233 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
6234 .map_err(session_control_error)?;
6235 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
6236}
6237
6238enum OpenRuntime {
6241 Hosted {
6244 runtime: Box<dyn RuntimeConnection>,
6245 capabilities: crate::RuntimeCapabilities,
6246 workspace: PathBuf,
6247 fresh: bool,
6249 },
6250 Joined { runtime: Box<dyn RuntimeConnection> },
6253}
6254
6255async fn open_runtime(
6260 method: &str,
6261 params: Value,
6262) -> std::result::Result<OpenRuntime, ServiceError> {
6263 match tokio::time::timeout(
6264 RUNTIME_OPEN_DEADLINE,
6265 open_runtime_unbounded(method, params),
6266 )
6267 .await
6268 {
6269 Ok(result) => result,
6270 Err(_) => Err(ServiceError::Operation(format!(
6271 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
6272 RUNTIME_OPEN_DEADLINE.as_secs()
6273 ))),
6274 }
6275}
6276
6277async fn open_runtime_unbounded(
6278 method: &str,
6279 params: Value,
6280) -> std::result::Result<OpenRuntime, ServiceError> {
6281 match method {
6282 "harness.v1.runtimes.start" => {
6283 let params = decode::<RuntimeStartParams>(params)?;
6284 let backend = runtime_backend(¶ms.backend)?;
6285 let capabilities = backend.capabilities();
6286 let workspace = params.cwd.clone();
6287 let runtime = backend
6288 .start(RuntimeStartRequest {
6289 cwd: params.cwd,
6290 launch: runtime_launch(¶ms.backend),
6291 mcp_servers: params.mcp_servers,
6292 approval_policy: params.approval_policy,
6293 })
6294 .await
6295 .map_err(operation)?;
6296 Ok(OpenRuntime::Hosted {
6297 runtime,
6298 capabilities,
6299 workspace,
6300 fresh: true,
6301 })
6302 }
6303 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
6304 let params = decode::<RuntimeAttachParams>(params)?;
6305 let backend = runtime_backend(¶ms.backend)?;
6306 let capabilities = backend.capabilities();
6307 let workspace = params
6308 .cwd
6309 .clone()
6310 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
6311 let runtime = backend
6312 .attach(RuntimeAttachRequest {
6313 runtime_id: params.runtime_id,
6314 cwd: params.cwd,
6315 launch: runtime_launch(¶ms.backend),
6316 mcp_servers: params.mcp_servers,
6317 approval_policy: params.approval_policy,
6318 })
6319 .await
6320 .map_err(operation)?;
6321 Ok(OpenRuntime::Hosted {
6322 runtime,
6323 capabilities,
6324 workspace,
6325 fresh: false,
6326 })
6327 }
6328 "harness.v1.runtimes.attach_existing" => {
6329 let params = decode::<RuntimeAttachParams>(params)?;
6330 let backend: Box<dyn RuntimeBackend> = match params
6331 .backend
6332 .base_url
6333 .as_deref()
6334 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
6335 {
6336 Some(endpoint) => {
6337 #[cfg(not(feature = "adapter-api"))]
6338 {
6339 let _ = endpoint;
6340 return Err(ServiceError::UnsupportedAction(
6341 "live HTTP attachment adapter is not compiled".into(),
6342 ));
6343 }
6344 #[cfg(feature = "adapter-api")]
6345 {
6346 let workspace = params.cwd.clone().ok_or_else(|| {
6347 ServiceError::InvalidParams(
6348 "Volter Harness live attach requires the project cwd".into(),
6349 )
6350 })?;
6351 let source = LiveRuntimeSource {
6352 harness: params.backend.harness.as_str().to_string(),
6353 session_id: params.runtime_id.clone(),
6354 workspace,
6355 };
6356 let receipt = resolve_live_runtime(&endpoint, &source)
6357 .map_err(|error| ServiceError::Operation(error.to_string()))?;
6358 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
6359 }
6360 }
6361 None => runtime_backend(¶ms.backend)?,
6362 };
6363 let capabilities = backend.capabilities();
6364 if !capabilities.attach_existing_process {
6365 return Err(ServiceError::Operation(format!(
6366 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
6367 backend.harness().as_str()
6368 )));
6369 }
6370 let runtime = backend
6371 .attach_existing(RuntimeAttachRequest {
6372 runtime_id: params.runtime_id,
6373 cwd: params.cwd,
6374 launch: runtime_launch(¶ms.backend),
6375 mcp_servers: params.mcp_servers,
6376 approval_policy: params.approval_policy,
6377 })
6378 .await
6379 .map_err(operation)?;
6380 Ok(OpenRuntime::Joined { runtime })
6381 }
6382 _ => Err(ServiceError::MethodNotFound),
6383 }
6384}
6385
6386fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
6388 match result {
6389 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
6390 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
6391 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
6392 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
6393 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
6394 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
6395 }
6396}
6397
6398fn runtime_backend(
6399 params: &RuntimeBackendParams,
6400) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
6401 if let Some(descriptor) = registry_connect_descriptor(params) {
6402 return open_connect_descriptor(&descriptor, &service_home()?);
6403 }
6404 if params.protocol.as_deref() == Some("acp") {
6405 let launch = params
6406 .launch
6407 .clone()
6408 .or_else(|| {
6409 harness_support_registry()
6410 .harnesses
6411 .into_iter()
6412 .find(|harness| harness.id == params.harness)
6413 .filter(|harness| {
6414 harness.runtime.implementation == ImplementationKind::GenericProtocol
6415 && harness.runtime.protocol.starts_with("acp")
6416 })
6417 .and_then(|harness| harness.runtime.default_launch)
6418 })
6419 .ok_or_else(|| {
6420 ServiceError::InvalidParams(
6421 "an ACP runtime requires `launch` unless the harness has a registered default"
6422 .into(),
6423 )
6424 })?;
6425 let resume_session = harness_support_registry()
6426 .harnesses
6427 .into_iter()
6428 .find(|harness| harness.id == params.harness)
6429 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
6430 return Ok(Box::new(
6431 AcpRuntimeBackend::new(params.harness.clone(), launch)
6432 .with_resume_support(resume_session),
6433 ));
6434 }
6435 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
6436 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
6437 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
6438 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
6439 HarnessId::OPENCODE => match ¶ms.base_url {
6440 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
6441 None => Box::new(OpenCodeRuntimeBackend::new()),
6442 },
6443 harness => {
6444 let descriptor = harness_support_registry()
6445 .harnesses
6446 .into_iter()
6447 .find(|descriptor| descriptor.id.as_str() == harness)
6448 .filter(|descriptor| {
6449 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
6450 && descriptor.runtime.protocol.starts_with("acp")
6451 });
6452 let Some(descriptor) = descriptor else {
6453 return Err(ServiceError::InvalidParams(format!(
6454 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
6455 )));
6456 };
6457 let resume = descriptor.runtime.capabilities.resume_session;
6458 Box::new(
6459 AcpRuntimeBackend::new(
6460 descriptor.id,
6461 descriptor
6462 .runtime
6463 .default_launch
6464 .expect("generic ACP registry entry includes its launch"),
6465 )
6466 .with_resume_support(resume),
6467 )
6468 }
6469 };
6470 Ok(backend)
6471}
6472
6473fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
6474 if let Some(launch) = ¶ms.launch {
6475 return Some(launch.clone());
6476 }
6477 if !matches!(params.policy, RuntimePolicy::Yolo) {
6478 return None;
6479 }
6480 let launch = match params.harness.as_str() {
6481 HarnessId::GROK => RuntimeLaunch {
6482 program: "grok".into(),
6483 arguments: {
6484 let mut arguments: Vec<String> = Vec::new();
6485 if crate::support::self_sandbox_supported() {
6486 arguments.extend(["--sandbox".into(), "workspace".into()]);
6487 }
6488 arguments.extend([
6489 "--always-approve".into(),
6490 "agent".into(),
6491 "--no-leader".into(),
6492 "stdio".into(),
6493 ]);
6494 arguments
6495 },
6496 env: crate::support::grok_env(),
6497 },
6498 HarnessId::CODEX => RuntimeLaunch {
6499 program: "codex".into(),
6500 arguments: vec![
6501 "--dangerously-bypass-approvals-and-sandbox".into(),
6502 "--dangerously-bypass-hook-trust".into(),
6503 "app-server".into(),
6504 ],
6505 env: BTreeMap::new(),
6506 },
6507 HarnessId::CLAUDE_CODE => RuntimeLaunch {
6508 program: "claude".into(),
6509 arguments: vec![
6510 "--dangerously-skip-permissions".into(),
6511 "--print".into(),
6512 "--input-format".into(),
6513 "stream-json".into(),
6514 "--output-format".into(),
6515 "stream-json".into(),
6516 "--verbose".into(),
6517 ],
6518 env: BTreeMap::new(),
6519 },
6520 HarnessId::PI => RuntimeLaunch {
6521 program: "pi".into(),
6522 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
6523 env: BTreeMap::new(),
6524 },
6525 HarnessId::OPENCODE => RuntimeLaunch {
6526 program: "opencode".into(),
6527 arguments: vec!["serve".into()],
6528 env: BTreeMap::new(),
6529 },
6530 HarnessId::GEMINI => RuntimeLaunch {
6531 program: "gemini".into(),
6532 arguments: vec!["--acp".into(), "--yolo".into()],
6533 env: BTreeMap::new(),
6534 },
6535 HarnessId::GOOSE => RuntimeLaunch {
6536 program: "goose".into(),
6537 arguments: vec!["acp".into()],
6538 env: BTreeMap::new(),
6539 },
6540 HarnessId::SUPERCODE => RuntimeLaunch {
6541 program: "supercode".into(),
6542 arguments: vec!["acp".into(), "--dangerous".into()],
6543 env: BTreeMap::new(),
6544 },
6545 _ => return None,
6546 };
6547 Some(launch)
6548}
6549
6550struct IsolatedProbeHome {
6556 launch: RuntimeLaunch,
6557 root: PathBuf,
6558}
6559
6560impl IsolatedProbeHome {
6561 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6562 let root = std::env::temp_dir().join(format!(
6563 "supercode-harness-probe-{harness}-{}",
6564 generated_session_id()
6565 ));
6566 std::fs::create_dir_all(&root)?;
6567 set_private_dir_permissions(&root)?;
6568
6569 if let Some(source_home) = supercode_interchange::user_home()
6570 .map(std::path::PathBuf::into_os_string)
6571 .map(PathBuf::from)
6572 {
6573 for relative in probe_auth_files(harness) {
6574 copy_probe_file(&source_home, &root, relative)?;
6575 }
6576 }
6577 if harness == HarnessId::SUPERCODE {
6582 let config_home = crate::agent::global_instructions_dir();
6583 for file in ["config.toml", "credentials.toml"] {
6584 copy_probe_path(
6585 &config_home.join(file),
6586 &root.join(".config/supercode").join(file),
6587 )?;
6588 }
6589 }
6590 configure_isolated_probe_auth(harness, &root)?;
6591
6592 let root_text = root.to_string_lossy().into_owned();
6593 for (key, value) in [
6594 ("HOME", root_text.clone()),
6595 (
6596 "XDG_CACHE_HOME",
6597 root.join(".cache").to_string_lossy().into_owned(),
6598 ),
6599 (
6600 "XDG_CONFIG_HOME",
6601 root.join(".config").to_string_lossy().into_owned(),
6602 ),
6603 (
6604 "XDG_DATA_HOME",
6605 root.join(".local/share").to_string_lossy().into_owned(),
6606 ),
6607 ] {
6608 launch.env.insert(key.into(), value);
6609 }
6610 let scoped = match harness {
6611 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6612 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6613 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6614 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6615 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6616 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6617 _ => None,
6618 };
6619 if let Some((key, value)) = scoped {
6620 launch
6621 .env
6622 .insert(key.into(), value.to_string_lossy().into_owned());
6623 }
6624 Ok(Self { launch, root })
6625 }
6626
6627 fn cleanup(&self) -> std::io::Result<()> {
6628 match std::fs::remove_dir_all(&self.root) {
6629 Ok(()) => Ok(()),
6630 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6631 Err(error) => Err(error),
6632 }
6633 }
6634}
6635
6636impl Drop for IsolatedProbeHome {
6637 fn drop(&mut self) {
6638 let _ = self.cleanup();
6639 }
6640}
6641
6642fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6643 match harness {
6644 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6645 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6649 HarnessId::CODEX => &[".codex/auth.json"],
6650 HarnessId::GEMINI => &[
6651 ".gemini/google_accounts.json",
6652 ".gemini/oauth_creds.json",
6653 ".gemini/settings.json",
6654 ],
6655 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6656 HarnessId::OPENCODE => &[
6657 ".config/opencode/auth.json",
6658 ".local/share/opencode/auth.json",
6659 ],
6660 HarnessId::PI => &[".pi/agent/auth.json"],
6661 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6666 _ => &[],
6667 }
6668}
6669
6670fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6671 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6672}
6673
6674fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6675 if !source.is_file() {
6676 return Ok(());
6677 }
6678 if let Some(parent) = destination.parent() {
6679 std::fs::create_dir_all(parent)?;
6680 set_private_dir_permissions(parent)?;
6681 }
6682 std::fs::copy(source, destination)?;
6683 set_private_file_permissions(destination)
6684}
6685
6686fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6687 if harness != HarnessId::GEMINI {
6688 return Ok(());
6689 }
6690 let oauth = probe_home.join(".gemini/oauth_creds.json");
6691 if !oauth.is_file() {
6692 return Ok(());
6693 }
6694 let settings_path = probe_home.join(".gemini/settings.json");
6695 let mut settings = std::fs::read_to_string(&settings_path)
6696 .ok()
6697 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6698 .unwrap_or_else(|| json!({}));
6699 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6700 std::fs::write(
6701 &settings_path,
6702 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6703 )?;
6704 set_private_file_permissions(&settings_path)
6705}
6706
6707#[cfg(unix)]
6708fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6709 use std::os::unix::fs::PermissionsExt;
6710 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6711}
6712
6713#[cfg(not(unix))]
6714fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6715 Ok(())
6716}
6717
6718#[cfg(unix)]
6719fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6720 use std::os::unix::fs::PermissionsExt;
6721 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6722}
6723
6724#[cfg(not(unix))]
6725fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6726 Ok(())
6727}
6728
6729fn find_executable(program: &str) -> Option<PathBuf> {
6730 let candidate = PathBuf::from(program);
6731 if candidate.components().count() > 1 {
6732 return candidate.is_file().then_some(candidate);
6733 }
6734 let path = std::env::var_os("PATH")?;
6735 for directory in std::env::split_paths(&path) {
6736 let candidate = directory.join(program);
6737 if candidate.is_file() {
6738 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6739 }
6740 #[cfg(windows)]
6741 {
6742 for extension in ["exe", "cmd", "bat"] {
6743 let candidate = directory.join(format!("{program}.{extension}"));
6744 if candidate.is_file() {
6745 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6746 }
6747 }
6748 }
6749 }
6750 None
6751}
6752
6753async fn executable_version(executable: &Path) -> Option<String> {
6754 let mut command = tokio::process::Command::new(executable);
6755 command
6756 .arg("--version")
6757 .stdin(std::process::Stdio::null())
6758 .stdout(std::process::Stdio::piped())
6759 .stderr(std::process::Stdio::piped())
6760 .kill_on_drop(true);
6761 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6762 .await
6763 .ok()?
6764 .ok()?;
6765 let stdout = String::from_utf8_lossy(&output.stdout);
6766 let stderr = String::from_utf8_lossy(&output.stderr);
6767 stdout
6768 .lines()
6769 .chain(stderr.lines())
6770 .map(str::trim)
6771 .find(|line| !line.is_empty())
6772 .map(|line| truncate_text(line, 200))
6773}
6774
6775pub(crate) fn auth_evidence(harness: &str) -> bool {
6776 let env_names: &[&str] = match harness {
6777 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6778 HarnessId::CODEX => &["OPENAI_API_KEY"],
6779 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6780 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6781 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6782 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6783 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6784 _ => &[],
6785 };
6786 if env_names
6787 .iter()
6788 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6789 {
6790 return true;
6791 }
6792 let Some(home) = supercode_interchange::user_home()
6793 .map(std::path::PathBuf::into_os_string)
6794 .map(PathBuf::from)
6795 else {
6796 return false;
6797 };
6798 let files: Vec<PathBuf> = match harness {
6799 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6800 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6801 HarnessId::OPENCODE => vec![
6802 home.join(".local/share/opencode/auth.json"),
6803 home.join(".config/opencode/auth.json"),
6804 ],
6805 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6806 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6807 HarnessId::GEMINI => vec![
6808 home.join(".gemini/oauth_creds.json"),
6809 home.join(".gemini/google_accounts.json"),
6810 ],
6811 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6812 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6813 _ => Vec::new(),
6814 };
6815 if files.into_iter().any(|path| {
6816 std::fs::metadata(path)
6817 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6818 .unwrap_or(false)
6819 }) {
6820 return true;
6821 }
6822 if harness == HarnessId::CLAUDE_CODE {
6829 return std::fs::read_to_string(home.join(".claude.json"))
6830 .map(|text| text.contains("\"oauthAccount\""))
6831 .unwrap_or(false);
6832 }
6833 false
6834}
6835
6836fn looks_like_auth_error(message: &str) -> bool {
6837 let message = message.to_ascii_lowercase();
6838 [
6839 "auth",
6840 "login",
6841 "sign in",
6842 "sign-in",
6843 "credential",
6844 "unauthorized",
6845 "forbidden",
6846 "token",
6847 ]
6848 .iter()
6849 .any(|needle| message.contains(needle))
6850}
6851
6852fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6853 crate::RuntimeCapabilities {
6854 start_session: false,
6855 resume_session: false,
6856 attach_existing_process: false,
6857 send_input: false,
6858 stream_events: false,
6859 interrupt: false,
6860 steer: false,
6861 respond_to_requests: false,
6862 }
6863}
6864
6865fn truncate_text(text: &str, max_chars: usize) -> String {
6866 let mut chars = text.chars();
6867 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6868 if chars.next().is_some() {
6869 format!("{truncated}…")
6870 } else {
6871 truncated
6872 }
6873}
6874
6875fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6882 match &handle.endpoint {
6883 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6884 crate::RuntimeEndpoint::Http { .. } => None,
6885 }
6886}
6887
6888fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6894 match process_group {
6895 #[cfg(unix)]
6896 Some(pid) => {
6897 crate::lsp::kill_process_group(pid);
6898 true
6899 }
6900 #[cfg(not(unix))]
6901 Some(_) => false,
6902 None => false,
6903 }
6904}
6905
6906fn error_message(error: ServiceError) -> String {
6907 match error {
6908 ServiceError::InvalidParams(message)
6909 | ServiceError::Operation(message)
6910 | ServiceError::UnsupportedAction(message) => message,
6911 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6912 ServiceError::Sdk(error) => error.to_string(),
6913 }
6914}
6915
6916#[derive(Debug)]
6917enum ServiceError {
6918 InvalidParams(String),
6919 MethodNotFound,
6920 UnsupportedAction(String),
6921 Operation(String),
6922 Sdk(SdkError),
6923}
6924
6925fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6926 match error {
6927 ServiceError::InvalidParams(message) => {
6928 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6929 }
6930 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6931 SdkError::unsupported(operation)
6932 }
6933 ServiceError::Operation(message) => {
6934 let code = if message.contains("already in progress") {
6935 SdkErrorCode::Busy
6936 } else if message.contains("not supported by this runtime") {
6937 SdkErrorCode::UnsupportedAction
6938 } else if message.contains("unknown runtime connection") {
6939 SdkErrorCode::NotFound
6940 } else {
6941 SdkErrorCode::Execution
6942 };
6943 SdkError::new(code, operation, message)
6944 }
6945 ServiceError::Sdk(error) => error,
6946 }
6947}
6948
6949fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6950 let error_code = error.code();
6951 let code = match error_code {
6952 SdkErrorCode::Unauthenticated => -32030,
6953 SdkErrorCode::Unauthorized => -32031,
6954 SdkErrorCode::ControllerRequired => -32032,
6955 SdkErrorCode::LeaseExpired => -32033,
6956 SdkErrorCode::InvalidArgument => -32602,
6957 SdkErrorCode::NotFound => -32004,
6958 SdkErrorCode::Busy => -32000,
6959 SdkErrorCode::UnsupportedAction => -32020,
6960 SdkErrorCode::Execution => -32002,
6961 SdkErrorCode::Transport => -32003,
6962 };
6963 json!({
6964 "jsonrpc": "2.0",
6965 "id": id,
6966 "error": {
6967 "code": code,
6968 "name": error_code,
6969 "operation": error.operation(),
6970 "message": error.to_string(),
6971 },
6972 })
6973}
6974
6975fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6976 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6977}
6978
6979fn operation(error: impl Into<crate::Error>) -> ServiceError {
6980 let error = error.into();
6981 match error {
6982 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6983 error => ServiceError::Operation(error.to_string()),
6984 }
6985}
6986
6987#[derive(Debug, Clone, Deserialize, Default)]
6991#[serde(default)]
6992struct MemoryRequest {
6993 harness: Option<String>,
6995 query: Option<String>,
6997 profile: Option<String>,
6999 session: Option<String>,
7001 full: bool,
7003 regex: bool,
7005 cwd: Option<std::path::PathBuf>,
7007 homes: crate::HarnessHomes,
7009}
7010
7011fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7014 let request = decode::<MemoryRequest>(params)?;
7015 let harness = request
7016 .harness
7017 .clone()
7018 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7019 let to_service = |error: crate::memory::MemoryError| match error {
7020 crate::memory::MemoryError::UnsupportedHarness { .. }
7021 | crate::memory::MemoryError::SessionNotScoped { .. } => {
7022 ServiceError::UnsupportedAction(error.to_string())
7023 }
7024 other => ServiceError::InvalidParams(other.to_string()),
7025 };
7026 match method {
7027 "harness.v1.memory.show" => {
7028 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
7029 harness,
7030 profile: request.profile,
7031 session: request.session,
7032 full: request.full,
7033 cwd: request.cwd,
7034 homes: request.homes,
7035 })
7036 .map_err(to_service)?;
7037 Ok(json!({
7038 "schema": crate::memory::MEMORY_SCHEMA,
7039 "documents": documents,
7040 }))
7041 }
7042 "harness.v1.memory.search" => {
7043 let query = request
7044 .query
7045 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
7046 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
7047 harness,
7048 query,
7049 profile: request.profile,
7050 regex: request.regex,
7051 cwd: request.cwd,
7052 homes: request.homes,
7053 })
7054 .map_err(to_service)?;
7055 Ok(json!({
7056 "schema": crate::memory::MEMORY_SCHEMA,
7057 "matches": matches,
7058 }))
7059 }
7060 _ => Err(ServiceError::MethodNotFound),
7061 }
7062}
7063
7064#[derive(Debug, Clone, Deserialize)]
7068#[serde(default)]
7069struct ProfilesQuery {
7070 harness: Option<String>,
7072 name: Option<String>,
7074 homes: crate::HarnessHomes,
7076}
7077
7078impl Default for ProfilesQuery {
7079 fn default() -> Self {
7080 Self {
7081 harness: None,
7082 name: None,
7083 homes: crate::HarnessHomes::default(),
7084 }
7085 }
7086}
7087
7088fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7091 let query = decode::<ProfilesQuery>(params)?;
7092 let to_service = |error: crate::profiles::ProfileError| match error {
7093 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
7094 ServiceError::UnsupportedAction(error.to_string())
7095 }
7096 crate::profiles::ProfileError::NotFound { .. } => {
7097 ServiceError::InvalidParams(error.to_string())
7098 }
7099 };
7100 match method {
7101 "harness.v1.profiles.list" => {
7102 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
7103 .map_err(to_service)?;
7104 Ok(json!({
7105 "schema": crate::profiles::PROFILES_SCHEMA,
7106 "profiles": profiles,
7107 }))
7108 }
7109 "harness.v1.profiles.get" => {
7110 let harness = query
7111 .harness
7112 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7113 let name = query
7114 .name
7115 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7116 let profile =
7117 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
7118 Ok(json!({
7119 "schema": crate::profiles::PROFILES_SCHEMA,
7120 "profile": profile,
7121 }))
7122 }
7123 _ => Err(ServiceError::MethodNotFound),
7124 }
7125}
7126
7127#[derive(Debug, Clone, Deserialize)]
7131#[serde(default)]
7132struct ChannelsQuery {
7133 harness: Option<String>,
7135 name: Option<String>,
7137 homes: crate::HarnessHomes,
7139}
7140
7141impl Default for ChannelsQuery {
7142 fn default() -> Self {
7143 Self {
7144 harness: None,
7145 name: None,
7146 homes: crate::HarnessHomes::default(),
7147 }
7148 }
7149}
7150
7151#[derive(Debug, Clone, Deserialize)]
7155#[serde(default)]
7156struct RoutesQuery {
7157 harness: Option<String>,
7158 profile: Option<String>,
7160 homes: crate::HarnessHomes,
7161}
7162
7163impl Default for RoutesQuery {
7164 fn default() -> Self {
7165 Self {
7166 harness: None,
7167 profile: None,
7168 homes: crate::HarnessHomes::default(),
7169 }
7170 }
7171}
7172
7173#[derive(Debug, Clone, Deserialize)]
7174#[serde(default)]
7175struct TriggersQuery {
7176 harness: Option<String>,
7177 homes: crate::HarnessHomes,
7178}
7179
7180impl Default for TriggersQuery {
7181 fn default() -> Self {
7182 Self {
7183 harness: None,
7184 homes: crate::HarnessHomes::default(),
7185 }
7186 }
7187}
7188
7189fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
7190 let query = decode::<TriggersQuery>(params)?;
7191 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
7192 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7193 Ok(json!({
7194 "schema": crate::triggers::TRIGGERS_SCHEMA,
7195 "triggers": triggers,
7196 }))
7197}
7198
7199fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
7200 let query = decode::<RoutesQuery>(params)?;
7201 let routes = crate::routes::list_routes(
7202 &query.homes,
7203 query.harness.as_deref(),
7204 query.profile.as_deref(),
7205 )
7206 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7207 Ok(json!({
7208 "schema": crate::routes::ROUTES_SCHEMA,
7209 "routes": routes,
7210 }))
7211}
7212
7213fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7214 let query = decode::<ChannelsQuery>(params)?;
7215 let to_service = |error: crate::channels::ChannelError| match error {
7216 crate::channels::ChannelError::UnsupportedHarness { .. } => {
7217 ServiceError::UnsupportedAction(error.to_string())
7218 }
7219 crate::channels::ChannelError::NotFound { .. } => {
7220 ServiceError::InvalidParams(error.to_string())
7221 }
7222 };
7223 match method {
7224 "harness.v1.channels.list" => {
7225 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
7226 .map_err(to_service)?;
7227 Ok(json!({
7228 "schema": crate::channels::CHANNELS_SCHEMA,
7229 "channels": channels,
7230 }))
7231 }
7232 "harness.v1.channels.status" => {
7233 let harness = query
7234 .harness
7235 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7236 let name = query
7237 .name
7238 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7239 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
7240 .map_err(to_service)?;
7241 Ok(json!({
7242 "schema": crate::channels::CHANNELS_SCHEMA,
7243 "channel": channel,
7244 }))
7245 }
7246 _ => Err(ServiceError::MethodNotFound),
7247 }
7248}
7249
7250fn rpc_error(id: Value, code: i64, message: &str) -> Value {
7251 json!({
7252 "jsonrpc": "2.0",
7253 "id": id,
7254 "error": {"code": code, "message": message},
7255 })
7256}