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