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, resolve_live_runtime, LiveRuntimeRegistration};
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.load",
52 "harness.v1.sessions.follow",
53 "harness.v1.sessions.unfollow",
54 "harness.v1.sessions.activity.subscribe",
55 "harness.v1.sessions.activity.unsubscribe",
56 "harness.v1.sessions.activity_under",
57 "harness.v1.sessions.index.subscribe",
58 "harness.v1.sessions.index.resize",
59 "harness.v1.sessions.index.unsubscribe",
60 "harness.v1.sessions.message",
61 "harness.v1.sessions.inbox",
62 "harness.v1.sessions.import",
63 "harness.v1.sessions.export",
64 "harness.v1.sessions.translate",
65 "harness.v1.sessions.reduce",
66 "harness.v1.sessions.branch",
67 "harness.v1.sessions.handoff",
68 "harness.v1.sessions.materialize",
69 "harness.v1.sessions.resume_instructions",
70 "harness.v1.skills.list",
71 "harness.v1.skills.install",
72 "harness.v1.skills.remove",
73 "harness.v1.memory.show",
74 "harness.v1.memory.search",
75 "harness.v1.jobs.list",
76 "harness.v1.jobs.get",
77 "harness.v1.jobs.create",
78 "harness.v1.jobs.update",
79 "harness.v1.jobs.pause",
80 "harness.v1.jobs.resume",
81 "harness.v1.jobs.run",
82 "harness.v1.jobs.delete",
83 "harness.v1.jobs.notepad",
84 "harness.v1.jobs.notepad_set",
85 "harness.v1.jobs.notepad_delete",
86 "harness.v1.sessions.new",
87 "harness.v1.sessions.reset",
88 "harness.v1.sessions.archive",
89 "harness.v1.sessions.delete",
90 "harness.v1.runs.list",
91 "harness.v1.runs.get",
92 "harness.v1.approvals.list",
93 "harness.v1.approvals.resolve",
94 "harness.v1.runtimes.capabilities",
95 "harness.v1.runtimes.start",
96 "harness.v1.runtimes.resume",
97 "harness.v1.runtimes.attach_existing",
98 "harness.v1.runtimes.attach",
99 "harness.v1.runtimes.send_input",
100 "harness.v1.runtimes.interrupt",
101 "harness.v1.runtimes.steer",
102 "harness.v1.runtimes.respond",
103 "harness.v1.runtimes.terminal_instructions",
104 "harness.v1.runtimes.acquire_control",
105 "harness.v1.runtimes.heartbeat",
106 "harness.v1.runtimes.detach",
107 "harness.v1.runtimes.close",
108 "harness.v1.profiles.list",
109 "harness.v1.profiles.get",
110 "harness.v1.profiles.create",
111 "harness.v1.profiles.delete",
112 "harness.v1.channels.list",
113 "harness.v1.routes.list",
114 "harness.v1.triggers.list",
115 "harness.v1.channels.status",
116 "harness.v1.orchestration.load",
117 "harness.v1.orchestration.save",
118 "harness.v1.orchestration.compile",
119 "harness.v1.orchestration.decompile",
120 "harness.v1.orchestration.import",
121 "harness.v1.orchestration.export",
122 "harness.v1.workflow.load",
123];
124
125pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
127pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
129pub const SESSION_ACTIVITY_EVENT_METHOD: &str = "harness.v1.sessions.activity_event";
131pub const SESSION_INDEX_EVENT_METHOD: &str = "harness.v1.sessions.index_event";
133pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
135
136pub struct HarnessSessionService {
139 catalog: HarnessCatalog,
140 followers: BTreeMap<String, SessionFollower>,
141 followed_sources: BTreeMap<String, FollowedSource>,
142 activity_subscriptions: BTreeMap<String, ActivitySubscription>,
143 index_subscriptions: BTreeMap<String, crate::session_index::SessionIndexSubscription>,
144 index_notifier: Arc<Notify>,
145 #[cfg(feature = "adapter-api")]
146 activity_monitor: crate::session_activity::SessionActivityMonitor,
147 next_subscription: u64,
148 runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
149 runtimes_in_flight: BTreeSet<String>,
154 terminal_launches: BTreeMap<String, StructuredLaunch>,
155 runtime_sequences: BTreeMap<String, u64>,
156 next_runtime: u64,
157 reduction_store_root: Option<PathBuf>,
158 approvals: crate::approvals::ApprovalRegistry,
162 subagent_approvals: Option<Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>>,
165}
166
167impl Default for HarnessSessionService {
168 fn default() -> Self {
169 Self::new()
170 }
171}
172
173impl HarnessSessionService {
174 pub fn new() -> Self {
176 Self {
177 catalog: HarnessCatalog::new(),
178 followers: BTreeMap::new(),
179 followed_sources: BTreeMap::new(),
180 activity_subscriptions: BTreeMap::new(),
181 index_subscriptions: BTreeMap::new(),
182 index_notifier: Arc::new(Notify::new()),
183 #[cfg(feature = "adapter-api")]
184 activity_monitor: Default::default(),
185 next_subscription: 1,
186 runtimes: BTreeMap::new(),
187 runtimes_in_flight: BTreeSet::new(),
188 terminal_launches: BTreeMap::new(),
189 runtime_sequences: BTreeMap::new(),
190 next_runtime: 1,
191 reduction_store_root: None,
192 approvals: crate::approvals::ApprovalRegistry::new(),
193 subagent_approvals: None,
194 }
195 }
196
197 pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
202 self.reduction_store_root = Some(root.into());
203 self
204 }
205
206 pub fn observe_subagent_approvals(
214 &mut self,
215 queue: Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>,
216 ) {
217 self.subagent_approvals = Some(queue);
218 }
219
220 pub fn approvals(&self, query: &crate::approvals::ApprovalsQuery) -> Vec<crate::ApprovalRow> {
227 let now = crate::approvals::now_ms();
228 let mut rows = self.approvals.rows(now);
229 if let Some(queue) = self.subagent_approvals.as_ref() {
230 let queued = queue
231 .lock()
232 .unwrap_or_else(std::sync::PoisonError::into_inner)
233 .clone();
234 rows.extend(crate::approvals::subagent_rows(&queued, now));
235 }
236 rows.retain(|row| query.matches(row));
237 rows.sort_by(|left, right| {
238 left.requested_at_ms
239 .cmp(&right.requested_at_ms)
240 .then_with(|| left.id.cmp(&right.id))
241 });
242 rows
243 }
244
245 async fn approvals_resolve(
255 &mut self,
256 params: Value,
257 ) -> std::result::Result<Value, ServiceError> {
258 let params = decode::<crate::approvals::ApprovalsResolveParams>(params)?;
259 if params.id.trim().is_empty() {
260 return Err(ServiceError::InvalidParams(
261 "approvals resolve requires the `id` of a listed approval row".into(),
262 ));
263 }
264 let choice = match (params.decision, params.option_id.as_deref()) {
265 (Some(_), Some(_)) => {
266 return Err(ServiceError::InvalidParams(
267 "approvals resolve takes either `decision` or `option_id`, not both".into(),
268 ))
269 }
270 (Some(decision), None) => crate::approvals::ApprovalChoice::Decision(decision),
271 (None, Some(option)) => crate::approvals::ApprovalChoice::Option(option.to_string()),
272 (None, None) => {
273 return Err(ServiceError::InvalidParams(format!(
274 "approvals resolve requires `decision` ({}) or an explicit `option_id`",
275 crate::approvals::ApprovalDecision::ALL
276 .map(|decision| decision.as_str())
277 .join(" | "),
278 )))
279 }
280 };
281 let resolution = self
282 .approvals
283 .resolution(¶ms.id, &choice)
284 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?;
285 self.runtime_call(
289 "harness.v1.runtimes.respond",
290 json!({
291 "connection": resolution.connection,
292 "request_id": resolution.request_id,
293 "response": resolution.response,
294 }),
295 )
296 .await?;
297 Ok(json!({
298 "id": params.id,
299 "decision": params.decision.map(|decision| decision.as_str()),
300 "option_id": resolution.option_id,
301 "resolved": true,
302 }))
303 }
304
305 #[cfg(feature = "adapter-api")]
308 pub fn session_index_notifier(&self) -> Arc<Notify> {
309 Arc::clone(&self.index_notifier)
310 }
311
312 #[cfg(feature = "adapter-api")]
314 pub fn handle(&mut self, request: Value) -> Value {
315 let id = request.get("id").cloned().unwrap_or(Value::Null);
316 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
317 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
318 }
319 let Some(method) = request.get("method").and_then(Value::as_str) else {
320 return rpc_error(id, -32600, "request is missing `method`");
321 };
322 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
323 match self.call(method, params) {
324 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
325 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
326 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
327 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
328 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
329 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
330 }
331 }
332
333 #[cfg(feature = "adapter-api")]
336 pub async fn handle_async(&mut self, request: Value) -> Value {
337 let method = request
338 .get("method")
339 .and_then(Value::as_str)
340 .unwrap_or_default();
341 if matches!(
342 method,
343 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
344 ) {
345 let id = request.get("id").cloned().unwrap_or(Value::Null);
346 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
347 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
348 }
349 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
350 return match self.inventory_call(method, params).await {
351 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
352 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
353 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
354 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
355 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
356 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
357 };
358 }
359 if matches!(
360 method,
361 "harness.v1.harnesses.auth.methods"
362 | "harness.v1.harnesses.auth.begin"
363 | "harness.v1.harnesses.auth.verify"
364 ) {
365 let id = request.get("id").cloned().unwrap_or(Value::Null);
366 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
367 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
368 }
369 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
370 return match self.harness_authentication_call(method, params).await {
371 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
372 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
373 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
374 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
375 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
376 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
377 };
378 }
379 if matches!(
385 method,
386 "harness.v1.sessions.new"
387 | "harness.v1.sessions.reset"
388 | "harness.v1.sessions.archive"
389 | "harness.v1.sessions.delete"
390 ) {
391 let id = request.get("id").cloned().unwrap_or(Value::Null);
392 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
393 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
394 }
395 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
396 let verb = match method {
397 "harness.v1.sessions.new" => crate::SessionVerb::New,
398 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
399 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
400 _ => crate::SessionVerb::Delete,
401 };
402 return match self.mutate_session(verb, params).await {
403 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
404 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
405 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
406 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
407 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
408 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
409 };
410 }
411 if method == "harness.v1.sessions.message" {
412 let id = request.get("id").cloned().unwrap_or(Value::Null);
413 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
414 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
415 }
416 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
417 return match self.message_call(params).await {
418 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
419 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
420 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
421 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
422 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
423 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
424 };
425 }
426 if matches!(
427 method,
428 "harness.v1.harnesses.settings" | "harness.v1.harnesses.configure"
429 ) {
430 let id = request.get("id").cloned().unwrap_or(Value::Null);
431 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
432 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
433 }
434 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
435 return match self.harness_settings_call(method, params) {
436 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
437 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
438 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
439 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
440 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
441 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
442 };
443 }
444 if method == "harness.v1.sessions.activity_under" {
445 let id = request.get("id").cloned().unwrap_or(Value::Null);
446 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
447 return match activity_under_call(params).await {
448 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
449 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
450 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
451 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
452 Err(_) => rpc_error(id, -32000, "sessions.activity_under failed"),
453 };
454 }
455 if method == "harness.v1.sessions.activity.subscribe" {
456 let id = request.get("id").cloned().unwrap_or(Value::Null);
457 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
458 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
459 }
460 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
461 return match self.subscribe_session_activity(params).await {
462 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
463 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
464 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
465 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
466 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
467 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
468 };
469 }
470 if let Some(operation) = SdkOperation::from_method(method) {
471 let id = request.get("id").cloned().unwrap_or(Value::Null);
472 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
473 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
474 }
475 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
476 return match self.execute(SdkRequest { operation, params }).await {
477 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
478 Err(error) => sdk_rpc_error(id, &error),
479 };
480 }
481 if !method.starts_with("harness.v1.runtimes.") {
482 return self.handle(request);
483 }
484 let id = request.get("id").cloned().unwrap_or(Value::Null);
485 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
486 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
487 }
488 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
489 match self.runtime_call(method, params).await {
490 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
491 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
492 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
493 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
494 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
495 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
496 }
497 }
498
499 #[cfg(feature = "adapter-api")]
502 pub fn poll(&mut self) -> Vec<Value> {
503 let mut notifications = Vec::new();
504 for (subscription, follower) in &mut self.followers {
505 match follower.poll() {
506 Ok(Some(event)) => notifications.push(json!({
507 "jsonrpc": "2.0",
508 "method": SESSION_EVENT_METHOD,
509 "params": {
510 "subscription": subscription,
511 "event": event.to_json(),
512 }
513 })),
514 Ok(None) => {}
515 Err(error) => notifications.push(json!({
516 "jsonrpc": "2.0",
517 "method": SESSION_EVENT_METHOD,
518 "params": {
519 "subscription": subscription,
520 "event": {
521 "type": "watch_error",
522 "recoverable": true,
523 "message": error.to_string(),
524 },
525 }
526 })),
527 }
528 }
529 notifications
530 }
531
532 #[cfg(feature = "adapter-api")]
543 pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
544 let registry = crate::LocalRuntimeRegistry::new();
545 let authorization = crate::RuntimeAuthorization::observer();
546 let mut notifications = Vec::new();
547 for (subscription, source) in &mut self.followed_sources {
548 let state = match registry
549 .source_state(&source.harness, &source.session_id, &authorization)
550 .await
551 {
552 Ok(Some(state)) => state,
553 Ok(None) => crate::RuntimeRegistryState::Persisted,
554 Err(_) => continue,
556 };
557 if source.reported.as_deref() == Some(state.as_str()) {
558 continue;
559 }
560 source.reported = Some(state.as_str().to_string());
561 notifications.push(json!({
562 "jsonrpc": "2.0",
563 "method": SESSION_EVENT_METHOD,
564 "params": {
565 "subscription": subscription,
566 "event": {"type": "runtime_state", "state": state.as_str()},
567 },
568 }));
569 }
570 notifications
571 }
572
573 #[cfg(feature = "adapter-api")]
577 pub async fn poll_session_activities(&mut self) -> Vec<Value> {
578 let subscriptions = self
579 .activity_subscriptions
580 .iter()
581 .map(|(id, subscription)| {
582 (
583 id.clone(),
584 subscription.locators.clone(),
585 subscription.homes.clone(),
586 )
587 })
588 .collect::<Vec<_>>();
589 let mut notifications = Vec::new();
590 for (subscription_id, locators, homes) in subscriptions {
591 let Ok(activities) = self.activity_monitor.resolve(&locators, &homes).await else {
592 continue;
595 };
596 let Some(subscription) = self.activity_subscriptions.get_mut(&subscription_id) else {
597 continue;
598 };
599 let mut changed = Vec::new();
600 for activity in activities {
601 let key = activity.key();
602 if subscription
603 .reported
604 .get(&key)
605 .is_some_and(|previous| previous.same_state(&activity))
606 {
607 continue;
608 }
609 subscription.reported.insert(key, activity.clone());
610 changed.push(activity);
611 }
612 if !changed.is_empty() {
613 notifications.push(json!({
614 "jsonrpc": "2.0",
615 "method": SESSION_ACTIVITY_EVENT_METHOD,
616 "params": {
617 "subscription": subscription_id,
618 "activities": changed,
619 },
620 }));
621 }
622 }
623 notifications
624 }
625
626 #[cfg(feature = "adapter-api")]
630 pub fn poll_session_indexes(&mut self) -> Vec<Value> {
631 let mut notifications = Vec::new();
632 for (subscription, index) in &mut self.index_subscriptions {
633 let homes = index.homes().clone();
634 match index.poll() {
635 Ok(Some(delta)) => match live_index_changes(delta.changes, &homes) {
636 Ok(changes) => notifications.push(json!({
637 "jsonrpc": "2.0",
638 "method": SESSION_INDEX_EVENT_METHOD,
639 "params": {
640 "subscription": subscription,
641 "revision": delta.revision,
642 "changes": changes,
643 },
644 })),
645 Err(error) => notifications.push(json!({
646 "jsonrpc": "2.0",
647 "method": SESSION_INDEX_EVENT_METHOD,
648 "params": {
649 "subscription": subscription,
650 "error": {"recoverable": true, "message": error_message(error)},
651 },
652 })),
653 },
654 Ok(None) => {}
655 Err(error) => notifications.push(json!({
656 "jsonrpc": "2.0",
657 "method": SESSION_INDEX_EVENT_METHOD,
658 "params": {
659 "subscription": subscription,
660 "error": {"recoverable": true, "message": error},
661 },
662 })),
663 }
664 }
665 notifications
666 }
667
668 #[cfg(feature = "adapter-api")]
669 async fn subscribe_session_activity(
670 &mut self,
671 params: Value,
672 ) -> std::result::Result<Value, ServiceError> {
673 let params = decode::<ActivitySubscribeParams>(params)?;
674 if params.locators.is_empty() {
675 return Err(ServiceError::InvalidParams(
676 "sessions.activity.subscribe requires at least one locator".into(),
677 ));
678 }
679 if params.locators.len() > 2_048 {
680 return Err(ServiceError::InvalidParams(
681 "sessions.activity.subscribe accepts at most 2048 locators".into(),
682 ));
683 }
684 let initial = self
685 .activity_monitor
686 .resolve(¶ms.locators, ¶ms.homes)
687 .await
688 .map_err(ServiceError::Sdk)?;
689 let subscription = format!("activity-sub-{}", self.next_subscription);
690 self.next_subscription += 1;
691 let reported = initial
692 .iter()
693 .cloned()
694 .map(|activity| (activity.key(), activity))
695 .collect();
696 self.activity_subscriptions.insert(
697 subscription.clone(),
698 ActivitySubscription {
699 locators: params.locators,
700 homes: params.homes,
701 reported,
702 },
703 );
704 Ok(json!({"subscription": subscription, "initial": initial}))
705 }
706
707 #[cfg(feature = "adapter-api")]
709 pub async fn poll_runtimes(&mut self) -> Vec<Value> {
710 self.poll_sdk_events()
711 .await
712 .into_iter()
713 .map(|(connection, runtime_event)| {
714 json!({
715 "jsonrpc": "2.0",
716 "method": RUNTIME_EVENT_METHOD,
717 "params": {
718 "connection": connection,
719 "session_id": runtime_event.session_id,
720 "sequence": runtime_event.event.sequence,
721 "event": {
722 "kind": runtime_event.event.kind,
723 "payload": runtime_event.event.payload,
724 },
725 },
726 })
727 })
728 .collect()
729 }
730
731 async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
732 let mut events = Vec::new();
733 let mut closed = Vec::new();
734 let now_ms = crate::approvals::now_ms();
735 for (connection, runtime) in &mut self.runtimes {
736 let session_id = runtime.handle().runtime_id.clone();
737 let harness = runtime.handle().harness.clone();
738 for _ in 0..256 {
743 match tokio::time::timeout(Duration::ZERO, runtime.next_event()).await {
744 Ok(Ok(Some(event))) => {
745 let terminal = event.kind == "transport_closed";
746 self.approvals
750 .observe(connection, &harness, &session_id, &event, now_ms);
751 let next_sequence = self
752 .runtime_sequences
753 .entry(session_id.clone())
754 .or_insert(0);
755 let sequence = event.sequence.unwrap_or_else(|| {
756 *next_sequence = next_sequence.saturating_add(1);
757 *next_sequence
758 });
759 *next_sequence = (*next_sequence).max(sequence);
760 events.push((
761 connection.clone(),
762 SdkRuntimeEvent {
763 session_id: session_id.clone(),
764 event: SdkEvent {
765 sequence,
766 kind: event.kind,
767 payload: event.payload,
768 },
769 },
770 ));
771 if terminal {
772 closed.push(connection.clone());
773 break;
774 }
775 }
776 Ok(Ok(None)) => {
777 let sequence = self
778 .runtime_sequences
779 .entry(session_id.clone())
780 .or_insert(0);
781 *sequence = sequence.saturating_add(1);
782 events.push((
783 connection.clone(),
784 SdkRuntimeEvent {
785 session_id,
786 event: SdkEvent {
787 sequence: *sequence,
788 kind: "transport_closed".into(),
789 payload: json!({"message": "Harness runtime transport closed."}),
790 },
791 },
792 ));
793 closed.push(connection.clone());
794 break;
795 }
796 Err(_) => break,
797 Ok(Err(error)) => {
798 let sequence = self
799 .runtime_sequences
800 .entry(session_id.clone())
801 .or_insert(0);
802 *sequence = sequence.saturating_add(1);
803 events.push((
804 connection.clone(),
805 SdkRuntimeEvent {
806 session_id,
807 event: SdkEvent {
808 sequence: *sequence,
809 kind: "transport_error".into(),
810 payload: json!({"message": error.to_string(), "terminal": true}),
811 },
812 },
813 ));
814 closed.push(connection.clone());
815 break;
816 }
817 }
818 }
819 }
820 for connection in closed {
821 if let Some(runtime) = self.runtimes.remove(&connection) {
822 self.runtime_sequences.remove(&runtime.handle().runtime_id);
823 }
824 self.terminal_launches.remove(&connection);
825 self.approvals.forget(&connection);
828 }
829 events
830 }
831
832 fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
833 match method {
834 "harness.v1.capabilities" => Ok(json!({
835 "version": HARNESS_SERVICE_VERSION,
836 "sdk": self.capabilities(),
837 "methods": HARNESS_SERVICE_METHODS,
838 "notifications": [
839 SESSION_EVENT_METHOD,
840 SESSION_ACTIVITY_EVENT_METHOD,
841 SESSION_INDEX_EVENT_METHOD,
842 RUNTIME_EVENT_METHOD
843 ],
844 "harnesses": harness_support_registry()
845 .harnesses
846 .into_iter()
847 .map(|harness| harness.id)
848 .collect::<Vec<_>>(),
849 })),
850 "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
851 .map_err(|error| ServiceError::Operation(error.to_string())),
852 "harness.v1.profiles.list" | "harness.v1.profiles.get" => profiles_call(method, params),
853 "harness.v1.profiles.create" => {
859 mutate_profile(crate::profiles_control::ProfileVerb::Create, params)
860 }
861 "harness.v1.profiles.delete" => {
862 mutate_profile(crate::profiles_control::ProfileVerb::Delete, params)
863 }
864 "harness.v1.channels.list" | "harness.v1.channels.status" => {
865 channels_call(method, params)
866 }
867 "harness.v1.routes.list" => routes_call(params),
870 "harness.v1.triggers.list" => triggers_call(params),
872 "harness.v1.workflow.load" => {
882 let params = decode::<WorkflowLoadParams>(params)?;
883 let read =
884 crate::workflow_doors::load(params.from, ¶ms.home).map_err(operation)?;
885 serde_json::to_value(read)
886 .map_err(|error| ServiceError::Operation(error.to_string()))
887 }
888 "harness.v1.orchestration.load" => {
889 let params = decode::<OrchestrationLoadParams>(params)?;
890 let read = crate::orchestration_doors::load(¶ms.root, params.flavor)
891 .map_err(operation)?;
892 serde_json::to_value(read)
893 .map_err(|error| ServiceError::Operation(error.to_string()))
894 }
895 "harness.v1.orchestration.save" => {
896 let params = decode::<OrchestrationSaveParams>(params)?;
897 let saved = crate::orchestration_doors::save(
898 ¶ms.root,
899 params.orchestration,
900 params.vault,
901 )
902 .map_err(operation)?;
903 serde_json::to_value(saved)
904 .map_err(|error| ServiceError::Operation(error.to_string()))
905 }
906 "harness.v1.orchestration.compile" => {
907 let params = decode::<OrchestrationCompileParams>(params)?;
908 let read = crate::orchestration_doors::compile(params.from, ¶ms.home)
909 .map_err(operation)?;
910 serde_json::to_value(read)
911 .map_err(|error| ServiceError::Operation(error.to_string()))
912 }
913 "harness.v1.orchestration.decompile" => {
914 let params = decode::<OrchestrationDecompileParams>(params)?;
915 let report = crate::orchestration_doors::decompile(
916 params.to,
917 params.orchestration,
918 ¶ms.source,
919 params.source_flavor,
920 ¶ms.dest,
921 params.vault,
922 )
923 .map_err(operation)?;
924 serde_json::to_value(report)
925 .map_err(|error| ServiceError::Operation(error.to_string()))
926 }
927 "harness.v1.orchestration.import" => {
932 let params = decode::<OrchestrationImportParams>(params)?;
933 let imported =
934 crate::orchestration_doors::import(params.from, ¶ms.home, ¶ms.into)
935 .map_err(operation)?;
936 serde_json::to_value(imported)
937 .map_err(|error| ServiceError::Operation(error.to_string()))
938 }
939 "harness.v1.orchestration.export" => {
940 let params = decode::<OrchestrationExportParams>(params)?;
941 let report =
942 crate::orchestration_doors::export(params.to, ¶ms.root, ¶ms.dest)
943 .map_err(operation)?;
944 serde_json::to_value(report)
945 .map_err(|error| ServiceError::Operation(error.to_string()))
946 }
947 "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
953 "harness.v1.skills.list" => {
958 let query = decode::<crate::skills::SkillsQuery>(params)?;
959 if let Some(harness) = query.harness.as_deref() {
960 if !crate::skills::SKILL_HARNESSES.contains(&harness) {
961 return Err(ServiceError::UnsupportedAction(format!(
962 "`{harness}` has no skills root Volter Harness reads"
963 )));
964 }
965 }
966 serde_json::to_value(crate::skills::list_skills(&query))
967 .map_err(|error| ServiceError::Operation(error.to_string()))
968 }
969 "harness.v1.skills.install" => {
977 mutate_skill(crate::skills_control::SkillVerb::Install, params)
978 }
979 "harness.v1.skills.remove" => {
980 mutate_skill(crate::skills_control::SkillVerb::Remove, params)
981 }
982 "harness.v1.approvals.list" => {
990 let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
991 if let Some(harness) = query.harness.as_deref() {
992 if !crate::approvals::lists_approvals(harness) {
993 return Err(ServiceError::UnsupportedAction(format!(
994 "`{harness}` has no runtime door that carries an approval request"
995 )));
996 }
997 }
998 serde_json::to_value(self.approvals(&query))
999 .map_err(|error| ServiceError::Operation(error.to_string()))
1000 }
1001 "harness.v1.sessions.discover" => {
1002 let query = decode::<DiscoveryQuery>(params)?;
1003 let mut page = discover_session_page(&query).map_err(operation)?;
1004 let doors = crate::mail_route::LiveSessions::read(&query.homes);
1009 if let Some(text) = query
1013 .query
1014 .as_deref()
1015 .map(str::trim)
1016 .filter(|text| !text.is_empty())
1017 {
1018 let needle = text.to_lowercase();
1019 for live in doors.all() {
1020 let name = live
1021 .name
1022 .split('@')
1023 .next()
1024 .unwrap_or(&live.name)
1025 .to_lowercase();
1026 if !name.contains(&needle)
1027 || page.sessions.iter().any(|session| {
1028 session.locator.session_id == live.address.session_id
1029 })
1030 {
1031 continue;
1032 }
1033 if query
1034 .limit
1035 .is_some_and(|limit| page.sessions.len() >= limit)
1036 {
1037 break;
1038 }
1039 let mut by_id = query.clone();
1041 by_id.query = Some(live.address.session_id.clone());
1042 by_id.cursor = None;
1043 if let Some(cwd) = &live.cwd {
1044 by_id.workspace = Some(cwd.clone());
1045 by_id.workspace_subtree = false;
1046 by_id.workspace_family = None;
1047 }
1048 if let Ok(found) = discover_session_page(&by_id) {
1049 page.sessions
1050 .extend(found.sessions.into_iter().filter(|session| {
1051 session.locator.session_id == live.address.session_id
1052 }));
1053 }
1054 }
1055 }
1056 let activities = crate::session_activity::resolve_stock_session_activities(
1057 &page
1058 .sessions
1059 .iter()
1060 .map(|session| session.locator.clone())
1061 .collect::<Vec<_>>(),
1062 &query.homes,
1063 )
1064 .into_iter()
1065 .map(|activity| (activity.key(), activity))
1066 .collect::<BTreeMap<_, _>>();
1067 let sessions = page
1068 .sessions
1069 .into_iter()
1070 .map(|session| {
1071 let mut value = live_descriptor_value(&session, &doors)?;
1072 let activity_key = (
1073 session.locator.harness.as_str().to_string(),
1074 session.locator.session_id.clone(),
1075 );
1076 if let Some(activity) = activities.get(&activity_key) {
1077 let mut activity = activity.clone();
1080 let waiting = doors.all().iter().any(|live| {
1081 live.address.harness == activity_key.0
1082 && live.address.session_id == activity_key.1
1083 && live.status == "waiting"
1084 });
1085 if waiting && activity.turn != crate::SessionTurnState::NeedsInput {
1086 activity.turn = crate::SessionTurnState::NeedsInput;
1087 activity.evidence.source = "pane_prompt".into();
1088 }
1089 let activity = &activity;
1090 value["activity"] = serde_json::to_value(activity)
1091 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1092 if let Some(status) = legacy_live_status(activity) {
1093 value["live_status"] = json!(status);
1094 }
1095 }
1096 Ok(value)
1097 })
1098 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1099 let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1100 if query.search_previews {
1103 result["receipt"] = serde_json::to_value(page.receipt)
1104 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1105 }
1106 Ok(result)
1107 }
1108 "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1109 "harness.v1.sessions.load" => {
1110 let params = decode::<LoadSessionParams>(params)?;
1111 if let Some(options) = ¶ms.options {
1112 options.validate()?;
1113 if let Some(result) = indexed_claude_window(¶ms.read.locator, options)? {
1114 return Ok(result);
1115 }
1116 return load_session(¶ms.read.locator)
1117 .map(|session| projected_session_result(&session, options))
1118 .map_err(operation);
1119 }
1120 let mut session = if params.read.display_history() {
1121 self.catalog
1122 .load_display_view(
1123 ¶ms.read.locator,
1124 params.read.read_fidelity(),
1125 params.read.tail_messages().unwrap_or(500),
1126 )
1127 .map_err(crate::Error::from)
1128 } else if params.read.include_subagents() {
1129 load_session_with_fidelity(¶ms.read.locator, params.read.read_fidelity())
1130 } else {
1131 self.catalog
1132 .load_parent_with_fidelity(
1133 ¶ms.read.locator,
1134 params.read.read_fidelity(),
1135 )
1136 .map_err(crate::Error::from)
1137 }
1138 .map_err(operation)?;
1139 params.read.bound_session(&mut session);
1140 Ok(json!({"session": normalized_session_json(&session)}))
1141 }
1142 "harness.v1.sessions.follow" => {
1143 let params = decode::<LocatorParams>(params)?;
1144 let mut follower = self
1145 .catalog
1146 .follow_read_view(
1147 ¶ms.locator,
1148 params.read_fidelity(),
1149 params.include_subagents(),
1150 params.tail_messages(),
1151 params.max_message_chars(),
1152 params.display_history(),
1153 )
1154 .map_err(operation)?;
1155 let initial = follower
1156 .poll()
1157 .map_err(operation)?
1158 .map(|event| event.to_json());
1159 let subscription = format!("sub-{}", self.next_subscription);
1160 self.next_subscription += 1;
1161 self.followers.insert(subscription.clone(), follower);
1162 self.followed_sources.insert(
1163 subscription.clone(),
1164 FollowedSource {
1165 harness: params.locator.harness.as_str().to_string(),
1166 session_id: params.locator.session_id.clone(),
1167 reported: None,
1168 },
1169 );
1170 Ok(json!({"subscription": subscription, "initial": initial}))
1171 }
1172 "harness.v1.sessions.unfollow" => {
1173 let params = decode::<UnfollowParams>(params)?;
1174 self.followed_sources.remove(¶ms.subscription);
1175 Ok(json!({
1176 "removed": self.followers.remove(¶ms.subscription).is_some()
1177 }))
1178 }
1179 "harness.v1.sessions.activity.unsubscribe" => {
1180 let params = decode::<UnfollowParams>(params)?;
1181 Ok(json!({
1182 "removed": self.activity_subscriptions.remove(¶ms.subscription).is_some()
1183 }))
1184 }
1185 "harness.v1.sessions.index.subscribe" => {
1186 let query = decode::<DiscoveryQuery>(params)?;
1187 crate::session_index::validate_query(&query)
1188 .map_err(ServiceError::InvalidParams)?;
1189 let homes = query.homes.clone();
1190 let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1191 query,
1192 Arc::clone(&self.index_notifier),
1193 )
1194 .map_err(ServiceError::Operation)?;
1195 let doors = crate::mail_route::LiveSessions::read(&homes);
1196 let initial = initial
1197 .iter()
1198 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1199 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1200 let subscription = format!("index-sub-{}", self.next_subscription);
1201 self.next_subscription += 1;
1202 self.index_subscriptions.insert(subscription.clone(), index);
1203 Ok(json!({
1204 "subscription": subscription,
1205 "revision": 1,
1206 "initial": initial,
1207 }))
1208 }
1209 "harness.v1.sessions.index.resize" => {
1210 let params = decode::<IndexResizeParams>(params)?;
1211 crate::session_index::validate_limit(params.limit)
1212 .map_err(ServiceError::InvalidParams)?;
1213 let index = self
1214 .index_subscriptions
1215 .get_mut(¶ms.subscription)
1216 .ok_or_else(|| {
1217 ServiceError::InvalidParams("unknown session index subscription".into())
1218 })?;
1219 let prepared = index
1220 .prepare_resize(params.limit)
1221 .map_err(ServiceError::Operation)?;
1222 let doors = crate::mail_route::LiveSessions::read(index.homes());
1223 let initial = prepared
1224 .page
1225 .sessions
1226 .iter()
1227 .map(|descriptor| live_descriptor_value(descriptor, &doors))
1228 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1229 let response = json!({
1230 "subscription": params.subscription,
1231 "revision": prepared.revision,
1232 "initial": initial,
1233 "receipt": prepared.page.receipt,
1234 });
1235 index.commit_resize(prepared);
1236 Ok(response)
1237 }
1238 "harness.v1.sessions.index.unsubscribe" => {
1239 let params = decode::<UnfollowParams>(params)?;
1240 Ok(json!({
1241 "removed": self.index_subscriptions.remove(¶ms.subscription).is_some()
1242 }))
1243 }
1244 "harness.v1.sessions.import" => {
1245 let params = decode::<ImportSessionParams>(params)?;
1246 let session = Session::load_str(¶ms.content, params.source_harness.into())
1247 .map_err(operation)?;
1248 Ok(json!({"session": normalized_session_json(&session)}))
1249 }
1250 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1251 let params = decode::<ExportSessionParams>(params)?;
1252 let session = load_session(¶ms.locator).map_err(operation)?;
1253 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
1254 if method == "harness.v1.sessions.export"
1255 && params.target_harness == TransferFormat::Hermes
1256 {
1257 let imported = crate::hermes_import::import_into_hermes(&session, None)
1259 .map_err(operation)?;
1260 return Ok(json!({"artifact": artifact, "imported": imported}));
1261 }
1262 Ok(json!({"artifact": artifact}))
1263 }
1264 "harness.v1.sessions.reduce" => {
1265 let params = decode::<ReduceSessionParams>(params)?;
1266 self.reduce_session(params)
1267 }
1268 "harness.v1.sessions.branch" => {
1269 let params = decode::<BranchSessionParams>(params)?;
1270 let session = load_session(¶ms.locator).map_err(operation)?;
1271 let storage = params.locator.storage.path().display().to_string();
1272 let bootstrap_prompt = format!(
1273 "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.",
1274 params.locator.harness.as_str(), params.locator.session_id, storage
1275 );
1276 let artifact = params
1277 .target_harness
1278 .map(|target| session_artifact(¶ms.locator, &session, target))
1279 .transpose()?;
1280 Ok(json!({
1281 "parent": params.locator,
1282 "session": normalized_session_json(&session),
1283 "bootstrap_prompt": bootstrap_prompt,
1284 "artifact": artifact,
1285 }))
1286 }
1287 "harness.v1.sessions.handoff" => {
1288 let params = decode::<HandoffSessionParams>(params)?;
1289 let session = load_session(¶ms.locator).map_err(operation)?;
1290 let cwd = params
1291 .cwd
1292 .or_else(|| session.meta.cwd.clone())
1293 .unwrap_or_else(|| PathBuf::from("."));
1294 let artifact = handoff_artifact(¶ms.locator, &session, params.target_harness)?;
1295 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1296 ServiceError::Operation(
1297 "handoff artifact omitted target session identity".into(),
1298 )
1299 })?;
1300 let instructions =
1301 handoff_instructions(params.target_harness, target_session_id, &cwd);
1302 Ok(json!({
1303 "artifact": artifact,
1304 "launch": instructions.launch,
1305 "materialize": instructions.materialize,
1306 "requires_materialization": instructions.requires_materialization,
1307 "note": instructions.note,
1308 }))
1309 }
1310 "harness.v1.sessions.materialize" => {
1311 let params = decode::<MaterializeSessionParams>(params)?;
1312 for file in ¶ms.artifact.files {
1316 if file.role == "source_recovery"
1317 && file.path == "recovery/source.supercode.jsonl"
1318 {
1319 if let Ok(source) = Session::from_native_str(&file.content) {
1320 crate::residue_store::store_segments(&source);
1321 }
1322 }
1323 }
1324 let locator = crate::native_materialize::materialize_native_artifact(
1325 params.artifact,
1326 ¶ms.cwd,
1327 ¶ms.homes,
1328 )
1329 .map_err(ServiceError::Operation)?;
1330 Ok(json!({"locator": locator}))
1331 }
1332 "harness.v1.jobs.list" => {
1336 let query = decode::<crate::jobs::JobsQuery>(params)?;
1337 if let Some(harness) = query.harness.as_deref() {
1338 refuse_harness_without_jobs(harness, "jobs.list")?;
1339 }
1340 let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1341 serde_json::to_value(listing)
1342 .map_err(|error| ServiceError::Operation(error.to_string()))
1343 }
1344 "harness.v1.jobs.get" => {
1345 let params = decode::<JobsGetParams>(params)?;
1346 refuse_harness_without_jobs(¶ms.harness, "jobs.get")?;
1347 match crate::jobs::get_job(¶ms.harness, ¶ms.id, ¶ms.homes)
1348 .map_err(operation)?
1349 {
1350 Some((job, source)) => Ok(json!({"job": job, "source": source})),
1351 None => Err(ServiceError::Operation(format!(
1352 "`{}` has no scheduled job `{}`",
1353 params.harness, params.id
1354 ))),
1355 }
1356 }
1357 "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1363 "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1364 "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1365 "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1366 "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1367 "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1368 "harness.v1.jobs.notepad"
1369 | "harness.v1.jobs.notepad_set"
1370 | "harness.v1.jobs.notepad_delete" => {
1371 let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1372 refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1373 let answer = match method {
1374 "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1375 "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1376 _ => crate::jobs_notepad::read(&request),
1377 }
1378 .map_err(job_control_error)?;
1379 serde_json::to_value(answer)
1380 .map_err(|error| ServiceError::Operation(error.to_string()))
1381 }
1382 "harness.v1.runs.list" => {
1386 let query = decode::<crate::runs::RunsQuery>(params)?;
1387 if let Some(harness) = query.harness.as_deref() {
1388 refuse_harness_without_runs(harness, "runs.list")?;
1389 }
1390 let listing = crate::runs::list_runs(&query).map_err(operation)?;
1391 serde_json::to_value(listing)
1392 .map_err(|error| ServiceError::Operation(error.to_string()))
1393 }
1394 "harness.v1.runs.get" => {
1395 let params = decode::<RunsGetParams>(params)?;
1396 refuse_harness_without_runs(¶ms.harness, "runs.get")?;
1397 match crate::runs::get_run(¶ms.harness, ¶ms.id, ¶ms.homes)
1398 .map_err(operation)?
1399 {
1400 Some((run, source)) => Ok(json!({"run": run, "source": source})),
1401 None => Err(ServiceError::Operation(format!(
1402 "`{}` has no run `{}`",
1403 params.harness, params.id
1404 ))),
1405 }
1406 }
1407 "harness.v1.sessions.resume_instructions" => {
1408 let params = decode::<ResumeInstructionsParams>(params)?;
1409 let session = load_session(¶ms.locator).map_err(operation)?;
1410 let cwd = params
1411 .cwd
1412 .or(session.meta.cwd)
1413 .unwrap_or_else(|| PathBuf::from("."));
1414 let launch = resume_launch(
1415 params.locator.harness.as_str(),
1416 ¶ms.locator.session_id,
1417 &cwd,
1418 params.policy,
1419 )?;
1420 Ok(json!({"launch": launch}))
1421 }
1422 _ => Err(ServiceError::MethodNotFound),
1423 }
1424 }
1425
1426 fn reduce_session(
1427 &self,
1428 params: ReduceSessionParams,
1429 ) -> std::result::Result<Value, ServiceError> {
1430 let session = load_session(¶ms.locator).map_err(operation)?;
1431 if session.messages.is_empty() {
1432 return Err(ServiceError::InvalidParams(
1433 "cannot reduce an empty session".into(),
1434 ));
1435 }
1436 let keep_last = params.keep_last.clamp(1, 128);
1437 let policy = reduce::ReductionPolicy {
1438 clear_turns_older_than: Some(keep_last),
1439 ..Default::default()
1440 };
1441 let (view, log) =
1442 reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1443 if log.reductions.is_empty() {
1444 return Err(ServiceError::UnsupportedAction(format!(
1445 "session `{}` is already too small for a meaningful reversible reduction",
1446 params.locator.session_id
1447 )));
1448 }
1449 let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1450 let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1451 if reduced_tokens >= source_tokens {
1452 return Err(ServiceError::UnsupportedAction(format!(
1453 "session `{}` has no token-reducing reversible projection",
1454 params.locator.session_id
1455 )));
1456 }
1457
1458 let store_root = self
1459 .reduction_store_root
1460 .clone()
1461 .unwrap_or_else(default_reduction_store_root);
1462 let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1463 let rescue_id = format!("rescue-{}", generated_session_id());
1464 let imported = session
1465 .imported_message_count
1466 .unwrap_or(session.messages.len())
1467 .min(session.messages.len());
1468 let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1469 let view_jsonl = messages_jsonl(&view)?;
1470 let title = format!(
1471 "Reduced {} continuation from {}",
1472 params.target_harness.id(),
1473 params.locator.session_id
1474 );
1475
1476 store
1481 .save_sidecar(&rescue_id, &sidecar_jsonl)
1482 .map_err(operation)?;
1483 store
1484 .save_reduction_log(&rescue_id, &log)
1485 .map_err(operation)?;
1486 store
1487 .save(&rescue_id, &title, &view_jsonl)
1488 .map_err(operation)?;
1489
1490 let source_bytes = serde_json::to_vec(&session.messages)
1491 .map_err(|error| ServiceError::Operation(error.to_string()))?
1492 .len() as u64;
1493 let reduced_bytes = serde_json::to_vec(&view)
1494 .map_err(|error| ServiceError::Operation(error.to_string()))?
1495 .len() as u64;
1496 store
1497 .set_reduction_stats(
1498 &rescue_id,
1499 &title,
1500 source_bytes,
1501 reduced_bytes,
1502 log.reductions.len() as u32,
1503 )
1504 .map_err(operation)?;
1505
1506 let reloaded_sidecar = store
1510 .load_sidecar(&rescue_id)
1511 .map_err(operation)?
1512 .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1513 let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1514 let reloaded_log = store
1515 .load_reduction_log(&rescue_id)
1516 .map_err(operation)?
1517 .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1518 let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1519 reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1520 let (restamped_view, restamped_log) =
1527 reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1528 if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1529 return Err(ServiceError::Operation(
1530 "persisted reduction view does not match its durable log and sidecar".into(),
1531 ));
1532 }
1533 if restamped_log != reloaded_log {
1534 return Err(ServiceError::Operation(
1535 "reapplying the durable reduction log changed its identity".into(),
1536 ));
1537 }
1538 let inverted =
1539 reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1540 if inverted != session.messages {
1541 return Err(ServiceError::Operation(
1542 "reduction inversion did not restore the source messages byte-exactly".into(),
1543 ));
1544 }
1545
1546 let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1547 let sidecar_path = store.sidecar_path(&rescue_id);
1548 let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1549 let bootstrap_prompt = reduced_bootstrap_prompt(
1550 ¶ms.locator,
1551 params.target_harness,
1552 &view_jsonl,
1553 &sidecar_path,
1554 &reduction_log_path,
1555 );
1556 let mut reduced_session = session.clone();
1557 reduced_session.meta.session_id = Some(rescue_id.clone());
1558 reduced_session.messages = view;
1559
1560 Ok(json!({
1561 "session": normalized_session_json(&reduced_session),
1562 "bootstrap_prompt": bootstrap_prompt,
1563 "receipt": {
1564 "id": rescue_id,
1565 "sidecar_id": rescue_id,
1566 "source_harness": params.locator.harness,
1567 "target_harness": params.target_harness.id(),
1568 "source_tokens": source_tokens,
1569 "reduced_tokens": reduced_tokens,
1570 "ratio": ratio,
1571 "source_bytes": source_bytes,
1572 "reduced_bytes": reduced_bytes,
1573 "reductions": reloaded_log.reductions.len(),
1574 "sidecar_path": sidecar_path,
1575 "reduction_log_path": reduction_log_path,
1576 "verified": true,
1577 "reversible": true,
1578 }
1579 }))
1580 }
1581
1582 pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1600 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1601 return None;
1602 }
1603 let method = request.get("method").and_then(Value::as_str)?;
1604 if !RUNTIME_OPEN_METHODS.contains(&method) {
1605 return None;
1606 }
1607 Some(RuntimeOpen {
1608 id: request.get("id").cloned().unwrap_or(Value::Null),
1609 method: method.to_string(),
1610 params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1611 })
1612 }
1613
1614 pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1636 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1637 return None;
1638 }
1639 let method = request.get("method").and_then(Value::as_str)?;
1640 if !DETACHED_METHODS.contains(&method) {
1641 return None;
1642 }
1643 let id = request.get("id").cloned().unwrap_or(Value::Null);
1644 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1645 let work = match method {
1646 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1647 .inventory_work(method, params)
1648 .map(DetachedWork::Inventory),
1649 "harness.v1.sessions.message" => {
1650 decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1651 }
1652 _ => {
1653 let verb = match method {
1654 "harness.v1.sessions.new" => crate::SessionVerb::New,
1655 "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1656 "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1657 _ => crate::SessionVerb::Delete,
1658 };
1659 match decode::<crate::SessionMutation>(params) {
1660 Ok(mutation) => {
1661 match crate::sessions_control::door(&mutation.harness, verb) {
1662 Ok(crate::SessionDoor::Live(_)) => return None,
1665 Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1666 Err(error) => Err(session_control_error(error)),
1667 }
1668 }
1669 Err(error) => Err(error),
1670 }
1671 }
1672 };
1673 Some(DetachedCall {
1674 id,
1675 method: method.to_string(),
1676 work: work.map(Work::Free),
1677 })
1678 }
1679
1680 pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1694 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1695 return None;
1696 }
1697 let method = request.get("method").and_then(Value::as_str)?;
1698 let id = request.get("id").cloned().unwrap_or(Value::Null);
1699 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1700 let work = match method {
1701 "harness.v1.runtimes.close" => decode::<RuntimeConnectionParams>(params)
1702 .and_then(|params| self.surrender_runtime(¶ms.connection))
1703 .map(|(runtime, process_group)| {
1704 Work::Runtime(RuntimeWork::Close {
1705 runtime,
1706 process_group,
1707 })
1708 }),
1709 "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1710 let verb = if method == "harness.v1.sessions.new" {
1711 crate::SessionVerb::New
1712 } else {
1713 crate::SessionVerb::Reset
1714 };
1715 let mutation = decode::<crate::SessionMutation>(params).ok()?;
1716 let Ok(crate::SessionDoor::Live(command)) =
1720 crate::sessions_control::door(&mutation.harness, verb)
1721 else {
1722 return None;
1723 };
1724 let connection = mutation
1725 .connection
1726 .clone()
1727 .filter(|value| !value.trim().is_empty())?;
1728 self.lend_runtime(&connection).map(|runtime| {
1729 let session = live_session_name(runtime.as_ref(), &mutation);
1730 Work::Runtime(RuntimeWork::LiveCommand {
1731 connection,
1732 runtime,
1733 verb,
1734 mutation,
1735 command,
1736 session,
1737 })
1738 })
1739 }
1740 _ => return None,
1741 };
1742 Some(DetachedCall {
1743 id,
1744 method: method.to_string(),
1745 work,
1746 })
1747 }
1748
1749 pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1754 let DetachedAnswer { response, returned } = answer;
1755 if let Some(ReturnedRuntime {
1756 connection,
1757 runtime,
1758 }) = returned
1759 {
1760 self.runtimes_in_flight.remove(&connection);
1761 self.runtimes.insert(connection, runtime);
1762 }
1763 response
1764 }
1765
1766 pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1770 let OpenedRuntime { id, outcome } = opened;
1771 let result = match outcome {
1772 Ok(open) => self.register_open_runtime(open).await,
1773 Err(error) => Err(error),
1774 };
1775 service_response(id, result)
1776 }
1777
1778 async fn register_open_runtime(
1780 &mut self,
1781 open: OpenRuntime,
1782 ) -> std::result::Result<Value, ServiceError> {
1783 match open {
1784 OpenRuntime::Hosted {
1785 runtime,
1786 capabilities,
1787 workspace,
1788 fresh,
1789 } => {
1790 self.insert_hosted_runtime(runtime, capabilities, workspace, fresh)
1791 .await
1792 }
1793 OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1794 }
1795 }
1796
1797 async fn runtime_call(
1798 &mut self,
1799 method: &str,
1800 params: Value,
1801 ) -> std::result::Result<Value, ServiceError> {
1802 match method {
1803 "harness.v1.runtimes.capabilities" => {
1804 let params = decode::<RuntimeBackendParams>(params)?;
1805 #[cfg(feature = "adapter-api")]
1809 if params
1810 .base_url
1811 .as_deref()
1812 .is_some_and(|value| LiveRuntimeEndpoint::parse(value).is_ok())
1813 {
1814 return Ok(json!({
1815 "harness": ¶ms.harness,
1816 "capabilities": SupercodeHttpRuntimeBackend::live_capabilities(),
1817 }));
1818 }
1819 let backend = runtime_backend(¶ms)?;
1820 Ok(json!({
1821 "harness": backend.harness(),
1822 "capabilities": backend.capabilities(),
1823 }))
1824 }
1825 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1826 self.register_open_runtime(open_runtime(method, params).await?)
1827 .await
1828 }
1829 "harness.v1.runtimes.send_input" => {
1830 let params = decode::<RuntimeInputParams>(params)?;
1831 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1832 let runtime = self.runtime_mut(¶ms.connection)?;
1833 let turn_id = within_control_deadline(
1834 method,
1835 runtime.send_input(RuntimeInput {
1836 text: params.text,
1837 image_urls,
1838 }),
1839 )
1840 .await?
1841 .map_err(operation)?;
1842 Ok(json!({"turn_id": turn_id}))
1843 }
1844 "harness.v1.runtimes.interrupt" => {
1845 let params = decode::<RuntimeConnectionParams>(params)?;
1846 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1847 .await?
1848 .map_err(operation)?;
1849 Ok(json!({}))
1850 }
1851 "harness.v1.runtimes.steer" => {
1852 let params = decode::<RuntimeInputParams>(params)?;
1853 if !params.image_urls.is_empty() {
1854 return Err(ServiceError::InvalidParams(
1855 "runtime steering accepts text only".into(),
1856 ));
1857 }
1858 let text = params.text.trim();
1859 if text.is_empty() || text.chars().count() > 50_000 {
1860 return Err(ServiceError::InvalidParams(
1861 "runtime steering requires 1 to 50,000 text characters".into(),
1862 ));
1863 }
1864 within_control_deadline(
1865 method,
1866 self.runtime_mut(¶ms.connection)?
1867 .steer(text.to_string()),
1868 )
1869 .await?
1870 .map_err(operation)?;
1871 Ok(json!({}))
1872 }
1873 "harness.v1.runtimes.respond" => {
1874 let params = decode::<RuntimeRespondParams>(params)?;
1875 let request_id = params.request_id.clone();
1876 within_control_deadline(
1877 method,
1878 self.runtime_mut(¶ms.connection)?
1879 .respond(params.request_id, params.response),
1880 )
1881 .await?
1882 .map_err(operation)?;
1883 self.approvals.answered(¶ms.connection, &request_id);
1885 Ok(json!({}))
1886 }
1887 "harness.v1.runtimes.acquire_control" => {
1888 let params = decode::<RuntimeConnectionParams>(params)?;
1889 let snapshot = within_control_deadline(
1890 method,
1891 self.runtime_mut(¶ms.connection)?.acquire_control(),
1892 )
1893 .await?
1894 .map_err(operation)?;
1895 serde_json::to_value(snapshot)
1896 .map_err(|error| ServiceError::Operation(error.to_string()))
1897 }
1898 "harness.v1.runtimes.heartbeat" => {
1899 let params = decode::<RuntimeConnectionParams>(params)?;
1900 let snapshot = within_control_deadline(
1901 method,
1902 self.runtime_mut(¶ms.connection)?.heartbeat(),
1903 )
1904 .await?
1905 .map_err(operation)?;
1906 serde_json::to_value(snapshot)
1907 .map_err(|error| ServiceError::Operation(error.to_string()))
1908 }
1909 "harness.v1.runtimes.detach" => {
1910 let params = decode::<RuntimeConnectionParams>(params)?;
1911 let snapshot =
1912 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1913 .await?
1914 .map_err(operation)?;
1915 serde_json::to_value(snapshot)
1916 .map_err(|error| ServiceError::Operation(error.to_string()))
1917 }
1918 "harness.v1.runtimes.terminal_instructions" => {
1919 let params = decode::<RuntimeConnectionParams>(params)?;
1920 let launch = self
1921 .terminal_launches
1922 .get(¶ms.connection)
1923 .ok_or_else(|| {
1924 ServiceError::Operation(
1925 "this runtime is not hosted for terminal attachment".into(),
1926 )
1927 })?;
1928 Ok(json!({"launch":launch}))
1929 }
1930 "harness.v1.runtimes.close" => {
1931 let params = decode::<RuntimeConnectionParams>(params)?;
1932 let (runtime, process_group) = self.surrender_runtime(¶ms.connection)?;
1933 close_runtime(runtime, process_group).await
1934 }
1935 _ => Err(ServiceError::MethodNotFound),
1936 }
1937 }
1938
1939 #[cfg(feature = "adapter-api")]
1941 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1942 let params = decode::<MessageSessionParams>(params)?;
1943 Ok(message_live_session(¶ms).await)
1944 }
1945
1946 #[cfg(feature = "adapter-api")]
1947 fn harness_settings_call(
1948 &self,
1949 method: &str,
1950 params: Value,
1951 ) -> std::result::Result<Value, ServiceError> {
1952 let homes = crate::HarnessHomes::default();
1953 match method {
1954 "harness.v1.harnesses.settings" => {
1955 let params = decode::<HarnessSettingsParams>(params)?;
1956 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1957 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1958 serde_json::to_value(report)
1959 .map_err(|error| ServiceError::Operation(error.to_string()))
1960 }
1961 "harness.v1.harnesses.configure" => {
1962 let params = decode::<ConfigureHarnessParams>(params)?;
1963 let report = crate::configure_harness_interop_settings(
1964 &homes,
1965 ¶ms.harness,
1966 ¶ms.changes,
1967 params.expected_revision.as_deref(),
1968 )
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 _ => Err(ServiceError::MethodNotFound),
1974 }
1975 }
1976
1977 fn insert_runtime(
1978 &mut self,
1979 runtime: Box<dyn RuntimeConnection>,
1980 ) -> std::result::Result<Value, ServiceError> {
1981 let connection = format!("runtime-{}", self.next_runtime);
1982 self.next_runtime += 1;
1983 let handle = runtime.handle().clone();
1984 self.runtime_sequences
1985 .entry(handle.runtime_id.clone())
1986 .or_insert(0);
1987 self.runtimes.insert(connection.clone(), runtime);
1988 Ok(json!({"connection": connection, "handle": handle}))
1989 }
1990
1991 #[cfg(feature = "adapter-api")]
1992 async fn insert_hosted_runtime(
1993 &mut self,
1994 runtime: Box<dyn RuntimeConnection>,
1995 capabilities: crate::RuntimeCapabilities,
1996 workspace: PathBuf,
1997 fresh: bool,
1998 ) -> std::result::Result<Value, ServiceError> {
1999 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
2000 let token: std::sync::Arc<str> = crate::server::generate_token().into();
2001 let server = crate::server::run_frontend_http(
2002 host.clone(),
2003 host.frontend_sender(),
2004 "127.0.0.1:0",
2005 token.clone(),
2006 connection.handle().runtime_id.clone(),
2007 )
2008 .await
2009 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2010 let source = LiveRuntimeSource {
2011 harness: connection.handle().harness.as_str().to_string(),
2012 session_id: connection.handle().runtime_id.clone(),
2013 workspace: workspace.clone(),
2014 };
2015 let registration = register_live_runtime(
2016 connection.handle().runtime_id.clone(),
2017 source.clone(),
2018 format!("http://{}", server.address()),
2019 token.to_string(),
2020 )
2021 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2022 let endpoint = registration.endpoint().to_string();
2023 let launch = StructuredLaunch {
2024 cwd: workspace,
2025 program: std::env::current_exe()
2029 .ok()
2030 .map(|path| path.to_string_lossy().into_owned())
2031 .unwrap_or_else(|| "supercode".into()),
2032 arguments: vec![
2033 "open".into(),
2034 endpoint,
2035 "--harness".into(),
2036 source.harness,
2037 "--session".into(),
2038 source.session_id,
2039 ],
2040 env: BTreeMap::new(),
2041 };
2042 let lease = HostedRuntimeLease {
2043 connection,
2044 _host: host,
2045 _registration: registration,
2046 _server: server,
2047 };
2048 let opened = self.insert_runtime(Box::new(lease))?;
2049 let connection_id = opened["connection"]
2050 .as_str()
2051 .expect("insert_runtime returns a connection id")
2052 .to_string();
2053 self.terminal_launches.insert(connection_id, launch);
2054 Ok(opened)
2055 }
2056
2057 #[cfg(not(feature = "adapter-api"))]
2058 async fn insert_hosted_runtime(
2059 &mut self,
2060 runtime: Box<dyn RuntimeConnection>,
2061 _capabilities: crate::RuntimeCapabilities,
2062 _workspace: PathBuf,
2063 _fresh: bool,
2064 ) -> std::result::Result<Value, ServiceError> {
2065 self.insert_runtime(runtime)
2066 }
2067
2068 fn runtime_mut(
2069 &mut self,
2070 connection: &str,
2071 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2072 if self.runtimes_in_flight.contains(connection) {
2073 return Err(self.lent_out(connection));
2074 }
2075 self.runtimes.get_mut(connection).ok_or_else(|| {
2076 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2077 })
2078 }
2079
2080 fn lent_out(&self, connection: &str) -> ServiceError {
2084 ServiceError::Operation(format!(
2085 "runtime connection `{connection}`: a harness turn is already in progress"
2086 ))
2087 }
2088
2089 fn lend_runtime(
2092 &mut self,
2093 connection: &str,
2094 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2095 if self.runtimes_in_flight.contains(connection) {
2096 return Err(self.lent_out(connection));
2097 }
2098 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2099 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2100 })?;
2101 self.runtimes_in_flight.insert(connection.to_string());
2102 Ok(runtime)
2103 }
2104
2105 fn surrender_runtime(
2117 &mut self,
2118 connection: &str,
2119 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2120 if self.runtimes_in_flight.contains(connection) {
2121 return Err(self.lent_out(connection));
2122 }
2123 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2124 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2125 })?;
2126 let process_group = runtime_process_group(runtime.handle());
2127 let runtime_id = runtime.handle().runtime_id.clone();
2128 self.terminal_launches.remove(connection);
2129 self.runtime_sequences.remove(&runtime_id);
2130 self.approvals.forget(connection);
2131 Ok((runtime, process_group))
2132 }
2133
2134 pub fn kill_all_runtime_groups(&self) -> usize {
2145 self.runtimes
2146 .values()
2147 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2148 .count()
2149 }
2150
2151 async fn mutate_session(
2162 &mut self,
2163 verb: crate::SessionVerb,
2164 params: Value,
2165 ) -> std::result::Result<Value, ServiceError> {
2166 let mutation = decode::<crate::SessionMutation>(params)?;
2167 let door = crate::sessions_control::door(&mutation.harness, verb)
2168 .map_err(session_control_error)?;
2169 let outcome = match door {
2170 #[cfg(not(feature = "adapter-api"))]
2174 crate::SessionDoor::Live(command) => {
2175 return Err(ServiceError::Operation(format!(
2176 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2177 session, which needs this build's `adapter-api` feature",
2178 mutation.harness,
2179 verb.as_str()
2180 )));
2181 }
2182 #[cfg(feature = "adapter-api")]
2183 crate::SessionDoor::Live(command) => {
2184 let connection = mutation
2185 .connection
2186 .clone()
2187 .filter(|value| !value.trim().is_empty())
2188 .ok_or_else(|| {
2189 ServiceError::InvalidParams(format!(
2190 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2191 driven session: pass the `connection` of an open runtime \
2192 (`harness.v1.runtimes.start`)",
2193 mutation.harness,
2194 verb.as_str()
2195 ))
2196 })?;
2197 let runtime = self.runtime_mut(&connection)?;
2198 let session = live_session_name(runtime.as_ref(), &mutation);
2199 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2205 .await;
2206 }
2207 _ => run_session_mutation(verb, &mutation).await?,
2208 };
2209 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2210 }
2211
2212 async fn inventory_call(
2216 &self,
2217 method: &str,
2218 params: Value,
2219 ) -> std::result::Result<Value, ServiceError> {
2220 run_inventory(self.inventory_work(method, params)?).await
2221 }
2222
2223 fn inventory_work(
2229 &self,
2230 method: &str,
2231 params: Value,
2232 ) -> std::result::Result<InventoryWork, ServiceError> {
2233 let mut params = decode::<HarnessInventoryParams>(params)?;
2234 if method == "harness.v1.harnesses.probe" {
2235 let harness = params.harness.take().ok_or_else(|| {
2236 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2237 })?;
2238 params.harnesses = vec![harness];
2239 }
2240 let selected = params
2241 .harnesses
2242 .iter()
2243 .map(HarnessId::as_str)
2244 .collect::<std::collections::BTreeSet<_>>();
2245 let supported = harness_support_registry()
2246 .harnesses
2247 .into_iter()
2248 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2249 .collect::<Vec<_>>();
2250 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2251 let known = supported
2252 .iter()
2253 .map(|harness| harness.id.as_str())
2254 .collect::<std::collections::BTreeSet<_>>();
2255 let missing = params
2256 .harnesses
2257 .iter()
2258 .filter(|id| !known.contains(id.as_str()))
2259 .map(HarnessId::as_str)
2260 .collect::<Vec<_>>();
2261 return Err(ServiceError::InvalidParams(format!(
2262 "unknown harness(es): {}",
2263 missing.join(", ")
2264 )));
2265 }
2266 let global_counts = params
2267 .include_sessions
2268 .then(|| self.session_counts(None, ¶ms.harnesses));
2269 let workspace_counts = params
2270 .include_sessions
2271 .then(|| {
2272 params
2273 .workspace
2274 .as_deref()
2275 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2276 })
2277 .flatten();
2278 Ok(InventoryWork {
2279 params,
2280 supported,
2281 global_counts,
2282 workspace_counts,
2283 })
2284 }
2285
2286 #[cfg(feature = "adapter-api")]
2287 async fn harness_authentication_call(
2288 &self,
2289 method: &str,
2290 params: Value,
2291 ) -> std::result::Result<Value, ServiceError> {
2292 match method {
2293 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2294 let params = decode::<HarnessAuthenticationParams>(params)?;
2295 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2296 .map_err(|error| ServiceError::Operation(error.to_string()))
2297 }
2298 "harness.v1.harnesses.auth.begin" => {
2299 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2300 let cwd = params
2301 .cwd
2302 .or_else(|| std::env::current_dir().ok())
2303 .unwrap_or_else(|| PathBuf::from("."));
2304 let plan = crate::harness_authentication_plan(
2305 ¶ms.harness,
2306 params.environment,
2307 params.method,
2308 &cwd,
2309 )
2310 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2311 serde_json::to_value(plan)
2312 .map_err(|error| ServiceError::Operation(error.to_string()))
2313 }
2314 _ => Err(ServiceError::MethodNotFound),
2315 }
2316 }
2317
2318 fn session_counts(
2319 &self,
2320 workspace: Option<&Path>,
2321 harnesses: &[HarnessId],
2322 ) -> BTreeMap<String, usize> {
2323 let mut counts = BTreeMap::new();
2324 for session in self
2325 .catalog
2326 .discover(&DiscoveryQuery {
2327 workspace: workspace.map(Path::to_path_buf),
2328 harnesses: harnesses.to_vec(),
2329 ..DiscoveryQuery::default()
2330 })
2331 .unwrap_or_default()
2332 {
2333 *counts
2334 .entry(session.locator.harness.as_str().to_string())
2335 .or_insert(0) += 1;
2336 }
2337 counts
2338 }
2339}
2340
2341#[async_trait::async_trait]
2342impl SdkService for HarnessSessionService {
2343 fn capabilities(&self) -> SdkCapabilities {
2344 SdkCapabilities::default()
2345 }
2346
2347 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2348 if request.operation == SdkOperation::Events {
2349 let events = self
2350 .poll_sdk_events()
2351 .await
2352 .into_iter()
2353 .map(|(_, event)| event)
2354 .collect::<Vec<_>>();
2355 return serde_json::to_value(events).map_err(|error| {
2356 SdkError::new(
2357 SdkErrorCode::Execution,
2358 request.operation,
2359 error.to_string(),
2360 )
2361 });
2362 }
2363 if self.runtimes.is_empty()
2364 && matches!(
2365 request.operation,
2366 SdkOperation::Input
2367 | SdkOperation::Interrupt
2368 | SdkOperation::Steer
2369 | SdkOperation::Respond
2370 | SdkOperation::Close
2371 )
2372 {
2373 return Err(SdkError::unsupported(request.operation));
2374 }
2375 let method = request
2376 .operation
2377 .method()
2378 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2379 let result = match request.operation {
2380 SdkOperation::Discover
2381 | SdkOperation::Load
2382 | SdkOperation::Export
2383 | SdkOperation::ProfilesList
2384 | SdkOperation::ProfilesGet
2385 | SdkOperation::ProfilesCreate
2386 | SdkOperation::ProfilesDelete
2387 | SdkOperation::SkillsList
2388 | SdkOperation::SkillsInstall
2389 | SdkOperation::SkillsRemove
2390 | SdkOperation::ChannelsList
2391 | SdkOperation::RoutesList
2392 | SdkOperation::TriggersList
2393 | SdkOperation::ChannelsStatus
2394 | SdkOperation::MemoryShow
2395 | SdkOperation::MemorySearch
2396 | SdkOperation::JobsList
2397 | SdkOperation::JobsGet
2398 | SdkOperation::JobsCreate
2399 | SdkOperation::JobsUpdate
2400 | SdkOperation::JobsPause
2401 | SdkOperation::JobsResume
2402 | SdkOperation::JobsRun
2403 | SdkOperation::JobsDelete
2404 | SdkOperation::JobsNotepad
2405 | SdkOperation::JobsNotepadSet
2406 | SdkOperation::JobsNotepadDelete
2407 | SdkOperation::RunsList
2408 | SdkOperation::RunsGet
2409 | SdkOperation::ApprovalsList
2410 | SdkOperation::OrchestrationLoad
2411 | SdkOperation::OrchestrationSave
2412 | SdkOperation::OrchestrationCompile
2413 | SdkOperation::OrchestrationDecompile
2414 | SdkOperation::OrchestrationImport
2415 | SdkOperation::OrchestrationExport
2416 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2417 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2420 SdkOperation::Start
2421 | SdkOperation::Resume
2422 | SdkOperation::Input
2423 | SdkOperation::Interrupt
2424 | SdkOperation::Steer
2425 | SdkOperation::Respond
2426 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2427 SdkOperation::SessionsNew => {
2432 self.mutate_session(crate::SessionVerb::New, request.params)
2433 .await
2434 }
2435 SdkOperation::SessionsReset => {
2436 self.mutate_session(crate::SessionVerb::Reset, request.params)
2437 .await
2438 }
2439 SdkOperation::SessionsArchive => {
2440 self.mutate_session(crate::SessionVerb::Archive, request.params)
2441 .await
2442 }
2443 SdkOperation::SessionsDelete => {
2444 self.mutate_session(crate::SessionVerb::Delete, request.params)
2445 .await
2446 }
2447 SdkOperation::Events => unreachable!("handled before method dispatch"),
2448 };
2449 result.map_err(|error| sdk_error(request.operation, error))
2450 }
2451
2452 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2453 Ok(self
2454 .poll_sdk_events()
2455 .await
2456 .into_iter()
2457 .map(|(_, event)| event)
2458 .collect())
2459 }
2460}
2461
2462#[cfg(feature = "adapter-api")]
2463struct HostedRuntimeLease {
2464 connection: HostedHarnessConnection,
2465 _host: std::sync::Arc<HostedHarnessRuntime>,
2466 _registration: LiveRuntimeRegistration,
2467 _server: crate::server::FrontendHttpServer,
2468}
2469
2470#[async_trait::async_trait]
2471#[cfg(feature = "adapter-api")]
2472impl RuntimeConnection for HostedRuntimeLease {
2473 fn handle(&self) -> &crate::RuntimeHandle {
2474 self.connection.handle()
2475 }
2476
2477 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2478 self.connection.send_input(input).await
2479 }
2480
2481 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2482 self.connection.next_event().await
2483 }
2484
2485 async fn interrupt(&mut self) -> crate::Result<()> {
2486 self.connection.interrupt().await
2487 }
2488
2489 async fn steer(&mut self, text: String) -> crate::Result<()> {
2492 self.connection.steer(text).await
2493 }
2494
2495 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2496 self.connection.respond(request_id, response).await
2497 }
2498
2499 async fn close(&mut self) -> crate::Result<()> {
2500 self.connection.close().await
2501 }
2502}
2503
2504struct InventoryWork {
2507 params: HarnessInventoryParams,
2508 supported: Vec<crate::HarnessSupportDescriptor>,
2509 global_counts: Option<BTreeMap<String, usize>>,
2510 workspace_counts: Option<BTreeMap<String, usize>>,
2511}
2512
2513async fn run_session_mutation(
2520 verb: crate::SessionVerb,
2521 mutation: &crate::SessionMutation,
2522) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2523 let door =
2530 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2531 if let crate::SessionDoor::Http = door {
2532 return crate::sessions_control::mutate(verb, mutation)
2533 .await
2534 .map_err(session_control_error);
2535 }
2536 let mutation = mutation.clone();
2537 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2538 .await
2539 .map_err(|error| {
2540 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2541 })?
2542 .map_err(session_control_error)
2543}
2544
2545async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2548 let InventoryWork {
2549 params,
2550 supported,
2551 global_counts,
2552 workspace_counts,
2553 } = work;
2554 let probes = supported.into_iter().map(|descriptor| {
2555 let global = global_counts
2556 .as_ref()
2557 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2558 let workspace = workspace_counts
2559 .as_ref()
2560 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2561 probe_harness(descriptor, ¶ms, global, workspace)
2562 });
2563 let harnesses = futures::future::join_all(probes).await;
2564 serde_json::to_value(HarnessInventoryReport {
2565 probe: params.probe,
2566 workspace: params.workspace,
2567 harnesses,
2568 })
2569 .map_err(|error| ServiceError::Operation(error.to_string()))
2570}
2571
2572async fn probe_harness(
2573 descriptor: crate::HarnessSupportDescriptor,
2574 params: &HarnessInventoryParams,
2575 global: Option<usize>,
2576 workspace: Option<usize>,
2577) -> LocalHarness {
2578 let launch = descriptor.runtime.default_launch.as_ref();
2579 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2584 .then(crate::orchestrator::daemon_entry)
2585 .and_then(Result::ok);
2586 let executable = match &orchestrator_entry {
2587 Some(entry) => Some(entry.clone()),
2588 None => launch.and_then(|launch| find_executable(&launch.program)),
2589 };
2590 let installed = executable.is_some();
2591 let version = if params.skip_versions || orchestrator_entry.is_some() {
2592 None
2595 } else {
2596 match executable.as_deref() {
2597 Some(path) => executable_version(path).await,
2598 None => None,
2599 }
2600 };
2601 let configured = auth_evidence(descriptor.id.as_str());
2602 let mut auth = if configured {
2603 HarnessAuthState::Configured
2604 } else if matches!(
2605 descriptor.id.as_str(),
2606 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2607 ) {
2608 HarnessAuthState::Required
2613 } else {
2614 HarnessAuthState::Unknown
2615 };
2616 let mut runtime = if installed {
2617 HarnessRuntimeState::Degraded
2618 } else {
2619 HarnessRuntimeState::Unavailable
2620 };
2621 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2622 let mut reason = (!installed).then(|| {
2623 if is_orchestrator {
2624 format!(
2625 "{} is supported but its daemon entry `{}` was not found",
2626 descriptor.display_name,
2627 crate::orchestrator::DAEMON_ENTRY
2628 )
2629 } else {
2630 format!(
2631 "{} is supported but `{}` was not found on PATH",
2632 descriptor.display_name,
2633 launch
2634 .map(|launch| launch.program.as_str())
2635 .unwrap_or("executable")
2636 )
2637 }
2638 });
2639 let mut repair = (!installed).then(|| {
2640 if is_orchestrator {
2641 format!(
2642 "Install the `supercode-orchestrator` package so `{}` resolves.",
2643 crate::orchestrator::DAEMON_ENTRY
2644 )
2645 } else {
2646 format!(
2647 "Install {} and ensure `{}` is on PATH.",
2648 descriptor.display_name,
2649 launch
2650 .map(|launch| launch.program.as_str())
2651 .unwrap_or("its executable")
2652 )
2653 }
2654 });
2655
2656 if installed && params.probe == HarnessProbeLevel::Handshake {
2657 let backend_params = RuntimeBackendParams {
2658 harness: descriptor.id.clone(),
2659 protocol: None,
2660 launch: None,
2661 base_url: None,
2662 policy: RuntimePolicy::Default,
2663 };
2664 match runtime_backend(&backend_params) {
2665 Ok(backend) => {
2666 let cwd = params
2667 .workspace
2668 .clone()
2669 .or_else(|| std::env::current_dir().ok())
2670 .unwrap_or_else(|| PathBuf::from("."));
2671 let isolated = descriptor
2672 .runtime
2673 .default_launch
2674 .clone()
2675 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2676 let Some(isolated) = isolated else {
2677 reason = Some(
2678 "No-prompt runtime handshake could not create its isolated harness home."
2679 .into(),
2680 );
2681 repair = Some(
2682 "Check temporary-directory permissions, then run the handshake probe again."
2683 .into(),
2684 );
2685 let running = probe_running_instance(descriptor.id.as_str());
2686 return LocalHarness {
2687 gateway: gateway_health(
2688 descriptor.id.as_str(),
2689 installed,
2690 running.as_ref(),
2691 version.as_deref(),
2692 ),
2693 id: descriptor.id,
2694 display_name: descriptor.display_name,
2695 supported: true,
2696 installed,
2697 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2698 version,
2699 auth,
2700 runtime,
2701 protocol: descriptor.runtime.protocol,
2702 capabilities: descriptor.runtime.capabilities.clone(),
2703 effective_capabilities: descriptor.runtime.capabilities,
2704 sessions: HarnessSessionCounts { global, workspace },
2705 running,
2706 reason,
2707 repair,
2708 };
2709 };
2710 match tokio::time::timeout(
2711 Duration::from_secs(30),
2712 backend.start(RuntimeStartRequest {
2713 cwd,
2714 launch: Some(isolated.launch.clone()),
2715 mcp_servers: Vec::new(),
2716 approval_policy: None,
2717 }),
2718 )
2719 .await
2720 {
2721 Ok(Ok(mut connection)) => {
2722 match stabilize_handshake(connection.as_mut()).await {
2723 Ok(()) => {
2724 auth = HarnessAuthState::Ready;
2725 runtime = HarnessRuntimeState::Ready;
2726 reason = Some(
2727 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2728 .into(),
2729 );
2730 repair = None;
2731 }
2732 Err(message) => {
2733 auth = if looks_like_auth_error(&message) {
2734 HarnessAuthState::Required
2735 } else if configured {
2736 HarnessAuthState::Configured
2737 } else {
2738 HarnessAuthState::Unknown
2739 };
2740 reason = Some(format!(
2741 "No-prompt runtime handshake became unhealthy during startup: {message}"
2742 ));
2743 repair = Some(if auth == HarnessAuthState::Required {
2744 format!(
2745 "Run `{}` interactively once and complete sign-in, then probe again.",
2746 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2747 )
2748 } else {
2749 "Run the harness directly to inspect its startup failure, then probe again."
2750 .into()
2751 });
2752 }
2753 }
2754 let _ =
2755 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2756 }
2757 Ok(Err(error)) => {
2758 let message = truncate_text(&error.to_string(), 500);
2759 auth = if looks_like_auth_error(&message) {
2760 HarnessAuthState::Required
2761 } else if configured {
2762 HarnessAuthState::Configured
2763 } else {
2764 HarnessAuthState::Unknown
2765 };
2766 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2767 repair = Some(if auth == HarnessAuthState::Required {
2768 format!(
2769 "Run `{}` interactively once and complete sign-in, then probe again.",
2770 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2771 )
2772 } else {
2773 "Check the harness installation and run the handshake probe again."
2774 .into()
2775 });
2776 }
2777 Err(_) => {
2778 reason =
2779 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2780 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2781 }
2782 }
2783 let _ = isolated.cleanup();
2792 tokio::time::sleep(Duration::from_millis(250)).await;
2793 if let Err(error) = isolated.cleanup() {
2794 auth = if configured {
2795 HarnessAuthState::Configured
2796 } else {
2797 HarnessAuthState::Unknown
2798 };
2799 runtime = HarnessRuntimeState::Degraded;
2800 reason = Some(format!(
2801 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2802 ));
2803 repair = Some(
2804 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2805 .into(),
2806 );
2807 }
2808 }
2809 Err(error) => {
2810 reason = Some(error_message(error));
2811 }
2812 }
2813 } else if installed && configured {
2814 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2815 } else if installed && auth == HarnessAuthState::Required {
2816 reason = Some("Executable found, but no native authentication evidence is present.".into());
2817 repair = Some(format!(
2818 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2819 descriptor.id.as_str()
2820 ));
2821 } else if installed {
2822 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2823 repair = Some(format!(
2824 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2825 launch
2826 .map(|launch| launch.program.as_str())
2827 .unwrap_or("the harness")
2828 ));
2829 }
2830
2831 let effective_capabilities = if installed {
2832 descriptor.runtime.capabilities.clone()
2833 } else {
2834 unavailable_capabilities()
2835 };
2836 let running = probe_running_instance(descriptor.id.as_str());
2837 LocalHarness {
2838 gateway: gateway_health(
2839 descriptor.id.as_str(),
2840 installed,
2841 running.as_ref(),
2842 version.as_deref(),
2843 ),
2844 id: descriptor.id,
2845 display_name: descriptor.display_name,
2846 supported: true,
2847 installed,
2848 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2849 version,
2850 auth,
2851 runtime,
2852 protocol: descriptor.runtime.protocol,
2853 capabilities: descriptor.runtime.capabilities,
2854 effective_capabilities,
2855 sessions: HarnessSessionCounts { global, workspace },
2856 running,
2857 reason,
2858 repair,
2859 }
2860}
2861
2862async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2863 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2864 loop {
2865 let now = tokio::time::Instant::now();
2866 if now >= deadline {
2867 return Ok(());
2868 }
2869 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2870 Err(_) => return Ok(()),
2871 Ok(Ok(Some(event))) => {
2872 if let Some(message) = handshake_event_failure(&event) {
2873 return Err(truncate_text(&message, 500));
2874 }
2875 }
2876 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2877 Ok(Err(error)) => return Err(error.to_string()),
2878 }
2879 }
2880}
2881
2882fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2883 let detail = event
2884 .payload
2885 .get("message")
2886 .or_else(|| event.payload.get("line"))
2887 .and_then(Value::as_str)
2888 .unwrap_or(event.kind.as_str());
2889 match event.kind.as_str() {
2890 "transport_closed" => Some("runtime transport closed during startup".into()),
2891 "transport_error" => Some(format!("runtime transport error: {detail}")),
2892 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2893 _ => None,
2898 }
2899}
2900
2901fn indexed_claude_window(
2902 locator: &SessionLocator,
2903 options: &SessionLoadOptions,
2904) -> std::result::Result<Option<Value>, ServiceError> {
2905 use supercode_interchange::session::ClaudeReadIndex;
2906 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2909 || options.include_subagents != Some(false)
2910 {
2911 return Ok(None);
2912 }
2913 let crate::StorageLocator::File { path } = &locator.storage else {
2914 return Ok(None);
2915 };
2916 if !ClaudeReadIndex::supports(path)
2917 .map_err(|error| ServiceError::Operation(error.to_string()))?
2918 {
2919 return Ok(None);
2920 }
2921 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2922 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2923 let total = index.len();
2924 let (offset, end) = projected_message_window(total, options);
2925 let session = index
2926 .read_messages(offset..end)
2927 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2928 let summary = index
2929 .read_summary()
2930 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2931 let selected_options = SessionLoadOptions {
2932 message_offset: None,
2933 message_limit: None,
2934 message_tail: None,
2935 ..options.clone()
2936 };
2937 let mut selected = projected_session_json(&session, &selected_options);
2938 selected["raw_record_count"] = json!(index.raw_record_count());
2939 Ok(Some(json!({
2940 "session": selected,
2941 "summary": projected_session_summary(&summary, options),
2942 "window": {
2943 "has_more": offset > 0 || end < total, "has_newer": end < total,
2944 "has_older": offset > 0, "newer_items": index.item_count(end..total),
2945 "offset": offset, "older_items": index.item_count(0..offset),
2946 "returned": end - offset, "total_messages": total,
2947 }
2948 })))
2949}
2950
2951fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
2952 let total_messages = session.messages.len();
2953 let (offset, end) = projected_message_window(total_messages, options);
2954 json!({
2955 "session": projected_session_json(session, options),
2956 "summary": projected_session_summary(session, options),
2957 "window": {
2958 "has_more": offset > 0 || end < total_messages,
2959 "has_newer": end < total_messages,
2960 "has_older": offset > 0,
2961 "newer_items": normalized_item_count(&session.messages[end..]),
2962 "offset": offset,
2963 "older_items": normalized_item_count(&session.messages[..offset]),
2964 "returned": end.saturating_sub(offset),
2965 "total_messages": total_messages,
2966 }
2967 })
2968}
2969
2970fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
2971 messages
2972 .iter()
2973 .map(|message| {
2974 let conversation = usize::from(
2975 matches!(message.role, Role::Assistant | Role::User)
2976 && message_has_content(message),
2977 );
2978 let tool_result =
2979 usize::from(message.role == Role::Tool && message_has_content(message));
2980 conversation + tool_result + message.tool_calls().len()
2981 })
2982 .sum()
2983}
2984
2985fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
2986 let mut conversational = session.messages.iter().filter(|message| {
2987 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
2988 });
2989 let first_message = conversational.clone().next();
2990 let last_message = conversational.next_back();
2991 let mut assistant = session
2992 .messages
2993 .iter()
2994 .filter(|message| message.role == Role::Assistant && message_has_content(message));
2995 let first_assistant_message = assistant.clone().next();
2996 let last_assistant_message = assistant.next_back();
2997 let end_of_turn = session
2998 .messages
2999 .iter()
3000 .rev()
3001 .find(|message| message.role != Role::System)
3002 .is_some_and(|message| {
3003 message.role == Role::Assistant
3004 && message_has_content(message)
3005 && message.tool_calls().is_empty()
3006 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
3008 });
3009 let project = |message: Option<&crate::ChatMessage>| {
3010 message.map(|message| project_inline_media(message_json(message), options))
3011 };
3012 json!({
3013 "end_of_turn": end_of_turn,
3014 "first_assistant_message": project(first_assistant_message),
3015 "first_message": project(first_message),
3016 "last_assistant_message": project(last_assistant_message),
3017 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3018 "last_message": project(last_message),
3019 })
3020}
3021
3022fn message_has_content(message: &crate::ChatMessage) -> bool {
3023 message
3024 .content
3025 .as_deref()
3026 .is_some_and(|content| !content.trim().is_empty())
3027 || message
3028 .content_parts
3029 .as_ref()
3030 .is_some_and(|parts| !parts.is_empty())
3031}
3032
3033fn message_text(message: &crate::ChatMessage) -> String {
3034 if let Some(content) = &message.content {
3035 return content.clone();
3036 }
3037 message
3038 .content_parts
3039 .as_ref()
3040 .into_iter()
3041 .flatten()
3042 .filter_map(|part| part.get("text").and_then(Value::as_str))
3043 .collect::<Vec<_>>()
3044 .join("\n")
3045}
3046
3047fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3048 let (offset, end) = projected_message_window(session.messages.len(), options);
3049 let messages = session.messages[offset..end]
3050 .iter()
3051 .map(|message| project_inline_media(message_json(message), options))
3052 .collect::<Vec<_>>();
3053 let subagents = if options.include_subagents.unwrap_or(true) {
3054 let subagent_options = SessionLoadOptions {
3059 message_limit: None,
3060 message_offset: None,
3061 message_tail: None,
3062 ..options.clone()
3063 };
3064 session
3065 .subagents
3066 .iter()
3067 .map(|subagent| projected_session_json(subagent, &subagent_options))
3068 .collect::<Vec<_>>()
3069 } else {
3070 Vec::new()
3071 };
3072 json!({
3073 "source": match session.meta.source {
3074 SessionSource::ClaudeCode => "claude_code",
3075 SessionSource::Codex => "codex",
3076 SessionSource::Gemini => "gemini",
3077 SessionSource::Goose => "goose",
3078 SessionSource::Grok => "grok",
3079 SessionSource::Native => "native",
3080 SessionSource::OpenClaw => "openclaw",
3081 SessionSource::Hermes => "hermes",
3082 SessionSource::OpenCode => "opencode",
3083 SessionSource::Pi => "pi",
3084 },
3085 "session_id": session.meta.session_id,
3086 "ended_at": session.meta.ended_at,
3087 "end_reason": session.meta.end_reason,
3088 "model": session.meta.model,
3089 "cwd": session.meta.cwd,
3090 "system_prompt": session.meta.system_prompt,
3091 "agent_id": session.meta.agent_id,
3092 "parent_tool_use_id": session.meta.parent_tool_use_id,
3093 "lineage": session.meta.lineage,
3094 "messages": messages,
3095 "subagents": subagents,
3096 "raw_record_count": session.raw.len(),
3097 "parse_error_lines": session.parse_error_lines,
3098 })
3099}
3100
3101fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3102 if let Some(tail) = options.message_tail {
3103 return (total.saturating_sub(tail), total);
3104 }
3105 let offset = options.message_offset.unwrap_or(0).min(total);
3106 let end = options
3107 .message_limit
3108 .map(|limit| offset.saturating_add(limit).min(total))
3109 .unwrap_or(total);
3110 (offset, end)
3111}
3112
3113fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3114 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3115 return message;
3116 };
3117 for part in parts {
3118 let Some(url) = part
3119 .get("image_url")
3120 .and_then(|image| image.get("url"))
3121 .and_then(Value::as_str)
3122 else {
3123 continue;
3124 };
3125 let Some(rest) = url.strip_prefix("data:") else {
3126 continue;
3127 };
3128 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3129 continue;
3130 };
3131 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3132 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3133 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3134 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3135 || options
3136 .max_inline_media_bytes
3137 .is_some_and(|limit| decoded_bytes > limit);
3138 if should_elide {
3139 *part = json!({
3140 "type": "media_reference",
3141 "media_type": media_type,
3142 "encoding": "base64",
3143 "encoded_bytes": encoded.len(),
3144 "decoded_bytes": decoded_bytes,
3145 "omitted": true,
3146 });
3147 }
3148 }
3149 message
3150}
3151
3152#[derive(Deserialize)]
3153struct LocatorParams {
3154 locator: SessionLocator,
3155 #[serde(default)]
3168 fidelity: Option<Fidelity>,
3169 #[serde(default)]
3172 view: Option<SessionReadView>,
3173}
3174
3175#[derive(Deserialize)]
3176struct SessionReadView {
3177 #[serde(default)]
3180 tail_messages: Option<usize>,
3181 #[serde(default)]
3184 include_subagents: bool,
3185 #[serde(default)]
3187 display_history: bool,
3188 #[serde(default)]
3191 max_message_chars: Option<usize>,
3192}
3193
3194impl LocatorParams {
3195 fn read_fidelity(&self) -> Fidelity {
3196 self.fidelity.unwrap_or(Fidelity::Semantic)
3197 }
3198
3199 fn include_subagents(&self) -> bool {
3200 self.view
3201 .as_ref()
3202 .map(|view| view.include_subagents)
3203 .unwrap_or(true)
3204 }
3205
3206 fn tail_messages(&self) -> Option<usize> {
3207 self.view
3208 .as_ref()
3209 .and_then(|view| view.tail_messages)
3210 .map(|limit| limit.clamp(1, 5_000))
3211 }
3212
3213 fn display_history(&self) -> bool {
3214 self.view.as_ref().is_some_and(|view| view.display_history)
3215 }
3216
3217 fn max_message_chars(&self) -> Option<usize> {
3218 self.view
3219 .as_ref()
3220 .and_then(|view| view.max_message_chars)
3221 .map(|limit| limit.clamp(256, 64_000))
3222 }
3223
3224 fn bound_session(&self, session: &mut Session) {
3225 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3226 }
3227}
3228
3229#[derive(Debug, Clone, Copy, Default, Deserialize)]
3230#[serde(rename_all = "snake_case")]
3231enum InlineMediaMode {
3232 #[default]
3233 Full,
3234 Metadata,
3235}
3236
3237#[derive(Debug, Clone, Default, Deserialize)]
3238#[serde(default)]
3239struct SessionLoadOptions {
3240 include_subagents: Option<bool>,
3241 inline_media: InlineMediaMode,
3242 max_inline_media_bytes: Option<usize>,
3243 message_limit: Option<usize>,
3244 message_offset: Option<usize>,
3245 message_tail: Option<usize>,
3246}
3247
3248impl SessionLoadOptions {
3249 fn validate(&self) -> std::result::Result<(), ServiceError> {
3250 if self.message_tail.is_some()
3251 && (self.message_limit.is_some() || self.message_offset.is_some())
3252 {
3253 return Err(ServiceError::InvalidParams(
3254 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3255 .into(),
3256 ));
3257 }
3258 Ok(())
3259 }
3260}
3261
3262#[derive(Deserialize)]
3263struct LoadSessionParams {
3264 #[serde(flatten)]
3265 read: LocatorParams,
3266 #[serde(default)]
3267 options: Option<SessionLoadOptions>,
3268}
3269
3270#[derive(Deserialize)]
3271struct UnfollowParams {
3272 subscription: String,
3273}
3274
3275#[derive(Debug, Deserialize)]
3276#[serde(deny_unknown_fields)]
3277struct IndexResizeParams {
3278 subscription: String,
3279 limit: usize,
3280}
3281
3282#[derive(Deserialize)]
3283struct ActivitySubscribeParams {
3284 locators: Vec<SessionLocator>,
3285 #[serde(default)]
3286 homes: crate::HarnessHomes,
3287}
3288
3289#[derive(Deserialize)]
3290struct MessageSessionParams {
3291 locator: SessionLocator,
3292 text: String,
3293 #[serde(default)]
3294 subject: Option<String>,
3295 #[serde(default)]
3297 idempotency_key: Option<String>,
3298 #[serde(default)]
3300 channel: bool,
3301 #[serde(default)]
3305 from_name: Option<String>,
3306 #[serde(default)]
3308 voice_for: Option<crate::mailbox::MailAddress>,
3309 #[serde(default)]
3311 sender_name: Option<String>,
3312 #[serde(default)]
3314 in_reply_to: Option<String>,
3315 #[serde(default)]
3317 notify_when_idle: bool,
3318 #[serde(default)]
3323 as_user: bool,
3324 #[serde(default)]
3327 homes: crate::HarnessHomes,
3328}
3329
3330#[derive(Deserialize)]
3331#[serde(deny_unknown_fields)]
3332struct InboxParams {
3333 #[serde(default)]
3335 from_name: Option<String>,
3336 #[serde(default)]
3338 address: Option<String>,
3339 #[serde(default)]
3341 all: bool,
3342}
3343
3344#[derive(Deserialize)]
3345#[serde(deny_unknown_fields)]
3346struct HarnessSettingsParams {
3347 harness: String,
3348}
3349
3350#[derive(Deserialize)]
3351#[serde(deny_unknown_fields)]
3352struct ConfigureHarnessParams {
3353 harness: String,
3354 #[serde(default)]
3355 changes: Vec<crate::HarnessSettingChange>,
3356 #[serde(default)]
3357 expected_revision: Option<String>,
3358}
3359
3360fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3361 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3362 Ok(report) => (
3363 serde_json::to_value(report).unwrap_or(Value::Null),
3364 Value::Null,
3365 ),
3366 Err(error) => (
3367 Value::Null,
3368 Value::String(format!(
3369 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3370 )),
3371 ),
3372 }
3373}
3374
3375async fn message_live_session(params: &MessageSessionParams) -> Value {
3389 use crate::mail_route::{Delivered, NoDoor, Refused};
3390 let (inbound_controls, inbound_controls_error) =
3391 claude_inbound_controls_or_error(¶ms.homes);
3392 let refused = |reason: &str, message: String| {
3393 json!({
3394 "delivered_to_bus": false,
3395 "refusal": {"reason": reason, "message": message},
3396 "inbound_controls": inbound_controls,
3397 "inbound_controls_error": inbound_controls_error,
3398 })
3399 };
3400 if params.text.trim().is_empty() {
3401 return refused(
3402 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3403 "refusing to deliver an empty message".into(),
3404 );
3405 }
3406 let sender = match operator_address(params.from_name.as_deref()) {
3407 Ok(sender) => sender,
3408 Err(message) => return refused("invalid_sender", message),
3409 };
3410 let receiver = match crate::mailbox::MailAddress::new(
3411 crate::mailbox::local_machine_name(),
3412 params.locator.harness.as_str(),
3413 ¶ms.locator.session_id,
3414 ) {
3415 Ok(receiver) => receiver,
3416 Err(error) => return refused("delivery_failed", error.to_string()),
3417 };
3418 if params
3419 .voice_for
3420 .as_ref()
3421 .is_some_and(|voice_for| voice_for != &receiver)
3422 {
3423 return refused(
3424 "invalid_sender",
3425 "a voice front must represent the receiving session".into(),
3426 );
3427 }
3428 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3430 && crate::runtime_mail::controlled_runtime(
3431 HarnessId::CLAUDE_CODE,
3432 ¶ms.locator.session_id,
3433 )
3434 .is_none()
3435 {
3436 if let Err(refusal) =
3437 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3438 {
3439 return refused(refusal.reason.as_str(), refusal.message);
3440 }
3441 }
3442 if params.as_user {
3443 return message_as_user(
3444 params,
3445 sender,
3446 receiver,
3447 inbound_controls,
3448 inbound_controls_error,
3449 )
3450 .await;
3451 }
3452 let door = match crate::mail_route::door_for(¶ms.homes, &receiver) {
3453 Ok(door) => door,
3454 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3455 return refused(
3456 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3457 format!(
3458 "no running `{}` session `{}` is reachable; its transcript is persisted only",
3459 params.locator.harness.as_str(),
3460 params.locator.session_id
3461 ),
3462 )
3463 }
3464 };
3465 let mut envelope = match crate::mailbox::Envelope::new(
3466 sender.clone(),
3467 params
3468 .sender_name
3469 .clone()
3470 .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3471 if params.channel {
3472 crate::mailbox::MailKind::Channel
3473 } else {
3474 crate::mailbox::MailKind::Peer
3475 },
3476 crate::mailbox::ReplyVia::Command,
3477 params.text.clone(),
3478 ) {
3479 Ok(envelope) => envelope,
3480 Err(error) => return refused("delivery_failed", error.to_string()),
3481 };
3482 if let Some(key) = ¶ms.idempotency_key {
3483 if key.is_empty() || key.len() > 256 {
3484 return refused("invalid_key", "idempotency_key needs 1–256 bytes".into());
3485 }
3486 let identity = serde_json::to_vec(&(sender.to_string(), receiver.to_string(), key))
3487 .expect("string tuple serializes");
3488 envelope.id = format!("m-{}", &blake3::hash(&identity).to_hex()[..24]);
3489 let mailbox = match crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &receiver) {
3490 Ok(mailbox) => mailbox,
3491 Err(error) => return refused("delivery_failed", error.to_string()),
3492 };
3493 match mailbox.find(&envelope.id) {
3494 Ok(Some(previous)) => {
3495 let saved = previous.envelope;
3496 let subject = params.subject.as_deref().map(|value| {
3497 value
3498 .lines()
3499 .next()
3500 .unwrap_or("")
3501 .chars()
3502 .take(200)
3503 .collect::<String>()
3504 });
3505 if saved.body != envelope.body
3506 || saved.subject != subject
3507 || saved.kind != envelope.kind
3508 || saved.from_name != envelope.from_name
3509 || saved.in_reply_to != params.in_reply_to
3510 || saved.voice_for != params.voice_for
3511 {
3512 return refused(
3513 "idempotency_conflict",
3514 "this key already names different mail".into(),
3515 );
3516 }
3517 }
3518 Ok(None) => {}
3519 Err(error) => return refused("delivery_failed", error.to_string()),
3520 }
3521 }
3522 envelope.in_reply_to = params.in_reply_to.clone();
3523 envelope.voice_for = params.voice_for.clone();
3524 envelope.subject = params.subject.as_deref().map(|value| {
3525 value
3526 .lines()
3527 .next()
3528 .unwrap_or("")
3529 .chars()
3530 .take(200)
3531 .collect()
3532 });
3533 let how = match crate::mail_route::deliver(
3534 &envelope,
3535 &receiver,
3536 &door,
3537 true,
3538 params.notify_when_idle,
3539 )
3540 .await
3541 {
3542 Err(detail) => {
3543 return refused(
3544 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3545 detail,
3546 )
3547 }
3548 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3549 Ok(Err(Refused::TooLong(bytes))) => {
3550 return refused(
3551 "too_long",
3552 format!(
3553 "the message is {bytes} bytes; the limit is {}",
3554 crate::mail_route::MAX_RELAYED_BYTES
3555 ),
3556 )
3557 }
3558 Ok(Ok(delivered)) => match delivered {
3559 Delivered::Steered => "steered",
3560 Delivered::Started => "started",
3561 Delivered::Native { busy: true } => "next_tool_call",
3562 Delivered::Native { busy: false } => "started",
3563 Delivered::Hooked => "hook",
3564 Delivered::HookWoken => "woken",
3565 Delivered::Queued => "queued",
3566 Delivered::Stored => "stored",
3567 Delivered::Operator => "filed",
3568 },
3569 };
3570 json!({
3571 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3572 "message_id": envelope.id,
3573 "reply_to": sender.to_string(),
3574 "target": {
3575 "harness": params.locator.harness.as_str(),
3576 "session_id": params.locator.session_id,
3577 "name": match &door {
3578 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
3579 _ => None,
3580 },
3581 },
3582 "delivery": {"door": door.name(), "how": how},
3583 "inbound_controls": inbound_controls,
3584 "inbound_controls_error": inbound_controls_error,
3585 })
3586}
3587
3588async fn message_as_user(
3592 params: &MessageSessionParams,
3593 sender: crate::mailbox::MailAddress,
3594 receiver: crate::mailbox::MailAddress,
3595 inbound_controls: Value,
3596 inbound_controls_error: Value,
3597) -> Value {
3598 let envelope = match crate::mailbox::Envelope::new(
3599 sender.clone(),
3600 format!("{}@{}", sender.session_id, sender.machine),
3601 crate::mailbox::MailKind::User,
3602 crate::mailbox::ReplyVia::None,
3603 params.text.clone(),
3604 ) {
3605 Ok(envelope) => envelope,
3606 Err(error) => {
3607 return json!({
3608 "delivered_to_bus": false,
3609 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
3610 "inbound_controls": inbound_controls,
3611 "inbound_controls_error": inbound_controls_error,
3612 })
3613 }
3614 };
3615 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
3616 Ok(turn) => json!({
3617 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
3618 "message_id": envelope.id,
3619 "target": {
3620 "harness": params.locator.harness.as_str(),
3621 "session_id": params.locator.session_id,
3622 },
3623 "delivery": {
3624 "door": match turn {
3625 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
3626 _ => "pane",
3627 },
3628 "how": turn.as_str(),
3629 },
3630 "inbound_controls": inbound_controls,
3631 "inbound_controls_error": inbound_controls_error,
3632 }),
3633 Err(message) => json!({
3634 "delivered_to_bus": false,
3635 "refusal": {"reason": "no_user_door", "message": message},
3636 "inbound_controls": inbound_controls,
3637 "inbound_controls_error": inbound_controls_error,
3638 }),
3639 }
3640}
3641
3642#[derive(Deserialize)]
3643#[serde(deny_unknown_fields)]
3644struct ActivityUnderParams {
3645 pids: Vec<u32>,
3647 #[serde(default)]
3648 homes: crate::HarnessHomes,
3649}
3650
3651async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
3655 let params = decode::<ActivityUnderParams>(params)?;
3656 if params.pids.len() > 1024 {
3657 return Err(ServiceError::InvalidParams(
3658 "sessions.activity_under accepts at most 1024 pids".into(),
3659 ));
3660 }
3661 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
3662 .await
3663 .map_err(ServiceError::Sdk)?;
3664 Ok(json!({
3665 "activities": found
3666 .into_iter()
3667 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
3668 .collect::<Vec<_>>(),
3669 }))
3670}
3671
3672fn operator_address(
3675 name: Option<&str>,
3676) -> std::result::Result<crate::mailbox::MailAddress, String> {
3677 let name: String = name
3678 .unwrap_or("supercode")
3679 .trim()
3680 .chars()
3681 .map(|character| {
3682 if character.is_whitespace() || character == '@' {
3683 '-'
3684 } else {
3685 character
3686 }
3687 })
3688 .collect();
3689 if name.is_empty() {
3690 return Err("from_name must not be empty".into());
3691 }
3692 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
3693 .map_err(|error| error.to_string())
3694}
3695
3696fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
3700 let address = match (¶ms.address, ¶ms.from_name) {
3701 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
3702 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
3703 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
3704 (Some(_), Some(_)) => {
3705 return Err(ServiceError::InvalidParams(
3706 "sessions.inbox takes from_name or address, not both".into(),
3707 ))
3708 }
3709 };
3710 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
3711 let mailbox =
3712 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
3713 let claimed = mailbox.claim_unread().map_err(operation)?;
3714 let mut messages: Vec<Value> = Vec::new();
3715 if params.all {
3716 for stored in mailbox.list().map_err(operation)? {
3717 if stored.state == crate::mailbox::MailState::Read {
3718 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3719 }
3720 }
3721 }
3722 for stored in &claimed {
3723 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3724 }
3725 for stored in &claimed {
3726 mailbox.acknowledge(stored).map_err(operation)?;
3727 }
3728 Ok(json!({"address": address.to_string(), "messages": messages}))
3729}
3730
3731#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3736struct FollowedSource {
3737 harness: String,
3738 session_id: String,
3739 reported: Option<String>,
3740}
3741
3742#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3743struct ActivitySubscription {
3744 locators: Vec<SessionLocator>,
3745 homes: crate::HarnessHomes,
3746 reported: BTreeMap<(String, String), crate::SessionActivity>,
3747}
3748
3749fn live_descriptor_value(
3756 session: &SessionDescriptor,
3757 doors: &crate::mail_route::LiveSessions,
3758) -> std::result::Result<Value, ServiceError> {
3759 let mut value = serde_json::to_value(session)
3760 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3761 if value.get("title").is_none_or(Value::is_null) {
3763 if let Some(live) = doors.all().iter().find(|live| {
3764 live.address.harness == session.locator.harness.as_str()
3765 && live.address.session_id == session.locator.session_id
3766 }) {
3767 let name = live.name.split('@').next().unwrap_or(&live.name);
3768 if !name.is_empty() {
3769 value["title"] = json!(name);
3770 }
3771 }
3772 }
3773 if let Some(workspace) = &session.cwd {
3774 let source = LiveRuntimeSource {
3775 harness: session.locator.harness.as_str().to_string(),
3776 session_id: session.locator.session_id.clone(),
3777 workspace: workspace.clone(),
3778 };
3779 if let Some(endpoint) = discover_live_runtime(&source)
3780 .map_err(|error| ServiceError::Operation(error.to_string()))?
3781 {
3782 value["live_endpoint"] = json!(endpoint.as_str());
3783 }
3784 }
3785 if let Some(door) = doors.door(
3789 session.locator.harness.as_str(),
3790 &session.locator.session_id,
3791 ) {
3792 value["delivery"] = json!(door);
3793 }
3794 if let Some(live) = doors.all().iter().find(|live| {
3797 live.address.harness == session.locator.harness.as_str()
3798 && live.address.session_id == session.locator.session_id
3799 && live.status == "waiting"
3800 }) {
3801 if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
3802 value["pending_request"] = request;
3803 }
3804 }
3805 Ok(value)
3806}
3807
3808fn live_index_changes(
3809 changes: Vec<crate::session_index::SessionIndexChange>,
3810 homes: &HarnessHomes,
3811) -> std::result::Result<Vec<Value>, ServiceError> {
3812 use crate::session_index::SessionIndexChange;
3813 let doors = crate::mail_route::LiveSessions::read(homes);
3814 changes
3815 .into_iter()
3816 .map(|change| match change {
3817 SessionIndexChange::Added { descriptor } => Ok(json!({
3818 "kind": "added",
3819 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3820 })),
3821 SessionIndexChange::Updated { descriptor } => Ok(json!({
3822 "kind": "updated",
3823 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3824 })),
3825 SessionIndexChange::Removed { key } => Ok(json!({
3826 "kind": "removed",
3827 "key": key,
3828 })),
3829 })
3830 .collect()
3831}
3832
3833fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
3834 use crate::{SessionPresence, SessionTurnState};
3835 match (activity.presence, activity.turn) {
3836 (SessionPresence::Persisted, _) => None,
3837 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
3838 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
3839 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
3841 (SessionPresence::Running, SessionTurnState::Unknown)
3845 if activity.evidence.native_state.is_none() =>
3846 {
3847 None
3848 }
3849 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
3850 }
3851}
3852
3853#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
3854#[serde(rename_all = "kebab-case")]
3855enum TransferFormat {
3856 ClaudeCode,
3857 Codex,
3858 #[serde(rename = "opencode", alias = "open-code")]
3859 OpenCode,
3860 Pi,
3861 Grok,
3862 Gemini,
3863 Goose,
3864 Hermes,
3868}
3869
3870impl TransferFormat {
3871 fn id(self) -> &'static str {
3872 match self {
3873 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
3874 Self::Codex => HarnessId::CODEX,
3875 Self::OpenCode => HarnessId::OPENCODE,
3876 Self::Pi => HarnessId::PI,
3877 Self::Grok => HarnessId::GROK,
3878 Self::Gemini => HarnessId::GEMINI,
3879 Self::Goose => HarnessId::GOOSE,
3880 Self::Hermes => HarnessId::HERMES,
3881 }
3882 }
3883}
3884
3885impl From<TransferFormat> for SessionFormat {
3886 fn from(value: TransferFormat) -> Self {
3887 match value {
3888 TransferFormat::ClaudeCode => Self::ClaudeCode,
3889 TransferFormat::Codex => Self::Codex,
3890 TransferFormat::OpenCode => Self::OpenCode,
3891 TransferFormat::Pi => Self::Pi,
3892 TransferFormat::Grok => Self::Grok,
3893 TransferFormat::Gemini => Self::Gemini,
3894 TransferFormat::Goose => Self::Goose,
3895 TransferFormat::Hermes => Self::Codex,
3897 }
3898 }
3899}
3900
3901#[derive(Deserialize)]
3902struct ImportSessionParams {
3903 source_harness: TransferFormat,
3904 content: String,
3905}
3906
3907#[derive(Deserialize)]
3908struct ExportSessionParams {
3909 locator: SessionLocator,
3910 target_harness: TransferFormat,
3911}
3912
3913#[derive(Deserialize)]
3914struct ReduceSessionParams {
3915 locator: SessionLocator,
3916 target_harness: TransferFormat,
3917 #[serde(default = "default_keep_last")]
3918 keep_last: usize,
3919}
3920
3921fn default_keep_last() -> usize {
3922 6
3923}
3924
3925#[derive(Deserialize)]
3926struct BranchSessionParams {
3927 locator: SessionLocator,
3928 #[serde(default)]
3929 target_harness: Option<TransferFormat>,
3930}
3931
3932#[derive(Deserialize)]
3933struct HandoffSessionParams {
3934 locator: SessionLocator,
3935 target_harness: TransferFormat,
3936 #[serde(default)]
3937 cwd: Option<PathBuf>,
3938}
3939
3940#[derive(Deserialize)]
3941struct MaterializeSessionParams {
3942 artifact: crate::native_materialize::MaterializeArtifact,
3943 cwd: PathBuf,
3944 #[serde(default)]
3946 homes: HarnessHomes,
3947}
3948
3949#[derive(Debug, Clone, Copy, Default, Deserialize)]
3952#[serde(rename_all = "snake_case")]
3953enum ResumePolicy {
3954 Default,
3955 #[default]
3956 Yolo,
3957}
3958
3959#[derive(Deserialize)]
3960struct ResumeInstructionsParams {
3961 locator: SessionLocator,
3962 #[serde(default)]
3963 cwd: Option<PathBuf>,
3964 #[serde(default)]
3965 policy: ResumePolicy,
3966}
3967
3968#[derive(Deserialize)]
3970struct WorkflowLoadParams {
3971 from: crate::workflow_doors::WorkflowHarness,
3972 home: PathBuf,
3973}
3974
3975#[derive(Deserialize)]
3978struct OrchestrationLoadParams {
3979 root: PathBuf,
3980 #[serde(default)]
3981 flavor: crate::orchestration_doors::HomeFlavor,
3982}
3983
3984#[derive(Deserialize)]
3987struct OrchestrationSaveParams {
3988 root: PathBuf,
3989 orchestration: crate::orchestration::Orchestration,
3990 #[serde(default)]
3991 vault: BTreeMap<String, String>,
3992}
3993
3994#[derive(Deserialize)]
3996struct OrchestrationCompileParams {
3997 from: crate::orchestration_doors::OrchestrationHarness,
3998 home: PathBuf,
3999}
4000
4001#[derive(Deserialize)]
4005struct OrchestrationDecompileParams {
4006 to: crate::orchestration_doors::OrchestrationHarness,
4007 orchestration: crate::orchestration::Orchestration,
4008 source: PathBuf,
4009 #[serde(default)]
4010 source_flavor: crate::orchestration_doors::SourceFlavor,
4011 dest: PathBuf,
4012 #[serde(default)]
4013 vault: BTreeMap<String, String>,
4014}
4015
4016#[derive(Deserialize)]
4019struct OrchestrationImportParams {
4020 from: crate::orchestration_doors::OrchestrationHarness,
4021 home: PathBuf,
4022 into: PathBuf,
4023}
4024
4025#[derive(Deserialize)]
4028struct OrchestrationExportParams {
4029 to: crate::orchestration_doors::OrchestrationHarness,
4030 root: PathBuf,
4031 dest: PathBuf,
4032}
4033
4034#[derive(Deserialize)]
4036struct JobsGetParams {
4037 harness: String,
4038 id: String,
4039 #[serde(default)]
4040 homes: crate::HarnessHomes,
4041}
4042
4043fn mutate_job(
4051 verb: crate::jobs_control::JobVerb,
4052 params: Value,
4053) -> std::result::Result<Value, ServiceError> {
4054 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4055 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4056 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4057 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4058}
4059
4060fn mutate_skill(
4067 verb: crate::skills_control::SkillVerb,
4068 params: Value,
4069) -> std::result::Result<Value, ServiceError> {
4070 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4071 if !crate::skills_control::supports_skill_control(&mutation.harness) {
4072 return Err(ServiceError::UnsupportedAction(format!(
4073 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4074 mutation.harness,
4075 verb.as_str(),
4076 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4077 )));
4078 }
4079 let outcome =
4080 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4081 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4082}
4083
4084fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4086 match error {
4087 crate::skills_control::SkillControlError::Unsupported(message) => {
4088 ServiceError::UnsupportedAction(message)
4089 }
4090 crate::skills_control::SkillControlError::Invalid(message) => {
4091 ServiceError::InvalidParams(message)
4092 }
4093 crate::skills_control::SkillControlError::Failed(message) => {
4094 ServiceError::Operation(message)
4095 }
4096 }
4097}
4098
4099fn mutate_profile(
4107 verb: crate::profiles_control::ProfileVerb,
4108 params: Value,
4109) -> std::result::Result<Value, ServiceError> {
4110 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4111 let outcome =
4112 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4113 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4114}
4115
4116fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4118 match error {
4119 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4120 ServiceError::UnsupportedAction(message)
4121 }
4122 crate::profiles_control::ProfileControlError::Invalid(message) => {
4123 ServiceError::InvalidParams(message)
4124 }
4125 crate::profiles_control::ProfileControlError::Failed(message) => {
4126 ServiceError::Operation(message)
4127 }
4128 }
4129}
4130
4131fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4135 match error {
4136 crate::jobs_control::JobControlError::Unsupported(message) => {
4137 ServiceError::UnsupportedAction(message)
4138 }
4139 crate::jobs_control::JobControlError::Invalid(message) => {
4140 ServiceError::InvalidParams(message)
4141 }
4142 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4143 }
4144}
4145
4146fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4151 match error {
4152 crate::SessionControlError::Unsupported(message) => {
4153 ServiceError::UnsupportedAction(message)
4154 }
4155 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4156 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4157 }
4158}
4159
4160fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4165 if crate::jobs::supports_jobs(harness) {
4166 return Ok(());
4167 }
4168 Err(ServiceError::UnsupportedAction(format!(
4169 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4170 crate::jobs::JOB_HARNESSES.join(", ")
4171 )))
4172}
4173
4174#[derive(Deserialize)]
4176struct RunsGetParams {
4177 harness: String,
4178 id: String,
4179 #[serde(default)]
4180 homes: crate::HarnessHomes,
4181}
4182
4183fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4188 if crate::runs::supports_runs(harness) {
4189 return Ok(());
4190 }
4191 Err(ServiceError::UnsupportedAction(format!(
4192 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4193 crate::runs::RUN_HARNESSES.join(", ")
4194 )))
4195}
4196
4197#[derive(Serialize)]
4198struct SessionArtifact {
4199 source_harness: HarnessId,
4200 target_harness: &'static str,
4201 session_id: Option<String>,
4202 content: String,
4203 suggested_filename: String,
4204 files: Vec<SessionArtifactFile>,
4205 fidelity: Fidelity,
4206 residue: Vec<String>,
4207}
4208
4209#[derive(Serialize)]
4210struct SessionArtifactFile {
4211 path: String,
4212 content: String,
4213 role: ArtifactFileRole,
4214}
4215
4216#[derive(Serialize)]
4217#[serde(rename_all = "snake_case")]
4218enum ArtifactFileRole {
4219 Primary,
4220 Subagent,
4221 Bundle,
4222 SourceRecovery,
4223}
4224
4225#[derive(Serialize)]
4226struct StructuredLaunch {
4227 cwd: PathBuf,
4228 program: String,
4229 arguments: Vec<String>,
4230 env: BTreeMap<String, String>,
4231}
4232
4233struct HandoffInstructions {
4234 launch: StructuredLaunch,
4235 materialize: Option<StructuredLaunch>,
4236 requires_materialization: bool,
4237 note: String,
4238}
4239
4240#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4241#[serde(rename_all = "snake_case")]
4242enum HarnessProbeLevel {
4243 #[default]
4244 Passive,
4245 Handshake,
4246}
4247
4248#[derive(Default, Deserialize)]
4249#[serde(default)]
4250struct HarnessInventoryParams {
4251 harness: Option<HarnessId>,
4252 harnesses: Vec<HarnessId>,
4253 workspace: Option<PathBuf>,
4254 probe: HarnessProbeLevel,
4255 include_sessions: bool,
4256 skip_versions: bool,
4258}
4259
4260#[derive(Deserialize)]
4261struct HarnessAuthenticationParams {
4262 harness: HarnessId,
4263}
4264
4265#[derive(Deserialize)]
4266struct BeginHarnessAuthenticationParams {
4267 harness: HarnessId,
4268 #[serde(default = "local_browser_authentication_environment")]
4269 environment: crate::HarnessAuthenticationEnvironment,
4270 #[serde(default)]
4271 method: Option<crate::HarnessAuthenticationMethodId>,
4272 #[serde(default)]
4273 cwd: Option<PathBuf>,
4274}
4275
4276fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4277 crate::HarnessAuthenticationEnvironment::LocalBrowser
4278}
4279
4280#[derive(Serialize)]
4281struct HarnessInventoryReport {
4282 probe: HarnessProbeLevel,
4283 workspace: Option<PathBuf>,
4284 harnesses: Vec<LocalHarness>,
4285}
4286
4287#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4288#[serde(rename_all = "snake_case")]
4289enum HarnessAuthState {
4290 Ready,
4291 Configured,
4292 Required,
4293 Unknown,
4294}
4295
4296#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4297#[serde(rename_all = "snake_case")]
4298enum HarnessRuntimeState {
4299 Ready,
4300 Degraded,
4301 Unavailable,
4302}
4303
4304#[derive(Serialize)]
4305struct HarnessSessionCounts {
4306 global: Option<usize>,
4307 workspace: Option<usize>,
4308}
4309
4310#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4321#[serde(rename_all = "snake_case")]
4322pub enum GatewayState {
4323 Up,
4324 Down,
4325 Unknown,
4326}
4327
4328#[derive(Debug, Clone, Serialize)]
4330pub struct GatewayHealth {
4331 pub state: GatewayState,
4332 #[serde(skip_serializing_if = "Option::is_none")]
4337 pub endpoint: Option<String>,
4338 #[serde(skip_serializing_if = "Option::is_none")]
4339 pub version: Option<String>,
4340 pub evidence: String,
4342 pub checked_at_ms: u64,
4343}
4344
4345fn openclaw_gateway_endpoint(home: &Path) -> String {
4349 let config_path = home.join(".openclaw/openclaw.json");
4350 let gateway = std::fs::read_to_string(&config_path)
4351 .ok()
4352 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4353 .and_then(|config| config.get("gateway").cloned());
4354 if let Some(url) = gateway
4355 .as_ref()
4356 .and_then(|gateway| gateway.get("url"))
4357 .and_then(serde_json::Value::as_str)
4358 {
4359 return url.to_string();
4360 }
4361 let port = gateway
4362 .as_ref()
4363 .and_then(|gateway| gateway.get("port"))
4364 .and_then(serde_json::Value::as_u64)
4365 .unwrap_or(18789);
4366 format!("ws://127.0.0.1:{port}")
4367}
4368
4369fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4376 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4377 let output = std::process::Command::new(&program)
4378 .args(["gateway", "status"])
4379 .stdin(std::process::Stdio::null())
4380 .output()
4381 .ok()?;
4382 let text = format!(
4383 "{}{}",
4384 String::from_utf8_lossy(&output.stdout),
4385 String::from_utf8_lossy(&output.stderr)
4386 );
4387 let verdict = text.lines().find_map(|line| {
4388 let l = line.trim();
4389 if l.contains("supervised by launchd (PID")
4390 || l.contains("supervised by systemd (PID")
4391 || l.contains("Gateway is running")
4392 || l.contains("process is running")
4393 {
4394 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4395 } else if l.contains("not running") || l.contains("not installed") {
4396 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4397 } else {
4398 None
4399 }
4400 });
4401 verdict
4402}
4403
4404fn gateway_health(
4405 id: &str,
4406 installed: bool,
4407 running: Option<&RunningInstance>,
4408 version: Option<&str>,
4409) -> GatewayHealth {
4410 let checked_at_ms = now_epoch_ms();
4411 let home = supercode_interchange::user_home()
4412 .map(std::path::PathBuf::into_os_string)
4413 .map(PathBuf::from);
4414 match id {
4415 HarnessId::HERMES | HarnessId::OPENCLAW => {
4416 let endpoint = (id == HarnessId::OPENCLAW)
4417 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4418 .flatten();
4419 let (state, evidence) = match running {
4420 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4421 None if !installed => (
4422 GatewayState::Unknown,
4423 format!("`{id}` is not installed; no gateway to probe"),
4424 ),
4425 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4426 Some((state, evidence)) => (state, evidence),
4429 None => (
4430 GatewayState::Down,
4431 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4432 ),
4433 },
4434 None => (
4435 GatewayState::Down,
4436 format!(
4437 "no TCP listener at {}",
4438 endpoint.as_deref().unwrap_or("the gateway endpoint")
4439 ),
4440 ),
4441 };
4442 GatewayHealth {
4443 state,
4444 endpoint,
4445 version: version.map(str::to_string),
4446 evidence,
4447 checked_at_ms,
4448 }
4449 }
4450 HarnessId::ORCHESTRATOR => {
4457 let root = crate::HarnessHomes::default().orchestrator;
4458 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4459 Some(lease) if lease.is_live() => (
4460 GatewayState::Up,
4461 format!(
4462 "`{}` names pid {} (started {}), which is live",
4463 crate::orchestrator::lock_path(&root).display(),
4464 lease.pid,
4465 lease.started_at
4466 ),
4467 ),
4468 Some(lease) => (
4469 GatewayState::Down,
4470 format!(
4471 "stale lease `{}`: pid {} is gone",
4472 crate::orchestrator::lock_path(&root).display(),
4473 lease.pid
4474 ),
4475 ),
4476 None => (
4477 GatewayState::Down,
4478 format!(
4479 "no lease at `{}`; `supercode orchestrator start` writes one",
4480 crate::orchestrator::lock_path(&root).display()
4481 ),
4482 ),
4483 };
4484 GatewayHealth {
4485 state,
4486 endpoint: None,
4487 version: version.map(str::to_string),
4488 evidence,
4489 checked_at_ms,
4490 }
4491 }
4492 _ => GatewayHealth {
4493 state: GatewayState::Unknown,
4494 endpoint: None,
4495 version: version.map(str::to_string),
4496 evidence: format!("`{id}` runs per session, not as a gateway"),
4497 checked_at_ms,
4498 },
4499 }
4500}
4501
4502#[derive(Debug, Clone, Serialize)]
4503struct RunningInstance {
4504 method: RunningInstanceMethod,
4506 evidence: String,
4508 checked_at_ms: u64,
4510}
4511
4512#[derive(Debug, Clone, Copy, Serialize)]
4513#[serde(rename_all = "snake_case")]
4514enum RunningInstanceMethod {
4515 GatewayConnect,
4518 StoreWalActivity,
4521}
4522
4523fn now_epoch_ms() -> u64 {
4524 std::time::SystemTime::now()
4525 .duration_since(std::time::UNIX_EPOCH)
4526 .map(|elapsed| elapsed.as_millis() as u64)
4527 .unwrap_or(0)
4528}
4529
4530fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4534 let config_path = home.join(".openclaw/openclaw.json");
4535 let text = std::fs::read_to_string(&config_path).ok();
4536 let gateway = text
4537 .as_deref()
4538 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4539 .and_then(|config| config.get("gateway").cloned());
4540 let address = gateway
4541 .as_ref()
4542 .and_then(|gateway| gateway.get("url"))
4543 .and_then(serde_json::Value::as_str)
4544 .and_then(|url| {
4545 url.split("://").nth(1).map(|rest| {
4546 rest.trim_end_matches('/')
4547 .split('/')
4548 .next()
4549 .unwrap_or(rest)
4550 .to_string()
4551 })
4552 })
4553 .unwrap_or_else(|| {
4554 let port = gateway
4555 .as_ref()
4556 .and_then(|gateway| gateway.get("port"))
4557 .and_then(serde_json::Value::as_u64)
4558 .unwrap_or(18789);
4559 format!("127.0.0.1:{port}")
4560 });
4561 let reachable = std::net::TcpStream::connect_timeout(
4562 &address.parse().ok()?,
4563 std::time::Duration::from_millis(400),
4564 )
4565 .is_ok();
4566 reachable.then(|| RunningInstance {
4567 method: RunningInstanceMethod::GatewayConnect,
4568 evidence: format!(
4569 "gateway endpoint {address} accepted a TCP connect (from {})",
4570 config_path.display()
4571 ),
4572 checked_at_ms: now_epoch_ms(),
4573 })
4574}
4575
4576fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4581 let wal = home.join(".hermes/state.db-wal");
4582 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4583 let age_ms = std::time::SystemTime::now()
4584 .duration_since(modified)
4585 .map(|age| age.as_millis() as u64)
4586 .unwrap_or(u64::MAX);
4587 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4588 method: RunningInstanceMethod::StoreWalActivity,
4589 evidence: format!(
4590 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4591 wal.display()
4592 ),
4593 checked_at_ms: now_epoch_ms(),
4594 })
4595}
4596
4597fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4599 let home = supercode_interchange::user_home()
4600 .map(std::path::PathBuf::into_os_string)
4601 .map(PathBuf::from)?;
4602 match id {
4603 HarnessId::OPENCLAW => probe_openclaw_running(&home),
4604 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4605 _ => None,
4606 }
4607}
4608
4609#[derive(Serialize)]
4610struct LocalHarness {
4611 id: HarnessId,
4612 display_name: String,
4613 supported: bool,
4614 installed: bool,
4615 executable: Option<String>,
4616 version: Option<String>,
4617 auth: HarnessAuthState,
4618 runtime: HarnessRuntimeState,
4619 protocol: String,
4620 capabilities: crate::RuntimeCapabilities,
4621 effective_capabilities: crate::RuntimeCapabilities,
4622 sessions: HarnessSessionCounts,
4623 #[serde(skip_serializing_if = "Option::is_none")]
4626 running: Option<RunningInstance>,
4627 gateway: GatewayHealth,
4629 reason: Option<String>,
4630 repair: Option<String>,
4631}
4632
4633#[derive(Clone, Deserialize)]
4634struct RuntimeBackendParams {
4635 harness: HarnessId,
4636 #[serde(default)]
4637 protocol: Option<String>,
4638 #[serde(default)]
4639 launch: Option<RuntimeLaunch>,
4640 #[serde(default)]
4641 base_url: Option<String>,
4642 #[serde(default)]
4643 policy: RuntimePolicy,
4644}
4645
4646#[derive(Debug, Clone, Copy, Default, Deserialize)]
4649#[serde(rename_all = "snake_case")]
4650enum RuntimePolicy {
4651 Default,
4652 #[default]
4653 Yolo,
4654}
4655
4656#[derive(Deserialize)]
4657struct RuntimeStartParams {
4658 #[serde(flatten)]
4659 backend: RuntimeBackendParams,
4660 cwd: PathBuf,
4661 #[serde(default)]
4664 mcp_servers: Vec<crate::McpServerLaunch>,
4665 #[serde(default)]
4667 approval_policy: Option<String>,
4668}
4669
4670#[derive(Deserialize)]
4671struct RuntimeAttachParams {
4672 #[serde(flatten)]
4673 backend: RuntimeBackendParams,
4674 runtime_id: String,
4675 #[serde(default)]
4676 cwd: Option<PathBuf>,
4677 #[serde(default)]
4680 mcp_servers: Vec<crate::McpServerLaunch>,
4681 #[serde(default)]
4683 approval_policy: Option<String>,
4684}
4685
4686#[derive(Deserialize)]
4687struct RuntimeConnectionParams {
4688 connection: String,
4689}
4690
4691#[derive(Deserialize)]
4692struct RuntimeInputParams {
4693 connection: String,
4694 text: String,
4695 #[serde(default)]
4696 image_urls: Vec<String>,
4697}
4698
4699const MAX_RUNTIME_IMAGES: usize = 4;
4700const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
4701const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
4702
4703fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
4704 if image_urls.len() > MAX_RUNTIME_IMAGES {
4705 return Err(ServiceError::InvalidParams(format!(
4706 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
4707 )));
4708 }
4709 let mut total = 0usize;
4710 for url in &image_urls {
4711 if !(url.starts_with("data:image/")
4712 || url.starts_with("https://")
4713 || url.starts_with("http://"))
4714 {
4715 return Err(ServiceError::InvalidParams(
4716 "runtime images must be image data URLs or HTTP(S) URLs".into(),
4717 ));
4718 }
4719 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
4720 return Err(ServiceError::InvalidParams(format!(
4721 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
4722 )));
4723 }
4724 total = total.saturating_add(url.len());
4725 }
4726 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
4727 return Err(ServiceError::InvalidParams(format!(
4728 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
4729 )));
4730 }
4731 Ok(image_urls)
4732}
4733
4734#[derive(Deserialize)]
4735struct RuntimeRespondParams {
4736 connection: String,
4737 request_id: Value,
4738 response: Value,
4739}
4740
4741fn default_reduction_store_root() -> PathBuf {
4742 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
4743 return PathBuf::from(root).join("sessions");
4744 }
4745 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
4746 return PathBuf::from(home).join(".supercode").join("sessions");
4747 }
4748 PathBuf::from(".supercode").join("sessions")
4749}
4750
4751fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
4752 let mut output = String::new();
4753 for message in messages {
4754 output.push_str(
4755 &serde_json::to_string(message)
4756 .map_err(|error| ServiceError::Operation(error.to_string()))?,
4757 );
4758 output.push('\n');
4759 }
4760 Ok(output)
4761}
4762
4763fn parse_messages_jsonl(
4764 content: &str,
4765) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
4766 content
4767 .lines()
4768 .enumerate()
4769 .filter(|(_, line)| !line.trim().is_empty())
4770 .map(|(index, line)| {
4771 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
4772 ServiceError::Operation(format!(
4773 "reduced transcript line {} is invalid: {error}",
4774 index + 1
4775 ))
4776 })
4777 })
4778 .collect()
4779}
4780
4781fn reduced_bootstrap_prompt(
4782 source: &SessionLocator,
4783 target: TransferFormat,
4784 view_jsonl: &str,
4785 sidecar_path: &Path,
4786 reduction_log_path: &Path,
4787) -> String {
4788 format!(
4789 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
4790 \n\
4791 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\
4792 \n\
4793 <supercode-reduced-session source-session=\"{source_id}\">\n\
4794 {view_jsonl}\
4795 </supercode-reduced-session>\n\
4796 \n\
4797 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
4798 source_harness = source.harness.as_str(),
4799 target_harness = target.id(),
4800 sidecar = sidecar_path.display(),
4801 log = reduction_log_path.display(),
4802 source_id = source.session_id,
4803 )
4804}
4805
4806fn session_artifact(
4807 locator: &SessionLocator,
4808 session: &Session,
4809 target: TransferFormat,
4810) -> std::result::Result<SessionArtifact, ServiceError> {
4811 session_artifact_with_id(locator, session, target, None)
4812}
4813
4814fn session_artifact_with_id(
4815 locator: &SessionLocator,
4816 session: &Session,
4817 target: TransferFormat,
4818 target_session_id: Option<&str>,
4819) -> std::result::Result<SessionArtifact, ServiceError> {
4820 let format: SessionFormat = target.into();
4821 let diagonal = format.source() == session.meta.source;
4822 crate::residue_store::store_segments(session);
4823 let has_appended_turns = session
4824 .imported_message_count
4825 .is_some_and(|imported| imported < session.messages.len());
4826 let mut restoration = None;
4827 let content = if let Some(id) = target_session_id {
4828 if diagonal && format != SessionFormat::OpenCode {
4829 session
4830 .to_jsonl_spliced(format, Some(id))
4831 .map_err(operation)?
4832 } else {
4833 let mut rewritten = session.clone();
4834 rewritten.meta.session_id = Some(id.to_string());
4835 rewritten.to_jsonl(format).map_err(operation)?
4836 }
4837 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
4838 session.raw_verbatim()
4839 } else if diagonal {
4840 session.to_jsonl_spliced(format, None).map_err(operation)?
4841 } else {
4842 match session
4845 .restore_residue(format, crate::residue_store::lookup)
4846 .map_err(operation)?
4847 {
4848 Some((content, report)) => {
4849 restoration = Some(report);
4850 content
4851 }
4852 None => session.to_jsonl(format).map_err(operation)?,
4853 }
4854 };
4855 let stem = sanitize_filename(
4856 target_session_id
4857 .or(session.meta.session_id.as_deref())
4858 .unwrap_or(&locator.session_id),
4859 );
4860 let suggested_filename = if diagonal && target == TransferFormat::Grok {
4861 "chat_history.jsonl".to_string()
4862 } else if target == TransferFormat::Goose {
4863 format!("{stem}.goose.json")
4864 } else {
4865 format!("{stem}.{}.jsonl", target.id())
4866 };
4867 let mut files = vec![SessionArtifactFile {
4868 path: suggested_filename.clone(),
4869 content: content.clone(),
4870 role: ArtifactFileRole::Primary,
4871 }];
4872 if target == TransferFormat::ClaudeCode {
4873 let bundle_stem = Path::new(&suggested_filename)
4874 .file_stem()
4875 .and_then(|stem| stem.to_str())
4876 .unwrap_or(&stem);
4877 let mut child_paths = BTreeSet::new();
4878 for (index, subagent) in session.subagents.iter().enumerate() {
4879 let agent_id = subagent
4880 .meta
4881 .agent_id
4882 .as_deref()
4883 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
4884 .map(sanitize_filename)
4885 .filter(|id| !id.is_empty())
4886 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4887 let child_has_appended_turns = subagent
4888 .imported_message_count
4889 .is_some_and(|imported| imported < subagent.messages.len());
4890 let child_content = if target_session_id.is_none()
4891 && subagent.meta.source == SessionSource::ClaudeCode
4892 && subagent.raw_is_verbatim
4893 && !child_has_appended_turns
4894 {
4895 subagent.raw_verbatim()
4896 } else if subagent.meta.source == SessionSource::ClaudeCode {
4897 subagent
4898 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
4899 .map_err(operation)?
4900 } else {
4901 let mut child = subagent.clone();
4902 if let Some(id) = target_session_id {
4903 child.meta.session_id = Some(id.to_string());
4904 }
4905 child
4906 .to_jsonl(SessionFormat::ClaudeCode)
4907 .map_err(operation)?
4908 };
4909 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
4910 if !child_paths.insert(path.clone()) {
4911 return Err(ServiceError::Operation(format!(
4912 "Claude subagent ids collide at artifact path `{path}`"
4913 )));
4914 }
4915 files.push(SessionArtifactFile {
4916 path,
4917 content: child_content,
4918 role: ArtifactFileRole::Subagent,
4919 });
4920 }
4921 }
4922 if diagonal && target == TransferFormat::Grok {
4923 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
4924 }
4925 if !diagonal || !session.raw_is_verbatim {
4926 files.push(SessionArtifactFile {
4927 path: "recovery/source.supercode.jsonl".into(),
4928 content: session.to_native_jsonl(),
4929 role: ArtifactFileRole::SourceRecovery,
4930 });
4931 for (index, subagent) in session.subagents.iter().enumerate() {
4932 let id = subagent
4933 .meta
4934 .agent_id
4935 .as_deref()
4936 .map(sanitize_filename)
4937 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4938 files.push(SessionArtifactFile {
4939 path: format!("recovery/subagents/{id}.supercode.jsonl"),
4940 content: subagent.to_native_jsonl(),
4941 role: ArtifactFileRole::SourceRecovery,
4942 });
4943 }
4944 }
4945 if !diagonal && session.meta.source == SessionSource::Grok {
4946 append_grok_bundle_files(
4947 locator,
4948 "recovery/grok/",
4949 ArtifactFileRole::SourceRecovery,
4950 &mut files,
4951 )?;
4952 }
4953 let (fidelity, residue) = if diagonal
4954 && target_session_id.is_none()
4955 && session.raw_is_verbatim
4956 && !has_appended_turns
4957 {
4958 (Fidelity::ByteLossless, Vec::new())
4959 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
4960 (
4961 Fidelity::ValueLossless,
4962 vec![if target_session_id.is_some() {
4963 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
4964 } else {
4965 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
4966 }],
4967 )
4968 } else {
4969 match restoration {
4970 Some(report) if report.rendered_messages == 0 => (
4971 Fidelity::ByteLossless,
4972 vec![format!(
4973 "restored verbatim from this conversation's {} source records in the residue store",
4974 target.id()
4975 )],
4976 ),
4977 Some(report) => (
4978 Fidelity::Semantic,
4979 vec![format!(
4980 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
4981 report.restored_messages,
4982 report.restored_messages + report.rendered_messages,
4983 report.rendered_messages,
4984 target.id()
4985 )],
4986 ),
4987 None => (
4988 Fidelity::Semantic,
4989 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
4990 ),
4991 }
4992 };
4993 Ok(SessionArtifact {
4994 source_harness: locator.harness.clone(),
4995 target_harness: target.id(),
4996 session_id: target_session_id
4997 .map(str::to_string)
4998 .or_else(|| session.meta.session_id.clone()),
4999 content,
5000 suggested_filename,
5001 files,
5002 fidelity,
5003 residue,
5004 })
5005}
5006
5007fn append_grok_bundle_files(
5008 locator: &SessionLocator,
5009 prefix: &str,
5010 role: ArtifactFileRole,
5011 files: &mut Vec<SessionArtifactFile>,
5012) -> std::result::Result<(), ServiceError> {
5013 let primary = locator.storage.path();
5014 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5015 return Err(ServiceError::Operation(format!(
5016 "Grok bundle locator must name chat_history.jsonl, got {}",
5017 primary.display()
5018 )));
5019 }
5020 let parent = primary.parent().ok_or_else(|| {
5021 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5022 })?;
5023 for name in ["summary.json", "updates.jsonl"] {
5024 let path = parent.join(name);
5025 let metadata = match std::fs::symlink_metadata(&path) {
5026 Ok(metadata) => metadata,
5027 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5028 Err(error) => return Err(ServiceError::Operation(error.to_string())),
5029 };
5030 if metadata.file_type().is_symlink() || !metadata.is_file() {
5031 return Err(ServiceError::Operation(format!(
5032 "refusing non-regular Grok bundle member {}",
5033 path.display()
5034 )));
5035 }
5036 let content = std::fs::read_to_string(&path).map_err(|error| {
5037 ServiceError::Operation(format!(
5038 "Grok bundle member {} is not representable as UTF-8: {error}",
5039 path.display()
5040 ))
5041 })?;
5042 files.push(SessionArtifactFile {
5043 path: format!("{prefix}{name}"),
5044 content,
5045 role: match role {
5046 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5047 _ => ArtifactFileRole::SourceRecovery,
5048 },
5049 });
5050 }
5051 Ok(())
5052}
5053
5054fn handoff_artifact(
5055 locator: &SessionLocator,
5056 session: &Session,
5057 target: TransferFormat,
5058) -> std::result::Result<SessionArtifact, ServiceError> {
5059 let target_session_id = target_session_id(target);
5060 session_artifact_with_id(locator, session, target, Some(&target_session_id))
5061}
5062
5063fn target_session_id(target: TransferFormat) -> String {
5064 let uuid = generated_session_id();
5065 match target {
5066 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5067 TransferFormat::ClaudeCode
5068 | TransferFormat::Codex
5069 | TransferFormat::Pi
5070 | TransferFormat::Grok
5071 | TransferFormat::Gemini
5072 | TransferFormat::Goose
5073 | TransferFormat::Hermes => uuid,
5074 }
5075}
5076
5077fn sanitize_filename(value: &str) -> String {
5078 let value = value
5079 .chars()
5080 .map(|character| {
5081 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5082 character
5083 } else {
5084 '-'
5085 }
5086 })
5087 .collect::<String>();
5088 let value = value.trim_matches('-');
5089 if value.is_empty() {
5090 "session".into()
5091 } else {
5092 value.chars().take(100).collect()
5093 }
5094}
5095
5096fn handoff_instructions(
5097 target: TransferFormat,
5098 session_id: &str,
5099 cwd: &Path,
5100) -> HandoffInstructions {
5101 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5102 cwd: cwd.to_path_buf(),
5103 program: program.into(),
5104 arguments,
5105 env: BTreeMap::new(),
5106 };
5107 match target {
5108 TransferFormat::ClaudeCode => HandoffInstructions {
5109 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5110 materialize: None,
5111 requires_materialization: true,
5112 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(),
5113 },
5114 TransferFormat::Hermes => HandoffInstructions {
5115 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5116 materialize: None,
5117 requires_materialization: true,
5118 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(),
5119 },
5120 TransferFormat::Codex => HandoffInstructions {
5121 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5122 materialize: None,
5123 requires_materialization: true,
5124 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5125 },
5126 TransferFormat::OpenCode => HandoffInstructions {
5127 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5128 materialize: Some(launch(
5129 "opencode",
5130 vec!["import".into(), "{artifact_path}".into()],
5131 )),
5132 requires_materialization: true,
5133 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5134 },
5135 TransferFormat::Pi => HandoffInstructions {
5136 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5137 materialize: None,
5138 requires_materialization: true,
5139 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5140 },
5141 TransferFormat::Grok => HandoffInstructions {
5142 launch: launch(
5143 "grok",
5144 vec!["--resume".into(), "{materialized_session_id}".into()],
5145 ),
5146 materialize: None,
5147 requires_materialization: true,
5148 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(),
5149 },
5150 TransferFormat::Gemini => HandoffInstructions {
5151 launch: launch(
5152 "gemini",
5153 vec!["--session-file".into(), "{artifact_path}".into()],
5154 ),
5155 materialize: None,
5156 requires_materialization: true,
5157 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(),
5158 },
5159 TransferFormat::Goose => HandoffInstructions {
5160 launch: launch(
5161 "goose",
5162 vec![
5163 "session".into(),
5164 "--resume".into(),
5165 "--session-id".into(),
5166 "{imported_session_id}".into(),
5167 ],
5168 ),
5169 materialize: Some(launch(
5170 "goose",
5171 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5172 )),
5173 requires_materialization: true,
5174 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(),
5175 },
5176 }
5177}
5178
5179fn resume_launch(
5180 harness: &str,
5181 session_id: &str,
5182 cwd: &Path,
5183 policy: ResumePolicy,
5184) -> std::result::Result<StructuredLaunch, ServiceError> {
5185 let mut arguments = Vec::new();
5186 let program = match harness {
5187 HarnessId::GROK => {
5188 if matches!(policy, ResumePolicy::Yolo) {
5189 if crate::support::self_sandbox_supported() {
5190 arguments.extend(["--sandbox".into(), "workspace".into()]);
5191 }
5192 arguments.push("--always-approve".into());
5193 }
5194 arguments.extend(["--resume".into(), session_id.into()]);
5195 "grok"
5196 }
5197 HarnessId::CODEX => {
5198 arguments.extend(crate::startup_prompts::startup_arguments(
5199 harness,
5200 Some(cwd),
5201 &[],
5202 matches!(policy, ResumePolicy::Yolo),
5203 ));
5204 arguments.extend(["resume".into(), session_id.into()]);
5205 "codex"
5206 }
5207 HarnessId::CLAUDE_CODE => {
5208 arguments.extend(crate::startup_prompts::startup_arguments(
5209 harness,
5210 Some(cwd),
5211 &[],
5212 matches!(policy, ResumePolicy::Yolo),
5213 ));
5214 arguments.extend(["--resume".into(), session_id.into()]);
5215 "claude"
5216 }
5217 HarnessId::GEMINI => {
5218 arguments.extend(crate::startup_prompts::startup_arguments(
5219 harness,
5220 Some(cwd),
5221 &[],
5222 matches!(policy, ResumePolicy::Yolo),
5223 ));
5224 arguments.extend(["--resume".into(), session_id.into()]);
5225 "gemini"
5226 }
5227 HarnessId::GOOSE => {
5228 arguments.extend([
5229 "session".into(),
5230 "--resume".into(),
5231 "--session-id".into(),
5232 session_id.into(),
5233 ]);
5234 "goose"
5235 }
5236 HarnessId::PI => {
5237 arguments.extend(crate::startup_prompts::startup_arguments(
5238 harness,
5239 Some(cwd),
5240 &[],
5241 matches!(policy, ResumePolicy::Yolo),
5242 ));
5243 arguments.extend(["--session".into(), session_id.into()]);
5244 "pi"
5245 }
5246 HarnessId::OPENCODE => {
5247 arguments.extend(["--session".into(), session_id.into()]);
5248 "opencode"
5249 }
5250 HarnessId::SUPERCODE => {
5251 if matches!(policy, ResumePolicy::Yolo) {
5252 arguments.push("--dangerous".into());
5253 }
5254 arguments.extend(["resume".into(), session_id.into()]);
5255 "supercode"
5256 }
5257 other => {
5258 return Err(ServiceError::InvalidParams(format!(
5259 "no structured resume launch is registered for harness `{other}`"
5260 )))
5261 }
5262 };
5263 Ok(StructuredLaunch {
5264 cwd: cwd.to_path_buf(),
5265 env: if program == "grok" {
5266 crate::support::grok_home_env()
5267 } else {
5268 BTreeMap::new()
5269 },
5270 program: program.into(),
5271 arguments,
5272 })
5273}
5274
5275fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5281 let digest = blake3::hash(address.as_bytes()).to_hex();
5282 let path = std::env::temp_dir().join(format!(
5283 "supercode-openclaw-gateway-token-{}",
5284 &digest.as_str()[..16]
5285 ));
5286 #[cfg(unix)]
5287 {
5288 use std::io::Write;
5289 use std::os::unix::fs::OpenOptionsExt;
5290 let mut file = std::fs::OpenOptions::new()
5291 .write(true)
5292 .create(true)
5293 .truncate(true)
5294 .mode(0o600)
5295 .open(&path)?;
5296 file.write_all(secret.as_bytes())?;
5297 }
5298 #[cfg(not(unix))]
5299 std::fs::write(&path, secret)?;
5300 Ok(path)
5301}
5302
5303fn open_connect_descriptor(
5309 descriptor: &crate::HarnessSupportDescriptor,
5310 home: &Path,
5311) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5312 let Some(connect) = &descriptor.runtime.connect_launch else {
5313 return Err(ServiceError::InvalidParams(format!(
5314 "harness `{}` has no registered connect-mode launch",
5315 descriptor.id.as_str()
5316 )));
5317 };
5318 let resolved = connect
5319 .resolve(home)
5320 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5321 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5322 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5323 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5324 if let Some(token) = resolved.auth {
5325 backend = backend.with_bearer(token);
5326 }
5327 Ok(Box::new(backend))
5328 }
5329 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5330 let mut env = BTreeMap::new();
5341 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5342 if let Some(token) = resolved.auth {
5343 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5344 .map_err(|error| {
5345 ServiceError::UnsupportedAction(format!(
5346 "could not stage the gateway credential for the bridge: {error}"
5347 ))
5348 })?;
5349 arguments.push("--token-file".into());
5350 arguments.push(token_path.to_string_lossy().into_owned());
5351 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5352 }
5353 let program = descriptor
5358 .runtime
5359 .default_launch
5360 .as_ref()
5361 .map(|launch| launch.program.clone())
5362 .unwrap_or_else(|| "openclaw".into());
5363 let launch = RuntimeLaunch {
5364 program,
5365 arguments,
5366 env,
5367 };
5368 Ok(Box::new(
5369 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5370 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5371 ))
5372 }
5373 _ => Err(ServiceError::UnsupportedAction(format!(
5374 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5375 descriptor.id.as_str(),
5376 connect.protocol
5377 ))),
5378 }
5379}
5380
5381fn registry_connect_descriptor(
5384 params: &RuntimeBackendParams,
5385) -> Option<crate::HarnessSupportDescriptor> {
5386 if params.launch.is_some() || params.base_url.is_some() {
5387 return None;
5388 }
5389 harness_support_registry()
5390 .harnesses
5391 .into_iter()
5392 .find(|descriptor| descriptor.id == params.harness)
5393 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5394}
5395
5396fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5397 supercode_interchange::user_home()
5398 .map(std::path::PathBuf::into_os_string)
5399 .map(PathBuf::from)
5400 .ok_or_else(|| {
5401 ServiceError::UnsupportedAction(
5402 "connect-mode launches need HOME to locate the harness config".into(),
5403 )
5404 })
5405}
5406
5407pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5410 "harness.v1.runtimes.start",
5411 "harness.v1.runtimes.resume",
5412 "harness.v1.runtimes.attach",
5413 "harness.v1.runtimes.attach_existing",
5414];
5415
5416pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5421
5422pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5428
5429pub const DETACHED_METHODS: &[&str] = &[
5435 "harness.v1.harnesses.list",
5436 "harness.v1.harnesses.probe",
5437 "harness.v1.sessions.message",
5438 "harness.v1.sessions.new",
5439 "harness.v1.sessions.reset",
5440 "harness.v1.sessions.archive",
5441 "harness.v1.sessions.delete",
5442];
5443
5444pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5450
5451pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5459
5460async fn within_control_deadline<F: std::future::Future>(
5463 method: &str,
5464 call: F,
5465) -> std::result::Result<F::Output, ServiceError> {
5466 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5467 .await
5468 .map_err(|_| {
5469 ServiceError::Operation(format!(
5470 "`{method}` gave up after {}s: the runtime did not answer",
5471 RUNTIME_CONTROL_DEADLINE.as_secs()
5472 ))
5473 })
5474}
5475
5476pub struct RuntimeOpen {
5480 id: Value,
5481 method: String,
5482 params: Value,
5483}
5484
5485impl RuntimeOpen {
5486 pub async fn open(self) -> OpenedRuntime {
5490 let Self { id, method, params } = self;
5491 let outcome = open_runtime(&method, params).await;
5492 OpenedRuntime { id, outcome }
5493 }
5494}
5495
5496pub struct OpenedRuntime {
5499 id: Value,
5500 outcome: std::result::Result<OpenRuntime, ServiceError>,
5501}
5502
5503pub struct DetachedCall {
5508 id: Value,
5509 method: String,
5510 work: std::result::Result<Work, ServiceError>,
5511}
5512
5513impl DetachedCall {
5514 pub async fn run(self) -> DetachedAnswer {
5517 let Self { id, method, work } = self;
5518 match work {
5519 Ok(Work::Runtime(work)) => {
5524 let (result, returned) = work.run().await;
5525 DetachedAnswer {
5526 response: service_response(id, result),
5527 returned,
5528 }
5529 }
5530 Ok(Work::Free(work)) => {
5531 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5532 Ok(result) => result,
5533 Err(_) => Err(ServiceError::Operation(format!(
5534 "`{method}` gave up after {}s: the harness it waits on did not answer",
5535 DETACHED_CALL_DEADLINE.as_secs()
5536 ))),
5537 };
5538 DetachedAnswer {
5539 response: service_response(id, result),
5540 returned: None,
5541 }
5542 }
5543 Err(error) => DetachedAnswer {
5544 response: service_response(id, Err(error)),
5545 returned: None,
5546 },
5547 }
5548 }
5549}
5550
5551pub struct DetachedAnswer {
5555 response: Value,
5556 returned: Option<ReturnedRuntime>,
5557}
5558
5559impl DetachedAnswer {
5560 pub fn into_response(self) -> Value {
5563 self.response
5564 }
5565}
5566
5567pub struct ReturnedRuntime {
5570 connection: String,
5571 runtime: Box<dyn RuntimeConnection>,
5572}
5573
5574enum Work {
5577 Free(DetachedWork),
5578 Runtime(RuntimeWork),
5579}
5580
5581enum DetachedWork {
5584 Inventory(InventoryWork),
5588 Message(MessageSessionParams),
5590 SessionMutation {
5593 verb: crate::SessionVerb,
5594 mutation: crate::SessionMutation,
5595 },
5596}
5597
5598impl DetachedWork {
5599 async fn run(self) -> std::result::Result<Value, ServiceError> {
5600 match self {
5601 Self::Inventory(work) => run_inventory(work).await,
5602 Self::Message(params) => Ok(message_live_session(¶ms).await),
5603 Self::SessionMutation { verb, mutation } => {
5604 let outcome = run_session_mutation(verb, &mutation).await?;
5605 serde_json::to_value(outcome)
5606 .map_err(|error| ServiceError::Operation(error.to_string()))
5607 }
5608 }
5609 }
5610}
5611
5612enum RuntimeWork {
5614 Close {
5616 runtime: Box<dyn RuntimeConnection>,
5617 process_group: Option<u32>,
5618 },
5619 LiveCommand {
5622 connection: String,
5623 runtime: Box<dyn RuntimeConnection>,
5624 verb: crate::SessionVerb,
5625 mutation: crate::SessionMutation,
5626 command: &'static str,
5627 session: String,
5628 },
5629}
5630
5631type RuntimeWorkAnswer = (
5634 std::result::Result<Value, ServiceError>,
5635 Option<ReturnedRuntime>,
5636);
5637
5638impl RuntimeWork {
5639 async fn run(self) -> RuntimeWorkAnswer {
5640 match self {
5641 Self::Close {
5642 runtime,
5643 process_group,
5644 } => (close_runtime(runtime, process_group).await, None),
5645 Self::LiveCommand {
5646 connection,
5647 mut runtime,
5648 verb,
5649 mutation,
5650 command,
5651 session,
5652 } => {
5653 let result =
5654 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5655 (
5656 result,
5657 Some(ReturnedRuntime {
5658 connection,
5659 runtime,
5660 }),
5661 )
5662 }
5663 }
5664 }
5665}
5666
5667async fn close_runtime(
5670 mut runtime: Box<dyn RuntimeConnection>,
5671 process_group: Option<u32>,
5672) -> std::result::Result<Value, ServiceError> {
5673 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
5674 Ok(result) => {
5675 result.map_err(operation)?;
5676 Ok(json!({"closed": true}))
5677 }
5678 Err(deadline) => {
5679 let killed = kill_runtime_process_group(process_group);
5684 drop(runtime);
5685 Ok(json!({
5686 "closed": true,
5687 "killed": killed,
5688 "detail": error_message(deadline),
5689 }))
5690 }
5691 }
5692}
5693
5694fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
5697 mutation
5698 .session
5699 .clone()
5700 .filter(|value| !value.trim().is_empty())
5701 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
5702}
5703
5704async fn type_live_command(
5708 runtime: &mut dyn RuntimeConnection,
5709 verb: crate::SessionVerb,
5710 mutation: &crate::SessionMutation,
5711 command: &str,
5712 session: String,
5713) -> std::result::Result<Value, ServiceError> {
5714 within_control_deadline(
5715 &format!("sessions.{}", verb.as_str()),
5716 runtime.send_input(RuntimeInput {
5717 text: command.to_string(),
5718 image_urls: Vec::new(),
5719 }),
5720 )
5721 .await?
5722 .map_err(operation)?;
5723 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
5724 .map_err(session_control_error)?;
5725 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
5726}
5727
5728enum OpenRuntime {
5731 Hosted {
5734 runtime: Box<dyn RuntimeConnection>,
5735 capabilities: crate::RuntimeCapabilities,
5736 workspace: PathBuf,
5737 fresh: bool,
5739 },
5740 Joined { runtime: Box<dyn RuntimeConnection> },
5743}
5744
5745async fn open_runtime(
5750 method: &str,
5751 params: Value,
5752) -> std::result::Result<OpenRuntime, ServiceError> {
5753 match tokio::time::timeout(
5754 RUNTIME_OPEN_DEADLINE,
5755 open_runtime_unbounded(method, params),
5756 )
5757 .await
5758 {
5759 Ok(result) => result,
5760 Err(_) => Err(ServiceError::Operation(format!(
5761 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
5762 RUNTIME_OPEN_DEADLINE.as_secs()
5763 ))),
5764 }
5765}
5766
5767async fn open_runtime_unbounded(
5768 method: &str,
5769 params: Value,
5770) -> std::result::Result<OpenRuntime, ServiceError> {
5771 match method {
5772 "harness.v1.runtimes.start" => {
5773 let params = decode::<RuntimeStartParams>(params)?;
5774 let backend = runtime_backend(¶ms.backend)?;
5775 let capabilities = backend.capabilities();
5776 let workspace = params.cwd.clone();
5777 let runtime = backend
5778 .start(RuntimeStartRequest {
5779 cwd: params.cwd,
5780 launch: runtime_launch(¶ms.backend),
5781 mcp_servers: params.mcp_servers,
5782 approval_policy: params.approval_policy,
5783 })
5784 .await
5785 .map_err(operation)?;
5786 Ok(OpenRuntime::Hosted {
5787 runtime,
5788 capabilities,
5789 workspace,
5790 fresh: true,
5791 })
5792 }
5793 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
5794 let params = decode::<RuntimeAttachParams>(params)?;
5795 let backend = runtime_backend(¶ms.backend)?;
5796 let capabilities = backend.capabilities();
5797 let workspace = params
5798 .cwd
5799 .clone()
5800 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
5801 let runtime = backend
5802 .attach(RuntimeAttachRequest {
5803 runtime_id: params.runtime_id,
5804 cwd: params.cwd,
5805 launch: runtime_launch(¶ms.backend),
5806 mcp_servers: params.mcp_servers,
5807 approval_policy: params.approval_policy,
5808 })
5809 .await
5810 .map_err(operation)?;
5811 Ok(OpenRuntime::Hosted {
5812 runtime,
5813 capabilities,
5814 workspace,
5815 fresh: false,
5816 })
5817 }
5818 "harness.v1.runtimes.attach_existing" => {
5819 let params = decode::<RuntimeAttachParams>(params)?;
5820 let backend: Box<dyn RuntimeBackend> = match params
5821 .backend
5822 .base_url
5823 .as_deref()
5824 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
5825 {
5826 Some(endpoint) => {
5827 #[cfg(not(feature = "adapter-api"))]
5828 {
5829 let _ = endpoint;
5830 return Err(ServiceError::UnsupportedAction(
5831 "live HTTP attachment adapter is not compiled".into(),
5832 ));
5833 }
5834 #[cfg(feature = "adapter-api")]
5835 {
5836 let workspace = params.cwd.clone().ok_or_else(|| {
5837 ServiceError::InvalidParams(
5838 "Volter Harness live attach requires the project cwd".into(),
5839 )
5840 })?;
5841 let source = LiveRuntimeSource {
5842 harness: params.backend.harness.as_str().to_string(),
5843 session_id: params.runtime_id.clone(),
5844 workspace,
5845 };
5846 let receipt = resolve_live_runtime(&endpoint, &source)
5847 .map_err(|error| ServiceError::Operation(error.to_string()))?;
5848 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
5849 }
5850 }
5851 None => runtime_backend(¶ms.backend)?,
5852 };
5853 let capabilities = backend.capabilities();
5854 if !capabilities.attach_existing_process {
5855 return Err(ServiceError::Operation(format!(
5856 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
5857 backend.harness().as_str()
5858 )));
5859 }
5860 let runtime = backend
5861 .attach_existing(RuntimeAttachRequest {
5862 runtime_id: params.runtime_id,
5863 cwd: params.cwd,
5864 launch: runtime_launch(¶ms.backend),
5865 mcp_servers: params.mcp_servers,
5866 approval_policy: params.approval_policy,
5867 })
5868 .await
5869 .map_err(operation)?;
5870 Ok(OpenRuntime::Joined { runtime })
5871 }
5872 _ => Err(ServiceError::MethodNotFound),
5873 }
5874}
5875
5876fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
5878 match result {
5879 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
5880 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
5881 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
5882 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
5883 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
5884 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
5885 }
5886}
5887
5888fn runtime_backend(
5889 params: &RuntimeBackendParams,
5890) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5891 if let Some(descriptor) = registry_connect_descriptor(params) {
5892 return open_connect_descriptor(&descriptor, &service_home()?);
5893 }
5894 if params.protocol.as_deref() == Some("acp") {
5895 let launch = params
5896 .launch
5897 .clone()
5898 .or_else(|| {
5899 harness_support_registry()
5900 .harnesses
5901 .into_iter()
5902 .find(|harness| harness.id == params.harness)
5903 .filter(|harness| {
5904 harness.runtime.implementation == ImplementationKind::GenericProtocol
5905 && harness.runtime.protocol.starts_with("acp")
5906 })
5907 .and_then(|harness| harness.runtime.default_launch)
5908 })
5909 .ok_or_else(|| {
5910 ServiceError::InvalidParams(
5911 "an ACP runtime requires `launch` unless the harness has a registered default"
5912 .into(),
5913 )
5914 })?;
5915 let resume_session = harness_support_registry()
5916 .harnesses
5917 .into_iter()
5918 .find(|harness| harness.id == params.harness)
5919 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
5920 return Ok(Box::new(
5921 AcpRuntimeBackend::new(params.harness.clone(), launch)
5922 .with_resume_support(resume_session),
5923 ));
5924 }
5925 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
5926 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
5927 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
5928 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
5929 HarnessId::OPENCODE => match ¶ms.base_url {
5930 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
5931 None => Box::new(OpenCodeRuntimeBackend::new()),
5932 },
5933 harness => {
5934 let descriptor = harness_support_registry()
5935 .harnesses
5936 .into_iter()
5937 .find(|descriptor| descriptor.id.as_str() == harness)
5938 .filter(|descriptor| {
5939 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
5940 && descriptor.runtime.protocol.starts_with("acp")
5941 });
5942 let Some(descriptor) = descriptor else {
5943 return Err(ServiceError::InvalidParams(format!(
5944 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
5945 )));
5946 };
5947 let resume = descriptor.runtime.capabilities.resume_session;
5948 Box::new(
5949 AcpRuntimeBackend::new(
5950 descriptor.id,
5951 descriptor
5952 .runtime
5953 .default_launch
5954 .expect("generic ACP registry entry includes its launch"),
5955 )
5956 .with_resume_support(resume),
5957 )
5958 }
5959 };
5960 Ok(backend)
5961}
5962
5963fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
5964 if let Some(launch) = ¶ms.launch {
5965 return Some(launch.clone());
5966 }
5967 if !matches!(params.policy, RuntimePolicy::Yolo) {
5968 return None;
5969 }
5970 let launch = match params.harness.as_str() {
5971 HarnessId::GROK => RuntimeLaunch {
5972 program: "grok".into(),
5973 arguments: {
5974 let mut arguments: Vec<String> = Vec::new();
5975 if crate::support::self_sandbox_supported() {
5976 arguments.extend(["--sandbox".into(), "workspace".into()]);
5977 }
5978 arguments.extend([
5979 "--always-approve".into(),
5980 "agent".into(),
5981 "--no-leader".into(),
5982 "stdio".into(),
5983 ]);
5984 arguments
5985 },
5986 env: crate::support::grok_env(),
5987 },
5988 HarnessId::CODEX => RuntimeLaunch {
5989 program: "codex".into(),
5990 arguments: vec![
5991 "--dangerously-bypass-approvals-and-sandbox".into(),
5992 "--dangerously-bypass-hook-trust".into(),
5993 "app-server".into(),
5994 ],
5995 env: BTreeMap::new(),
5996 },
5997 HarnessId::CLAUDE_CODE => RuntimeLaunch {
5998 program: "claude".into(),
5999 arguments: vec![
6000 "--dangerously-skip-permissions".into(),
6001 "--print".into(),
6002 "--input-format".into(),
6003 "stream-json".into(),
6004 "--output-format".into(),
6005 "stream-json".into(),
6006 "--verbose".into(),
6007 ],
6008 env: BTreeMap::new(),
6009 },
6010 HarnessId::PI => RuntimeLaunch {
6011 program: "pi".into(),
6012 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
6013 env: BTreeMap::new(),
6014 },
6015 HarnessId::OPENCODE => RuntimeLaunch {
6016 program: "opencode".into(),
6017 arguments: vec!["serve".into()],
6018 env: BTreeMap::new(),
6019 },
6020 HarnessId::GEMINI => RuntimeLaunch {
6021 program: "gemini".into(),
6022 arguments: vec!["--acp".into(), "--yolo".into()],
6023 env: BTreeMap::new(),
6024 },
6025 HarnessId::GOOSE => RuntimeLaunch {
6026 program: "goose".into(),
6027 arguments: vec!["acp".into()],
6028 env: BTreeMap::new(),
6029 },
6030 HarnessId::SUPERCODE => RuntimeLaunch {
6031 program: "supercode".into(),
6032 arguments: vec!["acp".into(), "--dangerous".into()],
6033 env: BTreeMap::new(),
6034 },
6035 _ => return None,
6036 };
6037 Some(launch)
6038}
6039
6040struct IsolatedProbeHome {
6046 launch: RuntimeLaunch,
6047 root: PathBuf,
6048}
6049
6050impl IsolatedProbeHome {
6051 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6052 let root = std::env::temp_dir().join(format!(
6053 "supercode-harness-probe-{harness}-{}",
6054 generated_session_id()
6055 ));
6056 std::fs::create_dir_all(&root)?;
6057 set_private_dir_permissions(&root)?;
6058
6059 if let Some(source_home) = supercode_interchange::user_home()
6060 .map(std::path::PathBuf::into_os_string)
6061 .map(PathBuf::from)
6062 {
6063 for relative in probe_auth_files(harness) {
6064 copy_probe_file(&source_home, &root, relative)?;
6065 }
6066 }
6067 if harness == HarnessId::SUPERCODE {
6072 let config_home = crate::agent::global_instructions_dir();
6073 for file in ["config.toml", "credentials.toml"] {
6074 copy_probe_path(
6075 &config_home.join(file),
6076 &root.join(".config/supercode").join(file),
6077 )?;
6078 }
6079 }
6080 configure_isolated_probe_auth(harness, &root)?;
6081
6082 let root_text = root.to_string_lossy().into_owned();
6083 for (key, value) in [
6084 ("HOME", root_text.clone()),
6085 (
6086 "XDG_CACHE_HOME",
6087 root.join(".cache").to_string_lossy().into_owned(),
6088 ),
6089 (
6090 "XDG_CONFIG_HOME",
6091 root.join(".config").to_string_lossy().into_owned(),
6092 ),
6093 (
6094 "XDG_DATA_HOME",
6095 root.join(".local/share").to_string_lossy().into_owned(),
6096 ),
6097 ] {
6098 launch.env.insert(key.into(), value);
6099 }
6100 let scoped = match harness {
6101 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6102 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6103 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6104 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6105 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6106 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6107 _ => None,
6108 };
6109 if let Some((key, value)) = scoped {
6110 launch
6111 .env
6112 .insert(key.into(), value.to_string_lossy().into_owned());
6113 }
6114 Ok(Self { launch, root })
6115 }
6116
6117 fn cleanup(&self) -> std::io::Result<()> {
6118 match std::fs::remove_dir_all(&self.root) {
6119 Ok(()) => Ok(()),
6120 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6121 Err(error) => Err(error),
6122 }
6123 }
6124}
6125
6126impl Drop for IsolatedProbeHome {
6127 fn drop(&mut self) {
6128 let _ = self.cleanup();
6129 }
6130}
6131
6132fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6133 match harness {
6134 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6135 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6139 HarnessId::CODEX => &[".codex/auth.json"],
6140 HarnessId::GEMINI => &[
6141 ".gemini/google_accounts.json",
6142 ".gemini/oauth_creds.json",
6143 ".gemini/settings.json",
6144 ],
6145 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6146 HarnessId::OPENCODE => &[
6147 ".config/opencode/auth.json",
6148 ".local/share/opencode/auth.json",
6149 ],
6150 HarnessId::PI => &[".pi/agent/auth.json"],
6151 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6156 _ => &[],
6157 }
6158}
6159
6160fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6161 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6162}
6163
6164fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6165 if !source.is_file() {
6166 return Ok(());
6167 }
6168 if let Some(parent) = destination.parent() {
6169 std::fs::create_dir_all(parent)?;
6170 set_private_dir_permissions(parent)?;
6171 }
6172 std::fs::copy(source, destination)?;
6173 set_private_file_permissions(destination)
6174}
6175
6176fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6177 if harness != HarnessId::GEMINI {
6178 return Ok(());
6179 }
6180 let oauth = probe_home.join(".gemini/oauth_creds.json");
6181 if !oauth.is_file() {
6182 return Ok(());
6183 }
6184 let settings_path = probe_home.join(".gemini/settings.json");
6185 let mut settings = std::fs::read_to_string(&settings_path)
6186 .ok()
6187 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6188 .unwrap_or_else(|| json!({}));
6189 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6190 std::fs::write(
6191 &settings_path,
6192 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6193 )?;
6194 set_private_file_permissions(&settings_path)
6195}
6196
6197#[cfg(unix)]
6198fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6199 use std::os::unix::fs::PermissionsExt;
6200 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6201}
6202
6203#[cfg(not(unix))]
6204fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6205 Ok(())
6206}
6207
6208#[cfg(unix)]
6209fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6210 use std::os::unix::fs::PermissionsExt;
6211 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6212}
6213
6214#[cfg(not(unix))]
6215fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6216 Ok(())
6217}
6218
6219fn find_executable(program: &str) -> Option<PathBuf> {
6220 let candidate = PathBuf::from(program);
6221 if candidate.components().count() > 1 {
6222 return candidate.is_file().then_some(candidate);
6223 }
6224 let path = std::env::var_os("PATH")?;
6225 for directory in std::env::split_paths(&path) {
6226 let candidate = directory.join(program);
6227 if candidate.is_file() {
6228 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6229 }
6230 #[cfg(windows)]
6231 {
6232 for extension in ["exe", "cmd", "bat"] {
6233 let candidate = directory.join(format!("{program}.{extension}"));
6234 if candidate.is_file() {
6235 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6236 }
6237 }
6238 }
6239 }
6240 None
6241}
6242
6243async fn executable_version(executable: &Path) -> Option<String> {
6244 let mut command = tokio::process::Command::new(executable);
6245 command
6246 .arg("--version")
6247 .stdin(std::process::Stdio::null())
6248 .stdout(std::process::Stdio::piped())
6249 .stderr(std::process::Stdio::piped())
6250 .kill_on_drop(true);
6251 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6252 .await
6253 .ok()?
6254 .ok()?;
6255 let stdout = String::from_utf8_lossy(&output.stdout);
6256 let stderr = String::from_utf8_lossy(&output.stderr);
6257 stdout
6258 .lines()
6259 .chain(stderr.lines())
6260 .map(str::trim)
6261 .find(|line| !line.is_empty())
6262 .map(|line| truncate_text(line, 200))
6263}
6264
6265pub(crate) fn auth_evidence(harness: &str) -> bool {
6266 let env_names: &[&str] = match harness {
6267 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6268 HarnessId::CODEX => &["OPENAI_API_KEY"],
6269 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6270 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6271 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6272 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6273 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6274 _ => &[],
6275 };
6276 if env_names
6277 .iter()
6278 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6279 {
6280 return true;
6281 }
6282 let Some(home) = supercode_interchange::user_home()
6283 .map(std::path::PathBuf::into_os_string)
6284 .map(PathBuf::from)
6285 else {
6286 return false;
6287 };
6288 let files: Vec<PathBuf> = match harness {
6289 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6290 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6291 HarnessId::OPENCODE => vec![
6292 home.join(".local/share/opencode/auth.json"),
6293 home.join(".config/opencode/auth.json"),
6294 ],
6295 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6296 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6297 HarnessId::GEMINI => vec![
6298 home.join(".gemini/oauth_creds.json"),
6299 home.join(".gemini/google_accounts.json"),
6300 ],
6301 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6302 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6303 _ => Vec::new(),
6304 };
6305 if files.into_iter().any(|path| {
6306 std::fs::metadata(path)
6307 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6308 .unwrap_or(false)
6309 }) {
6310 return true;
6311 }
6312 if harness == HarnessId::CLAUDE_CODE {
6319 return std::fs::read_to_string(home.join(".claude.json"))
6320 .map(|text| text.contains("\"oauthAccount\""))
6321 .unwrap_or(false);
6322 }
6323 false
6324}
6325
6326fn looks_like_auth_error(message: &str) -> bool {
6327 let message = message.to_ascii_lowercase();
6328 [
6329 "auth",
6330 "login",
6331 "sign in",
6332 "sign-in",
6333 "credential",
6334 "unauthorized",
6335 "forbidden",
6336 "token",
6337 ]
6338 .iter()
6339 .any(|needle| message.contains(needle))
6340}
6341
6342fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6343 crate::RuntimeCapabilities {
6344 start_session: false,
6345 resume_session: false,
6346 attach_existing_process: false,
6347 send_input: false,
6348 stream_events: false,
6349 interrupt: false,
6350 steer: false,
6351 respond_to_requests: false,
6352 }
6353}
6354
6355fn truncate_text(text: &str, max_chars: usize) -> String {
6356 let mut chars = text.chars();
6357 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6358 if chars.next().is_some() {
6359 format!("{truncated}…")
6360 } else {
6361 truncated
6362 }
6363}
6364
6365fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6372 match &handle.endpoint {
6373 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6374 crate::RuntimeEndpoint::Http { .. } => None,
6375 }
6376}
6377
6378fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6384 match process_group {
6385 #[cfg(unix)]
6386 Some(pid) => {
6387 crate::lsp::kill_process_group(pid);
6388 true
6389 }
6390 #[cfg(not(unix))]
6391 Some(_) => false,
6392 None => false,
6393 }
6394}
6395
6396fn error_message(error: ServiceError) -> String {
6397 match error {
6398 ServiceError::InvalidParams(message)
6399 | ServiceError::Operation(message)
6400 | ServiceError::UnsupportedAction(message) => message,
6401 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6402 ServiceError::Sdk(error) => error.to_string(),
6403 }
6404}
6405
6406#[derive(Debug)]
6407enum ServiceError {
6408 InvalidParams(String),
6409 MethodNotFound,
6410 UnsupportedAction(String),
6411 Operation(String),
6412 Sdk(SdkError),
6413}
6414
6415fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6416 match error {
6417 ServiceError::InvalidParams(message) => {
6418 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6419 }
6420 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6421 SdkError::unsupported(operation)
6422 }
6423 ServiceError::Operation(message) => {
6424 let code = if message.contains("already in progress") {
6425 SdkErrorCode::Busy
6426 } else if message.contains("not supported by this runtime") {
6427 SdkErrorCode::UnsupportedAction
6428 } else if message.contains("unknown runtime connection") {
6429 SdkErrorCode::NotFound
6430 } else {
6431 SdkErrorCode::Execution
6432 };
6433 SdkError::new(code, operation, message)
6434 }
6435 ServiceError::Sdk(error) => error,
6436 }
6437}
6438
6439fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6440 let error_code = error.code();
6441 let code = match error_code {
6442 SdkErrorCode::Unauthenticated => -32030,
6443 SdkErrorCode::Unauthorized => -32031,
6444 SdkErrorCode::ControllerRequired => -32032,
6445 SdkErrorCode::LeaseExpired => -32033,
6446 SdkErrorCode::InvalidArgument => -32602,
6447 SdkErrorCode::NotFound => -32004,
6448 SdkErrorCode::Busy => -32000,
6449 SdkErrorCode::UnsupportedAction => -32020,
6450 SdkErrorCode::Execution => -32002,
6451 SdkErrorCode::Transport => -32003,
6452 };
6453 json!({
6454 "jsonrpc": "2.0",
6455 "id": id,
6456 "error": {
6457 "code": code,
6458 "name": error_code,
6459 "operation": error.operation(),
6460 "message": error.to_string(),
6461 },
6462 })
6463}
6464
6465fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6466 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6467}
6468
6469fn operation(error: impl Into<crate::Error>) -> ServiceError {
6470 let error = error.into();
6471 match error {
6472 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6473 error => ServiceError::Operation(error.to_string()),
6474 }
6475}
6476
6477#[derive(Debug, Clone, Deserialize, Default)]
6481#[serde(default)]
6482struct MemoryRequest {
6483 harness: Option<String>,
6485 query: Option<String>,
6487 profile: Option<String>,
6489 session: Option<String>,
6491 full: bool,
6493 regex: bool,
6495 cwd: Option<std::path::PathBuf>,
6497 homes: crate::HarnessHomes,
6499}
6500
6501fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6504 let request = decode::<MemoryRequest>(params)?;
6505 let harness = request
6506 .harness
6507 .clone()
6508 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6509 let to_service = |error: crate::memory::MemoryError| match error {
6510 crate::memory::MemoryError::UnsupportedHarness { .. }
6511 | crate::memory::MemoryError::SessionNotScoped { .. } => {
6512 ServiceError::UnsupportedAction(error.to_string())
6513 }
6514 other => ServiceError::InvalidParams(other.to_string()),
6515 };
6516 match method {
6517 "harness.v1.memory.show" => {
6518 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6519 harness,
6520 profile: request.profile,
6521 session: request.session,
6522 full: request.full,
6523 cwd: request.cwd,
6524 homes: request.homes,
6525 })
6526 .map_err(to_service)?;
6527 Ok(json!({
6528 "schema": crate::memory::MEMORY_SCHEMA,
6529 "documents": documents,
6530 }))
6531 }
6532 "harness.v1.memory.search" => {
6533 let query = request
6534 .query
6535 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6536 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6537 harness,
6538 query,
6539 profile: request.profile,
6540 regex: request.regex,
6541 cwd: request.cwd,
6542 homes: request.homes,
6543 })
6544 .map_err(to_service)?;
6545 Ok(json!({
6546 "schema": crate::memory::MEMORY_SCHEMA,
6547 "matches": matches,
6548 }))
6549 }
6550 _ => Err(ServiceError::MethodNotFound),
6551 }
6552}
6553
6554#[derive(Debug, Clone, Deserialize)]
6558#[serde(default)]
6559struct ProfilesQuery {
6560 harness: Option<String>,
6562 name: Option<String>,
6564 homes: crate::HarnessHomes,
6566}
6567
6568impl Default for ProfilesQuery {
6569 fn default() -> Self {
6570 Self {
6571 harness: None,
6572 name: None,
6573 homes: crate::HarnessHomes::default(),
6574 }
6575 }
6576}
6577
6578fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6581 let query = decode::<ProfilesQuery>(params)?;
6582 let to_service = |error: crate::profiles::ProfileError| match error {
6583 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6584 ServiceError::UnsupportedAction(error.to_string())
6585 }
6586 crate::profiles::ProfileError::NotFound { .. } => {
6587 ServiceError::InvalidParams(error.to_string())
6588 }
6589 };
6590 match method {
6591 "harness.v1.profiles.list" => {
6592 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6593 .map_err(to_service)?;
6594 Ok(json!({
6595 "schema": crate::profiles::PROFILES_SCHEMA,
6596 "profiles": profiles,
6597 }))
6598 }
6599 "harness.v1.profiles.get" => {
6600 let harness = query
6601 .harness
6602 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6603 let name = query
6604 .name
6605 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6606 let profile =
6607 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6608 Ok(json!({
6609 "schema": crate::profiles::PROFILES_SCHEMA,
6610 "profile": profile,
6611 }))
6612 }
6613 _ => Err(ServiceError::MethodNotFound),
6614 }
6615}
6616
6617#[derive(Debug, Clone, Deserialize)]
6621#[serde(default)]
6622struct ChannelsQuery {
6623 harness: Option<String>,
6625 name: Option<String>,
6627 homes: crate::HarnessHomes,
6629}
6630
6631impl Default for ChannelsQuery {
6632 fn default() -> Self {
6633 Self {
6634 harness: None,
6635 name: None,
6636 homes: crate::HarnessHomes::default(),
6637 }
6638 }
6639}
6640
6641#[derive(Debug, Clone, Deserialize)]
6645#[serde(default)]
6646struct RoutesQuery {
6647 harness: Option<String>,
6648 profile: Option<String>,
6650 homes: crate::HarnessHomes,
6651}
6652
6653impl Default for RoutesQuery {
6654 fn default() -> Self {
6655 Self {
6656 harness: None,
6657 profile: None,
6658 homes: crate::HarnessHomes::default(),
6659 }
6660 }
6661}
6662
6663#[derive(Debug, Clone, Deserialize)]
6664#[serde(default)]
6665struct TriggersQuery {
6666 harness: Option<String>,
6667 homes: crate::HarnessHomes,
6668}
6669
6670impl Default for TriggersQuery {
6671 fn default() -> Self {
6672 Self {
6673 harness: None,
6674 homes: crate::HarnessHomes::default(),
6675 }
6676 }
6677}
6678
6679fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
6680 let query = decode::<TriggersQuery>(params)?;
6681 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
6682 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6683 Ok(json!({
6684 "schema": crate::triggers::TRIGGERS_SCHEMA,
6685 "triggers": triggers,
6686 }))
6687}
6688
6689fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
6690 let query = decode::<RoutesQuery>(params)?;
6691 let routes = crate::routes::list_routes(
6692 &query.homes,
6693 query.harness.as_deref(),
6694 query.profile.as_deref(),
6695 )
6696 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6697 Ok(json!({
6698 "schema": crate::routes::ROUTES_SCHEMA,
6699 "routes": routes,
6700 }))
6701}
6702
6703fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6704 let query = decode::<ChannelsQuery>(params)?;
6705 let to_service = |error: crate::channels::ChannelError| match error {
6706 crate::channels::ChannelError::UnsupportedHarness { .. } => {
6707 ServiceError::UnsupportedAction(error.to_string())
6708 }
6709 crate::channels::ChannelError::NotFound { .. } => {
6710 ServiceError::InvalidParams(error.to_string())
6711 }
6712 };
6713 match method {
6714 "harness.v1.channels.list" => {
6715 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
6716 .map_err(to_service)?;
6717 Ok(json!({
6718 "schema": crate::channels::CHANNELS_SCHEMA,
6719 "channels": channels,
6720 }))
6721 }
6722 "harness.v1.channels.status" => {
6723 let harness = query
6724 .harness
6725 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6726 let name = query
6727 .name
6728 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6729 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
6730 .map_err(to_service)?;
6731 Ok(json!({
6732 "schema": crate::channels::CHANNELS_SCHEMA,
6733 "channel": channel,
6734 }))
6735 }
6736 _ => Err(ServiceError::MethodNotFound),
6737 }
6738}
6739
6740fn rpc_error(id: Value, code: i64, message: &str) -> Value {
6741 json!({
6742 "jsonrpc": "2.0",
6743 "id": id,
6744 "error": {"code": code, "message": message},
6745 })
6746}