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 },
4098 };
4099 json!({
4100 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
4101 "message_id": envelope.id,
4102 "reply_to": sender.to_string(),
4103 "target": {
4104 "harness": params.locator.harness.as_str(),
4105 "session_id": params.locator.session_id,
4106 "name": match &door {
4107 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
4108 _ => None,
4109 },
4110 },
4111 "delivery": {"door": door.name(), "how": how},
4112 "inbound_controls": inbound_controls,
4113 "inbound_controls_error": inbound_controls_error,
4114 })
4115}
4116
4117async fn message_as_user(
4121 params: &MessageSessionParams,
4122 sender: crate::mailbox::MailAddress,
4123 receiver: crate::mailbox::MailAddress,
4124 agent: Option<crate::mailbox::MailAddress>,
4125 inbound_controls: Value,
4126 inbound_controls_error: Value,
4127) -> Value {
4128 let mut envelope = match crate::mailbox::Envelope::new(
4129 sender.clone(),
4130 format!("{}@{}", sender.session_id, sender.machine),
4131 crate::mailbox::MailKind::User,
4132 crate::mailbox::ReplyVia::None,
4133 params.text.clone(),
4134 ) {
4135 Ok(envelope) => envelope,
4136 Err(error) => {
4137 return json!({
4138 "delivered_to_bus": false,
4139 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4140 "inbound_controls": inbound_controls,
4141 "inbound_controls_error": inbound_controls_error,
4142 })
4143 }
4144 };
4145 let mut receiver = receiver;
4150 if let Some(agent) = &agent {
4151 envelope.in_reply_to = params.in_reply_to.clone();
4152 envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
4153 let channel = crate::mail_agent::Channel::default();
4154 let planned = match crate::mail_agent::plan(&mut envelope, agent, &channel) {
4155 Ok(planned) => planned,
4156 Err(error) => {
4157 return json!({
4158 "delivered_to_bus": false,
4159 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4160 "inbound_controls": inbound_controls,
4161 "inbound_controls_error": inbound_controls_error,
4162 })
4163 }
4164 };
4165 if let Some(planned) = planned {
4166 let typed_to = planned
4167 .recipients
4168 .iter()
4169 .find(|r| r.wake && r.address.harness != "operator")
4170 .map(|r| r.address.clone());
4171 if let Some(typed_to) = typed_to {
4172 receiver = typed_to;
4173 }
4174 if let Some(agent) = &planned.agent {
4175 let _ = crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), agent)
4177 .and_then(|mailbox| mailbox.deliver_read(&envelope));
4178 }
4179 let mut copy = envelope.clone();
4182 copy.kind = crate::mailbox::MailKind::Typed;
4183 for other in planned
4184 .recipients
4185 .iter()
4186 .filter(|r| r.address != receiver && r.address != sender)
4187 {
4188 let _ = crate::mail_agent::file_unread(&other.address, ©);
4189 }
4190 if let Some(account_manager) = crate::mail_agent::owners_account_manager() {
4193 let main = account_manager.main_session.clone();
4194 let already =
4195 main == receiver || planned.recipients.iter().any(|r| r.address == main);
4196 if account_manager.name != agent.session_id && !already {
4197 let _ = crate::mail_agent::file_unread(&main, ©);
4198 }
4199 }
4200 }
4201 }
4202 if agent.is_some()
4205 && crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4206 && crate::mail_agent::resumable(&receiver)
4207 {
4208 if let Err(error) = crate::mail_agent::resume(&receiver) {
4209 return json!({
4210 "delivered_to_bus": false,
4211 "refusal": {"reason": "no_user_door", "message": format!("{receiver} was not resumed: {error}")},
4212 "inbound_controls": inbound_controls,
4213 "inbound_controls_error": inbound_controls_error,
4214 });
4215 }
4216 let started = std::time::Instant::now();
4217 while crate::mail_route::door_for(¶ms.homes, &receiver).is_err()
4218 && started.elapsed() < crate::mail_agent::RESUME_WAIT
4219 {
4220 tokio::time::sleep(std::time::Duration::from_secs(1)).await;
4221 }
4222 }
4223 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
4224 Ok(turn) => json!({
4225 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
4226 "message_id": envelope.id,
4227 "target": {
4228 "harness": params.locator.harness.as_str(),
4229 "session_id": params.locator.session_id,
4230 },
4231 "delivery": {
4232 "door": match turn {
4233 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
4234 _ => "pane",
4235 },
4236 "how": turn.as_str(),
4237 },
4238 "inbound_controls": inbound_controls,
4239 "inbound_controls_error": inbound_controls_error,
4240 }),
4241 Err(message) => json!({
4242 "delivered_to_bus": false,
4243 "refusal": {"reason": "no_user_door", "message": message},
4244 "inbound_controls": inbound_controls,
4245 "inbound_controls_error": inbound_controls_error,
4246 }),
4247 }
4248}
4249
4250#[derive(Deserialize)]
4251#[serde(deny_unknown_fields)]
4252struct ActivityUnderParams {
4253 pids: Vec<u32>,
4255 #[serde(default)]
4256 homes: crate::HarnessHomes,
4257}
4258
4259async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
4263 let params = decode::<ActivityUnderParams>(params)?;
4264 if params.pids.len() > 1024 {
4265 return Err(ServiceError::InvalidParams(
4266 "sessions.activity_under accepts at most 1024 pids".into(),
4267 ));
4268 }
4269 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
4270 .await
4271 .map_err(ServiceError::Sdk)?;
4272 Ok(json!({
4273 "activities": found
4274 .into_iter()
4275 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
4276 .collect::<Vec<_>>(),
4277 }))
4278}
4279
4280fn operator_address(
4283 name: Option<&str>,
4284) -> std::result::Result<crate::mailbox::MailAddress, String> {
4285 let name: String = name
4286 .unwrap_or("supercode")
4287 .trim()
4288 .chars()
4289 .map(|character| {
4290 if character.is_whitespace() || character == '@' {
4291 '-'
4292 } else {
4293 character
4294 }
4295 })
4296 .collect();
4297 if name.is_empty() {
4298 return Err("from_name must not be empty".into());
4299 }
4300 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
4301 .map_err(|error| error.to_string())
4302}
4303
4304fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
4308 let address = match (¶ms.address, ¶ms.from_name) {
4309 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
4310 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
4311 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
4312 (Some(_), Some(_)) => {
4313 return Err(ServiceError::InvalidParams(
4314 "sessions.inbox takes from_name or address, not both".into(),
4315 ))
4316 }
4317 };
4318 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
4319 let mailbox =
4320 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
4321 let claimed = mailbox.claim_unread().map_err(operation)?;
4322 let mut messages: Vec<Value> = Vec::new();
4323 if params.all {
4324 for stored in mailbox.list().map_err(operation)? {
4325 if stored.state == crate::mailbox::MailState::Read {
4326 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4327 }
4328 }
4329 }
4330 for stored in &claimed {
4331 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4332 }
4333 for stored in &claimed {
4334 mailbox.acknowledge(stored).map_err(operation)?;
4335 }
4336 Ok(json!({"address": address.to_string(), "messages": messages}))
4337}
4338
4339#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4344struct FollowedSource {
4345 harness: String,
4346 session_id: String,
4347 reported: Option<String>,
4348}
4349
4350#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4351struct ActivitySubscription {
4352 locators: Vec<SessionLocator>,
4353 homes: crate::HarnessHomes,
4354 reported: BTreeMap<(String, String), crate::SessionActivity>,
4355}
4356
4357fn live_descriptor_value(
4364 session: &SessionDescriptor,
4365 doors: &crate::mail_route::LiveSessions,
4366) -> std::result::Result<Value, ServiceError> {
4367 let mut value = serde_json::to_value(session)
4368 .map_err(|error| ServiceError::Operation(error.to_string()))?;
4369 if value.get("title").is_none_or(Value::is_null) {
4371 if let Some(live) = doors.all().iter().find(|live| {
4372 live.address.harness == session.locator.harness.as_str()
4373 && live.address.session_id == session.locator.session_id
4374 }) {
4375 let name = live.name.split('@').next().unwrap_or(&live.name);
4376 if !name.is_empty() {
4377 value["title"] = json!(name);
4378 }
4379 }
4380 }
4381 if let Some(workspace) = &session.cwd {
4382 let source = LiveRuntimeSource {
4383 harness: session.locator.harness.as_str().to_string(),
4384 session_id: session.locator.session_id.clone(),
4385 workspace: workspace.clone(),
4386 };
4387 if let Some(endpoint) = discover_live_runtime(&source)
4388 .map_err(|error| ServiceError::Operation(error.to_string()))?
4389 {
4390 value["live_endpoint"] = json!(endpoint.as_str());
4391 }
4392 }
4393 if let Some(door) = doors.door(
4397 session.locator.harness.as_str(),
4398 &session.locator.session_id,
4399 ) {
4400 value["delivery"] = json!(door);
4401 }
4402 if let Some(live) = doors.all().iter().find(|live| {
4405 live.address.harness == session.locator.harness.as_str()
4406 && live.address.session_id == session.locator.session_id
4407 && live.status == "waiting"
4408 }) {
4409 if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
4410 value["pending_request"] = request;
4411 }
4412 }
4413 Ok(value)
4414}
4415
4416fn live_index_changes(
4417 changes: Vec<crate::session_index::SessionIndexChange>,
4418 homes: &HarnessHomes,
4419) -> std::result::Result<Vec<Value>, ServiceError> {
4420 use crate::session_index::SessionIndexChange;
4421 let doors = crate::mail_route::LiveSessions::read(homes);
4422 changes
4423 .into_iter()
4424 .map(|change| match change {
4425 SessionIndexChange::Added { descriptor } => Ok(json!({
4426 "kind": "added",
4427 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4428 })),
4429 SessionIndexChange::Updated { descriptor } => Ok(json!({
4430 "kind": "updated",
4431 "descriptor": live_descriptor_value(&descriptor, &doors)?,
4432 })),
4433 SessionIndexChange::Removed { key } => Ok(json!({
4434 "kind": "removed",
4435 "key": key,
4436 })),
4437 })
4438 .collect()
4439}
4440
4441fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
4442 use crate::{SessionPresence, SessionTurnState};
4443 match (activity.presence, activity.turn) {
4444 (SessionPresence::Persisted, _) => None,
4445 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
4446 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
4447 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
4449 (SessionPresence::Running, SessionTurnState::Unknown)
4453 if activity.evidence.native_state.is_none() =>
4454 {
4455 None
4456 }
4457 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
4458 }
4459}
4460
4461#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
4462#[serde(rename_all = "kebab-case")]
4463enum TransferFormat {
4464 ClaudeCode,
4465 Codex,
4466 #[serde(rename = "opencode", alias = "open-code")]
4467 OpenCode,
4468 Pi,
4469 Grok,
4470 Gemini,
4471 Goose,
4472 Hermes,
4476}
4477
4478impl TransferFormat {
4479 fn id(self) -> &'static str {
4480 match self {
4481 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
4482 Self::Codex => HarnessId::CODEX,
4483 Self::OpenCode => HarnessId::OPENCODE,
4484 Self::Pi => HarnessId::PI,
4485 Self::Grok => HarnessId::GROK,
4486 Self::Gemini => HarnessId::GEMINI,
4487 Self::Goose => HarnessId::GOOSE,
4488 Self::Hermes => HarnessId::HERMES,
4489 }
4490 }
4491}
4492
4493impl From<TransferFormat> for SessionFormat {
4494 fn from(value: TransferFormat) -> Self {
4495 match value {
4496 TransferFormat::ClaudeCode => Self::ClaudeCode,
4497 TransferFormat::Codex => Self::Codex,
4498 TransferFormat::OpenCode => Self::OpenCode,
4499 TransferFormat::Pi => Self::Pi,
4500 TransferFormat::Grok => Self::Grok,
4501 TransferFormat::Gemini => Self::Gemini,
4502 TransferFormat::Goose => Self::Goose,
4503 TransferFormat::Hermes => Self::Codex,
4505 }
4506 }
4507}
4508
4509#[derive(Deserialize)]
4510struct ImportSessionParams {
4511 source_harness: TransferFormat,
4512 content: String,
4513}
4514
4515#[derive(Deserialize)]
4516struct ExportSessionParams {
4517 locator: SessionLocator,
4518 target_harness: TransferFormat,
4519}
4520
4521#[derive(Deserialize)]
4522struct ReduceSessionParams {
4523 locator: SessionLocator,
4524 target_harness: TransferFormat,
4525 #[serde(default = "default_keep_last")]
4526 keep_last: usize,
4527}
4528
4529fn default_keep_last() -> usize {
4530 6
4531}
4532
4533#[derive(Deserialize)]
4534struct BranchSessionParams {
4535 locator: SessionLocator,
4536 #[serde(default)]
4537 target_harness: Option<TransferFormat>,
4538 #[serde(default)]
4541 before_message_id: Option<MessageId>,
4542}
4543
4544#[derive(Deserialize)]
4550struct MessageId {
4551 key: String,
4552 value: String,
4553}
4554
4555const MESSAGE_ID_KEYS: &[&str] = &["claude_uuid"];
4557
4558#[derive(Deserialize)]
4559struct HandoffSessionParams {
4560 locator: SessionLocator,
4561 target_harness: TransferFormat,
4562 #[serde(default)]
4563 cwd: Option<PathBuf>,
4564 #[serde(default)]
4567 before_message_id: Option<MessageId>,
4568}
4569
4570fn message_place(session: &Session, id: &MessageId) -> std::result::Result<usize, ServiceError> {
4572 if !MESSAGE_ID_KEYS.contains(&id.key.as_str()) {
4573 return Err(ServiceError::UnsupportedAction(format!(
4574 "{} is not a per-message id supercode reads; a fork at a turn names a message by {}",
4575 id.key,
4576 MESSAGE_ID_KEYS.join(", ")
4577 )));
4578 }
4579 let mut places = session
4580 .messages
4581 .iter()
4582 .enumerate()
4583 .filter(|(_, message)| {
4584 matches!(message.role, Role::User)
4585 && message.metadata.get(&id.key).map(String::as_str) == Some(id.value.as_str())
4586 })
4587 .map(|(index, _)| index + 1);
4588 match (places.next(), places.next()) {
4589 (Some(place), None) => Ok(place),
4590 (None, _) => Err(ServiceError::Operation(format!(
4591 "no user message in this session has {} {}",
4592 id.key, id.value
4593 ))),
4594 (Some(_), Some(_)) => Err(ServiceError::Operation(format!(
4595 "more than one user message has {} {}; it does not name one turn",
4596 id.key, id.value
4597 ))),
4598 }
4599}
4600
4601fn session_before(
4604 session: Session,
4605 id: Option<&MessageId>,
4606) -> std::result::Result<(Session, Option<usize>), ServiceError> {
4607 let before_message = id.map(|id| message_place(&session, id)).transpose()?;
4608 match before_message {
4609 None => Ok((session, None)),
4610 Some(place) => Ok((session.cut_before(place - 1), Some(place))),
4611 }
4612}
4613
4614#[derive(Deserialize)]
4615struct MaterializeSessionParams {
4616 artifact: crate::native_materialize::MaterializeArtifact,
4617 cwd: PathBuf,
4618 #[serde(default)]
4620 homes: HarnessHomes,
4621}
4622
4623#[derive(Debug, Clone, Copy, Default, Deserialize)]
4626#[serde(rename_all = "snake_case")]
4627enum ResumePolicy {
4628 Default,
4629 #[default]
4630 Yolo,
4631}
4632
4633#[derive(Deserialize)]
4634struct ResumeInstructionsParams {
4635 locator: SessionLocator,
4636 #[serde(default)]
4637 cwd: Option<PathBuf>,
4638 #[serde(default)]
4639 policy: ResumePolicy,
4640}
4641
4642#[derive(Deserialize)]
4644struct WorkflowLoadParams {
4645 #[serde(default)]
4647 from: Option<crate::workflow_doors::WorkflowHarness>,
4648 home: PathBuf,
4649}
4650
4651#[derive(Deserialize)]
4654struct OrchestrationLoadParams {
4655 root: PathBuf,
4656 #[serde(default)]
4657 flavor: crate::orchestration_doors::HomeFlavor,
4658}
4659
4660#[derive(Deserialize)]
4663struct OrchestrationSaveParams {
4664 root: PathBuf,
4665 orchestration: crate::orchestration::Orchestration,
4666 #[serde(default)]
4667 vault: BTreeMap<String, String>,
4668}
4669
4670#[derive(Deserialize)]
4672struct OrchestrationCompileParams {
4673 from: crate::orchestration_doors::OrchestrationHarness,
4674 home: PathBuf,
4675}
4676
4677#[derive(Deserialize)]
4681struct OrchestrationDecompileParams {
4682 to: crate::orchestration_doors::OrchestrationHarness,
4683 orchestration: crate::orchestration::Orchestration,
4684 source: PathBuf,
4685 #[serde(default)]
4686 source_flavor: crate::orchestration_doors::SourceFlavor,
4687 dest: PathBuf,
4688 #[serde(default)]
4689 vault: BTreeMap<String, String>,
4690}
4691
4692#[derive(Deserialize)]
4695struct OrchestrationImportParams {
4696 from: crate::orchestration_doors::OrchestrationHarness,
4697 home: PathBuf,
4698 into: PathBuf,
4699}
4700
4701#[derive(Deserialize)]
4704struct OrchestrationExportParams {
4705 to: crate::orchestration_doors::OrchestrationHarness,
4706 root: PathBuf,
4707 dest: PathBuf,
4708}
4709
4710#[derive(Deserialize)]
4712struct JobsGetParams {
4713 harness: String,
4714 id: String,
4715 #[serde(default)]
4716 homes: crate::HarnessHomes,
4717}
4718
4719fn mutate_job(
4727 verb: crate::jobs_control::JobVerb,
4728 params: Value,
4729) -> std::result::Result<Value, ServiceError> {
4730 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4731 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4732 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4733 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4734}
4735
4736fn mutate_skill(
4743 verb: crate::skills_control::SkillVerb,
4744 params: Value,
4745) -> std::result::Result<Value, ServiceError> {
4746 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4747 if !crate::skills_control::supports_skill_control(&mutation.harness) {
4748 return Err(ServiceError::UnsupportedAction(format!(
4749 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4750 mutation.harness,
4751 verb.as_str(),
4752 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4753 )));
4754 }
4755 let outcome =
4756 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4757 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4758}
4759
4760fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4762 match error {
4763 crate::skills_control::SkillControlError::Unsupported(message) => {
4764 ServiceError::UnsupportedAction(message)
4765 }
4766 crate::skills_control::SkillControlError::Invalid(message) => {
4767 ServiceError::InvalidParams(message)
4768 }
4769 crate::skills_control::SkillControlError::Failed(message) => {
4770 ServiceError::Operation(message)
4771 }
4772 }
4773}
4774
4775fn mutate_profile(
4783 verb: crate::profiles_control::ProfileVerb,
4784 params: Value,
4785) -> std::result::Result<Value, ServiceError> {
4786 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4787 let outcome =
4788 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4789 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4790}
4791
4792fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4794 match error {
4795 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4796 ServiceError::UnsupportedAction(message)
4797 }
4798 crate::profiles_control::ProfileControlError::Invalid(message) => {
4799 ServiceError::InvalidParams(message)
4800 }
4801 crate::profiles_control::ProfileControlError::Failed(message) => {
4802 ServiceError::Operation(message)
4803 }
4804 }
4805}
4806
4807fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4811 match error {
4812 crate::jobs_control::JobControlError::Unsupported(message) => {
4813 ServiceError::UnsupportedAction(message)
4814 }
4815 crate::jobs_control::JobControlError::Invalid(message) => {
4816 ServiceError::InvalidParams(message)
4817 }
4818 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4819 }
4820}
4821
4822fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4827 match error {
4828 crate::SessionControlError::Unsupported(message) => {
4829 ServiceError::UnsupportedAction(message)
4830 }
4831 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4832 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4833 }
4834}
4835
4836fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4841 if crate::jobs::supports_jobs(harness) {
4842 return Ok(());
4843 }
4844 Err(ServiceError::UnsupportedAction(format!(
4845 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4846 crate::jobs::JOB_HARNESSES.join(", ")
4847 )))
4848}
4849
4850#[derive(Deserialize)]
4852struct RunsGetParams {
4853 harness: String,
4854 id: String,
4855 #[serde(default)]
4856 homes: crate::HarnessHomes,
4857}
4858
4859fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4864 if crate::runs::supports_runs(harness) {
4865 return Ok(());
4866 }
4867 Err(ServiceError::UnsupportedAction(format!(
4868 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4869 crate::runs::RUN_HARNESSES.join(", ")
4870 )))
4871}
4872
4873#[derive(Serialize)]
4874struct SessionArtifact {
4875 source_harness: HarnessId,
4876 target_harness: &'static str,
4877 session_id: Option<String>,
4878 content: String,
4879 suggested_filename: String,
4880 files: Vec<SessionArtifactFile>,
4881 fidelity: Fidelity,
4882 residue: Vec<String>,
4883}
4884
4885#[derive(Serialize)]
4886struct SessionArtifactFile {
4887 path: String,
4888 content: String,
4889 role: ArtifactFileRole,
4890}
4891
4892#[derive(Serialize)]
4893#[serde(rename_all = "snake_case")]
4894enum ArtifactFileRole {
4895 Primary,
4896 Subagent,
4897 Bundle,
4898 SourceRecovery,
4899}
4900
4901#[derive(Serialize)]
4902struct StructuredLaunch {
4903 cwd: PathBuf,
4904 program: String,
4905 arguments: Vec<String>,
4906 env: BTreeMap<String, String>,
4907}
4908
4909struct HandoffInstructions {
4910 launch: StructuredLaunch,
4911 materialize: Option<StructuredLaunch>,
4912 requires_materialization: bool,
4913 note: String,
4914}
4915
4916#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4917#[serde(rename_all = "snake_case")]
4918enum HarnessProbeLevel {
4919 #[default]
4920 Passive,
4921 Handshake,
4922}
4923
4924#[derive(Default, Deserialize)]
4925#[serde(default)]
4926struct HarnessInventoryParams {
4927 harness: Option<HarnessId>,
4928 harnesses: Vec<HarnessId>,
4929 workspace: Option<PathBuf>,
4930 probe: HarnessProbeLevel,
4931 include_sessions: bool,
4932 skip_versions: bool,
4934}
4935
4936#[derive(Deserialize)]
4937struct HarnessAuthenticationParams {
4938 harness: HarnessId,
4939}
4940
4941#[derive(Deserialize)]
4942struct BeginHarnessAuthenticationParams {
4943 harness: HarnessId,
4944 #[serde(default = "local_browser_authentication_environment")]
4945 environment: crate::HarnessAuthenticationEnvironment,
4946 #[serde(default)]
4947 method: Option<crate::HarnessAuthenticationMethodId>,
4948 #[serde(default)]
4949 cwd: Option<PathBuf>,
4950}
4951
4952fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4953 crate::HarnessAuthenticationEnvironment::LocalBrowser
4954}
4955
4956#[derive(Serialize)]
4957struct HarnessInventoryReport {
4958 probe: HarnessProbeLevel,
4959 workspace: Option<PathBuf>,
4960 harnesses: Vec<LocalHarness>,
4961}
4962
4963#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4964#[serde(rename_all = "snake_case")]
4965enum HarnessAuthState {
4966 Ready,
4967 Configured,
4968 Required,
4969 Unknown,
4970}
4971
4972#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4973#[serde(rename_all = "snake_case")]
4974enum HarnessRuntimeState {
4975 Ready,
4976 Degraded,
4977 Unavailable,
4978}
4979
4980#[derive(Serialize)]
4981struct HarnessSessionCounts {
4982 global: Option<usize>,
4983 workspace: Option<usize>,
4984}
4985
4986#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4997#[serde(rename_all = "snake_case")]
4998pub enum GatewayState {
4999 Up,
5000 Down,
5001 Unknown,
5002}
5003
5004#[derive(Debug, Clone, Serialize)]
5006pub struct GatewayHealth {
5007 pub state: GatewayState,
5008 #[serde(skip_serializing_if = "Option::is_none")]
5013 pub endpoint: Option<String>,
5014 #[serde(skip_serializing_if = "Option::is_none")]
5015 pub version: Option<String>,
5016 pub evidence: String,
5018 pub checked_at_ms: u64,
5019}
5020
5021fn openclaw_gateway_endpoint(home: &Path) -> String {
5025 let config_path = home.join(".openclaw/openclaw.json");
5026 let gateway = std::fs::read_to_string(&config_path)
5027 .ok()
5028 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
5029 .and_then(|config| config.get("gateway").cloned());
5030 if let Some(url) = gateway
5031 .as_ref()
5032 .and_then(|gateway| gateway.get("url"))
5033 .and_then(serde_json::Value::as_str)
5034 {
5035 return url.to_string();
5036 }
5037 let port = gateway
5038 .as_ref()
5039 .and_then(|gateway| gateway.get("port"))
5040 .and_then(serde_json::Value::as_u64)
5041 .unwrap_or(18789);
5042 format!("ws://127.0.0.1:{port}")
5043}
5044
5045fn hermes_gateway_status() -> Option<(GatewayState, String)> {
5052 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
5053 let output = std::process::Command::new(&program)
5054 .args(["gateway", "status"])
5055 .stdin(std::process::Stdio::null())
5056 .output()
5057 .ok()?;
5058 let text = format!(
5059 "{}{}",
5060 String::from_utf8_lossy(&output.stdout),
5061 String::from_utf8_lossy(&output.stderr)
5062 );
5063 let verdict = text.lines().find_map(|line| {
5064 let l = line.trim();
5065 if l.contains("supervised by launchd (PID")
5066 || l.contains("supervised by systemd (PID")
5067 || l.contains("Gateway is running")
5068 || l.contains("process is running")
5069 {
5070 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
5071 } else if l.contains("not running") || l.contains("not installed") {
5072 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
5073 } else {
5074 None
5075 }
5076 });
5077 verdict
5078}
5079
5080fn gateway_health(
5081 id: &str,
5082 installed: bool,
5083 running: Option<&RunningInstance>,
5084 version: Option<&str>,
5085) -> GatewayHealth {
5086 let checked_at_ms = now_epoch_ms();
5087 let home = supercode_interchange::user_home()
5088 .map(std::path::PathBuf::into_os_string)
5089 .map(PathBuf::from);
5090 match id {
5091 HarnessId::HERMES | HarnessId::OPENCLAW => {
5092 let endpoint = (id == HarnessId::OPENCLAW)
5093 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
5094 .flatten();
5095 let (state, evidence) = match running {
5096 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
5097 None if !installed => (
5098 GatewayState::Unknown,
5099 format!("`{id}` is not installed; no gateway to probe"),
5100 ),
5101 None if id == HarnessId::HERMES => match hermes_gateway_status() {
5102 Some((state, evidence)) => (state, evidence),
5105 None => (
5106 GatewayState::Down,
5107 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
5108 ),
5109 },
5110 None => (
5111 GatewayState::Down,
5112 format!(
5113 "no TCP listener at {}",
5114 endpoint.as_deref().unwrap_or("the gateway endpoint")
5115 ),
5116 ),
5117 };
5118 GatewayHealth {
5119 state,
5120 endpoint,
5121 version: version.map(str::to_string),
5122 evidence,
5123 checked_at_ms,
5124 }
5125 }
5126 HarnessId::ORCHESTRATOR => {
5133 let root = crate::HarnessHomes::default().orchestrator;
5134 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
5135 Some(lease) if lease.is_live() => (
5136 GatewayState::Up,
5137 format!(
5138 "`{}` names pid {} (started {}), which is live",
5139 crate::orchestrator::lock_path(&root).display(),
5140 lease.pid,
5141 lease.started_at
5142 ),
5143 ),
5144 Some(lease) => (
5145 GatewayState::Down,
5146 format!(
5147 "stale lease `{}`: pid {} is gone",
5148 crate::orchestrator::lock_path(&root).display(),
5149 lease.pid
5150 ),
5151 ),
5152 None => (
5153 GatewayState::Down,
5154 format!(
5155 "no lease at `{}`; `supercode orchestrator start` writes one",
5156 crate::orchestrator::lock_path(&root).display()
5157 ),
5158 ),
5159 };
5160 GatewayHealth {
5161 state,
5162 endpoint: None,
5163 version: version.map(str::to_string),
5164 evidence,
5165 checked_at_ms,
5166 }
5167 }
5168 _ => GatewayHealth {
5169 state: GatewayState::Unknown,
5170 endpoint: None,
5171 version: version.map(str::to_string),
5172 evidence: format!("`{id}` runs per session, not as a gateway"),
5173 checked_at_ms,
5174 },
5175 }
5176}
5177
5178#[derive(Debug, Clone, Serialize)]
5179struct RunningInstance {
5180 method: RunningInstanceMethod,
5182 evidence: String,
5184 checked_at_ms: u64,
5186}
5187
5188#[derive(Debug, Clone, Copy, Serialize)]
5189#[serde(rename_all = "snake_case")]
5190enum RunningInstanceMethod {
5191 GatewayConnect,
5194 StoreWalActivity,
5197}
5198
5199fn now_epoch_ms() -> u64 {
5200 std::time::SystemTime::now()
5201 .duration_since(std::time::UNIX_EPOCH)
5202 .map(|elapsed| elapsed.as_millis() as u64)
5203 .unwrap_or(0)
5204}
5205
5206fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
5210 let config_path = home.join(".openclaw/openclaw.json");
5211 let text = std::fs::read_to_string(&config_path).ok();
5212 let gateway = text
5213 .as_deref()
5214 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
5215 .and_then(|config| config.get("gateway").cloned());
5216 let address = gateway
5217 .as_ref()
5218 .and_then(|gateway| gateway.get("url"))
5219 .and_then(serde_json::Value::as_str)
5220 .and_then(|url| {
5221 url.split("://").nth(1).map(|rest| {
5222 rest.trim_end_matches('/')
5223 .split('/')
5224 .next()
5225 .unwrap_or(rest)
5226 .to_string()
5227 })
5228 })
5229 .unwrap_or_else(|| {
5230 let port = gateway
5231 .as_ref()
5232 .and_then(|gateway| gateway.get("port"))
5233 .and_then(serde_json::Value::as_u64)
5234 .unwrap_or(18789);
5235 format!("127.0.0.1:{port}")
5236 });
5237 let reachable = std::net::TcpStream::connect_timeout(
5238 &address.parse().ok()?,
5239 std::time::Duration::from_millis(400),
5240 )
5241 .is_ok();
5242 reachable.then(|| RunningInstance {
5243 method: RunningInstanceMethod::GatewayConnect,
5244 evidence: format!(
5245 "gateway endpoint {address} accepted a TCP connect (from {})",
5246 config_path.display()
5247 ),
5248 checked_at_ms: now_epoch_ms(),
5249 })
5250}
5251
5252fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
5257 let wal = home.join(".hermes/state.db-wal");
5258 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
5259 let age_ms = std::time::SystemTime::now()
5260 .duration_since(modified)
5261 .map(|age| age.as_millis() as u64)
5262 .unwrap_or(u64::MAX);
5263 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
5264 method: RunningInstanceMethod::StoreWalActivity,
5265 evidence: format!(
5266 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
5267 wal.display()
5268 ),
5269 checked_at_ms: now_epoch_ms(),
5270 })
5271}
5272
5273fn probe_running_instance(id: &str) -> Option<RunningInstance> {
5275 let home = supercode_interchange::user_home()
5276 .map(std::path::PathBuf::into_os_string)
5277 .map(PathBuf::from)?;
5278 match id {
5279 HarnessId::OPENCLAW => probe_openclaw_running(&home),
5280 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
5281 _ => None,
5282 }
5283}
5284
5285#[derive(Serialize)]
5286struct LocalHarness {
5287 id: HarnessId,
5288 display_name: String,
5289 supported: bool,
5290 installed: bool,
5291 executable: Option<String>,
5292 version: Option<String>,
5293 auth: HarnessAuthState,
5294 runtime: HarnessRuntimeState,
5295 protocol: String,
5296 capabilities: crate::RuntimeCapabilities,
5297 effective_capabilities: crate::RuntimeCapabilities,
5298 sessions: HarnessSessionCounts,
5299 #[serde(skip_serializing_if = "Option::is_none")]
5302 running: Option<RunningInstance>,
5303 gateway: GatewayHealth,
5305 reason: Option<String>,
5306 repair: Option<String>,
5307}
5308
5309#[derive(Clone, Deserialize)]
5310struct RuntimeBackendParams {
5311 harness: HarnessId,
5312 #[serde(default)]
5313 protocol: Option<String>,
5314 #[serde(default)]
5315 launch: Option<RuntimeLaunch>,
5316 #[serde(default)]
5317 base_url: Option<String>,
5318 #[serde(default)]
5319 policy: RuntimePolicy,
5320}
5321
5322#[derive(Debug, Clone, Copy, Default, Deserialize)]
5325#[serde(rename_all = "snake_case")]
5326enum RuntimePolicy {
5327 Default,
5328 #[default]
5329 Yolo,
5330}
5331
5332#[derive(Deserialize)]
5333struct RuntimeStartParams {
5334 #[serde(flatten)]
5335 backend: RuntimeBackendParams,
5336 cwd: PathBuf,
5337 #[serde(default)]
5340 mcp_servers: Vec<crate::McpServerLaunch>,
5341 #[serde(default)]
5343 approval_policy: Option<String>,
5344 #[serde(default)]
5347 model: Option<String>,
5348}
5349
5350#[derive(Deserialize)]
5351struct RuntimeAttachParams {
5352 #[serde(flatten)]
5353 backend: RuntimeBackendParams,
5354 runtime_id: String,
5355 #[serde(default)]
5356 cwd: Option<PathBuf>,
5357 #[serde(default)]
5360 mcp_servers: Vec<crate::McpServerLaunch>,
5361 #[serde(default)]
5363 approval_policy: Option<String>,
5364}
5365
5366#[derive(Deserialize)]
5367struct RuntimeConnectionParams {
5368 connection: String,
5369}
5370
5371#[derive(Deserialize)]
5374struct RuntimeCloseParams {
5375 #[serde(default)]
5376 connection: String,
5377 #[serde(default)]
5378 runtime_id: Option<String>,
5379}
5380
5381#[derive(Deserialize)]
5382struct RuntimeInputParams {
5383 connection: String,
5384 text: String,
5385 #[serde(default)]
5386 image_urls: Vec<String>,
5387}
5388
5389const MAX_RUNTIME_IMAGES: usize = 4;
5390const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
5391const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
5392
5393fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
5394 if image_urls.len() > MAX_RUNTIME_IMAGES {
5395 return Err(ServiceError::InvalidParams(format!(
5396 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
5397 )));
5398 }
5399 let mut total = 0usize;
5400 for url in &image_urls {
5401 if !(url.starts_with("data:image/")
5402 || url.starts_with("https://")
5403 || url.starts_with("http://"))
5404 {
5405 return Err(ServiceError::InvalidParams(
5406 "runtime images must be image data URLs or HTTP(S) URLs".into(),
5407 ));
5408 }
5409 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
5410 return Err(ServiceError::InvalidParams(format!(
5411 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
5412 )));
5413 }
5414 total = total.saturating_add(url.len());
5415 }
5416 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
5417 return Err(ServiceError::InvalidParams(format!(
5418 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
5419 )));
5420 }
5421 Ok(image_urls)
5422}
5423
5424#[derive(Deserialize)]
5425struct RuntimeRespondParams {
5426 connection: String,
5427 request_id: Value,
5428 response: Value,
5429}
5430
5431fn default_reduction_store_root() -> PathBuf {
5432 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
5433 return PathBuf::from(root).join("sessions");
5434 }
5435 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
5436 return PathBuf::from(home).join(".supercode").join("sessions");
5437 }
5438 PathBuf::from(".supercode").join("sessions")
5439}
5440
5441fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
5442 let mut output = String::new();
5443 for message in messages {
5444 output.push_str(
5445 &serde_json::to_string(message)
5446 .map_err(|error| ServiceError::Operation(error.to_string()))?,
5447 );
5448 output.push('\n');
5449 }
5450 Ok(output)
5451}
5452
5453fn parse_messages_jsonl(
5454 content: &str,
5455) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
5456 content
5457 .lines()
5458 .enumerate()
5459 .filter(|(_, line)| !line.trim().is_empty())
5460 .map(|(index, line)| {
5461 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
5462 ServiceError::Operation(format!(
5463 "reduced transcript line {} is invalid: {error}",
5464 index + 1
5465 ))
5466 })
5467 })
5468 .collect()
5469}
5470
5471fn reduced_bootstrap_prompt(
5472 source: &SessionLocator,
5473 target: TransferFormat,
5474 view_jsonl: &str,
5475 sidecar_path: &Path,
5476 reduction_log_path: &Path,
5477) -> String {
5478 format!(
5479 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
5480 \n\
5481 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\
5482 \n\
5483 <supercode-reduced-session source-session=\"{source_id}\">\n\
5484 {view_jsonl}\
5485 </supercode-reduced-session>\n\
5486 \n\
5487 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
5488 source_harness = source.harness.as_str(),
5489 target_harness = target.id(),
5490 sidecar = sidecar_path.display(),
5491 log = reduction_log_path.display(),
5492 source_id = source.session_id,
5493 )
5494}
5495
5496fn session_artifact(
5497 locator: &SessionLocator,
5498 session: &Session,
5499 target: TransferFormat,
5500) -> std::result::Result<SessionArtifact, ServiceError> {
5501 session_artifact_with_id(locator, session, target, None)
5502}
5503
5504fn session_artifact_with_id(
5505 locator: &SessionLocator,
5506 session: &Session,
5507 target: TransferFormat,
5508 target_session_id: Option<&str>,
5509) -> std::result::Result<SessionArtifact, ServiceError> {
5510 let format: SessionFormat = target.into();
5511 let diagonal = format.source() == session.meta.source;
5512 crate::residue_store::store_segments(session);
5513 let has_appended_turns = session
5514 .imported_message_count
5515 .is_some_and(|imported| imported < session.messages.len());
5516 let mut restoration = None;
5517 let content = if let Some(id) = target_session_id {
5518 if diagonal && format != SessionFormat::OpenCode {
5519 session
5520 .to_jsonl_spliced(format, Some(id))
5521 .map_err(operation)?
5522 } else {
5523 let mut rewritten = session.clone();
5524 rewritten.meta.session_id = Some(id.to_string());
5525 rewritten.to_jsonl(format).map_err(operation)?
5526 }
5527 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
5528 session.raw_verbatim()
5529 } else if diagonal {
5530 session.to_jsonl_spliced(format, None).map_err(operation)?
5531 } else {
5532 match session
5535 .restore_residue(format, crate::residue_store::lookup)
5536 .map_err(operation)?
5537 {
5538 Some((content, report)) => {
5539 restoration = Some(report);
5540 content
5541 }
5542 None => session.to_jsonl(format).map_err(operation)?,
5543 }
5544 };
5545 let stem = sanitize_filename(
5546 target_session_id
5547 .or(session.meta.session_id.as_deref())
5548 .unwrap_or(&locator.session_id),
5549 );
5550 let suggested_filename = if diagonal && target == TransferFormat::Grok {
5551 "chat_history.jsonl".to_string()
5552 } else if target == TransferFormat::Goose {
5553 format!("{stem}.goose.json")
5554 } else {
5555 format!("{stem}.{}.jsonl", target.id())
5556 };
5557 let mut files = vec![SessionArtifactFile {
5558 path: suggested_filename.clone(),
5559 content: content.clone(),
5560 role: ArtifactFileRole::Primary,
5561 }];
5562 if target == TransferFormat::ClaudeCode {
5563 let bundle_stem = Path::new(&suggested_filename)
5564 .file_stem()
5565 .and_then(|stem| stem.to_str())
5566 .unwrap_or(&stem);
5567 let mut child_paths = BTreeSet::new();
5568 for (index, subagent) in session.subagents.iter().enumerate() {
5569 let agent_id = subagent
5570 .meta
5571 .agent_id
5572 .as_deref()
5573 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
5574 .map(sanitize_filename)
5575 .filter(|id| !id.is_empty())
5576 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5577 let child_has_appended_turns = subagent
5578 .imported_message_count
5579 .is_some_and(|imported| imported < subagent.messages.len());
5580 let child_content = if target_session_id.is_none()
5581 && subagent.meta.source == SessionSource::ClaudeCode
5582 && subagent.raw_is_verbatim
5583 && !child_has_appended_turns
5584 {
5585 subagent.raw_verbatim()
5586 } else if subagent.meta.source == SessionSource::ClaudeCode {
5587 subagent
5588 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
5589 .map_err(operation)?
5590 } else {
5591 let mut child = subagent.clone();
5592 if let Some(id) = target_session_id {
5593 child.meta.session_id = Some(id.to_string());
5594 }
5595 child
5596 .to_jsonl(SessionFormat::ClaudeCode)
5597 .map_err(operation)?
5598 };
5599 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
5600 if !child_paths.insert(path.clone()) {
5601 return Err(ServiceError::Operation(format!(
5602 "Claude subagent ids collide at artifact path `{path}`"
5603 )));
5604 }
5605 files.push(SessionArtifactFile {
5606 path,
5607 content: child_content,
5608 role: ArtifactFileRole::Subagent,
5609 });
5610 }
5611 }
5612 if diagonal && target == TransferFormat::Grok {
5613 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
5614 }
5615 if !diagonal || !session.raw_is_verbatim {
5616 files.push(SessionArtifactFile {
5617 path: "recovery/source.supercode.jsonl".into(),
5618 content: session.to_native_jsonl(),
5619 role: ArtifactFileRole::SourceRecovery,
5620 });
5621 for (index, subagent) in session.subagents.iter().enumerate() {
5622 let id = subagent
5623 .meta
5624 .agent_id
5625 .as_deref()
5626 .map(sanitize_filename)
5627 .unwrap_or_else(|| format!("subagent-{}", index + 1));
5628 files.push(SessionArtifactFile {
5629 path: format!("recovery/subagents/{id}.supercode.jsonl"),
5630 content: subagent.to_native_jsonl(),
5631 role: ArtifactFileRole::SourceRecovery,
5632 });
5633 }
5634 }
5635 if !diagonal && session.meta.source == SessionSource::Grok {
5636 append_grok_bundle_files(
5637 locator,
5638 "recovery/grok/",
5639 ArtifactFileRole::SourceRecovery,
5640 &mut files,
5641 )?;
5642 }
5643 let (fidelity, residue) = if diagonal
5644 && target_session_id.is_none()
5645 && session.raw_is_verbatim
5646 && !has_appended_turns
5647 {
5648 (Fidelity::ByteLossless, Vec::new())
5649 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
5650 (
5651 Fidelity::ValueLossless,
5652 vec![if target_session_id.is_some() {
5653 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
5654 } else {
5655 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
5656 }],
5657 )
5658 } else {
5659 match restoration {
5660 Some(report) if report.rendered_messages == 0 => (
5661 Fidelity::ByteLossless,
5662 vec![format!(
5663 "restored verbatim from this conversation's {} source records in the residue store",
5664 target.id()
5665 )],
5666 ),
5667 Some(report) => (
5668 Fidelity::Semantic,
5669 vec![format!(
5670 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
5671 report.restored_messages,
5672 report.restored_messages + report.rendered_messages,
5673 report.rendered_messages,
5674 target.id()
5675 )],
5676 ),
5677 None => (
5678 Fidelity::Semantic,
5679 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
5680 ),
5681 }
5682 };
5683 Ok(SessionArtifact {
5684 source_harness: locator.harness.clone(),
5685 target_harness: target.id(),
5686 session_id: target_session_id
5687 .map(str::to_string)
5688 .or_else(|| session.meta.session_id.clone()),
5689 content,
5690 suggested_filename,
5691 files,
5692 fidelity,
5693 residue,
5694 })
5695}
5696
5697fn append_grok_bundle_files(
5698 locator: &SessionLocator,
5699 prefix: &str,
5700 role: ArtifactFileRole,
5701 files: &mut Vec<SessionArtifactFile>,
5702) -> std::result::Result<(), ServiceError> {
5703 let primary = locator.storage.path();
5704 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5705 return Err(ServiceError::Operation(format!(
5706 "Grok bundle locator must name chat_history.jsonl, got {}",
5707 primary.display()
5708 )));
5709 }
5710 let parent = primary.parent().ok_or_else(|| {
5711 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5712 })?;
5713 for name in ["summary.json", "updates.jsonl"] {
5714 let path = parent.join(name);
5715 let metadata = match std::fs::symlink_metadata(&path) {
5716 Ok(metadata) => metadata,
5717 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5718 Err(error) => return Err(ServiceError::Operation(error.to_string())),
5719 };
5720 if metadata.file_type().is_symlink() || !metadata.is_file() {
5721 return Err(ServiceError::Operation(format!(
5722 "refusing non-regular Grok bundle member {}",
5723 path.display()
5724 )));
5725 }
5726 let content = std::fs::read_to_string(&path).map_err(|error| {
5727 ServiceError::Operation(format!(
5728 "Grok bundle member {} is not representable as UTF-8: {error}",
5729 path.display()
5730 ))
5731 })?;
5732 files.push(SessionArtifactFile {
5733 path: format!("{prefix}{name}"),
5734 content,
5735 role: match role {
5736 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5737 _ => ArtifactFileRole::SourceRecovery,
5738 },
5739 });
5740 }
5741 Ok(())
5742}
5743
5744fn handoff_artifact(
5745 locator: &SessionLocator,
5746 session: &Session,
5747 target: TransferFormat,
5748) -> std::result::Result<SessionArtifact, ServiceError> {
5749 let target_session_id = target_session_id(target);
5750 session_artifact_with_id(locator, session, target, Some(&target_session_id))
5751}
5752
5753fn target_session_id(target: TransferFormat) -> String {
5754 let uuid = generated_session_id();
5755 match target {
5756 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5757 TransferFormat::ClaudeCode
5758 | TransferFormat::Codex
5759 | TransferFormat::Pi
5760 | TransferFormat::Grok
5761 | TransferFormat::Gemini
5762 | TransferFormat::Goose
5763 | TransferFormat::Hermes => uuid,
5764 }
5765}
5766
5767fn sanitize_filename(value: &str) -> String {
5768 let value = value
5769 .chars()
5770 .map(|character| {
5771 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5772 character
5773 } else {
5774 '-'
5775 }
5776 })
5777 .collect::<String>();
5778 let value = value.trim_matches('-');
5779 if value.is_empty() {
5780 "session".into()
5781 } else {
5782 value.chars().take(100).collect()
5783 }
5784}
5785
5786fn handoff_instructions(
5787 target: TransferFormat,
5788 session_id: &str,
5789 cwd: &Path,
5790) -> HandoffInstructions {
5791 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5792 cwd: cwd.to_path_buf(),
5793 program: program.into(),
5794 arguments,
5795 env: BTreeMap::new(),
5796 };
5797 match target {
5798 TransferFormat::ClaudeCode => HandoffInstructions {
5799 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5800 materialize: None,
5801 requires_materialization: true,
5802 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(),
5803 },
5804 TransferFormat::Hermes => HandoffInstructions {
5805 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5806 materialize: None,
5807 requires_materialization: true,
5808 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(),
5809 },
5810 TransferFormat::Codex => HandoffInstructions {
5811 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5812 materialize: None,
5813 requires_materialization: true,
5814 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5815 },
5816 TransferFormat::OpenCode => HandoffInstructions {
5817 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5818 materialize: Some(launch(
5819 "opencode",
5820 vec!["import".into(), "{artifact_path}".into()],
5821 )),
5822 requires_materialization: true,
5823 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5824 },
5825 TransferFormat::Pi => HandoffInstructions {
5826 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5827 materialize: None,
5828 requires_materialization: true,
5829 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5830 },
5831 TransferFormat::Grok => HandoffInstructions {
5832 launch: launch(
5833 "grok",
5834 vec!["--resume".into(), "{materialized_session_id}".into()],
5835 ),
5836 materialize: None,
5837 requires_materialization: true,
5838 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(),
5839 },
5840 TransferFormat::Gemini => HandoffInstructions {
5841 launch: launch(
5842 "gemini",
5843 vec!["--session-file".into(), "{artifact_path}".into()],
5844 ),
5845 materialize: None,
5846 requires_materialization: true,
5847 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(),
5848 },
5849 TransferFormat::Goose => HandoffInstructions {
5850 launch: launch(
5851 "goose",
5852 vec![
5853 "session".into(),
5854 "--resume".into(),
5855 "--session-id".into(),
5856 "{imported_session_id}".into(),
5857 ],
5858 ),
5859 materialize: Some(launch(
5860 "goose",
5861 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5862 )),
5863 requires_materialization: true,
5864 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(),
5865 },
5866 }
5867}
5868
5869fn resume_launch(
5870 harness: &str,
5871 session_id: &str,
5872 cwd: &Path,
5873 policy: ResumePolicy,
5874) -> std::result::Result<StructuredLaunch, ServiceError> {
5875 let mut arguments = Vec::new();
5876 let program = match harness {
5877 HarnessId::GROK => {
5878 if matches!(policy, ResumePolicy::Yolo) {
5879 if crate::support::self_sandbox_supported() {
5880 arguments.extend(["--sandbox".into(), "workspace".into()]);
5881 }
5882 arguments.push("--always-approve".into());
5883 }
5884 arguments.extend(["--resume".into(), session_id.into()]);
5885 "grok"
5886 }
5887 HarnessId::CODEX => {
5888 arguments.extend(crate::startup_prompts::startup_arguments(
5889 harness,
5890 Some(cwd),
5891 &[],
5892 matches!(policy, ResumePolicy::Yolo),
5893 ));
5894 arguments.extend(["resume".into(), session_id.into()]);
5895 "codex"
5896 }
5897 HarnessId::CLAUDE_CODE => {
5898 arguments.extend(crate::startup_prompts::startup_arguments(
5899 harness,
5900 Some(cwd),
5901 &[],
5902 matches!(policy, ResumePolicy::Yolo),
5903 ));
5904 arguments.extend(["--resume".into(), session_id.into()]);
5905 "claude"
5906 }
5907 HarnessId::GEMINI => {
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 "gemini"
5916 }
5917 HarnessId::GOOSE => {
5918 arguments.extend([
5919 "session".into(),
5920 "--resume".into(),
5921 "--session-id".into(),
5922 session_id.into(),
5923 ]);
5924 "goose"
5925 }
5926 HarnessId::PI => {
5927 arguments.extend(crate::startup_prompts::startup_arguments(
5928 harness,
5929 Some(cwd),
5930 &[],
5931 matches!(policy, ResumePolicy::Yolo),
5932 ));
5933 arguments.extend(["--session".into(), session_id.into()]);
5934 "pi"
5935 }
5936 HarnessId::OPENCODE => {
5937 arguments.extend(["--session".into(), session_id.into()]);
5938 "opencode"
5939 }
5940 HarnessId::SUPERCODE => {
5941 if matches!(policy, ResumePolicy::Yolo) {
5942 arguments.push("--dangerous".into());
5943 }
5944 arguments.extend(["resume".into(), session_id.into()]);
5945 "supercode"
5946 }
5947 other => {
5948 return Err(ServiceError::InvalidParams(format!(
5949 "no structured resume launch is registered for harness `{other}`"
5950 )))
5951 }
5952 };
5953 Ok(StructuredLaunch {
5954 cwd: cwd.to_path_buf(),
5955 env: if program == "grok" {
5956 crate::support::grok_home_env()
5957 } else {
5958 BTreeMap::new()
5959 },
5960 program: program.into(),
5961 arguments,
5962 })
5963}
5964
5965fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5971 let digest = blake3::hash(address.as_bytes()).to_hex();
5972 let path = std::env::temp_dir().join(format!(
5973 "supercode-openclaw-gateway-token-{}",
5974 &digest.as_str()[..16]
5975 ));
5976 #[cfg(unix)]
5977 {
5978 use std::io::Write;
5979 use std::os::unix::fs::OpenOptionsExt;
5980 let mut file = std::fs::OpenOptions::new()
5981 .write(true)
5982 .create(true)
5983 .truncate(true)
5984 .mode(0o600)
5985 .open(&path)?;
5986 file.write_all(secret.as_bytes())?;
5987 }
5988 #[cfg(not(unix))]
5989 std::fs::write(&path, secret)?;
5990 Ok(path)
5991}
5992
5993fn open_connect_descriptor(
5999 descriptor: &crate::HarnessSupportDescriptor,
6000 home: &Path,
6001) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
6002 let Some(connect) = &descriptor.runtime.connect_launch else {
6003 return Err(ServiceError::InvalidParams(format!(
6004 "harness `{}` has no registered connect-mode launch",
6005 descriptor.id.as_str()
6006 )));
6007 };
6008 let resolved = connect
6009 .resolve(home)
6010 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6011 match (descriptor.id.as_str(), connect.protocol.as_str()) {
6012 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
6013 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
6014 if let Some(token) = resolved.auth {
6015 backend = backend.with_bearer(token);
6016 }
6017 Ok(Box::new(backend))
6018 }
6019 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
6020 let mut env = BTreeMap::new();
6031 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
6032 if let Some(token) = resolved.auth {
6033 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
6034 .map_err(|error| {
6035 ServiceError::UnsupportedAction(format!(
6036 "could not stage the gateway credential for the bridge: {error}"
6037 ))
6038 })?;
6039 arguments.push("--token-file".into());
6040 arguments.push(token_path.to_string_lossy().into_owned());
6041 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
6042 }
6043 let program = descriptor
6048 .runtime
6049 .default_launch
6050 .as_ref()
6051 .map(|launch| launch.program.clone())
6052 .unwrap_or_else(|| "openclaw".into());
6053 let launch = RuntimeLaunch {
6054 program,
6055 arguments,
6056 env,
6057 };
6058 Ok(Box::new(
6059 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
6060 .with_resume_support(descriptor.runtime.capabilities.resume_session),
6061 ))
6062 }
6063 _ => Err(ServiceError::UnsupportedAction(format!(
6064 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
6065 descriptor.id.as_str(),
6066 connect.protocol
6067 ))),
6068 }
6069}
6070
6071fn registry_connect_descriptor(
6074 params: &RuntimeBackendParams,
6075) -> Option<crate::HarnessSupportDescriptor> {
6076 if params.launch.is_some() || params.base_url.is_some() {
6077 return None;
6078 }
6079 harness_support_registry()
6080 .harnesses
6081 .into_iter()
6082 .find(|descriptor| descriptor.id == params.harness)
6083 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
6084}
6085
6086fn service_home() -> std::result::Result<PathBuf, ServiceError> {
6087 supercode_interchange::user_home()
6088 .map(std::path::PathBuf::into_os_string)
6089 .map(PathBuf::from)
6090 .ok_or_else(|| {
6091 ServiceError::UnsupportedAction(
6092 "connect-mode launches need HOME to locate the harness config".into(),
6093 )
6094 })
6095}
6096
6097pub const RUNTIME_OPEN_METHODS: &[&str] = &[
6100 "harness.v1.runtimes.start",
6101 "harness.v1.runtimes.resume",
6102 "harness.v1.runtimes.attach",
6103 "harness.v1.runtimes.attach_existing",
6104];
6105
6106pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
6111
6112pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
6118
6119pub const DETACHED_METHODS: &[&str] = &[
6125 "harness.v1.harnesses.list",
6126 "harness.v1.sessions.load",
6127 "harness.v1.harnesses.probe",
6128 "harness.v1.sessions.message",
6129 "harness.v1.sessions.new",
6130 "harness.v1.sessions.reset",
6131 "harness.v1.sessions.archive",
6132 "harness.v1.sessions.delete",
6133];
6134
6135pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
6141
6142pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
6150
6151async fn within_control_deadline<F: std::future::Future>(
6154 method: &str,
6155 call: F,
6156) -> std::result::Result<F::Output, ServiceError> {
6157 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
6158 .await
6159 .map_err(|_| {
6160 ServiceError::Operation(format!(
6161 "`{method}` gave up after {}s: the runtime did not answer",
6162 RUNTIME_CONTROL_DEADLINE.as_secs()
6163 ))
6164 })
6165}
6166
6167pub struct RuntimeOpen {
6171 id: Value,
6172 method: String,
6173 params: Value,
6174}
6175
6176impl RuntimeOpen {
6177 pub async fn open(self) -> OpenedRuntime {
6181 let Self { id, method, params } = self;
6182 let outcome = open_runtime(&method, params).await;
6183 OpenedRuntime { id, outcome }
6184 }
6185}
6186
6187pub struct OpenedRuntime {
6190 id: Value,
6191 outcome: std::result::Result<OpenRuntime, ServiceError>,
6192}
6193
6194pub struct DetachedCall {
6199 id: Value,
6200 method: String,
6201 work: std::result::Result<Work, ServiceError>,
6202}
6203
6204impl DetachedCall {
6205 pub async fn run(self) -> DetachedAnswer {
6208 let Self { id, method, work } = self;
6209 match work {
6210 Ok(Work::Runtime(work)) => {
6215 let (result, returned) = work.run().await;
6216 DetachedAnswer {
6217 response: service_response(id, result),
6218 returned,
6219 }
6220 }
6221 Ok(Work::Free(work)) => {
6222 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
6223 Ok(result) => result,
6224 Err(_) => Err(ServiceError::Operation(format!(
6225 "`{method}` gave up after {}s: the harness it waits on did not answer",
6226 DETACHED_CALL_DEADLINE.as_secs()
6227 ))),
6228 };
6229 DetachedAnswer {
6230 response: service_response(id, result),
6231 returned: None,
6232 }
6233 }
6234 Err(error) => DetachedAnswer {
6235 response: service_response(id, Err(error)),
6236 returned: None,
6237 },
6238 }
6239 }
6240}
6241
6242pub struct DetachedAnswer {
6246 response: Value,
6247 returned: Option<ReturnedRuntime>,
6248}
6249
6250impl DetachedAnswer {
6251 pub fn into_response(self) -> Value {
6254 self.response
6255 }
6256}
6257
6258pub struct ReturnedRuntime {
6261 connection: String,
6262 runtime: Box<dyn RuntimeConnection>,
6263}
6264
6265enum Work {
6268 Free(DetachedWork),
6269 Runtime(RuntimeWork),
6270}
6271
6272enum DetachedWork {
6275 Inventory(InventoryWork),
6279 Message(MessageSessionParams),
6281 Load(LoadSessionParams),
6283 SessionMutation {
6286 verb: crate::SessionVerb,
6287 mutation: crate::SessionMutation,
6288 },
6289}
6290
6291impl DetachedWork {
6292 async fn run(self) -> std::result::Result<Value, ServiceError> {
6293 match self {
6294 Self::Inventory(work) => run_inventory(work).await,
6295 Self::Message(params) => Ok(message_live_session(¶ms).await),
6296 Self::Load(params) => tokio::task::spawn_blocking(move || load_session_door(params))
6297 .await
6298 .map_err(|error| {
6299 ServiceError::Operation(format!("the session read stopped: {error}"))
6300 })?,
6301 Self::SessionMutation { verb, mutation } => {
6302 let outcome = run_session_mutation(verb, &mutation).await?;
6303 serde_json::to_value(outcome)
6304 .map_err(|error| ServiceError::Operation(error.to_string()))
6305 }
6306 }
6307 }
6308}
6309
6310enum RuntimeWork {
6312 Close {
6314 runtime: Box<dyn RuntimeConnection>,
6315 process_group: Option<u32>,
6316 },
6317 LiveCommand {
6320 connection: String,
6321 runtime: Box<dyn RuntimeConnection>,
6322 verb: crate::SessionVerb,
6323 mutation: crate::SessionMutation,
6324 command: &'static str,
6325 session: String,
6326 },
6327}
6328
6329type RuntimeWorkAnswer = (
6332 std::result::Result<Value, ServiceError>,
6333 Option<ReturnedRuntime>,
6334);
6335
6336impl RuntimeWork {
6337 async fn run(self) -> RuntimeWorkAnswer {
6338 match self {
6339 Self::Close {
6340 runtime,
6341 process_group,
6342 } => (close_runtime(runtime, process_group).await, None),
6343 Self::LiveCommand {
6344 connection,
6345 mut runtime,
6346 verb,
6347 mutation,
6348 command,
6349 session,
6350 } => {
6351 let result =
6352 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
6353 (
6354 result,
6355 Some(ReturnedRuntime {
6356 connection,
6357 runtime,
6358 }),
6359 )
6360 }
6361 }
6362 }
6363}
6364
6365async fn close_runtime(
6368 mut runtime: Box<dyn RuntimeConnection>,
6369 process_group: Option<u32>,
6370) -> std::result::Result<Value, ServiceError> {
6371 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
6372 Ok(result) => {
6373 result.map_err(operation)?;
6374 Ok(json!({"closed": true}))
6375 }
6376 Err(deadline) => {
6377 let killed = kill_runtime_process_group(process_group);
6382 drop(runtime);
6383 Ok(json!({
6384 "closed": true,
6385 "killed": killed,
6386 "detail": error_message(deadline),
6387 }))
6388 }
6389 }
6390}
6391
6392fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
6395 mutation
6396 .session
6397 .clone()
6398 .filter(|value| !value.trim().is_empty())
6399 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
6400}
6401
6402async fn type_live_command(
6406 runtime: &mut dyn RuntimeConnection,
6407 verb: crate::SessionVerb,
6408 mutation: &crate::SessionMutation,
6409 command: &str,
6410 session: String,
6411) -> std::result::Result<Value, ServiceError> {
6412 within_control_deadline(
6413 &format!("sessions.{}", verb.as_str()),
6414 runtime.send_input(RuntimeInput {
6415 text: command.to_string(),
6416 image_urls: Vec::new(),
6417 }),
6418 )
6419 .await?
6420 .map_err(operation)?;
6421 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
6422 .map_err(session_control_error)?;
6423 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
6424}
6425
6426enum OpenRuntime {
6429 Hosted {
6432 runtime: Box<dyn RuntimeConnection>,
6433 capabilities: crate::RuntimeCapabilities,
6434 workspace: PathBuf,
6435 fresh: bool,
6437 },
6438 Joined { runtime: Box<dyn RuntimeConnection> },
6441}
6442
6443async fn open_runtime(
6448 method: &str,
6449 params: Value,
6450) -> std::result::Result<OpenRuntime, ServiceError> {
6451 match tokio::time::timeout(
6452 RUNTIME_OPEN_DEADLINE,
6453 open_runtime_unbounded(method, params),
6454 )
6455 .await
6456 {
6457 Ok(result) => result,
6458 Err(_) => Err(ServiceError::Operation(format!(
6459 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
6460 RUNTIME_OPEN_DEADLINE.as_secs()
6461 ))),
6462 }
6463}
6464
6465async fn open_runtime_unbounded(
6466 method: &str,
6467 params: Value,
6468) -> std::result::Result<OpenRuntime, ServiceError> {
6469 match method {
6470 "harness.v1.runtimes.start" => {
6471 let params = decode::<RuntimeStartParams>(params)?;
6472 if params.model.is_some()
6475 && (params.backend.protocol.is_some()
6476 || !matches!(
6477 params.backend.harness.as_str(),
6478 HarnessId::CLAUDE_CODE | HarnessId::CODEX
6479 ))
6480 {
6481 return Err(ServiceError::UnsupportedAction(format!(
6482 "{} takes no model at start through supercode",
6483 params.backend.harness.as_str()
6484 )));
6485 }
6486 let backend = runtime_backend(¶ms.backend)?;
6487 let capabilities = backend.capabilities();
6488 let workspace = params.cwd.clone();
6489 let runtime = backend
6490 .start(RuntimeStartRequest {
6491 cwd: params.cwd,
6492 launch: runtime_launch(¶ms.backend),
6493 mcp_servers: params.mcp_servers,
6494 approval_policy: params.approval_policy,
6495 model: params.model,
6496 })
6497 .await
6498 .map_err(operation)?;
6499 Ok(OpenRuntime::Hosted {
6500 runtime,
6501 capabilities,
6502 workspace,
6503 fresh: true,
6504 })
6505 }
6506 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
6507 let params = decode::<RuntimeAttachParams>(params)?;
6508 let backend = runtime_backend(¶ms.backend)?;
6509 let capabilities = backend.capabilities();
6510 let workspace = params
6511 .cwd
6512 .clone()
6513 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
6514 let runtime = backend
6515 .attach(RuntimeAttachRequest {
6516 runtime_id: params.runtime_id,
6517 cwd: params.cwd,
6518 launch: runtime_launch(¶ms.backend),
6519 mcp_servers: params.mcp_servers,
6520 approval_policy: params.approval_policy,
6521 })
6522 .await
6523 .map_err(operation)?;
6524 Ok(OpenRuntime::Hosted {
6525 runtime,
6526 capabilities,
6527 workspace,
6528 fresh: false,
6529 })
6530 }
6531 "harness.v1.runtimes.attach_existing" => {
6532 let params = decode::<RuntimeAttachParams>(params)?;
6533 let backend: Box<dyn RuntimeBackend> = match params
6534 .backend
6535 .base_url
6536 .as_deref()
6537 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
6538 {
6539 Some(endpoint) => {
6540 #[cfg(not(feature = "adapter-api"))]
6541 {
6542 let _ = endpoint;
6543 return Err(ServiceError::UnsupportedAction(
6544 "live HTTP attachment adapter is not compiled".into(),
6545 ));
6546 }
6547 #[cfg(feature = "adapter-api")]
6548 {
6549 let workspace = params.cwd.clone().ok_or_else(|| {
6550 ServiceError::InvalidParams(
6551 "Volter Harness live attach requires the project cwd".into(),
6552 )
6553 })?;
6554 let source = LiveRuntimeSource {
6555 harness: params.backend.harness.as_str().to_string(),
6556 session_id: params.runtime_id.clone(),
6557 workspace,
6558 };
6559 let receipt = resolve_live_runtime(&endpoint, &source)
6560 .map_err(|error| ServiceError::Operation(error.to_string()))?;
6561 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
6562 }
6563 }
6564 None => runtime_backend(¶ms.backend)?,
6565 };
6566 let capabilities = backend.capabilities();
6567 if !capabilities.attach_existing_process {
6568 return Err(ServiceError::Operation(format!(
6569 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
6570 backend.harness().as_str()
6571 )));
6572 }
6573 let runtime = backend
6574 .attach_existing(RuntimeAttachRequest {
6575 runtime_id: params.runtime_id,
6576 cwd: params.cwd,
6577 launch: runtime_launch(¶ms.backend),
6578 mcp_servers: params.mcp_servers,
6579 approval_policy: params.approval_policy,
6580 })
6581 .await
6582 .map_err(operation)?;
6583 Ok(OpenRuntime::Joined { runtime })
6584 }
6585 _ => Err(ServiceError::MethodNotFound),
6586 }
6587}
6588
6589fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
6591 match result {
6592 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
6593 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
6594 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
6595 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
6596 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
6597 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
6598 }
6599}
6600
6601fn runtime_backend(
6602 params: &RuntimeBackendParams,
6603) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
6604 if let Some(descriptor) = registry_connect_descriptor(params) {
6605 return open_connect_descriptor(&descriptor, &service_home()?);
6606 }
6607 if params.protocol.as_deref() == Some("acp") {
6608 let launch = params
6609 .launch
6610 .clone()
6611 .or_else(|| {
6612 harness_support_registry()
6613 .harnesses
6614 .into_iter()
6615 .find(|harness| harness.id == params.harness)
6616 .filter(|harness| {
6617 harness.runtime.implementation == ImplementationKind::GenericProtocol
6618 && harness.runtime.protocol.starts_with("acp")
6619 })
6620 .and_then(|harness| harness.runtime.default_launch)
6621 })
6622 .ok_or_else(|| {
6623 ServiceError::InvalidParams(
6624 "an ACP runtime requires `launch` unless the harness has a registered default"
6625 .into(),
6626 )
6627 })?;
6628 let resume_session = harness_support_registry()
6629 .harnesses
6630 .into_iter()
6631 .find(|harness| harness.id == params.harness)
6632 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
6633 return Ok(Box::new(
6634 AcpRuntimeBackend::new(params.harness.clone(), launch)
6635 .with_resume_support(resume_session),
6636 ));
6637 }
6638 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
6639 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
6640 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
6641 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
6642 HarnessId::OPENCODE => match ¶ms.base_url {
6643 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
6644 None => Box::new(OpenCodeRuntimeBackend::new()),
6645 },
6646 harness => {
6647 let descriptor = harness_support_registry()
6648 .harnesses
6649 .into_iter()
6650 .find(|descriptor| descriptor.id.as_str() == harness)
6651 .filter(|descriptor| {
6652 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
6653 && descriptor.runtime.protocol.starts_with("acp")
6654 });
6655 let Some(descriptor) = descriptor else {
6656 return Err(ServiceError::InvalidParams(format!(
6657 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
6658 )));
6659 };
6660 let resume = descriptor.runtime.capabilities.resume_session;
6661 Box::new(
6662 AcpRuntimeBackend::new(
6663 descriptor.id,
6664 descriptor
6665 .runtime
6666 .default_launch
6667 .expect("generic ACP registry entry includes its launch"),
6668 )
6669 .with_resume_support(resume),
6670 )
6671 }
6672 };
6673 Ok(backend)
6674}
6675
6676fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
6677 if let Some(launch) = ¶ms.launch {
6678 return Some(launch.clone());
6679 }
6680 if !matches!(params.policy, RuntimePolicy::Yolo) {
6681 return None;
6682 }
6683 let launch = match params.harness.as_str() {
6684 HarnessId::GROK => RuntimeLaunch {
6685 program: "grok".into(),
6686 arguments: {
6687 let mut arguments: Vec<String> = Vec::new();
6688 if crate::support::self_sandbox_supported() {
6689 arguments.extend(["--sandbox".into(), "workspace".into()]);
6690 }
6691 arguments.extend([
6692 "--always-approve".into(),
6693 "agent".into(),
6694 "--no-leader".into(),
6695 "stdio".into(),
6696 ]);
6697 arguments
6698 },
6699 env: crate::support::grok_env(),
6700 },
6701 HarnessId::CODEX => RuntimeLaunch {
6702 program: "codex".into(),
6703 arguments: vec![
6704 "--dangerously-bypass-approvals-and-sandbox".into(),
6705 "--dangerously-bypass-hook-trust".into(),
6706 "app-server".into(),
6707 ],
6708 env: BTreeMap::new(),
6709 },
6710 HarnessId::CLAUDE_CODE => RuntimeLaunch {
6711 program: "claude".into(),
6712 arguments: vec![
6713 "--dangerously-skip-permissions".into(),
6714 "--print".into(),
6715 "--input-format".into(),
6716 "stream-json".into(),
6717 "--output-format".into(),
6718 "stream-json".into(),
6719 "--verbose".into(),
6720 ],
6721 env: BTreeMap::new(),
6722 },
6723 HarnessId::PI => RuntimeLaunch {
6724 program: "pi".into(),
6725 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
6726 env: BTreeMap::new(),
6727 },
6728 HarnessId::OPENCODE => RuntimeLaunch {
6729 program: "opencode".into(),
6730 arguments: vec!["serve".into()],
6731 env: BTreeMap::new(),
6732 },
6733 HarnessId::GEMINI => RuntimeLaunch {
6734 program: "gemini".into(),
6735 arguments: vec!["--acp".into(), "--yolo".into()],
6736 env: BTreeMap::new(),
6737 },
6738 HarnessId::GOOSE => RuntimeLaunch {
6739 program: "goose".into(),
6740 arguments: vec!["acp".into()],
6741 env: BTreeMap::new(),
6742 },
6743 HarnessId::SUPERCODE => RuntimeLaunch {
6744 program: "supercode".into(),
6745 arguments: vec!["acp".into(), "--dangerous".into()],
6746 env: BTreeMap::new(),
6747 },
6748 _ => return None,
6749 };
6750 Some(launch)
6751}
6752
6753struct IsolatedProbeHome {
6759 launch: RuntimeLaunch,
6760 root: PathBuf,
6761}
6762
6763impl IsolatedProbeHome {
6764 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6765 let root = std::env::temp_dir().join(format!(
6766 "supercode-harness-probe-{harness}-{}",
6767 generated_session_id()
6768 ));
6769 std::fs::create_dir_all(&root)?;
6770 set_private_dir_permissions(&root)?;
6771
6772 if let Some(source_home) = supercode_interchange::user_home()
6773 .map(std::path::PathBuf::into_os_string)
6774 .map(PathBuf::from)
6775 {
6776 for relative in probe_auth_files(harness) {
6777 copy_probe_file(&source_home, &root, relative)?;
6778 }
6779 }
6780 if harness == HarnessId::SUPERCODE {
6785 let config_home = crate::agent::global_instructions_dir();
6786 for file in ["config.toml", "credentials.toml"] {
6787 copy_probe_path(
6788 &config_home.join(file),
6789 &root.join(".config/supercode").join(file),
6790 )?;
6791 }
6792 }
6793 configure_isolated_probe_auth(harness, &root)?;
6794
6795 let root_text = root.to_string_lossy().into_owned();
6796 for (key, value) in [
6797 ("HOME", root_text.clone()),
6798 (
6799 "XDG_CACHE_HOME",
6800 root.join(".cache").to_string_lossy().into_owned(),
6801 ),
6802 (
6803 "XDG_CONFIG_HOME",
6804 root.join(".config").to_string_lossy().into_owned(),
6805 ),
6806 (
6807 "XDG_DATA_HOME",
6808 root.join(".local/share").to_string_lossy().into_owned(),
6809 ),
6810 ] {
6811 launch.env.insert(key.into(), value);
6812 }
6813 let scoped = match harness {
6814 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6815 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6816 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6817 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6818 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6819 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6820 _ => None,
6821 };
6822 if let Some((key, value)) = scoped {
6823 launch
6824 .env
6825 .insert(key.into(), value.to_string_lossy().into_owned());
6826 }
6827 Ok(Self { launch, root })
6828 }
6829
6830 fn cleanup(&self) -> std::io::Result<()> {
6831 match std::fs::remove_dir_all(&self.root) {
6832 Ok(()) => Ok(()),
6833 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6834 Err(error) => Err(error),
6835 }
6836 }
6837}
6838
6839impl Drop for IsolatedProbeHome {
6840 fn drop(&mut self) {
6841 let _ = self.cleanup();
6842 }
6843}
6844
6845fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6846 match harness {
6847 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6848 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6852 HarnessId::CODEX => &[".codex/auth.json"],
6853 HarnessId::GEMINI => &[
6854 ".gemini/google_accounts.json",
6855 ".gemini/oauth_creds.json",
6856 ".gemini/settings.json",
6857 ],
6858 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6859 HarnessId::OPENCODE => &[
6860 ".config/opencode/auth.json",
6861 ".local/share/opencode/auth.json",
6862 ],
6863 HarnessId::PI => &[".pi/agent/auth.json"],
6864 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6869 _ => &[],
6870 }
6871}
6872
6873fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6874 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6875}
6876
6877fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6878 if !source.is_file() {
6879 return Ok(());
6880 }
6881 if let Some(parent) = destination.parent() {
6882 std::fs::create_dir_all(parent)?;
6883 set_private_dir_permissions(parent)?;
6884 }
6885 std::fs::copy(source, destination)?;
6886 set_private_file_permissions(destination)
6887}
6888
6889fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6890 if harness != HarnessId::GEMINI {
6891 return Ok(());
6892 }
6893 let oauth = probe_home.join(".gemini/oauth_creds.json");
6894 if !oauth.is_file() {
6895 return Ok(());
6896 }
6897 let settings_path = probe_home.join(".gemini/settings.json");
6898 let mut settings = std::fs::read_to_string(&settings_path)
6899 .ok()
6900 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6901 .unwrap_or_else(|| json!({}));
6902 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6903 std::fs::write(
6904 &settings_path,
6905 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6906 )?;
6907 set_private_file_permissions(&settings_path)
6908}
6909
6910#[cfg(unix)]
6911fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6912 use std::os::unix::fs::PermissionsExt;
6913 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6914}
6915
6916#[cfg(not(unix))]
6917fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6918 Ok(())
6919}
6920
6921#[cfg(unix)]
6922fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6923 use std::os::unix::fs::PermissionsExt;
6924 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6925}
6926
6927#[cfg(not(unix))]
6928fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6929 Ok(())
6930}
6931
6932fn find_executable(program: &str) -> Option<PathBuf> {
6933 let candidate = PathBuf::from(program);
6934 if candidate.components().count() > 1 {
6935 return candidate.is_file().then_some(candidate);
6936 }
6937 let path = std::env::var_os("PATH")?;
6938 for directory in std::env::split_paths(&path) {
6939 let candidate = directory.join(program);
6940 if candidate.is_file() {
6941 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6942 }
6943 #[cfg(windows)]
6944 {
6945 for extension in ["exe", "cmd", "bat"] {
6946 let candidate = directory.join(format!("{program}.{extension}"));
6947 if candidate.is_file() {
6948 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6949 }
6950 }
6951 }
6952 }
6953 None
6954}
6955
6956async fn executable_version(executable: &Path) -> Option<String> {
6957 let mut command = tokio::process::Command::new(executable);
6958 command
6959 .arg("--version")
6960 .stdin(std::process::Stdio::null())
6961 .stdout(std::process::Stdio::piped())
6962 .stderr(std::process::Stdio::piped())
6963 .kill_on_drop(true);
6964 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6965 .await
6966 .ok()?
6967 .ok()?;
6968 let stdout = String::from_utf8_lossy(&output.stdout);
6969 let stderr = String::from_utf8_lossy(&output.stderr);
6970 stdout
6971 .lines()
6972 .chain(stderr.lines())
6973 .map(str::trim)
6974 .find(|line| !line.is_empty())
6975 .map(|line| truncate_text(line, 200))
6976}
6977
6978pub(crate) fn auth_evidence(harness: &str) -> bool {
6979 let env_names: &[&str] = match harness {
6980 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6981 HarnessId::CODEX => &["OPENAI_API_KEY"],
6982 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6983 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6984 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6985 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6986 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6987 _ => &[],
6988 };
6989 if env_names
6990 .iter()
6991 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6992 {
6993 return true;
6994 }
6995 let Some(home) = supercode_interchange::user_home()
6996 .map(std::path::PathBuf::into_os_string)
6997 .map(PathBuf::from)
6998 else {
6999 return false;
7000 };
7001 let files: Vec<PathBuf> = match harness {
7002 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
7003 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
7004 HarnessId::OPENCODE => vec![
7005 home.join(".local/share/opencode/auth.json"),
7006 home.join(".config/opencode/auth.json"),
7007 ],
7008 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
7009 HarnessId::GROK => vec![home.join(".grok/auth.json")],
7010 HarnessId::GEMINI => vec![
7011 home.join(".gemini/oauth_creds.json"),
7012 home.join(".gemini/google_accounts.json"),
7013 ],
7014 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
7015 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
7016 _ => Vec::new(),
7017 };
7018 if files.into_iter().any(|path| {
7019 std::fs::metadata(path)
7020 .map(|metadata| metadata.is_file() && metadata.len() > 2)
7021 .unwrap_or(false)
7022 }) {
7023 return true;
7024 }
7025 if harness == HarnessId::CLAUDE_CODE {
7032 return std::fs::read_to_string(home.join(".claude.json"))
7033 .map(|text| text.contains("\"oauthAccount\""))
7034 .unwrap_or(false);
7035 }
7036 false
7037}
7038
7039fn looks_like_auth_error(message: &str) -> bool {
7040 let message = message.to_ascii_lowercase();
7041 [
7042 "auth",
7043 "login",
7044 "sign in",
7045 "sign-in",
7046 "credential",
7047 "unauthorized",
7048 "forbidden",
7049 "token",
7050 ]
7051 .iter()
7052 .any(|needle| message.contains(needle))
7053}
7054
7055fn unavailable_capabilities() -> crate::RuntimeCapabilities {
7056 crate::RuntimeCapabilities {
7057 start_session: false,
7058 resume_session: false,
7059 attach_existing_process: false,
7060 send_input: false,
7061 stream_events: false,
7062 interrupt: false,
7063 steer: false,
7064 respond_to_requests: false,
7065 }
7066}
7067
7068fn truncate_text(text: &str, max_chars: usize) -> String {
7069 let mut chars = text.chars();
7070 let truncated = chars.by_ref().take(max_chars).collect::<String>();
7071 if chars.next().is_some() {
7072 format!("{truncated}…")
7073 } else {
7074 truncated
7075 }
7076}
7077
7078fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
7085 match &handle.endpoint {
7086 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
7087 crate::RuntimeEndpoint::Http { .. } => None,
7088 }
7089}
7090
7091fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
7097 match process_group {
7098 #[cfg(unix)]
7099 Some(pid) => {
7100 crate::lsp::kill_process_group(pid);
7101 true
7102 }
7103 #[cfg(not(unix))]
7104 Some(_) => false,
7105 None => false,
7106 }
7107}
7108
7109fn error_message(error: ServiceError) -> String {
7110 match error {
7111 ServiceError::InvalidParams(message)
7112 | ServiceError::Operation(message)
7113 | ServiceError::UnsupportedAction(message) => message,
7114 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
7115 ServiceError::Sdk(error) => error.to_string(),
7116 }
7117}
7118
7119#[derive(Debug)]
7120enum ServiceError {
7121 InvalidParams(String),
7122 MethodNotFound,
7123 UnsupportedAction(String),
7124 Operation(String),
7125 Sdk(SdkError),
7126}
7127
7128fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
7129 match error {
7130 ServiceError::InvalidParams(message) => {
7131 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
7132 }
7133 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
7134 SdkError::unsupported(operation)
7135 }
7136 ServiceError::Operation(message) => {
7137 let code = if message.contains("already in progress") {
7138 SdkErrorCode::Busy
7139 } else if message.contains("not supported by this runtime") {
7140 SdkErrorCode::UnsupportedAction
7141 } else if message.contains("unknown runtime connection") {
7142 SdkErrorCode::NotFound
7143 } else {
7144 SdkErrorCode::Execution
7145 };
7146 SdkError::new(code, operation, message)
7147 }
7148 ServiceError::Sdk(error) => error,
7149 }
7150}
7151
7152fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
7153 let error_code = error.code();
7154 let code = match error_code {
7155 SdkErrorCode::Unauthenticated => -32030,
7156 SdkErrorCode::Unauthorized => -32031,
7157 SdkErrorCode::ControllerRequired => -32032,
7158 SdkErrorCode::LeaseExpired => -32033,
7159 SdkErrorCode::InvalidArgument => -32602,
7160 SdkErrorCode::NotFound => -32004,
7161 SdkErrorCode::Busy => -32000,
7162 SdkErrorCode::UnsupportedAction => -32020,
7163 SdkErrorCode::Execution => -32002,
7164 SdkErrorCode::Transport => -32003,
7165 };
7166 json!({
7167 "jsonrpc": "2.0",
7168 "id": id,
7169 "error": {
7170 "code": code,
7171 "name": error_code,
7172 "operation": error.operation(),
7173 "message": error.to_string(),
7174 },
7175 })
7176}
7177
7178fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
7179 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
7180}
7181
7182fn operation(error: impl Into<crate::Error>) -> ServiceError {
7183 let error = error.into();
7184 match error {
7185 crate::Error::Sdk(error) => ServiceError::Sdk(error),
7186 error => ServiceError::Operation(error.to_string()),
7187 }
7188}
7189
7190#[derive(Debug, Clone, Deserialize, Default)]
7194#[serde(default)]
7195struct MemoryRequest {
7196 harness: Option<String>,
7198 query: Option<String>,
7200 profile: Option<String>,
7202 session: Option<String>,
7204 full: bool,
7206 regex: bool,
7208 cwd: Option<std::path::PathBuf>,
7210 homes: crate::HarnessHomes,
7212}
7213
7214fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7217 let request = decode::<MemoryRequest>(params)?;
7218 let harness = request
7219 .harness
7220 .clone()
7221 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7222 let to_service = |error: crate::memory::MemoryError| match error {
7223 crate::memory::MemoryError::UnsupportedHarness { .. }
7224 | crate::memory::MemoryError::SessionNotScoped { .. } => {
7225 ServiceError::UnsupportedAction(error.to_string())
7226 }
7227 other => ServiceError::InvalidParams(other.to_string()),
7228 };
7229 match method {
7230 "harness.v1.memory.show" => {
7231 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
7232 harness,
7233 profile: request.profile,
7234 session: request.session,
7235 full: request.full,
7236 cwd: request.cwd,
7237 homes: request.homes,
7238 })
7239 .map_err(to_service)?;
7240 Ok(json!({
7241 "schema": crate::memory::MEMORY_SCHEMA,
7242 "documents": documents,
7243 }))
7244 }
7245 "harness.v1.memory.search" => {
7246 let query = request
7247 .query
7248 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
7249 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
7250 harness,
7251 query,
7252 profile: request.profile,
7253 regex: request.regex,
7254 cwd: request.cwd,
7255 homes: request.homes,
7256 })
7257 .map_err(to_service)?;
7258 Ok(json!({
7259 "schema": crate::memory::MEMORY_SCHEMA,
7260 "matches": matches,
7261 }))
7262 }
7263 _ => Err(ServiceError::MethodNotFound),
7264 }
7265}
7266
7267#[derive(Debug, Clone, Deserialize)]
7271#[serde(default)]
7272struct ProfilesQuery {
7273 harness: Option<String>,
7275 name: Option<String>,
7277 homes: crate::HarnessHomes,
7279}
7280
7281impl Default for ProfilesQuery {
7282 fn default() -> Self {
7283 Self {
7284 harness: None,
7285 name: None,
7286 homes: crate::HarnessHomes::default(),
7287 }
7288 }
7289}
7290
7291fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7294 let query = decode::<ProfilesQuery>(params)?;
7295 let to_service = |error: crate::profiles::ProfileError| match error {
7296 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
7297 ServiceError::UnsupportedAction(error.to_string())
7298 }
7299 crate::profiles::ProfileError::NotFound { .. } => {
7300 ServiceError::InvalidParams(error.to_string())
7301 }
7302 };
7303 match method {
7304 "harness.v1.profiles.list" => {
7305 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
7306 .map_err(to_service)?;
7307 Ok(json!({
7308 "schema": crate::profiles::PROFILES_SCHEMA,
7309 "profiles": profiles,
7310 }))
7311 }
7312 "harness.v1.profiles.get" => {
7313 let harness = query
7314 .harness
7315 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7316 let name = query
7317 .name
7318 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7319 let profile =
7320 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
7321 Ok(json!({
7322 "schema": crate::profiles::PROFILES_SCHEMA,
7323 "profile": profile,
7324 }))
7325 }
7326 _ => Err(ServiceError::MethodNotFound),
7327 }
7328}
7329
7330#[derive(Debug, Clone, Deserialize)]
7334#[serde(default)]
7335struct ChannelsQuery {
7336 harness: Option<String>,
7338 name: Option<String>,
7340 homes: crate::HarnessHomes,
7342}
7343
7344impl Default for ChannelsQuery {
7345 fn default() -> Self {
7346 Self {
7347 harness: None,
7348 name: None,
7349 homes: crate::HarnessHomes::default(),
7350 }
7351 }
7352}
7353
7354#[derive(Debug, Clone, Deserialize)]
7358#[serde(default)]
7359struct RoutesQuery {
7360 harness: Option<String>,
7361 profile: Option<String>,
7363 homes: crate::HarnessHomes,
7364}
7365
7366impl Default for RoutesQuery {
7367 fn default() -> Self {
7368 Self {
7369 harness: None,
7370 profile: None,
7371 homes: crate::HarnessHomes::default(),
7372 }
7373 }
7374}
7375
7376#[derive(Debug, Clone, Deserialize)]
7377#[serde(default)]
7378struct TriggersQuery {
7379 harness: Option<String>,
7380 homes: crate::HarnessHomes,
7381}
7382
7383impl Default for TriggersQuery {
7384 fn default() -> Self {
7385 Self {
7386 harness: None,
7387 homes: crate::HarnessHomes::default(),
7388 }
7389 }
7390}
7391
7392fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
7393 let query = decode::<TriggersQuery>(params)?;
7394 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
7395 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7396 Ok(json!({
7397 "schema": crate::triggers::TRIGGERS_SCHEMA,
7398 "triggers": triggers,
7399 }))
7400}
7401
7402fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
7403 let query = decode::<RoutesQuery>(params)?;
7404 let routes = crate::routes::list_routes(
7405 &query.homes,
7406 query.harness.as_deref(),
7407 query.profile.as_deref(),
7408 )
7409 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7410 Ok(json!({
7411 "schema": crate::routes::ROUTES_SCHEMA,
7412 "routes": routes,
7413 }))
7414}
7415
7416fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7417 let query = decode::<ChannelsQuery>(params)?;
7418 let to_service = |error: crate::channels::ChannelError| match error {
7419 crate::channels::ChannelError::UnsupportedHarness { .. } => {
7420 ServiceError::UnsupportedAction(error.to_string())
7421 }
7422 crate::channels::ChannelError::NotFound { .. } => {
7423 ServiceError::InvalidParams(error.to_string())
7424 }
7425 };
7426 match method {
7427 "harness.v1.channels.list" => {
7428 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
7429 .map_err(to_service)?;
7430 Ok(json!({
7431 "schema": crate::channels::CHANNELS_SCHEMA,
7432 "channels": channels,
7433 }))
7434 }
7435 "harness.v1.channels.status" => {
7436 let harness = query
7437 .harness
7438 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7439 let name = query
7440 .name
7441 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7442 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
7443 .map_err(to_service)?;
7444 Ok(json!({
7445 "schema": crate::channels::CHANNELS_SCHEMA,
7446 "channel": channel,
7447 }))
7448 }
7449 _ => Err(ServiceError::MethodNotFound),
7450 }
7451}
7452
7453fn rpc_error(id: Value, code: i64, message: &str) -> Value {
7454 json!({
7455 "jsonrpc": "2.0",
7456 "id": id,
7457 "error": {"code": code, "message": message},
7458 })
7459}