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