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