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 let backend = runtime_backend(¶ms)?;
1806 Ok(json!({
1807 "harness": backend.harness(),
1808 "capabilities": backend.capabilities(),
1809 }))
1810 }
1811 method if RUNTIME_OPEN_METHODS.contains(&method) => {
1812 self.register_open_runtime(open_runtime(method, params).await?)
1813 .await
1814 }
1815 "harness.v1.runtimes.send_input" => {
1816 let params = decode::<RuntimeInputParams>(params)?;
1817 let image_urls = validate_runtime_image_urls(params.image_urls)?;
1818 let runtime = self.runtime_mut(¶ms.connection)?;
1819 let turn_id = within_control_deadline(
1820 method,
1821 runtime.send_input(RuntimeInput {
1822 text: params.text,
1823 image_urls,
1824 }),
1825 )
1826 .await?
1827 .map_err(operation)?;
1828 Ok(json!({"turn_id": turn_id}))
1829 }
1830 "harness.v1.runtimes.interrupt" => {
1831 let params = decode::<RuntimeConnectionParams>(params)?;
1832 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.interrupt())
1833 .await?
1834 .map_err(operation)?;
1835 Ok(json!({}))
1836 }
1837 "harness.v1.runtimes.steer" => {
1838 let params = decode::<RuntimeInputParams>(params)?;
1839 if !params.image_urls.is_empty() {
1840 return Err(ServiceError::InvalidParams(
1841 "runtime steering accepts text only".into(),
1842 ));
1843 }
1844 let text = params.text.trim();
1845 if text.is_empty() || text.chars().count() > 50_000 {
1846 return Err(ServiceError::InvalidParams(
1847 "runtime steering requires 1 to 50,000 text characters".into(),
1848 ));
1849 }
1850 within_control_deadline(
1851 method,
1852 self.runtime_mut(¶ms.connection)?
1853 .steer(text.to_string()),
1854 )
1855 .await?
1856 .map_err(operation)?;
1857 Ok(json!({}))
1858 }
1859 "harness.v1.runtimes.respond" => {
1860 let params = decode::<RuntimeRespondParams>(params)?;
1861 let request_id = params.request_id.clone();
1862 within_control_deadline(
1863 method,
1864 self.runtime_mut(¶ms.connection)?
1865 .respond(params.request_id, params.response),
1866 )
1867 .await?
1868 .map_err(operation)?;
1869 self.approvals.answered(¶ms.connection, &request_id);
1871 Ok(json!({}))
1872 }
1873 "harness.v1.runtimes.acquire_control" => {
1874 let params = decode::<RuntimeConnectionParams>(params)?;
1875 let snapshot = within_control_deadline(
1876 method,
1877 self.runtime_mut(¶ms.connection)?.acquire_control(),
1878 )
1879 .await?
1880 .map_err(operation)?;
1881 serde_json::to_value(snapshot)
1882 .map_err(|error| ServiceError::Operation(error.to_string()))
1883 }
1884 "harness.v1.runtimes.heartbeat" => {
1885 let params = decode::<RuntimeConnectionParams>(params)?;
1886 let snapshot = within_control_deadline(
1887 method,
1888 self.runtime_mut(¶ms.connection)?.heartbeat(),
1889 )
1890 .await?
1891 .map_err(operation)?;
1892 serde_json::to_value(snapshot)
1893 .map_err(|error| ServiceError::Operation(error.to_string()))
1894 }
1895 "harness.v1.runtimes.detach" => {
1896 let params = decode::<RuntimeConnectionParams>(params)?;
1897 let snapshot =
1898 within_control_deadline(method, self.runtime_mut(¶ms.connection)?.detach())
1899 .await?
1900 .map_err(operation)?;
1901 serde_json::to_value(snapshot)
1902 .map_err(|error| ServiceError::Operation(error.to_string()))
1903 }
1904 "harness.v1.runtimes.terminal_instructions" => {
1905 let params = decode::<RuntimeConnectionParams>(params)?;
1906 let launch = self
1907 .terminal_launches
1908 .get(¶ms.connection)
1909 .ok_or_else(|| {
1910 ServiceError::Operation(
1911 "this runtime is not hosted for terminal attachment".into(),
1912 )
1913 })?;
1914 Ok(json!({"launch":launch}))
1915 }
1916 "harness.v1.runtimes.close" => {
1917 let params = decode::<RuntimeConnectionParams>(params)?;
1918 let (runtime, process_group) = self.surrender_runtime(¶ms.connection)?;
1919 close_runtime(runtime, process_group).await
1920 }
1921 _ => Err(ServiceError::MethodNotFound),
1922 }
1923 }
1924
1925 #[cfg(feature = "adapter-api")]
1927 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1928 let params = decode::<MessageSessionParams>(params)?;
1929 Ok(message_live_session(¶ms).await)
1930 }
1931
1932 #[cfg(feature = "adapter-api")]
1933 fn harness_settings_call(
1934 &self,
1935 method: &str,
1936 params: Value,
1937 ) -> std::result::Result<Value, ServiceError> {
1938 let homes = crate::HarnessHomes::default();
1939 match method {
1940 "harness.v1.harnesses.settings" => {
1941 let params = decode::<HarnessSettingsParams>(params)?;
1942 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1943 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1944 serde_json::to_value(report)
1945 .map_err(|error| ServiceError::Operation(error.to_string()))
1946 }
1947 "harness.v1.harnesses.configure" => {
1948 let params = decode::<ConfigureHarnessParams>(params)?;
1949 let report = crate::configure_harness_interop_settings(
1950 &homes,
1951 ¶ms.harness,
1952 ¶ms.changes,
1953 params.expected_revision.as_deref(),
1954 )
1955 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1956 serde_json::to_value(report)
1957 .map_err(|error| ServiceError::Operation(error.to_string()))
1958 }
1959 _ => Err(ServiceError::MethodNotFound),
1960 }
1961 }
1962
1963 fn insert_runtime(
1964 &mut self,
1965 runtime: Box<dyn RuntimeConnection>,
1966 ) -> std::result::Result<Value, ServiceError> {
1967 let connection = format!("runtime-{}", self.next_runtime);
1968 self.next_runtime += 1;
1969 let handle = runtime.handle().clone();
1970 self.runtime_sequences
1971 .entry(handle.runtime_id.clone())
1972 .or_insert(0);
1973 self.runtimes.insert(connection.clone(), runtime);
1974 Ok(json!({"connection": connection, "handle": handle}))
1975 }
1976
1977 #[cfg(feature = "adapter-api")]
1978 async fn insert_hosted_runtime(
1979 &mut self,
1980 runtime: Box<dyn RuntimeConnection>,
1981 capabilities: crate::RuntimeCapabilities,
1982 workspace: PathBuf,
1983 fresh: bool,
1984 ) -> std::result::Result<Value, ServiceError> {
1985 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
1986 let token: std::sync::Arc<str> = crate::server::generate_token().into();
1987 let server = crate::server::run_frontend_http(
1988 host.clone(),
1989 host.frontend_sender(),
1990 "127.0.0.1:0",
1991 token.clone(),
1992 connection.handle().runtime_id.clone(),
1993 )
1994 .await
1995 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1996 let source = LiveRuntimeSource {
1997 harness: connection.handle().harness.as_str().to_string(),
1998 session_id: connection.handle().runtime_id.clone(),
1999 workspace: workspace.clone(),
2000 };
2001 let registration = register_live_runtime(
2002 connection.handle().runtime_id.clone(),
2003 source.clone(),
2004 format!("http://{}", server.address()),
2005 token.to_string(),
2006 )
2007 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2008 let endpoint = registration.endpoint().to_string();
2009 let launch = StructuredLaunch {
2010 cwd: workspace,
2011 program: std::env::current_exe()
2015 .ok()
2016 .map(|path| path.to_string_lossy().into_owned())
2017 .unwrap_or_else(|| "supercode".into()),
2018 arguments: vec![
2019 "open".into(),
2020 endpoint,
2021 "--harness".into(),
2022 source.harness,
2023 "--session".into(),
2024 source.session_id,
2025 ],
2026 env: BTreeMap::new(),
2027 };
2028 let lease = HostedRuntimeLease {
2029 connection,
2030 _host: host,
2031 _registration: registration,
2032 _server: server,
2033 };
2034 let opened = self.insert_runtime(Box::new(lease))?;
2035 let connection_id = opened["connection"]
2036 .as_str()
2037 .expect("insert_runtime returns a connection id")
2038 .to_string();
2039 self.terminal_launches.insert(connection_id, launch);
2040 Ok(opened)
2041 }
2042
2043 #[cfg(not(feature = "adapter-api"))]
2044 async fn insert_hosted_runtime(
2045 &mut self,
2046 runtime: Box<dyn RuntimeConnection>,
2047 _capabilities: crate::RuntimeCapabilities,
2048 _workspace: PathBuf,
2049 _fresh: bool,
2050 ) -> std::result::Result<Value, ServiceError> {
2051 self.insert_runtime(runtime)
2052 }
2053
2054 fn runtime_mut(
2055 &mut self,
2056 connection: &str,
2057 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2058 if self.runtimes_in_flight.contains(connection) {
2059 return Err(self.lent_out(connection));
2060 }
2061 self.runtimes.get_mut(connection).ok_or_else(|| {
2062 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2063 })
2064 }
2065
2066 fn lent_out(&self, connection: &str) -> ServiceError {
2070 ServiceError::Operation(format!(
2071 "runtime connection `{connection}`: a harness turn is already in progress"
2072 ))
2073 }
2074
2075 fn lend_runtime(
2078 &mut self,
2079 connection: &str,
2080 ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2081 if self.runtimes_in_flight.contains(connection) {
2082 return Err(self.lent_out(connection));
2083 }
2084 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2085 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2086 })?;
2087 self.runtimes_in_flight.insert(connection.to_string());
2088 Ok(runtime)
2089 }
2090
2091 fn surrender_runtime(
2103 &mut self,
2104 connection: &str,
2105 ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2106 if self.runtimes_in_flight.contains(connection) {
2107 return Err(self.lent_out(connection));
2108 }
2109 let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2110 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2111 })?;
2112 let process_group = runtime_process_group(runtime.handle());
2113 let runtime_id = runtime.handle().runtime_id.clone();
2114 self.terminal_launches.remove(connection);
2115 self.runtime_sequences.remove(&runtime_id);
2116 self.approvals.forget(connection);
2117 Ok((runtime, process_group))
2118 }
2119
2120 pub fn kill_all_runtime_groups(&self) -> usize {
2131 self.runtimes
2132 .values()
2133 .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2134 .count()
2135 }
2136
2137 async fn mutate_session(
2148 &mut self,
2149 verb: crate::SessionVerb,
2150 params: Value,
2151 ) -> std::result::Result<Value, ServiceError> {
2152 let mutation = decode::<crate::SessionMutation>(params)?;
2153 let door = crate::sessions_control::door(&mutation.harness, verb)
2154 .map_err(session_control_error)?;
2155 let outcome = match door {
2156 #[cfg(not(feature = "adapter-api"))]
2160 crate::SessionDoor::Live(command) => {
2161 return Err(ServiceError::Operation(format!(
2162 "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2163 session, which needs this build's `adapter-api` feature",
2164 mutation.harness,
2165 verb.as_str()
2166 )));
2167 }
2168 #[cfg(feature = "adapter-api")]
2169 crate::SessionDoor::Live(command) => {
2170 let connection = mutation
2171 .connection
2172 .clone()
2173 .filter(|value| !value.trim().is_empty())
2174 .ok_or_else(|| {
2175 ServiceError::InvalidParams(format!(
2176 "`{}` performs `sessions.{}` by typing `{command}` into a live \
2177 driven session: pass the `connection` of an open runtime \
2178 (`harness.v1.runtimes.start`)",
2179 mutation.harness,
2180 verb.as_str()
2181 ))
2182 })?;
2183 let runtime = self.runtime_mut(&connection)?;
2184 let session = live_session_name(runtime.as_ref(), &mutation);
2185 return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2191 .await;
2192 }
2193 _ => run_session_mutation(verb, &mutation).await?,
2194 };
2195 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2196 }
2197
2198 async fn inventory_call(
2202 &self,
2203 method: &str,
2204 params: Value,
2205 ) -> std::result::Result<Value, ServiceError> {
2206 run_inventory(self.inventory_work(method, params)?).await
2207 }
2208
2209 fn inventory_work(
2215 &self,
2216 method: &str,
2217 params: Value,
2218 ) -> std::result::Result<InventoryWork, ServiceError> {
2219 let mut params = decode::<HarnessInventoryParams>(params)?;
2220 if method == "harness.v1.harnesses.probe" {
2221 let harness = params.harness.take().ok_or_else(|| {
2222 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2223 })?;
2224 params.harnesses = vec![harness];
2225 }
2226 let selected = params
2227 .harnesses
2228 .iter()
2229 .map(HarnessId::as_str)
2230 .collect::<std::collections::BTreeSet<_>>();
2231 let supported = harness_support_registry()
2232 .harnesses
2233 .into_iter()
2234 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2235 .collect::<Vec<_>>();
2236 if !params.harnesses.is_empty() && supported.len() != selected.len() {
2237 let known = supported
2238 .iter()
2239 .map(|harness| harness.id.as_str())
2240 .collect::<std::collections::BTreeSet<_>>();
2241 let missing = params
2242 .harnesses
2243 .iter()
2244 .filter(|id| !known.contains(id.as_str()))
2245 .map(HarnessId::as_str)
2246 .collect::<Vec<_>>();
2247 return Err(ServiceError::InvalidParams(format!(
2248 "unknown harness(es): {}",
2249 missing.join(", ")
2250 )));
2251 }
2252 let global_counts = params
2253 .include_sessions
2254 .then(|| self.session_counts(None, ¶ms.harnesses));
2255 let workspace_counts = params
2256 .include_sessions
2257 .then(|| {
2258 params
2259 .workspace
2260 .as_deref()
2261 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
2262 })
2263 .flatten();
2264 Ok(InventoryWork {
2265 params,
2266 supported,
2267 global_counts,
2268 workspace_counts,
2269 })
2270 }
2271
2272 #[cfg(feature = "adapter-api")]
2273 async fn harness_authentication_call(
2274 &self,
2275 method: &str,
2276 params: Value,
2277 ) -> std::result::Result<Value, ServiceError> {
2278 match method {
2279 "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2280 let params = decode::<HarnessAuthenticationParams>(params)?;
2281 serde_json::to_value(crate::inspect_harness_authentication(¶ms.harness).await)
2282 .map_err(|error| ServiceError::Operation(error.to_string()))
2283 }
2284 "harness.v1.harnesses.auth.begin" => {
2285 let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2286 let cwd = params
2287 .cwd
2288 .or_else(|| std::env::current_dir().ok())
2289 .unwrap_or_else(|| PathBuf::from("."));
2290 let plan = crate::harness_authentication_plan(
2291 ¶ms.harness,
2292 params.environment,
2293 params.method,
2294 &cwd,
2295 )
2296 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2297 serde_json::to_value(plan)
2298 .map_err(|error| ServiceError::Operation(error.to_string()))
2299 }
2300 _ => Err(ServiceError::MethodNotFound),
2301 }
2302 }
2303
2304 fn session_counts(
2305 &self,
2306 workspace: Option<&Path>,
2307 harnesses: &[HarnessId],
2308 ) -> BTreeMap<String, usize> {
2309 let mut counts = BTreeMap::new();
2310 for session in self
2311 .catalog
2312 .discover(&DiscoveryQuery {
2313 workspace: workspace.map(Path::to_path_buf),
2314 harnesses: harnesses.to_vec(),
2315 ..DiscoveryQuery::default()
2316 })
2317 .unwrap_or_default()
2318 {
2319 *counts
2320 .entry(session.locator.harness.as_str().to_string())
2321 .or_insert(0) += 1;
2322 }
2323 counts
2324 }
2325}
2326
2327#[async_trait::async_trait]
2328impl SdkService for HarnessSessionService {
2329 fn capabilities(&self) -> SdkCapabilities {
2330 SdkCapabilities::default()
2331 }
2332
2333 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2334 if request.operation == SdkOperation::Events {
2335 let events = self
2336 .poll_sdk_events()
2337 .await
2338 .into_iter()
2339 .map(|(_, event)| event)
2340 .collect::<Vec<_>>();
2341 return serde_json::to_value(events).map_err(|error| {
2342 SdkError::new(
2343 SdkErrorCode::Execution,
2344 request.operation,
2345 error.to_string(),
2346 )
2347 });
2348 }
2349 if self.runtimes.is_empty()
2350 && matches!(
2351 request.operation,
2352 SdkOperation::Input
2353 | SdkOperation::Interrupt
2354 | SdkOperation::Steer
2355 | SdkOperation::Respond
2356 | SdkOperation::Close
2357 )
2358 {
2359 return Err(SdkError::unsupported(request.operation));
2360 }
2361 let method = request
2362 .operation
2363 .method()
2364 .ok_or_else(|| SdkError::unsupported(request.operation))?;
2365 let result = match request.operation {
2366 SdkOperation::Discover
2367 | SdkOperation::Load
2368 | SdkOperation::Export
2369 | SdkOperation::ProfilesList
2370 | SdkOperation::ProfilesGet
2371 | SdkOperation::ProfilesCreate
2372 | SdkOperation::ProfilesDelete
2373 | SdkOperation::SkillsList
2374 | SdkOperation::SkillsInstall
2375 | SdkOperation::SkillsRemove
2376 | SdkOperation::ChannelsList
2377 | SdkOperation::RoutesList
2378 | SdkOperation::TriggersList
2379 | SdkOperation::ChannelsStatus
2380 | SdkOperation::MemoryShow
2381 | SdkOperation::MemorySearch
2382 | SdkOperation::JobsList
2383 | SdkOperation::JobsGet
2384 | SdkOperation::JobsCreate
2385 | SdkOperation::JobsUpdate
2386 | SdkOperation::JobsPause
2387 | SdkOperation::JobsResume
2388 | SdkOperation::JobsRun
2389 | SdkOperation::JobsDelete
2390 | SdkOperation::JobsNotepad
2391 | SdkOperation::JobsNotepadSet
2392 | SdkOperation::JobsNotepadDelete
2393 | SdkOperation::RunsList
2394 | SdkOperation::RunsGet
2395 | SdkOperation::ApprovalsList
2396 | SdkOperation::OrchestrationLoad
2397 | SdkOperation::OrchestrationSave
2398 | SdkOperation::OrchestrationCompile
2399 | SdkOperation::OrchestrationDecompile
2400 | SdkOperation::OrchestrationImport
2401 | SdkOperation::OrchestrationExport
2402 | SdkOperation::WorkflowLoad => self.call(method, request.params),
2403 SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2406 SdkOperation::Start
2407 | SdkOperation::Resume
2408 | SdkOperation::Input
2409 | SdkOperation::Interrupt
2410 | SdkOperation::Steer
2411 | SdkOperation::Respond
2412 | SdkOperation::Close => self.runtime_call(method, request.params).await,
2413 SdkOperation::SessionsNew => {
2418 self.mutate_session(crate::SessionVerb::New, request.params)
2419 .await
2420 }
2421 SdkOperation::SessionsReset => {
2422 self.mutate_session(crate::SessionVerb::Reset, request.params)
2423 .await
2424 }
2425 SdkOperation::SessionsArchive => {
2426 self.mutate_session(crate::SessionVerb::Archive, request.params)
2427 .await
2428 }
2429 SdkOperation::SessionsDelete => {
2430 self.mutate_session(crate::SessionVerb::Delete, request.params)
2431 .await
2432 }
2433 SdkOperation::Events => unreachable!("handled before method dispatch"),
2434 };
2435 result.map_err(|error| sdk_error(request.operation, error))
2436 }
2437
2438 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2439 Ok(self
2440 .poll_sdk_events()
2441 .await
2442 .into_iter()
2443 .map(|(_, event)| event)
2444 .collect())
2445 }
2446}
2447
2448#[cfg(feature = "adapter-api")]
2449struct HostedRuntimeLease {
2450 connection: HostedHarnessConnection,
2451 _host: std::sync::Arc<HostedHarnessRuntime>,
2452 _registration: LiveRuntimeRegistration,
2453 _server: crate::server::FrontendHttpServer,
2454}
2455
2456#[async_trait::async_trait]
2457#[cfg(feature = "adapter-api")]
2458impl RuntimeConnection for HostedRuntimeLease {
2459 fn handle(&self) -> &crate::RuntimeHandle {
2460 self.connection.handle()
2461 }
2462
2463 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2464 self.connection.send_input(input).await
2465 }
2466
2467 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2468 self.connection.next_event().await
2469 }
2470
2471 async fn interrupt(&mut self) -> crate::Result<()> {
2472 self.connection.interrupt().await
2473 }
2474
2475 async fn steer(&mut self, text: String) -> crate::Result<()> {
2478 self.connection.steer(text).await
2479 }
2480
2481 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2482 self.connection.respond(request_id, response).await
2483 }
2484
2485 async fn close(&mut self) -> crate::Result<()> {
2486 self.connection.close().await
2487 }
2488}
2489
2490struct InventoryWork {
2493 params: HarnessInventoryParams,
2494 supported: Vec<crate::HarnessSupportDescriptor>,
2495 global_counts: Option<BTreeMap<String, usize>>,
2496 workspace_counts: Option<BTreeMap<String, usize>>,
2497}
2498
2499async fn run_session_mutation(
2506 verb: crate::SessionVerb,
2507 mutation: &crate::SessionMutation,
2508) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2509 let door =
2516 crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2517 if let crate::SessionDoor::Http = door {
2518 return crate::sessions_control::mutate(verb, mutation)
2519 .await
2520 .map_err(session_control_error);
2521 }
2522 let mutation = mutation.clone();
2523 tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2524 .await
2525 .map_err(|error| {
2526 ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2527 })?
2528 .map_err(session_control_error)
2529}
2530
2531async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2534 let InventoryWork {
2535 params,
2536 supported,
2537 global_counts,
2538 workspace_counts,
2539 } = work;
2540 let probes = supported.into_iter().map(|descriptor| {
2541 let global = global_counts
2542 .as_ref()
2543 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2544 let workspace = workspace_counts
2545 .as_ref()
2546 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2547 probe_harness(descriptor, ¶ms, global, workspace)
2548 });
2549 let harnesses = futures::future::join_all(probes).await;
2550 serde_json::to_value(HarnessInventoryReport {
2551 probe: params.probe,
2552 workspace: params.workspace,
2553 harnesses,
2554 })
2555 .map_err(|error| ServiceError::Operation(error.to_string()))
2556}
2557
2558async fn probe_harness(
2559 descriptor: crate::HarnessSupportDescriptor,
2560 params: &HarnessInventoryParams,
2561 global: Option<usize>,
2562 workspace: Option<usize>,
2563) -> LocalHarness {
2564 let launch = descriptor.runtime.default_launch.as_ref();
2565 let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2570 .then(crate::orchestrator::daemon_entry)
2571 .and_then(Result::ok);
2572 let executable = match &orchestrator_entry {
2573 Some(entry) => Some(entry.clone()),
2574 None => launch.and_then(|launch| find_executable(&launch.program)),
2575 };
2576 let installed = executable.is_some();
2577 let version = if params.skip_versions || orchestrator_entry.is_some() {
2578 None
2581 } else {
2582 match executable.as_deref() {
2583 Some(path) => executable_version(path).await,
2584 None => None,
2585 }
2586 };
2587 let configured = auth_evidence(descriptor.id.as_str());
2588 let mut auth = if configured {
2589 HarnessAuthState::Configured
2590 } else if matches!(
2591 descriptor.id.as_str(),
2592 HarnessId::CLAUDE_CODE | HarnessId::CODEX
2593 ) {
2594 HarnessAuthState::Required
2599 } else {
2600 HarnessAuthState::Unknown
2601 };
2602 let mut runtime = if installed {
2603 HarnessRuntimeState::Degraded
2604 } else {
2605 HarnessRuntimeState::Unavailable
2606 };
2607 let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2608 let mut reason = (!installed).then(|| {
2609 if is_orchestrator {
2610 format!(
2611 "{} is supported but its daemon entry `{}` was not found",
2612 descriptor.display_name,
2613 crate::orchestrator::DAEMON_ENTRY
2614 )
2615 } else {
2616 format!(
2617 "{} is supported but `{}` was not found on PATH",
2618 descriptor.display_name,
2619 launch
2620 .map(|launch| launch.program.as_str())
2621 .unwrap_or("executable")
2622 )
2623 }
2624 });
2625 let mut repair = (!installed).then(|| {
2626 if is_orchestrator {
2627 format!(
2628 "Install the `supercode-orchestrator` package so `{}` resolves.",
2629 crate::orchestrator::DAEMON_ENTRY
2630 )
2631 } else {
2632 format!(
2633 "Install {} and ensure `{}` is on PATH.",
2634 descriptor.display_name,
2635 launch
2636 .map(|launch| launch.program.as_str())
2637 .unwrap_or("its executable")
2638 )
2639 }
2640 });
2641
2642 if installed && params.probe == HarnessProbeLevel::Handshake {
2643 let backend_params = RuntimeBackendParams {
2644 harness: descriptor.id.clone(),
2645 protocol: None,
2646 launch: None,
2647 base_url: None,
2648 policy: RuntimePolicy::Default,
2649 };
2650 match runtime_backend(&backend_params) {
2651 Ok(backend) => {
2652 let cwd = params
2653 .workspace
2654 .clone()
2655 .or_else(|| std::env::current_dir().ok())
2656 .unwrap_or_else(|| PathBuf::from("."));
2657 let isolated = descriptor
2658 .runtime
2659 .default_launch
2660 .clone()
2661 .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2662 let Some(isolated) = isolated else {
2663 reason = Some(
2664 "No-prompt runtime handshake could not create its isolated harness home."
2665 .into(),
2666 );
2667 repair = Some(
2668 "Check temporary-directory permissions, then run the handshake probe again."
2669 .into(),
2670 );
2671 let running = probe_running_instance(descriptor.id.as_str());
2672 return LocalHarness {
2673 gateway: gateway_health(
2674 descriptor.id.as_str(),
2675 installed,
2676 running.as_ref(),
2677 version.as_deref(),
2678 ),
2679 id: descriptor.id,
2680 display_name: descriptor.display_name,
2681 supported: true,
2682 installed,
2683 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2684 version,
2685 auth,
2686 runtime,
2687 protocol: descriptor.runtime.protocol,
2688 capabilities: descriptor.runtime.capabilities.clone(),
2689 effective_capabilities: descriptor.runtime.capabilities,
2690 sessions: HarnessSessionCounts { global, workspace },
2691 running,
2692 reason,
2693 repair,
2694 };
2695 };
2696 match tokio::time::timeout(
2697 Duration::from_secs(30),
2698 backend.start(RuntimeStartRequest {
2699 cwd,
2700 launch: Some(isolated.launch.clone()),
2701 mcp_servers: Vec::new(),
2702 approval_policy: None,
2703 }),
2704 )
2705 .await
2706 {
2707 Ok(Ok(mut connection)) => {
2708 match stabilize_handshake(connection.as_mut()).await {
2709 Ok(()) => {
2710 auth = HarnessAuthState::Ready;
2711 runtime = HarnessRuntimeState::Ready;
2712 reason = Some(
2713 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2714 .into(),
2715 );
2716 repair = None;
2717 }
2718 Err(message) => {
2719 auth = if looks_like_auth_error(&message) {
2720 HarnessAuthState::Required
2721 } else if configured {
2722 HarnessAuthState::Configured
2723 } else {
2724 HarnessAuthState::Unknown
2725 };
2726 reason = Some(format!(
2727 "No-prompt runtime handshake became unhealthy during startup: {message}"
2728 ));
2729 repair = Some(if auth == HarnessAuthState::Required {
2730 format!(
2731 "Run `{}` interactively once and complete sign-in, then probe again.",
2732 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2733 )
2734 } else {
2735 "Run the harness directly to inspect its startup failure, then probe again."
2736 .into()
2737 });
2738 }
2739 }
2740 let _ =
2741 tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2742 }
2743 Ok(Err(error)) => {
2744 let message = truncate_text(&error.to_string(), 500);
2745 auth = if looks_like_auth_error(&message) {
2746 HarnessAuthState::Required
2747 } else if configured {
2748 HarnessAuthState::Configured
2749 } else {
2750 HarnessAuthState::Unknown
2751 };
2752 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2753 repair = Some(if auth == HarnessAuthState::Required {
2754 format!(
2755 "Run `{}` interactively once and complete sign-in, then probe again.",
2756 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2757 )
2758 } else {
2759 "Check the harness installation and run the handshake probe again."
2760 .into()
2761 });
2762 }
2763 Err(_) => {
2764 reason =
2765 Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2766 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2767 }
2768 }
2769 let _ = isolated.cleanup();
2778 tokio::time::sleep(Duration::from_millis(250)).await;
2779 if let Err(error) = isolated.cleanup() {
2780 auth = if configured {
2781 HarnessAuthState::Configured
2782 } else {
2783 HarnessAuthState::Unknown
2784 };
2785 runtime = HarnessRuntimeState::Degraded;
2786 reason = Some(format!(
2787 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2788 ));
2789 repair = Some(
2790 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2791 .into(),
2792 );
2793 }
2794 }
2795 Err(error) => {
2796 reason = Some(error_message(error));
2797 }
2798 }
2799 } else if installed && configured {
2800 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2801 } else if installed && auth == HarnessAuthState::Required {
2802 reason = Some("Executable found, but no native authentication evidence is present.".into());
2803 repair = Some(format!(
2804 "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2805 descriptor.id.as_str()
2806 ));
2807 } else if installed {
2808 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2809 repair = Some(format!(
2810 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2811 launch
2812 .map(|launch| launch.program.as_str())
2813 .unwrap_or("the harness")
2814 ));
2815 }
2816
2817 let effective_capabilities = if installed {
2818 descriptor.runtime.capabilities.clone()
2819 } else {
2820 unavailable_capabilities()
2821 };
2822 let running = probe_running_instance(descriptor.id.as_str());
2823 LocalHarness {
2824 gateway: gateway_health(
2825 descriptor.id.as_str(),
2826 installed,
2827 running.as_ref(),
2828 version.as_deref(),
2829 ),
2830 id: descriptor.id,
2831 display_name: descriptor.display_name,
2832 supported: true,
2833 installed,
2834 executable: executable.map(|path| path.to_string_lossy().into_owned()),
2835 version,
2836 auth,
2837 runtime,
2838 protocol: descriptor.runtime.protocol,
2839 capabilities: descriptor.runtime.capabilities,
2840 effective_capabilities,
2841 sessions: HarnessSessionCounts { global, workspace },
2842 running,
2843 reason,
2844 repair,
2845 }
2846}
2847
2848async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2849 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2850 loop {
2851 let now = tokio::time::Instant::now();
2852 if now >= deadline {
2853 return Ok(());
2854 }
2855 match tokio::time::timeout(deadline - now, connection.next_event()).await {
2856 Err(_) => return Ok(()),
2857 Ok(Ok(Some(event))) => {
2858 if let Some(message) = handshake_event_failure(&event) {
2859 return Err(truncate_text(&message, 500));
2860 }
2861 }
2862 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2863 Ok(Err(error)) => return Err(error.to_string()),
2864 }
2865 }
2866}
2867
2868fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2869 let detail = event
2870 .payload
2871 .get("message")
2872 .or_else(|| event.payload.get("line"))
2873 .and_then(Value::as_str)
2874 .unwrap_or(event.kind.as_str());
2875 match event.kind.as_str() {
2876 "transport_closed" => Some("runtime transport closed during startup".into()),
2877 "transport_error" => Some(format!("runtime transport error: {detail}")),
2878 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2879 _ => None,
2884 }
2885}
2886
2887fn indexed_claude_window(
2888 locator: &SessionLocator,
2889 options: &SessionLoadOptions,
2890) -> std::result::Result<Option<Value>, ServiceError> {
2891 use supercode_interchange::session::ClaudeReadIndex;
2892 if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2895 || options.include_subagents != Some(false)
2896 {
2897 return Ok(None);
2898 }
2899 let crate::StorageLocator::File { path } = &locator.storage else {
2900 return Ok(None);
2901 };
2902 if !ClaudeReadIndex::supports(path)
2903 .map_err(|error| ServiceError::Operation(error.to_string()))?
2904 {
2905 return Ok(None);
2906 }
2907 let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2908 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2909 let total = index.len();
2910 let (offset, end) = projected_message_window(total, options);
2911 let session = index
2912 .read_messages(offset..end)
2913 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2914 let summary = index
2915 .read_summary()
2916 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2917 let selected_options = SessionLoadOptions {
2918 message_offset: None,
2919 message_limit: None,
2920 message_tail: None,
2921 ..options.clone()
2922 };
2923 let mut selected = projected_session_json(&session, &selected_options);
2924 selected["raw_record_count"] = json!(index.raw_record_count());
2925 Ok(Some(json!({
2926 "session": selected,
2927 "summary": projected_session_summary(&summary, options),
2928 "window": {
2929 "has_more": offset > 0 || end < total, "has_newer": end < total,
2930 "has_older": offset > 0, "newer_items": index.item_count(end..total),
2931 "offset": offset, "older_items": index.item_count(0..offset),
2932 "returned": end - offset, "total_messages": total,
2933 }
2934 })))
2935}
2936
2937fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
2938 let total_messages = session.messages.len();
2939 let (offset, end) = projected_message_window(total_messages, options);
2940 json!({
2941 "session": projected_session_json(session, options),
2942 "summary": projected_session_summary(session, options),
2943 "window": {
2944 "has_more": offset > 0 || end < total_messages,
2945 "has_newer": end < total_messages,
2946 "has_older": offset > 0,
2947 "newer_items": normalized_item_count(&session.messages[end..]),
2948 "offset": offset,
2949 "older_items": normalized_item_count(&session.messages[..offset]),
2950 "returned": end.saturating_sub(offset),
2951 "total_messages": total_messages,
2952 }
2953 })
2954}
2955
2956fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
2957 messages
2958 .iter()
2959 .map(|message| {
2960 let conversation = usize::from(
2961 matches!(message.role, Role::Assistant | Role::User)
2962 && message_has_content(message),
2963 );
2964 let tool_result =
2965 usize::from(message.role == Role::Tool && message_has_content(message));
2966 conversation + tool_result + message.tool_calls().len()
2967 })
2968 .sum()
2969}
2970
2971fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
2972 let mut conversational = session.messages.iter().filter(|message| {
2973 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
2974 });
2975 let first_message = conversational.clone().next();
2976 let last_message = conversational.next_back();
2977 let mut assistant = session
2978 .messages
2979 .iter()
2980 .filter(|message| message.role == Role::Assistant && message_has_content(message));
2981 let first_assistant_message = assistant.clone().next();
2982 let last_assistant_message = assistant.next_back();
2983 let end_of_turn = session
2984 .messages
2985 .iter()
2986 .rev()
2987 .find(|message| message.role != Role::System)
2988 .is_some_and(|message| {
2989 message.role == Role::Assistant
2990 && message_has_content(message)
2991 && message.tool_calls().is_empty()
2992 && message.metadata.get("phase").map(String::as_str) != Some("commentary")
2994 });
2995 let project = |message: Option<&crate::ChatMessage>| {
2996 message.map(|message| project_inline_media(message_json(message), options))
2997 };
2998 json!({
2999 "end_of_turn": end_of_turn,
3000 "first_assistant_message": project(first_assistant_message),
3001 "first_message": project(first_message),
3002 "last_assistant_message": project(last_assistant_message),
3003 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3004 "last_message": project(last_message),
3005 })
3006}
3007
3008fn message_has_content(message: &crate::ChatMessage) -> bool {
3009 message
3010 .content
3011 .as_deref()
3012 .is_some_and(|content| !content.trim().is_empty())
3013 || message
3014 .content_parts
3015 .as_ref()
3016 .is_some_and(|parts| !parts.is_empty())
3017}
3018
3019fn message_text(message: &crate::ChatMessage) -> String {
3020 if let Some(content) = &message.content {
3021 return content.clone();
3022 }
3023 message
3024 .content_parts
3025 .as_ref()
3026 .into_iter()
3027 .flatten()
3028 .filter_map(|part| part.get("text").and_then(Value::as_str))
3029 .collect::<Vec<_>>()
3030 .join("\n")
3031}
3032
3033fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3034 let (offset, end) = projected_message_window(session.messages.len(), options);
3035 let messages = session.messages[offset..end]
3036 .iter()
3037 .map(|message| project_inline_media(message_json(message), options))
3038 .collect::<Vec<_>>();
3039 let subagents = if options.include_subagents.unwrap_or(true) {
3040 let subagent_options = SessionLoadOptions {
3045 message_limit: None,
3046 message_offset: None,
3047 message_tail: None,
3048 ..options.clone()
3049 };
3050 session
3051 .subagents
3052 .iter()
3053 .map(|subagent| projected_session_json(subagent, &subagent_options))
3054 .collect::<Vec<_>>()
3055 } else {
3056 Vec::new()
3057 };
3058 json!({
3059 "source": match session.meta.source {
3060 SessionSource::ClaudeCode => "claude_code",
3061 SessionSource::Codex => "codex",
3062 SessionSource::Gemini => "gemini",
3063 SessionSource::Goose => "goose",
3064 SessionSource::Grok => "grok",
3065 SessionSource::Native => "native",
3066 SessionSource::OpenClaw => "openclaw",
3067 SessionSource::Hermes => "hermes",
3068 SessionSource::OpenCode => "opencode",
3069 SessionSource::Pi => "pi",
3070 },
3071 "session_id": session.meta.session_id,
3072 "ended_at": session.meta.ended_at,
3073 "end_reason": session.meta.end_reason,
3074 "model": session.meta.model,
3075 "cwd": session.meta.cwd,
3076 "system_prompt": session.meta.system_prompt,
3077 "agent_id": session.meta.agent_id,
3078 "parent_tool_use_id": session.meta.parent_tool_use_id,
3079 "lineage": session.meta.lineage,
3080 "messages": messages,
3081 "subagents": subagents,
3082 "raw_record_count": session.raw.len(),
3083 "parse_error_lines": session.parse_error_lines,
3084 })
3085}
3086
3087fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3088 if let Some(tail) = options.message_tail {
3089 return (total.saturating_sub(tail), total);
3090 }
3091 let offset = options.message_offset.unwrap_or(0).min(total);
3092 let end = options
3093 .message_limit
3094 .map(|limit| offset.saturating_add(limit).min(total))
3095 .unwrap_or(total);
3096 (offset, end)
3097}
3098
3099fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3100 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3101 return message;
3102 };
3103 for part in parts {
3104 let Some(url) = part
3105 .get("image_url")
3106 .and_then(|image| image.get("url"))
3107 .and_then(Value::as_str)
3108 else {
3109 continue;
3110 };
3111 let Some(rest) = url.strip_prefix("data:") else {
3112 continue;
3113 };
3114 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3115 continue;
3116 };
3117 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3118 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3119 let decoded_bytes = decoded_bytes.saturating_sub(padding);
3120 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3121 || options
3122 .max_inline_media_bytes
3123 .is_some_and(|limit| decoded_bytes > limit);
3124 if should_elide {
3125 *part = json!({
3126 "type": "media_reference",
3127 "media_type": media_type,
3128 "encoding": "base64",
3129 "encoded_bytes": encoded.len(),
3130 "decoded_bytes": decoded_bytes,
3131 "omitted": true,
3132 });
3133 }
3134 }
3135 message
3136}
3137
3138#[derive(Deserialize)]
3139struct LocatorParams {
3140 locator: SessionLocator,
3141 #[serde(default)]
3154 fidelity: Option<Fidelity>,
3155 #[serde(default)]
3158 view: Option<SessionReadView>,
3159}
3160
3161#[derive(Deserialize)]
3162struct SessionReadView {
3163 #[serde(default)]
3166 tail_messages: Option<usize>,
3167 #[serde(default)]
3170 include_subagents: bool,
3171 #[serde(default)]
3173 display_history: bool,
3174 #[serde(default)]
3177 max_message_chars: Option<usize>,
3178}
3179
3180impl LocatorParams {
3181 fn read_fidelity(&self) -> Fidelity {
3182 self.fidelity.unwrap_or(Fidelity::Semantic)
3183 }
3184
3185 fn include_subagents(&self) -> bool {
3186 self.view
3187 .as_ref()
3188 .map(|view| view.include_subagents)
3189 .unwrap_or(true)
3190 }
3191
3192 fn tail_messages(&self) -> Option<usize> {
3193 self.view
3194 .as_ref()
3195 .and_then(|view| view.tail_messages)
3196 .map(|limit| limit.clamp(1, 5_000))
3197 }
3198
3199 fn display_history(&self) -> bool {
3200 self.view.as_ref().is_some_and(|view| view.display_history)
3201 }
3202
3203 fn max_message_chars(&self) -> Option<usize> {
3204 self.view
3205 .as_ref()
3206 .and_then(|view| view.max_message_chars)
3207 .map(|limit| limit.clamp(256, 64_000))
3208 }
3209
3210 fn bound_session(&self, session: &mut Session) {
3211 bound_session_view(session, self.tail_messages(), self.max_message_chars());
3212 }
3213}
3214
3215#[derive(Debug, Clone, Copy, Default, Deserialize)]
3216#[serde(rename_all = "snake_case")]
3217enum InlineMediaMode {
3218 #[default]
3219 Full,
3220 Metadata,
3221}
3222
3223#[derive(Debug, Clone, Default, Deserialize)]
3224#[serde(default)]
3225struct SessionLoadOptions {
3226 include_subagents: Option<bool>,
3227 inline_media: InlineMediaMode,
3228 max_inline_media_bytes: Option<usize>,
3229 message_limit: Option<usize>,
3230 message_offset: Option<usize>,
3231 message_tail: Option<usize>,
3232}
3233
3234impl SessionLoadOptions {
3235 fn validate(&self) -> std::result::Result<(), ServiceError> {
3236 if self.message_tail.is_some()
3237 && (self.message_limit.is_some() || self.message_offset.is_some())
3238 {
3239 return Err(ServiceError::InvalidParams(
3240 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3241 .into(),
3242 ));
3243 }
3244 Ok(())
3245 }
3246}
3247
3248#[derive(Deserialize)]
3249struct LoadSessionParams {
3250 #[serde(flatten)]
3251 read: LocatorParams,
3252 #[serde(default)]
3253 options: Option<SessionLoadOptions>,
3254}
3255
3256#[derive(Deserialize)]
3257struct UnfollowParams {
3258 subscription: String,
3259}
3260
3261#[derive(Debug, Deserialize)]
3262#[serde(deny_unknown_fields)]
3263struct IndexResizeParams {
3264 subscription: String,
3265 limit: usize,
3266}
3267
3268#[derive(Deserialize)]
3269struct ActivitySubscribeParams {
3270 locators: Vec<SessionLocator>,
3271 #[serde(default)]
3272 homes: crate::HarnessHomes,
3273}
3274
3275#[derive(Deserialize)]
3276struct MessageSessionParams {
3277 locator: SessionLocator,
3278 text: String,
3279 #[serde(default)]
3280 subject: Option<String>,
3281 #[serde(default)]
3283 idempotency_key: Option<String>,
3284 #[serde(default)]
3286 channel: bool,
3287 #[serde(default)]
3291 from_name: Option<String>,
3292 #[serde(default)]
3294 voice_for: Option<crate::mailbox::MailAddress>,
3295 #[serde(default)]
3297 sender_name: Option<String>,
3298 #[serde(default)]
3300 in_reply_to: Option<String>,
3301 #[serde(default)]
3303 notify_when_idle: bool,
3304 #[serde(default)]
3309 as_user: bool,
3310 #[serde(default)]
3313 homes: crate::HarnessHomes,
3314}
3315
3316#[derive(Deserialize)]
3317#[serde(deny_unknown_fields)]
3318struct InboxParams {
3319 #[serde(default)]
3321 from_name: Option<String>,
3322 #[serde(default)]
3324 address: Option<String>,
3325 #[serde(default)]
3327 all: bool,
3328}
3329
3330#[derive(Deserialize)]
3331#[serde(deny_unknown_fields)]
3332struct HarnessSettingsParams {
3333 harness: String,
3334}
3335
3336#[derive(Deserialize)]
3337#[serde(deny_unknown_fields)]
3338struct ConfigureHarnessParams {
3339 harness: String,
3340 #[serde(default)]
3341 changes: Vec<crate::HarnessSettingChange>,
3342 #[serde(default)]
3343 expected_revision: Option<String>,
3344}
3345
3346fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3347 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3348 Ok(report) => (
3349 serde_json::to_value(report).unwrap_or(Value::Null),
3350 Value::Null,
3351 ),
3352 Err(error) => (
3353 Value::Null,
3354 Value::String(format!(
3355 "Volter Harness could not inspect Claude Code inbound controls: {error}"
3356 )),
3357 ),
3358 }
3359}
3360
3361async fn message_live_session(params: &MessageSessionParams) -> Value {
3375 use crate::mail_route::{Delivered, NoDoor, Refused};
3376 let (inbound_controls, inbound_controls_error) =
3377 claude_inbound_controls_or_error(¶ms.homes);
3378 let refused = |reason: &str, message: String| {
3379 json!({
3380 "delivered_to_bus": false,
3381 "refusal": {"reason": reason, "message": message},
3382 "inbound_controls": inbound_controls,
3383 "inbound_controls_error": inbound_controls_error,
3384 })
3385 };
3386 if params.text.trim().is_empty() {
3387 return refused(
3388 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3389 "refusing to deliver an empty message".into(),
3390 );
3391 }
3392 let sender = match operator_address(params.from_name.as_deref()) {
3393 Ok(sender) => sender,
3394 Err(message) => return refused("invalid_sender", message),
3395 };
3396 let receiver = match crate::mailbox::MailAddress::new(
3397 crate::mailbox::local_machine_name(),
3398 params.locator.harness.as_str(),
3399 ¶ms.locator.session_id,
3400 ) {
3401 Ok(receiver) => receiver,
3402 Err(error) => return refused("delivery_failed", error.to_string()),
3403 };
3404 if params
3405 .voice_for
3406 .as_ref()
3407 .is_some_and(|voice_for| voice_for != &receiver)
3408 {
3409 return refused(
3410 "invalid_sender",
3411 "a voice front must represent the receiving session".into(),
3412 );
3413 }
3414 if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3416 && crate::runtime_mail::controlled_runtime(
3417 HarnessId::CLAUDE_CODE,
3418 ¶ms.locator.session_id,
3419 )
3420 .is_none()
3421 {
3422 if let Err(refusal) =
3423 crate::claude_peer::resolve_live_session(¶ms.homes, ¶ms.locator.session_id)
3424 {
3425 return refused(refusal.reason.as_str(), refusal.message);
3426 }
3427 }
3428 if params.as_user {
3429 return message_as_user(
3430 params,
3431 sender,
3432 receiver,
3433 inbound_controls,
3434 inbound_controls_error,
3435 )
3436 .await;
3437 }
3438 let door = match crate::mail_route::door_for(¶ms.homes, &receiver) {
3439 Ok(door) => door,
3440 Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3441 return refused(
3442 crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3443 format!(
3444 "no running `{}` session `{}` is reachable; its transcript is persisted only",
3445 params.locator.harness.as_str(),
3446 params.locator.session_id
3447 ),
3448 )
3449 }
3450 };
3451 let mut envelope = match crate::mailbox::Envelope::new(
3452 sender.clone(),
3453 params
3454 .sender_name
3455 .clone()
3456 .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3457 if params.channel {
3458 crate::mailbox::MailKind::Channel
3459 } else {
3460 crate::mailbox::MailKind::Peer
3461 },
3462 crate::mailbox::ReplyVia::Command,
3463 params.text.clone(),
3464 ) {
3465 Ok(envelope) => envelope,
3466 Err(error) => return refused("delivery_failed", error.to_string()),
3467 };
3468 if let Some(key) = ¶ms.idempotency_key {
3469 if key.is_empty() || key.len() > 256 {
3470 return refused("invalid_key", "idempotency_key needs 1–256 bytes".into());
3471 }
3472 let identity = serde_json::to_vec(&(sender.to_string(), receiver.to_string(), key))
3473 .expect("string tuple serializes");
3474 envelope.id = format!("m-{}", &blake3::hash(&identity).to_hex()[..24]);
3475 let mailbox = match crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &receiver) {
3476 Ok(mailbox) => mailbox,
3477 Err(error) => return refused("delivery_failed", error.to_string()),
3478 };
3479 match mailbox.find(&envelope.id) {
3480 Ok(Some(previous)) => {
3481 let saved = previous.envelope;
3482 let subject = params.subject.as_deref().map(|value| {
3483 value
3484 .lines()
3485 .next()
3486 .unwrap_or("")
3487 .chars()
3488 .take(200)
3489 .collect::<String>()
3490 });
3491 if saved.body != envelope.body
3492 || saved.subject != subject
3493 || saved.kind != envelope.kind
3494 || saved.from_name != envelope.from_name
3495 || saved.in_reply_to != params.in_reply_to
3496 || saved.voice_for != params.voice_for
3497 {
3498 return refused(
3499 "idempotency_conflict",
3500 "this key already names different mail".into(),
3501 );
3502 }
3503 }
3504 Ok(None) => {}
3505 Err(error) => return refused("delivery_failed", error.to_string()),
3506 }
3507 }
3508 envelope.in_reply_to = params.in_reply_to.clone();
3509 envelope.voice_for = params.voice_for.clone();
3510 envelope.subject = params.subject.as_deref().map(|value| {
3511 value
3512 .lines()
3513 .next()
3514 .unwrap_or("")
3515 .chars()
3516 .take(200)
3517 .collect()
3518 });
3519 let how = match crate::mail_route::deliver(
3520 &envelope,
3521 &receiver,
3522 &door,
3523 true,
3524 params.notify_when_idle,
3525 )
3526 .await
3527 {
3528 Err(detail) => {
3529 return refused(
3530 crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3531 detail,
3532 )
3533 }
3534 Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3535 Ok(Err(Refused::TooLong(bytes))) => {
3536 return refused(
3537 "too_long",
3538 format!(
3539 "the message is {bytes} bytes; the limit is {}",
3540 crate::mail_route::MAX_RELAYED_BYTES
3541 ),
3542 )
3543 }
3544 Ok(Ok(delivered)) => match delivered {
3545 Delivered::Steered => "steered",
3546 Delivered::Started => "started",
3547 Delivered::Native { busy: true } => "next_tool_call",
3548 Delivered::Native { busy: false } => "started",
3549 Delivered::Hooked => "hook",
3550 Delivered::HookWoken => "woken",
3551 Delivered::Queued => "queued",
3552 Delivered::Stored => "stored",
3553 Delivered::Operator => "filed",
3554 },
3555 };
3556 json!({
3557 "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3558 "message_id": envelope.id,
3559 "reply_to": sender.to_string(),
3560 "target": {
3561 "harness": params.locator.harness.as_str(),
3562 "session_id": params.locator.session_id,
3563 "name": match &door {
3564 crate::mail_route::Door::Native(session) => Some(session.name.clone()),
3565 _ => None,
3566 },
3567 },
3568 "delivery": {"door": door.name(), "how": how},
3569 "inbound_controls": inbound_controls,
3570 "inbound_controls_error": inbound_controls_error,
3571 })
3572}
3573
3574async fn message_as_user(
3578 params: &MessageSessionParams,
3579 sender: crate::mailbox::MailAddress,
3580 receiver: crate::mailbox::MailAddress,
3581 inbound_controls: Value,
3582 inbound_controls_error: Value,
3583) -> Value {
3584 let envelope = match crate::mailbox::Envelope::new(
3585 sender.clone(),
3586 format!("{}@{}", sender.session_id, sender.machine),
3587 crate::mailbox::MailKind::User,
3588 crate::mailbox::ReplyVia::None,
3589 params.text.clone(),
3590 ) {
3591 Ok(envelope) => envelope,
3592 Err(error) => {
3593 return json!({
3594 "delivered_to_bus": false,
3595 "refusal": {"reason": "delivery_failed", "message": error.to_string()},
3596 "inbound_controls": inbound_controls,
3597 "inbound_controls_error": inbound_controls_error,
3598 })
3599 }
3600 };
3601 match crate::mail_route::deliver_user_turn(¶ms.homes, &envelope, &receiver).await {
3602 Ok(turn) => json!({
3603 "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
3604 "message_id": envelope.id,
3605 "target": {
3606 "harness": params.locator.harness.as_str(),
3607 "session_id": params.locator.session_id,
3608 },
3609 "delivery": {
3610 "door": match turn {
3611 crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
3612 _ => "pane",
3613 },
3614 "how": turn.as_str(),
3615 },
3616 "inbound_controls": inbound_controls,
3617 "inbound_controls_error": inbound_controls_error,
3618 }),
3619 Err(message) => json!({
3620 "delivered_to_bus": false,
3621 "refusal": {"reason": "no_user_door", "message": message},
3622 "inbound_controls": inbound_controls,
3623 "inbound_controls_error": inbound_controls_error,
3624 }),
3625 }
3626}
3627
3628#[derive(Deserialize)]
3629#[serde(deny_unknown_fields)]
3630struct ActivityUnderParams {
3631 pids: Vec<u32>,
3633 #[serde(default)]
3634 homes: crate::HarnessHomes,
3635}
3636
3637async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
3641 let params = decode::<ActivityUnderParams>(params)?;
3642 if params.pids.len() > 1024 {
3643 return Err(ServiceError::InvalidParams(
3644 "sessions.activity_under accepts at most 1024 pids".into(),
3645 ));
3646 }
3647 let found = crate::session_activity::activity_under(¶ms.pids, ¶ms.homes)
3648 .await
3649 .map_err(ServiceError::Sdk)?;
3650 Ok(json!({
3651 "activities": found
3652 .into_iter()
3653 .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
3654 .collect::<Vec<_>>(),
3655 }))
3656}
3657
3658fn operator_address(
3661 name: Option<&str>,
3662) -> std::result::Result<crate::mailbox::MailAddress, String> {
3663 let name: String = name
3664 .unwrap_or("supercode")
3665 .trim()
3666 .chars()
3667 .map(|character| {
3668 if character.is_whitespace() || character == '@' {
3669 '-'
3670 } else {
3671 character
3672 }
3673 })
3674 .collect();
3675 if name.is_empty() {
3676 return Err("from_name must not be empty".into());
3677 }
3678 crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
3679 .map_err(|error| error.to_string())
3680}
3681
3682fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
3686 let address = match (¶ms.address, ¶ms.from_name) {
3687 (Some(address), None) => crate::mailbox::MailAddress::parse(address)
3688 .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
3689 (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
3690 (Some(_), Some(_)) => {
3691 return Err(ServiceError::InvalidParams(
3692 "sessions.inbox takes from_name or address, not both".into(),
3693 ))
3694 }
3695 };
3696 let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
3697 let mailbox =
3698 crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
3699 let claimed = mailbox.claim_unread().map_err(operation)?;
3700 let mut messages: Vec<Value> = Vec::new();
3701 if params.all {
3702 for stored in mailbox.list().map_err(operation)? {
3703 if stored.state == crate::mailbox::MailState::Read {
3704 messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3705 }
3706 }
3707 }
3708 for stored in &claimed {
3709 messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3710 }
3711 for stored in &claimed {
3712 mailbox.acknowledge(stored).map_err(operation)?;
3713 }
3714 Ok(json!({"address": address.to_string(), "messages": messages}))
3715}
3716
3717#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3722struct FollowedSource {
3723 harness: String,
3724 session_id: String,
3725 reported: Option<String>,
3726}
3727
3728#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3729struct ActivitySubscription {
3730 locators: Vec<SessionLocator>,
3731 homes: crate::HarnessHomes,
3732 reported: BTreeMap<(String, String), crate::SessionActivity>,
3733}
3734
3735fn live_descriptor_value(
3742 session: &SessionDescriptor,
3743 doors: &crate::mail_route::LiveSessions,
3744) -> std::result::Result<Value, ServiceError> {
3745 let mut value = serde_json::to_value(session)
3746 .map_err(|error| ServiceError::Operation(error.to_string()))?;
3747 if value.get("title").is_none_or(Value::is_null) {
3749 if let Some(live) = doors.all().iter().find(|live| {
3750 live.address.harness == session.locator.harness.as_str()
3751 && live.address.session_id == session.locator.session_id
3752 }) {
3753 let name = live.name.split('@').next().unwrap_or(&live.name);
3754 if !name.is_empty() {
3755 value["title"] = json!(name);
3756 }
3757 }
3758 }
3759 if let Some(workspace) = &session.cwd {
3760 let source = LiveRuntimeSource {
3761 harness: session.locator.harness.as_str().to_string(),
3762 session_id: session.locator.session_id.clone(),
3763 workspace: workspace.clone(),
3764 };
3765 if let Some(endpoint) = discover_live_runtime(&source)
3766 .map_err(|error| ServiceError::Operation(error.to_string()))?
3767 {
3768 value["live_endpoint"] = json!(endpoint.as_str());
3769 }
3770 }
3771 if let Some(door) = doors.door(
3775 session.locator.harness.as_str(),
3776 &session.locator.session_id,
3777 ) {
3778 value["delivery"] = json!(door);
3779 }
3780 if let Some(live) = doors.all().iter().find(|live| {
3783 live.address.harness == session.locator.harness.as_str()
3784 && live.address.session_id == session.locator.session_id
3785 && live.status == "waiting"
3786 }) {
3787 if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
3788 value["pending_request"] = request;
3789 }
3790 }
3791 Ok(value)
3792}
3793
3794fn live_index_changes(
3795 changes: Vec<crate::session_index::SessionIndexChange>,
3796 homes: &HarnessHomes,
3797) -> std::result::Result<Vec<Value>, ServiceError> {
3798 use crate::session_index::SessionIndexChange;
3799 let doors = crate::mail_route::LiveSessions::read(homes);
3800 changes
3801 .into_iter()
3802 .map(|change| match change {
3803 SessionIndexChange::Added { descriptor } => Ok(json!({
3804 "kind": "added",
3805 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3806 })),
3807 SessionIndexChange::Updated { descriptor } => Ok(json!({
3808 "kind": "updated",
3809 "descriptor": live_descriptor_value(&descriptor, &doors)?,
3810 })),
3811 SessionIndexChange::Removed { key } => Ok(json!({
3812 "kind": "removed",
3813 "key": key,
3814 })),
3815 })
3816 .collect()
3817}
3818
3819fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
3820 use crate::{SessionPresence, SessionTurnState};
3821 match (activity.presence, activity.turn) {
3822 (SessionPresence::Persisted, _) => None,
3823 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
3824 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
3825 (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
3827 (SessionPresence::Running, SessionTurnState::Unknown)
3831 if activity.evidence.native_state.is_none() =>
3832 {
3833 None
3834 }
3835 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
3836 }
3837}
3838
3839#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
3840#[serde(rename_all = "kebab-case")]
3841enum TransferFormat {
3842 ClaudeCode,
3843 Codex,
3844 #[serde(rename = "opencode", alias = "open-code")]
3845 OpenCode,
3846 Pi,
3847 Grok,
3848 Gemini,
3849 Goose,
3850 Hermes,
3854}
3855
3856impl TransferFormat {
3857 fn id(self) -> &'static str {
3858 match self {
3859 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
3860 Self::Codex => HarnessId::CODEX,
3861 Self::OpenCode => HarnessId::OPENCODE,
3862 Self::Pi => HarnessId::PI,
3863 Self::Grok => HarnessId::GROK,
3864 Self::Gemini => HarnessId::GEMINI,
3865 Self::Goose => HarnessId::GOOSE,
3866 Self::Hermes => HarnessId::HERMES,
3867 }
3868 }
3869}
3870
3871impl From<TransferFormat> for SessionFormat {
3872 fn from(value: TransferFormat) -> Self {
3873 match value {
3874 TransferFormat::ClaudeCode => Self::ClaudeCode,
3875 TransferFormat::Codex => Self::Codex,
3876 TransferFormat::OpenCode => Self::OpenCode,
3877 TransferFormat::Pi => Self::Pi,
3878 TransferFormat::Grok => Self::Grok,
3879 TransferFormat::Gemini => Self::Gemini,
3880 TransferFormat::Goose => Self::Goose,
3881 TransferFormat::Hermes => Self::Codex,
3883 }
3884 }
3885}
3886
3887#[derive(Deserialize)]
3888struct ImportSessionParams {
3889 source_harness: TransferFormat,
3890 content: String,
3891}
3892
3893#[derive(Deserialize)]
3894struct ExportSessionParams {
3895 locator: SessionLocator,
3896 target_harness: TransferFormat,
3897}
3898
3899#[derive(Deserialize)]
3900struct ReduceSessionParams {
3901 locator: SessionLocator,
3902 target_harness: TransferFormat,
3903 #[serde(default = "default_keep_last")]
3904 keep_last: usize,
3905}
3906
3907fn default_keep_last() -> usize {
3908 6
3909}
3910
3911#[derive(Deserialize)]
3912struct BranchSessionParams {
3913 locator: SessionLocator,
3914 #[serde(default)]
3915 target_harness: Option<TransferFormat>,
3916}
3917
3918#[derive(Deserialize)]
3919struct HandoffSessionParams {
3920 locator: SessionLocator,
3921 target_harness: TransferFormat,
3922 #[serde(default)]
3923 cwd: Option<PathBuf>,
3924}
3925
3926#[derive(Deserialize)]
3927struct MaterializeSessionParams {
3928 artifact: crate::native_materialize::MaterializeArtifact,
3929 cwd: PathBuf,
3930 #[serde(default)]
3932 homes: HarnessHomes,
3933}
3934
3935#[derive(Debug, Clone, Copy, Default, Deserialize)]
3938#[serde(rename_all = "snake_case")]
3939enum ResumePolicy {
3940 Default,
3941 #[default]
3942 Yolo,
3943}
3944
3945#[derive(Deserialize)]
3946struct ResumeInstructionsParams {
3947 locator: SessionLocator,
3948 #[serde(default)]
3949 cwd: Option<PathBuf>,
3950 #[serde(default)]
3951 policy: ResumePolicy,
3952}
3953
3954#[derive(Deserialize)]
3956struct WorkflowLoadParams {
3957 from: crate::workflow_doors::WorkflowHarness,
3958 home: PathBuf,
3959}
3960
3961#[derive(Deserialize)]
3964struct OrchestrationLoadParams {
3965 root: PathBuf,
3966 #[serde(default)]
3967 flavor: crate::orchestration_doors::HomeFlavor,
3968}
3969
3970#[derive(Deserialize)]
3973struct OrchestrationSaveParams {
3974 root: PathBuf,
3975 orchestration: crate::orchestration::Orchestration,
3976 #[serde(default)]
3977 vault: BTreeMap<String, String>,
3978}
3979
3980#[derive(Deserialize)]
3982struct OrchestrationCompileParams {
3983 from: crate::orchestration_doors::OrchestrationHarness,
3984 home: PathBuf,
3985}
3986
3987#[derive(Deserialize)]
3991struct OrchestrationDecompileParams {
3992 to: crate::orchestration_doors::OrchestrationHarness,
3993 orchestration: crate::orchestration::Orchestration,
3994 source: PathBuf,
3995 #[serde(default)]
3996 source_flavor: crate::orchestration_doors::SourceFlavor,
3997 dest: PathBuf,
3998 #[serde(default)]
3999 vault: BTreeMap<String, String>,
4000}
4001
4002#[derive(Deserialize)]
4005struct OrchestrationImportParams {
4006 from: crate::orchestration_doors::OrchestrationHarness,
4007 home: PathBuf,
4008 into: PathBuf,
4009}
4010
4011#[derive(Deserialize)]
4014struct OrchestrationExportParams {
4015 to: crate::orchestration_doors::OrchestrationHarness,
4016 root: PathBuf,
4017 dest: PathBuf,
4018}
4019
4020#[derive(Deserialize)]
4022struct JobsGetParams {
4023 harness: String,
4024 id: String,
4025 #[serde(default)]
4026 homes: crate::HarnessHomes,
4027}
4028
4029fn mutate_job(
4037 verb: crate::jobs_control::JobVerb,
4038 params: Value,
4039) -> std::result::Result<Value, ServiceError> {
4040 let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4041 refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4042 let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4043 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4044}
4045
4046fn mutate_skill(
4053 verb: crate::skills_control::SkillVerb,
4054 params: Value,
4055) -> std::result::Result<Value, ServiceError> {
4056 let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4057 if !crate::skills_control::supports_skill_control(&mutation.harness) {
4058 return Err(ServiceError::UnsupportedAction(format!(
4059 "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4060 mutation.harness,
4061 verb.as_str(),
4062 crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4063 )));
4064 }
4065 let outcome =
4066 crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4067 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4068}
4069
4070fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4072 match error {
4073 crate::skills_control::SkillControlError::Unsupported(message) => {
4074 ServiceError::UnsupportedAction(message)
4075 }
4076 crate::skills_control::SkillControlError::Invalid(message) => {
4077 ServiceError::InvalidParams(message)
4078 }
4079 crate::skills_control::SkillControlError::Failed(message) => {
4080 ServiceError::Operation(message)
4081 }
4082 }
4083}
4084
4085fn mutate_profile(
4093 verb: crate::profiles_control::ProfileVerb,
4094 params: Value,
4095) -> std::result::Result<Value, ServiceError> {
4096 let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4097 let outcome =
4098 crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4099 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4100}
4101
4102fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4104 match error {
4105 crate::profiles_control::ProfileControlError::Unsupported(message) => {
4106 ServiceError::UnsupportedAction(message)
4107 }
4108 crate::profiles_control::ProfileControlError::Invalid(message) => {
4109 ServiceError::InvalidParams(message)
4110 }
4111 crate::profiles_control::ProfileControlError::Failed(message) => {
4112 ServiceError::Operation(message)
4113 }
4114 }
4115}
4116
4117fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4121 match error {
4122 crate::jobs_control::JobControlError::Unsupported(message) => {
4123 ServiceError::UnsupportedAction(message)
4124 }
4125 crate::jobs_control::JobControlError::Invalid(message) => {
4126 ServiceError::InvalidParams(message)
4127 }
4128 crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4129 }
4130}
4131
4132fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4137 match error {
4138 crate::SessionControlError::Unsupported(message) => {
4139 ServiceError::UnsupportedAction(message)
4140 }
4141 crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4142 crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4143 }
4144}
4145
4146fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4151 if crate::jobs::supports_jobs(harness) {
4152 return Ok(());
4153 }
4154 Err(ServiceError::UnsupportedAction(format!(
4155 "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4156 crate::jobs::JOB_HARNESSES.join(", ")
4157 )))
4158}
4159
4160#[derive(Deserialize)]
4162struct RunsGetParams {
4163 harness: String,
4164 id: String,
4165 #[serde(default)]
4166 homes: crate::HarnessHomes,
4167}
4168
4169fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4174 if crate::runs::supports_runs(harness) {
4175 return Ok(());
4176 }
4177 Err(ServiceError::UnsupportedAction(format!(
4178 "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4179 crate::runs::RUN_HARNESSES.join(", ")
4180 )))
4181}
4182
4183#[derive(Serialize)]
4184struct SessionArtifact {
4185 source_harness: HarnessId,
4186 target_harness: &'static str,
4187 session_id: Option<String>,
4188 content: String,
4189 suggested_filename: String,
4190 files: Vec<SessionArtifactFile>,
4191 fidelity: Fidelity,
4192 residue: Vec<String>,
4193}
4194
4195#[derive(Serialize)]
4196struct SessionArtifactFile {
4197 path: String,
4198 content: String,
4199 role: ArtifactFileRole,
4200}
4201
4202#[derive(Serialize)]
4203#[serde(rename_all = "snake_case")]
4204enum ArtifactFileRole {
4205 Primary,
4206 Subagent,
4207 Bundle,
4208 SourceRecovery,
4209}
4210
4211#[derive(Serialize)]
4212struct StructuredLaunch {
4213 cwd: PathBuf,
4214 program: String,
4215 arguments: Vec<String>,
4216 env: BTreeMap<String, String>,
4217}
4218
4219struct HandoffInstructions {
4220 launch: StructuredLaunch,
4221 materialize: Option<StructuredLaunch>,
4222 requires_materialization: bool,
4223 note: String,
4224}
4225
4226#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4227#[serde(rename_all = "snake_case")]
4228enum HarnessProbeLevel {
4229 #[default]
4230 Passive,
4231 Handshake,
4232}
4233
4234#[derive(Default, Deserialize)]
4235#[serde(default)]
4236struct HarnessInventoryParams {
4237 harness: Option<HarnessId>,
4238 harnesses: Vec<HarnessId>,
4239 workspace: Option<PathBuf>,
4240 probe: HarnessProbeLevel,
4241 include_sessions: bool,
4242 skip_versions: bool,
4244}
4245
4246#[derive(Deserialize)]
4247struct HarnessAuthenticationParams {
4248 harness: HarnessId,
4249}
4250
4251#[derive(Deserialize)]
4252struct BeginHarnessAuthenticationParams {
4253 harness: HarnessId,
4254 #[serde(default = "local_browser_authentication_environment")]
4255 environment: crate::HarnessAuthenticationEnvironment,
4256 #[serde(default)]
4257 method: Option<crate::HarnessAuthenticationMethodId>,
4258 #[serde(default)]
4259 cwd: Option<PathBuf>,
4260}
4261
4262fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4263 crate::HarnessAuthenticationEnvironment::LocalBrowser
4264}
4265
4266#[derive(Serialize)]
4267struct HarnessInventoryReport {
4268 probe: HarnessProbeLevel,
4269 workspace: Option<PathBuf>,
4270 harnesses: Vec<LocalHarness>,
4271}
4272
4273#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4274#[serde(rename_all = "snake_case")]
4275enum HarnessAuthState {
4276 Ready,
4277 Configured,
4278 Required,
4279 Unknown,
4280}
4281
4282#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4283#[serde(rename_all = "snake_case")]
4284enum HarnessRuntimeState {
4285 Ready,
4286 Degraded,
4287 Unavailable,
4288}
4289
4290#[derive(Serialize)]
4291struct HarnessSessionCounts {
4292 global: Option<usize>,
4293 workspace: Option<usize>,
4294}
4295
4296#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4307#[serde(rename_all = "snake_case")]
4308pub enum GatewayState {
4309 Up,
4310 Down,
4311 Unknown,
4312}
4313
4314#[derive(Debug, Clone, Serialize)]
4316pub struct GatewayHealth {
4317 pub state: GatewayState,
4318 #[serde(skip_serializing_if = "Option::is_none")]
4323 pub endpoint: Option<String>,
4324 #[serde(skip_serializing_if = "Option::is_none")]
4325 pub version: Option<String>,
4326 pub evidence: String,
4328 pub checked_at_ms: u64,
4329}
4330
4331fn openclaw_gateway_endpoint(home: &Path) -> String {
4335 let config_path = home.join(".openclaw/openclaw.json");
4336 let gateway = std::fs::read_to_string(&config_path)
4337 .ok()
4338 .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4339 .and_then(|config| config.get("gateway").cloned());
4340 if let Some(url) = gateway
4341 .as_ref()
4342 .and_then(|gateway| gateway.get("url"))
4343 .and_then(serde_json::Value::as_str)
4344 {
4345 return url.to_string();
4346 }
4347 let port = gateway
4348 .as_ref()
4349 .and_then(|gateway| gateway.get("port"))
4350 .and_then(serde_json::Value::as_u64)
4351 .unwrap_or(18789);
4352 format!("ws://127.0.0.1:{port}")
4353}
4354
4355fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4362 let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4363 let output = std::process::Command::new(&program)
4364 .args(["gateway", "status"])
4365 .stdin(std::process::Stdio::null())
4366 .output()
4367 .ok()?;
4368 let text = format!(
4369 "{}{}",
4370 String::from_utf8_lossy(&output.stdout),
4371 String::from_utf8_lossy(&output.stderr)
4372 );
4373 let verdict = text.lines().find_map(|line| {
4374 let l = line.trim();
4375 if l.contains("supervised by launchd (PID")
4376 || l.contains("supervised by systemd (PID")
4377 || l.contains("Gateway is running")
4378 || l.contains("process is running")
4379 {
4380 Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4381 } else if l.contains("not running") || l.contains("not installed") {
4382 Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4383 } else {
4384 None
4385 }
4386 });
4387 verdict
4388}
4389
4390fn gateway_health(
4391 id: &str,
4392 installed: bool,
4393 running: Option<&RunningInstance>,
4394 version: Option<&str>,
4395) -> GatewayHealth {
4396 let checked_at_ms = now_epoch_ms();
4397 let home = supercode_interchange::user_home()
4398 .map(std::path::PathBuf::into_os_string)
4399 .map(PathBuf::from);
4400 match id {
4401 HarnessId::HERMES | HarnessId::OPENCLAW => {
4402 let endpoint = (id == HarnessId::OPENCLAW)
4403 .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4404 .flatten();
4405 let (state, evidence) = match running {
4406 Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4407 None if !installed => (
4408 GatewayState::Unknown,
4409 format!("`{id}` is not installed; no gateway to probe"),
4410 ),
4411 None if id == HarnessId::HERMES => match hermes_gateway_status() {
4412 Some((state, evidence)) => (state, evidence),
4415 None => (
4416 GatewayState::Down,
4417 "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4418 ),
4419 },
4420 None => (
4421 GatewayState::Down,
4422 format!(
4423 "no TCP listener at {}",
4424 endpoint.as_deref().unwrap_or("the gateway endpoint")
4425 ),
4426 ),
4427 };
4428 GatewayHealth {
4429 state,
4430 endpoint,
4431 version: version.map(str::to_string),
4432 evidence,
4433 checked_at_ms,
4434 }
4435 }
4436 HarnessId::ORCHESTRATOR => {
4443 let root = crate::HarnessHomes::default().orchestrator;
4444 let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4445 Some(lease) if lease.is_live() => (
4446 GatewayState::Up,
4447 format!(
4448 "`{}` names pid {} (started {}), which is live",
4449 crate::orchestrator::lock_path(&root).display(),
4450 lease.pid,
4451 lease.started_at
4452 ),
4453 ),
4454 Some(lease) => (
4455 GatewayState::Down,
4456 format!(
4457 "stale lease `{}`: pid {} is gone",
4458 crate::orchestrator::lock_path(&root).display(),
4459 lease.pid
4460 ),
4461 ),
4462 None => (
4463 GatewayState::Down,
4464 format!(
4465 "no lease at `{}`; `supercode orchestrator start` writes one",
4466 crate::orchestrator::lock_path(&root).display()
4467 ),
4468 ),
4469 };
4470 GatewayHealth {
4471 state,
4472 endpoint: None,
4473 version: version.map(str::to_string),
4474 evidence,
4475 checked_at_ms,
4476 }
4477 }
4478 _ => GatewayHealth {
4479 state: GatewayState::Unknown,
4480 endpoint: None,
4481 version: version.map(str::to_string),
4482 evidence: format!("`{id}` runs per session, not as a gateway"),
4483 checked_at_ms,
4484 },
4485 }
4486}
4487
4488#[derive(Debug, Clone, Serialize)]
4489struct RunningInstance {
4490 method: RunningInstanceMethod,
4492 evidence: String,
4494 checked_at_ms: u64,
4496}
4497
4498#[derive(Debug, Clone, Copy, Serialize)]
4499#[serde(rename_all = "snake_case")]
4500enum RunningInstanceMethod {
4501 GatewayConnect,
4504 StoreWalActivity,
4507}
4508
4509fn now_epoch_ms() -> u64 {
4510 std::time::SystemTime::now()
4511 .duration_since(std::time::UNIX_EPOCH)
4512 .map(|elapsed| elapsed.as_millis() as u64)
4513 .unwrap_or(0)
4514}
4515
4516fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4520 let config_path = home.join(".openclaw/openclaw.json");
4521 let text = std::fs::read_to_string(&config_path).ok();
4522 let gateway = text
4523 .as_deref()
4524 .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4525 .and_then(|config| config.get("gateway").cloned());
4526 let address = gateway
4527 .as_ref()
4528 .and_then(|gateway| gateway.get("url"))
4529 .and_then(serde_json::Value::as_str)
4530 .and_then(|url| {
4531 url.split("://").nth(1).map(|rest| {
4532 rest.trim_end_matches('/')
4533 .split('/')
4534 .next()
4535 .unwrap_or(rest)
4536 .to_string()
4537 })
4538 })
4539 .unwrap_or_else(|| {
4540 let port = gateway
4541 .as_ref()
4542 .and_then(|gateway| gateway.get("port"))
4543 .and_then(serde_json::Value::as_u64)
4544 .unwrap_or(18789);
4545 format!("127.0.0.1:{port}")
4546 });
4547 let reachable = std::net::TcpStream::connect_timeout(
4548 &address.parse().ok()?,
4549 std::time::Duration::from_millis(400),
4550 )
4551 .is_ok();
4552 reachable.then(|| RunningInstance {
4553 method: RunningInstanceMethod::GatewayConnect,
4554 evidence: format!(
4555 "gateway endpoint {address} accepted a TCP connect (from {})",
4556 config_path.display()
4557 ),
4558 checked_at_ms: now_epoch_ms(),
4559 })
4560}
4561
4562fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4567 let wal = home.join(".hermes/state.db-wal");
4568 let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4569 let age_ms = std::time::SystemTime::now()
4570 .duration_since(modified)
4571 .map(|age| age.as_millis() as u64)
4572 .unwrap_or(u64::MAX);
4573 (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4574 method: RunningInstanceMethod::StoreWalActivity,
4575 evidence: format!(
4576 "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4577 wal.display()
4578 ),
4579 checked_at_ms: now_epoch_ms(),
4580 })
4581}
4582
4583fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4585 let home = supercode_interchange::user_home()
4586 .map(std::path::PathBuf::into_os_string)
4587 .map(PathBuf::from)?;
4588 match id {
4589 HarnessId::OPENCLAW => probe_openclaw_running(&home),
4590 HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4591 _ => None,
4592 }
4593}
4594
4595#[derive(Serialize)]
4596struct LocalHarness {
4597 id: HarnessId,
4598 display_name: String,
4599 supported: bool,
4600 installed: bool,
4601 executable: Option<String>,
4602 version: Option<String>,
4603 auth: HarnessAuthState,
4604 runtime: HarnessRuntimeState,
4605 protocol: String,
4606 capabilities: crate::RuntimeCapabilities,
4607 effective_capabilities: crate::RuntimeCapabilities,
4608 sessions: HarnessSessionCounts,
4609 #[serde(skip_serializing_if = "Option::is_none")]
4612 running: Option<RunningInstance>,
4613 gateway: GatewayHealth,
4615 reason: Option<String>,
4616 repair: Option<String>,
4617}
4618
4619#[derive(Clone, Deserialize)]
4620struct RuntimeBackendParams {
4621 harness: HarnessId,
4622 #[serde(default)]
4623 protocol: Option<String>,
4624 #[serde(default)]
4625 launch: Option<RuntimeLaunch>,
4626 #[serde(default)]
4627 base_url: Option<String>,
4628 #[serde(default)]
4629 policy: RuntimePolicy,
4630}
4631
4632#[derive(Debug, Clone, Copy, Default, Deserialize)]
4635#[serde(rename_all = "snake_case")]
4636enum RuntimePolicy {
4637 Default,
4638 #[default]
4639 Yolo,
4640}
4641
4642#[derive(Deserialize)]
4643struct RuntimeStartParams {
4644 #[serde(flatten)]
4645 backend: RuntimeBackendParams,
4646 cwd: PathBuf,
4647 #[serde(default)]
4650 mcp_servers: Vec<crate::McpServerLaunch>,
4651 #[serde(default)]
4653 approval_policy: Option<String>,
4654}
4655
4656#[derive(Deserialize)]
4657struct RuntimeAttachParams {
4658 #[serde(flatten)]
4659 backend: RuntimeBackendParams,
4660 runtime_id: String,
4661 #[serde(default)]
4662 cwd: Option<PathBuf>,
4663 #[serde(default)]
4666 mcp_servers: Vec<crate::McpServerLaunch>,
4667 #[serde(default)]
4669 approval_policy: Option<String>,
4670}
4671
4672#[derive(Deserialize)]
4673struct RuntimeConnectionParams {
4674 connection: String,
4675}
4676
4677#[derive(Deserialize)]
4678struct RuntimeInputParams {
4679 connection: String,
4680 text: String,
4681 #[serde(default)]
4682 image_urls: Vec<String>,
4683}
4684
4685const MAX_RUNTIME_IMAGES: usize = 4;
4686const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
4687const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
4688
4689fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
4690 if image_urls.len() > MAX_RUNTIME_IMAGES {
4691 return Err(ServiceError::InvalidParams(format!(
4692 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
4693 )));
4694 }
4695 let mut total = 0usize;
4696 for url in &image_urls {
4697 if !(url.starts_with("data:image/")
4698 || url.starts_with("https://")
4699 || url.starts_with("http://"))
4700 {
4701 return Err(ServiceError::InvalidParams(
4702 "runtime images must be image data URLs or HTTP(S) URLs".into(),
4703 ));
4704 }
4705 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
4706 return Err(ServiceError::InvalidParams(format!(
4707 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
4708 )));
4709 }
4710 total = total.saturating_add(url.len());
4711 }
4712 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
4713 return Err(ServiceError::InvalidParams(format!(
4714 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
4715 )));
4716 }
4717 Ok(image_urls)
4718}
4719
4720#[derive(Deserialize)]
4721struct RuntimeRespondParams {
4722 connection: String,
4723 request_id: Value,
4724 response: Value,
4725}
4726
4727fn default_reduction_store_root() -> PathBuf {
4728 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
4729 return PathBuf::from(root).join("sessions");
4730 }
4731 if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
4732 return PathBuf::from(home).join(".supercode").join("sessions");
4733 }
4734 PathBuf::from(".supercode").join("sessions")
4735}
4736
4737fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
4738 let mut output = String::new();
4739 for message in messages {
4740 output.push_str(
4741 &serde_json::to_string(message)
4742 .map_err(|error| ServiceError::Operation(error.to_string()))?,
4743 );
4744 output.push('\n');
4745 }
4746 Ok(output)
4747}
4748
4749fn parse_messages_jsonl(
4750 content: &str,
4751) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
4752 content
4753 .lines()
4754 .enumerate()
4755 .filter(|(_, line)| !line.trim().is_empty())
4756 .map(|(index, line)| {
4757 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
4758 ServiceError::Operation(format!(
4759 "reduced transcript line {} is invalid: {error}",
4760 index + 1
4761 ))
4762 })
4763 })
4764 .collect()
4765}
4766
4767fn reduced_bootstrap_prompt(
4768 source: &SessionLocator,
4769 target: TransferFormat,
4770 view_jsonl: &str,
4771 sidecar_path: &Path,
4772 reduction_log_path: &Path,
4773) -> String {
4774 format!(
4775 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
4776 \n\
4777 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\
4778 \n\
4779 <supercode-reduced-session source-session=\"{source_id}\">\n\
4780 {view_jsonl}\
4781 </supercode-reduced-session>\n\
4782 \n\
4783 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
4784 source_harness = source.harness.as_str(),
4785 target_harness = target.id(),
4786 sidecar = sidecar_path.display(),
4787 log = reduction_log_path.display(),
4788 source_id = source.session_id,
4789 )
4790}
4791
4792fn session_artifact(
4793 locator: &SessionLocator,
4794 session: &Session,
4795 target: TransferFormat,
4796) -> std::result::Result<SessionArtifact, ServiceError> {
4797 session_artifact_with_id(locator, session, target, None)
4798}
4799
4800fn session_artifact_with_id(
4801 locator: &SessionLocator,
4802 session: &Session,
4803 target: TransferFormat,
4804 target_session_id: Option<&str>,
4805) -> std::result::Result<SessionArtifact, ServiceError> {
4806 let format: SessionFormat = target.into();
4807 let diagonal = format.source() == session.meta.source;
4808 crate::residue_store::store_segments(session);
4809 let has_appended_turns = session
4810 .imported_message_count
4811 .is_some_and(|imported| imported < session.messages.len());
4812 let mut restoration = None;
4813 let content = if let Some(id) = target_session_id {
4814 if diagonal && format != SessionFormat::OpenCode {
4815 session
4816 .to_jsonl_spliced(format, Some(id))
4817 .map_err(operation)?
4818 } else {
4819 let mut rewritten = session.clone();
4820 rewritten.meta.session_id = Some(id.to_string());
4821 rewritten.to_jsonl(format).map_err(operation)?
4822 }
4823 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
4824 session.raw_verbatim()
4825 } else if diagonal {
4826 session.to_jsonl_spliced(format, None).map_err(operation)?
4827 } else {
4828 match session
4831 .restore_residue(format, crate::residue_store::lookup)
4832 .map_err(operation)?
4833 {
4834 Some((content, report)) => {
4835 restoration = Some(report);
4836 content
4837 }
4838 None => session.to_jsonl(format).map_err(operation)?,
4839 }
4840 };
4841 let stem = sanitize_filename(
4842 target_session_id
4843 .or(session.meta.session_id.as_deref())
4844 .unwrap_or(&locator.session_id),
4845 );
4846 let suggested_filename = if diagonal && target == TransferFormat::Grok {
4847 "chat_history.jsonl".to_string()
4848 } else if target == TransferFormat::Goose {
4849 format!("{stem}.goose.json")
4850 } else {
4851 format!("{stem}.{}.jsonl", target.id())
4852 };
4853 let mut files = vec![SessionArtifactFile {
4854 path: suggested_filename.clone(),
4855 content: content.clone(),
4856 role: ArtifactFileRole::Primary,
4857 }];
4858 if target == TransferFormat::ClaudeCode {
4859 let bundle_stem = Path::new(&suggested_filename)
4860 .file_stem()
4861 .and_then(|stem| stem.to_str())
4862 .unwrap_or(&stem);
4863 let mut child_paths = BTreeSet::new();
4864 for (index, subagent) in session.subagents.iter().enumerate() {
4865 let agent_id = subagent
4866 .meta
4867 .agent_id
4868 .as_deref()
4869 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
4870 .map(sanitize_filename)
4871 .filter(|id| !id.is_empty())
4872 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4873 let child_has_appended_turns = subagent
4874 .imported_message_count
4875 .is_some_and(|imported| imported < subagent.messages.len());
4876 let child_content = if target_session_id.is_none()
4877 && subagent.meta.source == SessionSource::ClaudeCode
4878 && subagent.raw_is_verbatim
4879 && !child_has_appended_turns
4880 {
4881 subagent.raw_verbatim()
4882 } else if subagent.meta.source == SessionSource::ClaudeCode {
4883 subagent
4884 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
4885 .map_err(operation)?
4886 } else {
4887 let mut child = subagent.clone();
4888 if let Some(id) = target_session_id {
4889 child.meta.session_id = Some(id.to_string());
4890 }
4891 child
4892 .to_jsonl(SessionFormat::ClaudeCode)
4893 .map_err(operation)?
4894 };
4895 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
4896 if !child_paths.insert(path.clone()) {
4897 return Err(ServiceError::Operation(format!(
4898 "Claude subagent ids collide at artifact path `{path}`"
4899 )));
4900 }
4901 files.push(SessionArtifactFile {
4902 path,
4903 content: child_content,
4904 role: ArtifactFileRole::Subagent,
4905 });
4906 }
4907 }
4908 if diagonal && target == TransferFormat::Grok {
4909 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
4910 }
4911 if !diagonal || !session.raw_is_verbatim {
4912 files.push(SessionArtifactFile {
4913 path: "recovery/source.supercode.jsonl".into(),
4914 content: session.to_native_jsonl(),
4915 role: ArtifactFileRole::SourceRecovery,
4916 });
4917 for (index, subagent) in session.subagents.iter().enumerate() {
4918 let id = subagent
4919 .meta
4920 .agent_id
4921 .as_deref()
4922 .map(sanitize_filename)
4923 .unwrap_or_else(|| format!("subagent-{}", index + 1));
4924 files.push(SessionArtifactFile {
4925 path: format!("recovery/subagents/{id}.supercode.jsonl"),
4926 content: subagent.to_native_jsonl(),
4927 role: ArtifactFileRole::SourceRecovery,
4928 });
4929 }
4930 }
4931 if !diagonal && session.meta.source == SessionSource::Grok {
4932 append_grok_bundle_files(
4933 locator,
4934 "recovery/grok/",
4935 ArtifactFileRole::SourceRecovery,
4936 &mut files,
4937 )?;
4938 }
4939 let (fidelity, residue) = if diagonal
4940 && target_session_id.is_none()
4941 && session.raw_is_verbatim
4942 && !has_appended_turns
4943 {
4944 (Fidelity::ByteLossless, Vec::new())
4945 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
4946 (
4947 Fidelity::ValueLossless,
4948 vec![if target_session_id.is_some() {
4949 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
4950 } else {
4951 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
4952 }],
4953 )
4954 } else {
4955 match restoration {
4956 Some(report) if report.rendered_messages == 0 => (
4957 Fidelity::ByteLossless,
4958 vec![format!(
4959 "restored verbatim from this conversation's {} source records in the residue store",
4960 target.id()
4961 )],
4962 ),
4963 Some(report) => (
4964 Fidelity::Semantic,
4965 vec![format!(
4966 "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
4967 report.restored_messages,
4968 report.restored_messages + report.rendered_messages,
4969 report.rendered_messages,
4970 target.id()
4971 )],
4972 ),
4973 None => (
4974 Fidelity::Semantic,
4975 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
4976 ),
4977 }
4978 };
4979 Ok(SessionArtifact {
4980 source_harness: locator.harness.clone(),
4981 target_harness: target.id(),
4982 session_id: target_session_id
4983 .map(str::to_string)
4984 .or_else(|| session.meta.session_id.clone()),
4985 content,
4986 suggested_filename,
4987 files,
4988 fidelity,
4989 residue,
4990 })
4991}
4992
4993fn append_grok_bundle_files(
4994 locator: &SessionLocator,
4995 prefix: &str,
4996 role: ArtifactFileRole,
4997 files: &mut Vec<SessionArtifactFile>,
4998) -> std::result::Result<(), ServiceError> {
4999 let primary = locator.storage.path();
5000 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5001 return Err(ServiceError::Operation(format!(
5002 "Grok bundle locator must name chat_history.jsonl, got {}",
5003 primary.display()
5004 )));
5005 }
5006 let parent = primary.parent().ok_or_else(|| {
5007 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5008 })?;
5009 for name in ["summary.json", "updates.jsonl"] {
5010 let path = parent.join(name);
5011 let metadata = match std::fs::symlink_metadata(&path) {
5012 Ok(metadata) => metadata,
5013 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5014 Err(error) => return Err(ServiceError::Operation(error.to_string())),
5015 };
5016 if metadata.file_type().is_symlink() || !metadata.is_file() {
5017 return Err(ServiceError::Operation(format!(
5018 "refusing non-regular Grok bundle member {}",
5019 path.display()
5020 )));
5021 }
5022 let content = std::fs::read_to_string(&path).map_err(|error| {
5023 ServiceError::Operation(format!(
5024 "Grok bundle member {} is not representable as UTF-8: {error}",
5025 path.display()
5026 ))
5027 })?;
5028 files.push(SessionArtifactFile {
5029 path: format!("{prefix}{name}"),
5030 content,
5031 role: match role {
5032 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5033 _ => ArtifactFileRole::SourceRecovery,
5034 },
5035 });
5036 }
5037 Ok(())
5038}
5039
5040fn handoff_artifact(
5041 locator: &SessionLocator,
5042 session: &Session,
5043 target: TransferFormat,
5044) -> std::result::Result<SessionArtifact, ServiceError> {
5045 let target_session_id = target_session_id(target);
5046 session_artifact_with_id(locator, session, target, Some(&target_session_id))
5047}
5048
5049fn target_session_id(target: TransferFormat) -> String {
5050 let uuid = generated_session_id();
5051 match target {
5052 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5053 TransferFormat::ClaudeCode
5054 | TransferFormat::Codex
5055 | TransferFormat::Pi
5056 | TransferFormat::Grok
5057 | TransferFormat::Gemini
5058 | TransferFormat::Goose
5059 | TransferFormat::Hermes => uuid,
5060 }
5061}
5062
5063fn sanitize_filename(value: &str) -> String {
5064 let value = value
5065 .chars()
5066 .map(|character| {
5067 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5068 character
5069 } else {
5070 '-'
5071 }
5072 })
5073 .collect::<String>();
5074 let value = value.trim_matches('-');
5075 if value.is_empty() {
5076 "session".into()
5077 } else {
5078 value.chars().take(100).collect()
5079 }
5080}
5081
5082fn handoff_instructions(
5083 target: TransferFormat,
5084 session_id: &str,
5085 cwd: &Path,
5086) -> HandoffInstructions {
5087 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5088 cwd: cwd.to_path_buf(),
5089 program: program.into(),
5090 arguments,
5091 env: BTreeMap::new(),
5092 };
5093 match target {
5094 TransferFormat::ClaudeCode => HandoffInstructions {
5095 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5096 materialize: None,
5097 requires_materialization: true,
5098 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(),
5099 },
5100 TransferFormat::Hermes => HandoffInstructions {
5101 launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5102 materialize: None,
5103 requires_materialization: true,
5104 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(),
5105 },
5106 TransferFormat::Codex => HandoffInstructions {
5107 launch: launch("codex", vec!["resume".into(), session_id.into()]),
5108 materialize: None,
5109 requires_materialization: true,
5110 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5111 },
5112 TransferFormat::OpenCode => HandoffInstructions {
5113 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5114 materialize: Some(launch(
5115 "opencode",
5116 vec!["import".into(), "{artifact_path}".into()],
5117 )),
5118 requires_materialization: true,
5119 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5120 },
5121 TransferFormat::Pi => HandoffInstructions {
5122 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5123 materialize: None,
5124 requires_materialization: true,
5125 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5126 },
5127 TransferFormat::Grok => HandoffInstructions {
5128 launch: launch(
5129 "grok",
5130 vec!["--resume".into(), "{materialized_session_id}".into()],
5131 ),
5132 materialize: None,
5133 requires_materialization: true,
5134 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(),
5135 },
5136 TransferFormat::Gemini => HandoffInstructions {
5137 launch: launch(
5138 "gemini",
5139 vec!["--session-file".into(), "{artifact_path}".into()],
5140 ),
5141 materialize: None,
5142 requires_materialization: true,
5143 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(),
5144 },
5145 TransferFormat::Goose => HandoffInstructions {
5146 launch: launch(
5147 "goose",
5148 vec![
5149 "session".into(),
5150 "--resume".into(),
5151 "--session-id".into(),
5152 "{imported_session_id}".into(),
5153 ],
5154 ),
5155 materialize: Some(launch(
5156 "goose",
5157 vec!["session".into(), "import".into(), "{artifact_path}".into()],
5158 )),
5159 requires_materialization: true,
5160 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(),
5161 },
5162 }
5163}
5164
5165fn resume_launch(
5166 harness: &str,
5167 session_id: &str,
5168 cwd: &Path,
5169 policy: ResumePolicy,
5170) -> std::result::Result<StructuredLaunch, ServiceError> {
5171 let mut arguments = Vec::new();
5172 let program = match harness {
5173 HarnessId::GROK => {
5174 if matches!(policy, ResumePolicy::Yolo) {
5175 if crate::support::self_sandbox_supported() {
5176 arguments.extend(["--sandbox".into(), "workspace".into()]);
5177 }
5178 arguments.push("--always-approve".into());
5179 }
5180 arguments.extend(["--resume".into(), session_id.into()]);
5181 "grok"
5182 }
5183 HarnessId::CODEX => {
5184 arguments.extend(crate::startup_prompts::startup_arguments(
5185 harness,
5186 Some(cwd),
5187 &[],
5188 matches!(policy, ResumePolicy::Yolo),
5189 ));
5190 arguments.extend(["resume".into(), session_id.into()]);
5191 "codex"
5192 }
5193 HarnessId::CLAUDE_CODE => {
5194 arguments.extend(crate::startup_prompts::startup_arguments(
5195 harness,
5196 Some(cwd),
5197 &[],
5198 matches!(policy, ResumePolicy::Yolo),
5199 ));
5200 arguments.extend(["--resume".into(), session_id.into()]);
5201 "claude"
5202 }
5203 HarnessId::GEMINI => {
5204 arguments.extend(crate::startup_prompts::startup_arguments(
5205 harness,
5206 Some(cwd),
5207 &[],
5208 matches!(policy, ResumePolicy::Yolo),
5209 ));
5210 arguments.extend(["--resume".into(), session_id.into()]);
5211 "gemini"
5212 }
5213 HarnessId::GOOSE => {
5214 arguments.extend([
5215 "session".into(),
5216 "--resume".into(),
5217 "--session-id".into(),
5218 session_id.into(),
5219 ]);
5220 "goose"
5221 }
5222 HarnessId::PI => {
5223 arguments.extend(crate::startup_prompts::startup_arguments(
5224 harness,
5225 Some(cwd),
5226 &[],
5227 matches!(policy, ResumePolicy::Yolo),
5228 ));
5229 arguments.extend(["--session".into(), session_id.into()]);
5230 "pi"
5231 }
5232 HarnessId::OPENCODE => {
5233 arguments.extend(["--session".into(), session_id.into()]);
5234 "opencode"
5235 }
5236 HarnessId::SUPERCODE => {
5237 if matches!(policy, ResumePolicy::Yolo) {
5238 arguments.push("--dangerous".into());
5239 }
5240 arguments.extend(["resume".into(), session_id.into()]);
5241 "supercode"
5242 }
5243 other => {
5244 return Err(ServiceError::InvalidParams(format!(
5245 "no structured resume launch is registered for harness `{other}`"
5246 )))
5247 }
5248 };
5249 Ok(StructuredLaunch {
5250 cwd: cwd.to_path_buf(),
5251 env: if program == "grok" {
5252 crate::support::grok_home_env()
5253 } else {
5254 BTreeMap::new()
5255 },
5256 program: program.into(),
5257 arguments,
5258 })
5259}
5260
5261fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5267 let digest = blake3::hash(address.as_bytes()).to_hex();
5268 let path = std::env::temp_dir().join(format!(
5269 "supercode-openclaw-gateway-token-{}",
5270 &digest.as_str()[..16]
5271 ));
5272 #[cfg(unix)]
5273 {
5274 use std::io::Write;
5275 use std::os::unix::fs::OpenOptionsExt;
5276 let mut file = std::fs::OpenOptions::new()
5277 .write(true)
5278 .create(true)
5279 .truncate(true)
5280 .mode(0o600)
5281 .open(&path)?;
5282 file.write_all(secret.as_bytes())?;
5283 }
5284 #[cfg(not(unix))]
5285 std::fs::write(&path, secret)?;
5286 Ok(path)
5287}
5288
5289fn open_connect_descriptor(
5295 descriptor: &crate::HarnessSupportDescriptor,
5296 home: &Path,
5297) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5298 let Some(connect) = &descriptor.runtime.connect_launch else {
5299 return Err(ServiceError::InvalidParams(format!(
5300 "harness `{}` has no registered connect-mode launch",
5301 descriptor.id.as_str()
5302 )));
5303 };
5304 let resolved = connect
5305 .resolve(home)
5306 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5307 match (descriptor.id.as_str(), connect.protocol.as_str()) {
5308 (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5309 let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5310 if let Some(token) = resolved.auth {
5311 backend = backend.with_bearer(token);
5312 }
5313 Ok(Box::new(backend))
5314 }
5315 (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5316 let mut env = BTreeMap::new();
5327 let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5328 if let Some(token) = resolved.auth {
5329 let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5330 .map_err(|error| {
5331 ServiceError::UnsupportedAction(format!(
5332 "could not stage the gateway credential for the bridge: {error}"
5333 ))
5334 })?;
5335 arguments.push("--token-file".into());
5336 arguments.push(token_path.to_string_lossy().into_owned());
5337 env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5338 }
5339 let program = descriptor
5344 .runtime
5345 .default_launch
5346 .as_ref()
5347 .map(|launch| launch.program.clone())
5348 .unwrap_or_else(|| "openclaw".into());
5349 let launch = RuntimeLaunch {
5350 program,
5351 arguments,
5352 env,
5353 };
5354 Ok(Box::new(
5355 crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5356 .with_resume_support(descriptor.runtime.capabilities.resume_session),
5357 ))
5358 }
5359 _ => Err(ServiceError::UnsupportedAction(format!(
5360 "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5361 descriptor.id.as_str(),
5362 connect.protocol
5363 ))),
5364 }
5365}
5366
5367fn registry_connect_descriptor(
5370 params: &RuntimeBackendParams,
5371) -> Option<crate::HarnessSupportDescriptor> {
5372 if params.launch.is_some() || params.base_url.is_some() {
5373 return None;
5374 }
5375 harness_support_registry()
5376 .harnesses
5377 .into_iter()
5378 .find(|descriptor| descriptor.id == params.harness)
5379 .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5380}
5381
5382fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5383 supercode_interchange::user_home()
5384 .map(std::path::PathBuf::into_os_string)
5385 .map(PathBuf::from)
5386 .ok_or_else(|| {
5387 ServiceError::UnsupportedAction(
5388 "connect-mode launches need HOME to locate the harness config".into(),
5389 )
5390 })
5391}
5392
5393pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5396 "harness.v1.runtimes.start",
5397 "harness.v1.runtimes.resume",
5398 "harness.v1.runtimes.attach",
5399 "harness.v1.runtimes.attach_existing",
5400];
5401
5402pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5407
5408pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5414
5415pub const DETACHED_METHODS: &[&str] = &[
5421 "harness.v1.harnesses.list",
5422 "harness.v1.harnesses.probe",
5423 "harness.v1.sessions.message",
5424 "harness.v1.sessions.new",
5425 "harness.v1.sessions.reset",
5426 "harness.v1.sessions.archive",
5427 "harness.v1.sessions.delete",
5428];
5429
5430pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5436
5437pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5445
5446async fn within_control_deadline<F: std::future::Future>(
5449 method: &str,
5450 call: F,
5451) -> std::result::Result<F::Output, ServiceError> {
5452 tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5453 .await
5454 .map_err(|_| {
5455 ServiceError::Operation(format!(
5456 "`{method}` gave up after {}s: the runtime did not answer",
5457 RUNTIME_CONTROL_DEADLINE.as_secs()
5458 ))
5459 })
5460}
5461
5462pub struct RuntimeOpen {
5466 id: Value,
5467 method: String,
5468 params: Value,
5469}
5470
5471impl RuntimeOpen {
5472 pub async fn open(self) -> OpenedRuntime {
5476 let Self { id, method, params } = self;
5477 let outcome = open_runtime(&method, params).await;
5478 OpenedRuntime { id, outcome }
5479 }
5480}
5481
5482pub struct OpenedRuntime {
5485 id: Value,
5486 outcome: std::result::Result<OpenRuntime, ServiceError>,
5487}
5488
5489pub struct DetachedCall {
5494 id: Value,
5495 method: String,
5496 work: std::result::Result<Work, ServiceError>,
5497}
5498
5499impl DetachedCall {
5500 pub async fn run(self) -> DetachedAnswer {
5503 let Self { id, method, work } = self;
5504 match work {
5505 Ok(Work::Runtime(work)) => {
5510 let (result, returned) = work.run().await;
5511 DetachedAnswer {
5512 response: service_response(id, result),
5513 returned,
5514 }
5515 }
5516 Ok(Work::Free(work)) => {
5517 let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5518 Ok(result) => result,
5519 Err(_) => Err(ServiceError::Operation(format!(
5520 "`{method}` gave up after {}s: the harness it waits on did not answer",
5521 DETACHED_CALL_DEADLINE.as_secs()
5522 ))),
5523 };
5524 DetachedAnswer {
5525 response: service_response(id, result),
5526 returned: None,
5527 }
5528 }
5529 Err(error) => DetachedAnswer {
5530 response: service_response(id, Err(error)),
5531 returned: None,
5532 },
5533 }
5534 }
5535}
5536
5537pub struct DetachedAnswer {
5541 response: Value,
5542 returned: Option<ReturnedRuntime>,
5543}
5544
5545impl DetachedAnswer {
5546 pub fn into_response(self) -> Value {
5549 self.response
5550 }
5551}
5552
5553pub struct ReturnedRuntime {
5556 connection: String,
5557 runtime: Box<dyn RuntimeConnection>,
5558}
5559
5560enum Work {
5563 Free(DetachedWork),
5564 Runtime(RuntimeWork),
5565}
5566
5567enum DetachedWork {
5570 Inventory(InventoryWork),
5574 Message(MessageSessionParams),
5576 SessionMutation {
5579 verb: crate::SessionVerb,
5580 mutation: crate::SessionMutation,
5581 },
5582}
5583
5584impl DetachedWork {
5585 async fn run(self) -> std::result::Result<Value, ServiceError> {
5586 match self {
5587 Self::Inventory(work) => run_inventory(work).await,
5588 Self::Message(params) => Ok(message_live_session(¶ms).await),
5589 Self::SessionMutation { verb, mutation } => {
5590 let outcome = run_session_mutation(verb, &mutation).await?;
5591 serde_json::to_value(outcome)
5592 .map_err(|error| ServiceError::Operation(error.to_string()))
5593 }
5594 }
5595 }
5596}
5597
5598enum RuntimeWork {
5600 Close {
5602 runtime: Box<dyn RuntimeConnection>,
5603 process_group: Option<u32>,
5604 },
5605 LiveCommand {
5608 connection: String,
5609 runtime: Box<dyn RuntimeConnection>,
5610 verb: crate::SessionVerb,
5611 mutation: crate::SessionMutation,
5612 command: &'static str,
5613 session: String,
5614 },
5615}
5616
5617type RuntimeWorkAnswer = (
5620 std::result::Result<Value, ServiceError>,
5621 Option<ReturnedRuntime>,
5622);
5623
5624impl RuntimeWork {
5625 async fn run(self) -> RuntimeWorkAnswer {
5626 match self {
5627 Self::Close {
5628 runtime,
5629 process_group,
5630 } => (close_runtime(runtime, process_group).await, None),
5631 Self::LiveCommand {
5632 connection,
5633 mut runtime,
5634 verb,
5635 mutation,
5636 command,
5637 session,
5638 } => {
5639 let result =
5640 type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5641 (
5642 result,
5643 Some(ReturnedRuntime {
5644 connection,
5645 runtime,
5646 }),
5647 )
5648 }
5649 }
5650 }
5651}
5652
5653async fn close_runtime(
5656 mut runtime: Box<dyn RuntimeConnection>,
5657 process_group: Option<u32>,
5658) -> std::result::Result<Value, ServiceError> {
5659 match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
5660 Ok(result) => {
5661 result.map_err(operation)?;
5662 Ok(json!({"closed": true}))
5663 }
5664 Err(deadline) => {
5665 let killed = kill_runtime_process_group(process_group);
5670 drop(runtime);
5671 Ok(json!({
5672 "closed": true,
5673 "killed": killed,
5674 "detail": error_message(deadline),
5675 }))
5676 }
5677 }
5678}
5679
5680fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
5683 mutation
5684 .session
5685 .clone()
5686 .filter(|value| !value.trim().is_empty())
5687 .unwrap_or_else(|| runtime.handle().runtime_id.clone())
5688}
5689
5690async fn type_live_command(
5694 runtime: &mut dyn RuntimeConnection,
5695 verb: crate::SessionVerb,
5696 mutation: &crate::SessionMutation,
5697 command: &str,
5698 session: String,
5699) -> std::result::Result<Value, ServiceError> {
5700 within_control_deadline(
5701 &format!("sessions.{}", verb.as_str()),
5702 runtime.send_input(RuntimeInput {
5703 text: command.to_string(),
5704 image_urls: Vec::new(),
5705 }),
5706 )
5707 .await?
5708 .map_err(operation)?;
5709 let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
5710 .map_err(session_control_error)?;
5711 serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
5712}
5713
5714enum OpenRuntime {
5717 Hosted {
5720 runtime: Box<dyn RuntimeConnection>,
5721 capabilities: crate::RuntimeCapabilities,
5722 workspace: PathBuf,
5723 fresh: bool,
5725 },
5726 Joined { runtime: Box<dyn RuntimeConnection> },
5729}
5730
5731async fn open_runtime(
5736 method: &str,
5737 params: Value,
5738) -> std::result::Result<OpenRuntime, ServiceError> {
5739 match tokio::time::timeout(
5740 RUNTIME_OPEN_DEADLINE,
5741 open_runtime_unbounded(method, params),
5742 )
5743 .await
5744 {
5745 Ok(result) => result,
5746 Err(_) => Err(ServiceError::Operation(format!(
5747 "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
5748 RUNTIME_OPEN_DEADLINE.as_secs()
5749 ))),
5750 }
5751}
5752
5753async fn open_runtime_unbounded(
5754 method: &str,
5755 params: Value,
5756) -> std::result::Result<OpenRuntime, ServiceError> {
5757 match method {
5758 "harness.v1.runtimes.start" => {
5759 let params = decode::<RuntimeStartParams>(params)?;
5760 let backend = runtime_backend(¶ms.backend)?;
5761 let capabilities = backend.capabilities();
5762 let workspace = params.cwd.clone();
5763 let runtime = backend
5764 .start(RuntimeStartRequest {
5765 cwd: params.cwd,
5766 launch: runtime_launch(¶ms.backend),
5767 mcp_servers: params.mcp_servers,
5768 approval_policy: params.approval_policy,
5769 })
5770 .await
5771 .map_err(operation)?;
5772 Ok(OpenRuntime::Hosted {
5773 runtime,
5774 capabilities,
5775 workspace,
5776 fresh: true,
5777 })
5778 }
5779 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
5780 let params = decode::<RuntimeAttachParams>(params)?;
5781 let backend = runtime_backend(¶ms.backend)?;
5782 let capabilities = backend.capabilities();
5783 let workspace = params
5784 .cwd
5785 .clone()
5786 .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
5787 let runtime = backend
5788 .attach(RuntimeAttachRequest {
5789 runtime_id: params.runtime_id,
5790 cwd: params.cwd,
5791 launch: runtime_launch(¶ms.backend),
5792 mcp_servers: params.mcp_servers,
5793 approval_policy: params.approval_policy,
5794 })
5795 .await
5796 .map_err(operation)?;
5797 Ok(OpenRuntime::Hosted {
5798 runtime,
5799 capabilities,
5800 workspace,
5801 fresh: false,
5802 })
5803 }
5804 "harness.v1.runtimes.attach_existing" => {
5805 let params = decode::<RuntimeAttachParams>(params)?;
5806 let backend: Box<dyn RuntimeBackend> = match params
5807 .backend
5808 .base_url
5809 .as_deref()
5810 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
5811 {
5812 Some(endpoint) => {
5813 #[cfg(not(feature = "adapter-api"))]
5814 {
5815 let _ = endpoint;
5816 return Err(ServiceError::UnsupportedAction(
5817 "live HTTP attachment adapter is not compiled".into(),
5818 ));
5819 }
5820 #[cfg(feature = "adapter-api")]
5821 {
5822 let workspace = params.cwd.clone().ok_or_else(|| {
5823 ServiceError::InvalidParams(
5824 "Volter Harness live attach requires the project cwd".into(),
5825 )
5826 })?;
5827 let source = LiveRuntimeSource {
5828 harness: params.backend.harness.as_str().to_string(),
5829 session_id: params.runtime_id.clone(),
5830 workspace,
5831 };
5832 let receipt = resolve_live_runtime(&endpoint, &source)
5833 .map_err(|error| ServiceError::Operation(error.to_string()))?;
5834 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
5835 }
5836 }
5837 None => runtime_backend(¶ms.backend)?,
5838 };
5839 let capabilities = backend.capabilities();
5840 if !capabilities.attach_existing_process {
5841 return Err(ServiceError::Operation(format!(
5842 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
5843 backend.harness().as_str()
5844 )));
5845 }
5846 let runtime = backend
5847 .attach_existing(RuntimeAttachRequest {
5848 runtime_id: params.runtime_id,
5849 cwd: params.cwd,
5850 launch: runtime_launch(¶ms.backend),
5851 mcp_servers: params.mcp_servers,
5852 approval_policy: params.approval_policy,
5853 })
5854 .await
5855 .map_err(operation)?;
5856 Ok(OpenRuntime::Joined { runtime })
5857 }
5858 _ => Err(ServiceError::MethodNotFound),
5859 }
5860}
5861
5862fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
5864 match result {
5865 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
5866 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
5867 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
5868 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
5869 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
5870 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
5871 }
5872}
5873
5874fn runtime_backend(
5875 params: &RuntimeBackendParams,
5876) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5877 if let Some(descriptor) = registry_connect_descriptor(params) {
5878 return open_connect_descriptor(&descriptor, &service_home()?);
5879 }
5880 if params.protocol.as_deref() == Some("acp") {
5881 let launch = params
5882 .launch
5883 .clone()
5884 .or_else(|| {
5885 harness_support_registry()
5886 .harnesses
5887 .into_iter()
5888 .find(|harness| harness.id == params.harness)
5889 .filter(|harness| {
5890 harness.runtime.implementation == ImplementationKind::GenericProtocol
5891 && harness.runtime.protocol.starts_with("acp")
5892 })
5893 .and_then(|harness| harness.runtime.default_launch)
5894 })
5895 .ok_or_else(|| {
5896 ServiceError::InvalidParams(
5897 "an ACP runtime requires `launch` unless the harness has a registered default"
5898 .into(),
5899 )
5900 })?;
5901 let resume_session = harness_support_registry()
5902 .harnesses
5903 .into_iter()
5904 .find(|harness| harness.id == params.harness)
5905 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
5906 return Ok(Box::new(
5907 AcpRuntimeBackend::new(params.harness.clone(), launch)
5908 .with_resume_support(resume_session),
5909 ));
5910 }
5911 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
5912 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
5913 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
5914 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
5915 HarnessId::OPENCODE => match ¶ms.base_url {
5916 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
5917 None => Box::new(OpenCodeRuntimeBackend::new()),
5918 },
5919 harness => {
5920 let descriptor = harness_support_registry()
5921 .harnesses
5922 .into_iter()
5923 .find(|descriptor| descriptor.id.as_str() == harness)
5924 .filter(|descriptor| {
5925 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
5926 && descriptor.runtime.protocol.starts_with("acp")
5927 });
5928 let Some(descriptor) = descriptor else {
5929 return Err(ServiceError::InvalidParams(format!(
5930 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
5931 )));
5932 };
5933 let resume = descriptor.runtime.capabilities.resume_session;
5934 Box::new(
5935 AcpRuntimeBackend::new(
5936 descriptor.id,
5937 descriptor
5938 .runtime
5939 .default_launch
5940 .expect("generic ACP registry entry includes its launch"),
5941 )
5942 .with_resume_support(resume),
5943 )
5944 }
5945 };
5946 Ok(backend)
5947}
5948
5949fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
5950 if let Some(launch) = ¶ms.launch {
5951 return Some(launch.clone());
5952 }
5953 if !matches!(params.policy, RuntimePolicy::Yolo) {
5954 return None;
5955 }
5956 let launch = match params.harness.as_str() {
5957 HarnessId::GROK => RuntimeLaunch {
5958 program: "grok".into(),
5959 arguments: {
5960 let mut arguments: Vec<String> = Vec::new();
5961 if crate::support::self_sandbox_supported() {
5962 arguments.extend(["--sandbox".into(), "workspace".into()]);
5963 }
5964 arguments.extend([
5965 "--always-approve".into(),
5966 "agent".into(),
5967 "--no-leader".into(),
5968 "stdio".into(),
5969 ]);
5970 arguments
5971 },
5972 env: crate::support::grok_env(),
5973 },
5974 HarnessId::CODEX => RuntimeLaunch {
5975 program: "codex".into(),
5976 arguments: vec![
5977 "--dangerously-bypass-approvals-and-sandbox".into(),
5978 "--dangerously-bypass-hook-trust".into(),
5979 "app-server".into(),
5980 ],
5981 env: BTreeMap::new(),
5982 },
5983 HarnessId::CLAUDE_CODE => RuntimeLaunch {
5984 program: "claude".into(),
5985 arguments: vec![
5986 "--dangerously-skip-permissions".into(),
5987 "--print".into(),
5988 "--input-format".into(),
5989 "stream-json".into(),
5990 "--output-format".into(),
5991 "stream-json".into(),
5992 "--verbose".into(),
5993 ],
5994 env: BTreeMap::new(),
5995 },
5996 HarnessId::PI => RuntimeLaunch {
5997 program: "pi".into(),
5998 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
5999 env: BTreeMap::new(),
6000 },
6001 HarnessId::OPENCODE => RuntimeLaunch {
6002 program: "opencode".into(),
6003 arguments: vec!["serve".into()],
6004 env: BTreeMap::new(),
6005 },
6006 HarnessId::GEMINI => RuntimeLaunch {
6007 program: "gemini".into(),
6008 arguments: vec!["--acp".into(), "--yolo".into()],
6009 env: BTreeMap::new(),
6010 },
6011 HarnessId::GOOSE => RuntimeLaunch {
6012 program: "goose".into(),
6013 arguments: vec!["acp".into()],
6014 env: BTreeMap::new(),
6015 },
6016 HarnessId::SUPERCODE => RuntimeLaunch {
6017 program: "supercode".into(),
6018 arguments: vec!["acp".into(), "--dangerous".into()],
6019 env: BTreeMap::new(),
6020 },
6021 _ => return None,
6022 };
6023 Some(launch)
6024}
6025
6026struct IsolatedProbeHome {
6032 launch: RuntimeLaunch,
6033 root: PathBuf,
6034}
6035
6036impl IsolatedProbeHome {
6037 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6038 let root = std::env::temp_dir().join(format!(
6039 "supercode-harness-probe-{harness}-{}",
6040 generated_session_id()
6041 ));
6042 std::fs::create_dir_all(&root)?;
6043 set_private_dir_permissions(&root)?;
6044
6045 if let Some(source_home) = supercode_interchange::user_home()
6046 .map(std::path::PathBuf::into_os_string)
6047 .map(PathBuf::from)
6048 {
6049 for relative in probe_auth_files(harness) {
6050 copy_probe_file(&source_home, &root, relative)?;
6051 }
6052 }
6053 if harness == HarnessId::SUPERCODE {
6058 let config_home = crate::agent::global_instructions_dir();
6059 for file in ["config.toml", "credentials.toml"] {
6060 copy_probe_path(
6061 &config_home.join(file),
6062 &root.join(".config/supercode").join(file),
6063 )?;
6064 }
6065 }
6066 configure_isolated_probe_auth(harness, &root)?;
6067
6068 let root_text = root.to_string_lossy().into_owned();
6069 for (key, value) in [
6070 ("HOME", root_text.clone()),
6071 (
6072 "XDG_CACHE_HOME",
6073 root.join(".cache").to_string_lossy().into_owned(),
6074 ),
6075 (
6076 "XDG_CONFIG_HOME",
6077 root.join(".config").to_string_lossy().into_owned(),
6078 ),
6079 (
6080 "XDG_DATA_HOME",
6081 root.join(".local/share").to_string_lossy().into_owned(),
6082 ),
6083 ] {
6084 launch.env.insert(key.into(), value);
6085 }
6086 let scoped = match harness {
6087 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6088 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6089 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6090 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6091 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6092 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6093 _ => None,
6094 };
6095 if let Some((key, value)) = scoped {
6096 launch
6097 .env
6098 .insert(key.into(), value.to_string_lossy().into_owned());
6099 }
6100 Ok(Self { launch, root })
6101 }
6102
6103 fn cleanup(&self) -> std::io::Result<()> {
6104 match std::fs::remove_dir_all(&self.root) {
6105 Ok(()) => Ok(()),
6106 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6107 Err(error) => Err(error),
6108 }
6109 }
6110}
6111
6112impl Drop for IsolatedProbeHome {
6113 fn drop(&mut self) {
6114 let _ = self.cleanup();
6115 }
6116}
6117
6118fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6119 match harness {
6120 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6121 HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6125 HarnessId::CODEX => &[".codex/auth.json"],
6126 HarnessId::GEMINI => &[
6127 ".gemini/google_accounts.json",
6128 ".gemini/oauth_creds.json",
6129 ".gemini/settings.json",
6130 ],
6131 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6132 HarnessId::OPENCODE => &[
6133 ".config/opencode/auth.json",
6134 ".local/share/opencode/auth.json",
6135 ],
6136 HarnessId::PI => &[".pi/agent/auth.json"],
6137 HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6142 _ => &[],
6143 }
6144}
6145
6146fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6147 copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6148}
6149
6150fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6151 if !source.is_file() {
6152 return Ok(());
6153 }
6154 if let Some(parent) = destination.parent() {
6155 std::fs::create_dir_all(parent)?;
6156 set_private_dir_permissions(parent)?;
6157 }
6158 std::fs::copy(source, destination)?;
6159 set_private_file_permissions(destination)
6160}
6161
6162fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6163 if harness != HarnessId::GEMINI {
6164 return Ok(());
6165 }
6166 let oauth = probe_home.join(".gemini/oauth_creds.json");
6167 if !oauth.is_file() {
6168 return Ok(());
6169 }
6170 let settings_path = probe_home.join(".gemini/settings.json");
6171 let mut settings = std::fs::read_to_string(&settings_path)
6172 .ok()
6173 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6174 .unwrap_or_else(|| json!({}));
6175 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6176 std::fs::write(
6177 &settings_path,
6178 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6179 )?;
6180 set_private_file_permissions(&settings_path)
6181}
6182
6183#[cfg(unix)]
6184fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6185 use std::os::unix::fs::PermissionsExt;
6186 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6187}
6188
6189#[cfg(not(unix))]
6190fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6191 Ok(())
6192}
6193
6194#[cfg(unix)]
6195fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6196 use std::os::unix::fs::PermissionsExt;
6197 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6198}
6199
6200#[cfg(not(unix))]
6201fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6202 Ok(())
6203}
6204
6205fn find_executable(program: &str) -> Option<PathBuf> {
6206 let candidate = PathBuf::from(program);
6207 if candidate.components().count() > 1 {
6208 return candidate.is_file().then_some(candidate);
6209 }
6210 let path = std::env::var_os("PATH")?;
6211 for directory in std::env::split_paths(&path) {
6212 let candidate = directory.join(program);
6213 if candidate.is_file() {
6214 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6215 }
6216 #[cfg(windows)]
6217 {
6218 for extension in ["exe", "cmd", "bat"] {
6219 let candidate = directory.join(format!("{program}.{extension}"));
6220 if candidate.is_file() {
6221 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6222 }
6223 }
6224 }
6225 }
6226 None
6227}
6228
6229async fn executable_version(executable: &Path) -> Option<String> {
6230 let mut command = tokio::process::Command::new(executable);
6231 command
6232 .arg("--version")
6233 .stdin(std::process::Stdio::null())
6234 .stdout(std::process::Stdio::piped())
6235 .stderr(std::process::Stdio::piped())
6236 .kill_on_drop(true);
6237 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6238 .await
6239 .ok()?
6240 .ok()?;
6241 let stdout = String::from_utf8_lossy(&output.stdout);
6242 let stderr = String::from_utf8_lossy(&output.stderr);
6243 stdout
6244 .lines()
6245 .chain(stderr.lines())
6246 .map(str::trim)
6247 .find(|line| !line.is_empty())
6248 .map(|line| truncate_text(line, 200))
6249}
6250
6251pub(crate) fn auth_evidence(harness: &str) -> bool {
6252 let env_names: &[&str] = match harness {
6253 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6254 HarnessId::CODEX => &["OPENAI_API_KEY"],
6255 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6256 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6257 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6258 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6259 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6260 _ => &[],
6261 };
6262 if env_names
6263 .iter()
6264 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6265 {
6266 return true;
6267 }
6268 let Some(home) = supercode_interchange::user_home()
6269 .map(std::path::PathBuf::into_os_string)
6270 .map(PathBuf::from)
6271 else {
6272 return false;
6273 };
6274 let files: Vec<PathBuf> = match harness {
6275 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6276 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6277 HarnessId::OPENCODE => vec![
6278 home.join(".local/share/opencode/auth.json"),
6279 home.join(".config/opencode/auth.json"),
6280 ],
6281 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6282 HarnessId::GROK => vec![home.join(".grok/auth.json")],
6283 HarnessId::GEMINI => vec![
6284 home.join(".gemini/oauth_creds.json"),
6285 home.join(".gemini/google_accounts.json"),
6286 ],
6287 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6288 HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6289 _ => Vec::new(),
6290 };
6291 if files.into_iter().any(|path| {
6292 std::fs::metadata(path)
6293 .map(|metadata| metadata.is_file() && metadata.len() > 2)
6294 .unwrap_or(false)
6295 }) {
6296 return true;
6297 }
6298 if harness == HarnessId::CLAUDE_CODE {
6305 return std::fs::read_to_string(home.join(".claude.json"))
6306 .map(|text| text.contains("\"oauthAccount\""))
6307 .unwrap_or(false);
6308 }
6309 false
6310}
6311
6312fn looks_like_auth_error(message: &str) -> bool {
6313 let message = message.to_ascii_lowercase();
6314 [
6315 "auth",
6316 "login",
6317 "sign in",
6318 "sign-in",
6319 "credential",
6320 "unauthorized",
6321 "forbidden",
6322 "token",
6323 ]
6324 .iter()
6325 .any(|needle| message.contains(needle))
6326}
6327
6328fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6329 crate::RuntimeCapabilities {
6330 start_session: false,
6331 resume_session: false,
6332 attach_existing_process: false,
6333 send_input: false,
6334 stream_events: false,
6335 interrupt: false,
6336 steer: false,
6337 respond_to_requests: false,
6338 }
6339}
6340
6341fn truncate_text(text: &str, max_chars: usize) -> String {
6342 let mut chars = text.chars();
6343 let truncated = chars.by_ref().take(max_chars).collect::<String>();
6344 if chars.next().is_some() {
6345 format!("{truncated}…")
6346 } else {
6347 truncated
6348 }
6349}
6350
6351fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6358 match &handle.endpoint {
6359 crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6360 crate::RuntimeEndpoint::Http { .. } => None,
6361 }
6362}
6363
6364fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6370 match process_group {
6371 #[cfg(unix)]
6372 Some(pid) => {
6373 crate::lsp::kill_process_group(pid);
6374 true
6375 }
6376 #[cfg(not(unix))]
6377 Some(_) => false,
6378 None => false,
6379 }
6380}
6381
6382fn error_message(error: ServiceError) -> String {
6383 match error {
6384 ServiceError::InvalidParams(message)
6385 | ServiceError::Operation(message)
6386 | ServiceError::UnsupportedAction(message) => message,
6387 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6388 ServiceError::Sdk(error) => error.to_string(),
6389 }
6390}
6391
6392#[derive(Debug)]
6393enum ServiceError {
6394 InvalidParams(String),
6395 MethodNotFound,
6396 UnsupportedAction(String),
6397 Operation(String),
6398 Sdk(SdkError),
6399}
6400
6401fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6402 match error {
6403 ServiceError::InvalidParams(message) => {
6404 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6405 }
6406 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6407 SdkError::unsupported(operation)
6408 }
6409 ServiceError::Operation(message) => {
6410 let code = if message.contains("already in progress") {
6411 SdkErrorCode::Busy
6412 } else if message.contains("not supported by this runtime") {
6413 SdkErrorCode::UnsupportedAction
6414 } else if message.contains("unknown runtime connection") {
6415 SdkErrorCode::NotFound
6416 } else {
6417 SdkErrorCode::Execution
6418 };
6419 SdkError::new(code, operation, message)
6420 }
6421 ServiceError::Sdk(error) => error,
6422 }
6423}
6424
6425fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6426 let error_code = error.code();
6427 let code = match error_code {
6428 SdkErrorCode::Unauthenticated => -32030,
6429 SdkErrorCode::Unauthorized => -32031,
6430 SdkErrorCode::ControllerRequired => -32032,
6431 SdkErrorCode::LeaseExpired => -32033,
6432 SdkErrorCode::InvalidArgument => -32602,
6433 SdkErrorCode::NotFound => -32004,
6434 SdkErrorCode::Busy => -32000,
6435 SdkErrorCode::UnsupportedAction => -32020,
6436 SdkErrorCode::Execution => -32002,
6437 SdkErrorCode::Transport => -32003,
6438 };
6439 json!({
6440 "jsonrpc": "2.0",
6441 "id": id,
6442 "error": {
6443 "code": code,
6444 "name": error_code,
6445 "operation": error.operation(),
6446 "message": error.to_string(),
6447 },
6448 })
6449}
6450
6451fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6452 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6453}
6454
6455fn operation(error: impl Into<crate::Error>) -> ServiceError {
6456 let error = error.into();
6457 match error {
6458 crate::Error::Sdk(error) => ServiceError::Sdk(error),
6459 error => ServiceError::Operation(error.to_string()),
6460 }
6461}
6462
6463#[derive(Debug, Clone, Deserialize, Default)]
6467#[serde(default)]
6468struct MemoryRequest {
6469 harness: Option<String>,
6471 query: Option<String>,
6473 profile: Option<String>,
6475 session: Option<String>,
6477 full: bool,
6479 regex: bool,
6481 cwd: Option<std::path::PathBuf>,
6483 homes: crate::HarnessHomes,
6485}
6486
6487fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6490 let request = decode::<MemoryRequest>(params)?;
6491 let harness = request
6492 .harness
6493 .clone()
6494 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6495 let to_service = |error: crate::memory::MemoryError| match error {
6496 crate::memory::MemoryError::UnsupportedHarness { .. }
6497 | crate::memory::MemoryError::SessionNotScoped { .. } => {
6498 ServiceError::UnsupportedAction(error.to_string())
6499 }
6500 other => ServiceError::InvalidParams(other.to_string()),
6501 };
6502 match method {
6503 "harness.v1.memory.show" => {
6504 let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6505 harness,
6506 profile: request.profile,
6507 session: request.session,
6508 full: request.full,
6509 cwd: request.cwd,
6510 homes: request.homes,
6511 })
6512 .map_err(to_service)?;
6513 Ok(json!({
6514 "schema": crate::memory::MEMORY_SCHEMA,
6515 "documents": documents,
6516 }))
6517 }
6518 "harness.v1.memory.search" => {
6519 let query = request
6520 .query
6521 .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6522 let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6523 harness,
6524 query,
6525 profile: request.profile,
6526 regex: request.regex,
6527 cwd: request.cwd,
6528 homes: request.homes,
6529 })
6530 .map_err(to_service)?;
6531 Ok(json!({
6532 "schema": crate::memory::MEMORY_SCHEMA,
6533 "matches": matches,
6534 }))
6535 }
6536 _ => Err(ServiceError::MethodNotFound),
6537 }
6538}
6539
6540#[derive(Debug, Clone, Deserialize)]
6544#[serde(default)]
6545struct ProfilesQuery {
6546 harness: Option<String>,
6548 name: Option<String>,
6550 homes: crate::HarnessHomes,
6552}
6553
6554impl Default for ProfilesQuery {
6555 fn default() -> Self {
6556 Self {
6557 harness: None,
6558 name: None,
6559 homes: crate::HarnessHomes::default(),
6560 }
6561 }
6562}
6563
6564fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6567 let query = decode::<ProfilesQuery>(params)?;
6568 let to_service = |error: crate::profiles::ProfileError| match error {
6569 crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6570 ServiceError::UnsupportedAction(error.to_string())
6571 }
6572 crate::profiles::ProfileError::NotFound { .. } => {
6573 ServiceError::InvalidParams(error.to_string())
6574 }
6575 };
6576 match method {
6577 "harness.v1.profiles.list" => {
6578 let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6579 .map_err(to_service)?;
6580 Ok(json!({
6581 "schema": crate::profiles::PROFILES_SCHEMA,
6582 "profiles": profiles,
6583 }))
6584 }
6585 "harness.v1.profiles.get" => {
6586 let harness = query
6587 .harness
6588 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6589 let name = query
6590 .name
6591 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6592 let profile =
6593 crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6594 Ok(json!({
6595 "schema": crate::profiles::PROFILES_SCHEMA,
6596 "profile": profile,
6597 }))
6598 }
6599 _ => Err(ServiceError::MethodNotFound),
6600 }
6601}
6602
6603#[derive(Debug, Clone, Deserialize)]
6607#[serde(default)]
6608struct ChannelsQuery {
6609 harness: Option<String>,
6611 name: Option<String>,
6613 homes: crate::HarnessHomes,
6615}
6616
6617impl Default for ChannelsQuery {
6618 fn default() -> Self {
6619 Self {
6620 harness: None,
6621 name: None,
6622 homes: crate::HarnessHomes::default(),
6623 }
6624 }
6625}
6626
6627#[derive(Debug, Clone, Deserialize)]
6631#[serde(default)]
6632struct RoutesQuery {
6633 harness: Option<String>,
6634 profile: Option<String>,
6636 homes: crate::HarnessHomes,
6637}
6638
6639impl Default for RoutesQuery {
6640 fn default() -> Self {
6641 Self {
6642 harness: None,
6643 profile: None,
6644 homes: crate::HarnessHomes::default(),
6645 }
6646 }
6647}
6648
6649#[derive(Debug, Clone, Deserialize)]
6650#[serde(default)]
6651struct TriggersQuery {
6652 harness: Option<String>,
6653 homes: crate::HarnessHomes,
6654}
6655
6656impl Default for TriggersQuery {
6657 fn default() -> Self {
6658 Self {
6659 harness: None,
6660 homes: crate::HarnessHomes::default(),
6661 }
6662 }
6663}
6664
6665fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
6666 let query = decode::<TriggersQuery>(params)?;
6667 let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
6668 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6669 Ok(json!({
6670 "schema": crate::triggers::TRIGGERS_SCHEMA,
6671 "triggers": triggers,
6672 }))
6673}
6674
6675fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
6676 let query = decode::<RoutesQuery>(params)?;
6677 let routes = crate::routes::list_routes(
6678 &query.homes,
6679 query.harness.as_deref(),
6680 query.profile.as_deref(),
6681 )
6682 .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6683 Ok(json!({
6684 "schema": crate::routes::ROUTES_SCHEMA,
6685 "routes": routes,
6686 }))
6687}
6688
6689fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6690 let query = decode::<ChannelsQuery>(params)?;
6691 let to_service = |error: crate::channels::ChannelError| match error {
6692 crate::channels::ChannelError::UnsupportedHarness { .. } => {
6693 ServiceError::UnsupportedAction(error.to_string())
6694 }
6695 crate::channels::ChannelError::NotFound { .. } => {
6696 ServiceError::InvalidParams(error.to_string())
6697 }
6698 };
6699 match method {
6700 "harness.v1.channels.list" => {
6701 let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
6702 .map_err(to_service)?;
6703 Ok(json!({
6704 "schema": crate::channels::CHANNELS_SCHEMA,
6705 "channels": channels,
6706 }))
6707 }
6708 "harness.v1.channels.status" => {
6709 let harness = query
6710 .harness
6711 .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6712 let name = query
6713 .name
6714 .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6715 let channel = crate::channels::channel_status(&query.homes, &harness, &name)
6716 .map_err(to_service)?;
6717 Ok(json!({
6718 "schema": crate::channels::CHANNELS_SCHEMA,
6719 "channel": channel,
6720 }))
6721 }
6722 _ => Err(ServiceError::MethodNotFound),
6723 }
6724}
6725
6726fn rpc_error(id: Value, code: i64, message: &str) -> Value {
6727 json!({
6728 "jsonrpc": "2.0",
6729 "id": id,
6730 "error": {"code": code, "message": message},
6731 })
6732}