Skip to main content

supercode_harness/
harness_service.rs

1//! Versioned, language-neutral service over persisted harness sessions.
2//!
3//! The service is transport-agnostic: [`HarnessSessionService::handle`] accepts
4//! one JSON-RPC value and [`HarnessSessionService::poll`] produces subscription
5//! notifications. The CLI exposes those primitives as NDJSON over stdio.
6
7use 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
39/// Every JSON-RPC method the harness service dispatches (`harness.v1.capabilities`
40/// reports it; ORCH-4 registry tiers must cite entries of it).
41pub 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
125/// Protocol namespace implemented by this service.
126pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
127/// Notification method emitted for followed-session changes.
128pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
129/// Notification method emitted for normalized session-activity transitions.
130pub const SESSION_ACTIVITY_EVENT_METHOD: &str = "harness.v1.sessions.activity_event";
131/// Notification method emitted for revisioned session-list changes.
132pub const SESSION_INDEX_EVENT_METHOD: &str = "harness.v1.sessions.index_event";
133/// Notification method emitted for live runtime events.
134pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
135
136/// Stateful persisted-session service. Each instance owns its follow
137/// subscriptions; discovery and loading remain read-only.
138pub 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    /// Connections lent to a detached call that is running right now. The
150    /// runtime itself is OUT of `runtimes` for that whole call, and these
151    /// names are how a second caller is told the connection is busy rather
152    /// than unknown.
153    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    /// ORCH-9: live permission/approval requests outstanding on the open
159    /// runtime connections above, fed by the same event pump that publishes
160    /// `harness.v1.runtimes.event`.
161    approvals: crate::approvals::ApprovalRegistry,
162    /// ORCH-9: supercode's own queued subagent approvals, when the host that
163    /// owns this service publishes its parent queue here.
164    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    /// Create an empty service instance.
175    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    /// Override the trusted, service-owned store used for durable reduction
198    /// bundles. Embedders and tests use this to keep all writes inside an
199    /// explicitly selected root; the CLI otherwise uses the normal
200    /// `$SUPERCODE_HOME/sessions` location.
201    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    /// ORCH-9: publish the parent's own subagent-approval queue into
207    /// `harness.v1.approvals.list`.
208    ///
209    /// This is the SAME `Arc` an [`crate::Agent`] pushes into
210    /// (`Agent::pending_child_approvals`), so a host that runs supercode's own
211    /// loop beside this service surfaces those requests through the uniform
212    /// door without copying them anywhere.
213    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    /// ORCH-9: every approval request this service can see, newest last.
221    ///
222    /// Two sources, both live: the requests outstanding on the open runtime
223    /// connections, and supercode's own queued subagent approvals. There is
224    /// no file or database source at the pinned harness versions (see
225    /// [`crate::approvals`]), so a stored or proposal row is never produced.
226    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    /// ORCH-20 (controlled tier): answer one listed approval request by its
246    /// row id and one uniform decision.
247    ///
248    /// The decision is translated into the option token and reply envelope
249    /// the door that raised the request already accepts
250    /// ([`crate::approvals::plan_reply`]), and the answer is then sent by
251    /// calling `harness.v1.runtimes.respond` itself — the same code path, the
252    /// same adapter, the same bookkeeping that drops the row. This verb adds
253    /// a translation and nothing else.
254    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(&params.id, &choice)
284            .map_err(|error| ServiceError::InvalidParams(error.to_string()))?;
285        // The harness's own door, unchanged: this is the identical call
286        // `harness.v1.runtimes.respond` performs for a caller who built the
287        // envelope by hand, including dropping the answered row.
288        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    /// Return the edge-triggered wakeup used by session-index filesystem
306    /// subscriptions. Transports can await this instead of polling indexes.
307    #[cfg(feature = "adapter-api")]
308    pub fn session_index_notifier(&self) -> Arc<Notify> {
309        Arc::clone(&self.index_notifier)
310    }
311
312    /// Handle one JSON-RPC 2.0 request and return one JSON-RPC response.
313    #[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    /// Handle either a persisted-session request or an asynchronous live
334    /// runtime request.
335    #[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        // ORCH-19 controlled tier. Answered here rather than through the SDK
380        // operation dispatch below so the harness's OWN refusal reaches the
381        // caller: `sdk_error` collapses every `UnsupportedAction` to one
382        // generic sentence, and the whole point of this tier is that a
383        // refusal names which door the harness does have.
384        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    /// Poll all active subscriptions once and return zero or more JSON-RPC
500    /// notifications. Recoverable follower errors are delivered as events.
501    #[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    /// Report each followed session's live-runtime lifecycle state on that
533    /// session's own subscription, emitting only when the state changes.
534    ///
535    /// A growing transcript is not evidence that an agent is working, so the
536    /// state comes from the live-runtime registry and nowhere else. A followed
537    /// session with no registered Supercode runtime — a harness running outside
538    /// Supercode — reports `persisted`, which says plainly that its activity is
539    /// unknown rather than guessing at it. These events carry no sequence
540    /// number and no transcript content; they never interleave with the
541    /// content follower's sequenced stream.
542    #[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                // A failed registry read is not evidence of a state change.
555                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    /// Poll normalized activity subscriptions, emitting only proven state
574    /// transitions. Every subscription is bulk-sampled so stock-harness
575    /// process and registry discovery happens once per UI, not once per row.
576    #[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                // A failed evidence read proves no transition. Retain the last
593                // good state instead of flashing every row to persisted.
594                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    /// Drain native-store invalidations and emit revisioned descriptor deltas.
627    /// An idle subscription performs no catalog or transcript reads between
628    /// its minute-scale recovery reconciliations.
629    #[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(&params.locators, &params.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    /// Non-blockingly sample one event from every connected live runtime.
708    #[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            // Drain what the runtime already has: a turn is several events
739            // (updates, then the protocol's completion), and delivering one
740            // per poll would cost a poll interval each. A zero timeout takes
741            // only what is ready — an idle runtime costs nothing.
742            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                        // ORCH-9: a permission/approval request arrives as an
747                        // ordinary event; it becomes listable here and stops
748                        // being listable when `runtimes.respond` answers it.
749                        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            // A connection that is gone cannot answer anything it was
826            // holding; those requests stop being listable with it.
827            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            // ORCH-21 controlled tier. Each verb translates to the HARNESS'S
854            // OWN profile verb and runs it (`crate::profiles_control`);
855            // supercode makes and removes nothing itself. The row returned is
856            // re-read through the ORCH-10 loader afterwards, and `ran`
857            // narrates the exact command.
858            "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            // ORCH-15 observed tier: which profile / agent a surface tuple
868            // resolves to, read from each gateway harness's own config.
869            "harness.v1.routes.list" => routes_call(params),
870            // ORCH-16 observed tier: inbound webhook routes / hook mappings.
871            "harness.v1.triggers.list" => triggers_call(params),
872            // ONT-4: the orchestration doors. One home folder in, one typed orchestration
873            // value out (and back). Every one of the four is
874            // `crate::orchestration_doors`, which the `supercode orchestration` verbs call
875            // too — the RPC adds nothing but the envelope. A vault VALUE
876            // never crosses this wire: a load or a compile answers with the
877            // `.env` KEY NAMES, and a caller that needs a value reads the
878            // home's own `.env`.
879            // the workflow layer's read door: a harness's board as one typed value,
880            // the same code the `supercode workflow load` verb calls
881            "harness.v1.workflow.load" => {
882                let params = decode::<WorkflowLoadParams>(params)?;
883                let read =
884                    crate::workflow_doors::load(params.from, &params.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(&params.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                    &params.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, &params.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                    &params.source,
919                    params.source_flavor,
920                    &params.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            // a migration keeps the credential in this process: a compile and
928            // a save (import), a load and a decompile (export), composed here
929            // because composed by a client the secret would have to cross
930            // the wire
931            "harness.v1.orchestration.import" => {
932                let params = decode::<OrchestrationImportParams>(params)?;
933                let imported =
934                    crate::orchestration_doors::import(params.from, &params.home, &params.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, &params.root, &params.dest)
943                        .map_err(operation)?;
944                serde_json::to_value(report)
945                    .map_err(|error| ServiceError::Operation(error.to_string()))
946            }
947            // ORCH-12 observed tier: read and search the persistent memory
948            // documents a harness keeps on disk. Read-only — every write
949            // (`hermes memory off`, `openclaw memory forget|reset`, Claude
950            // Code's `/memory`) stays the harness's own verb. A harness with
951            // no memory store is refused with UnsupportedAction.
952            "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
953            // ORCH-11 observed tier: read-only enumeration of every harness's
954            // installed skill packages. An unknown harness id is refused with
955            // UnsupportedAction — every harness supports skills, so a filter
956            // that matches nothing is a caller error, never an empty listing.
957            "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            // ORCH-22 controlled tier: each verb goes through the door the
970            // HARNESS publishes — `hermes skills install|uninstall`,
971            // `openclaw skills install`, and for the core four the loader's
972            // own directory, which is the only skills door those harnesses
973            // have. supercode resolves no registry and unpacks no archive.
974            // The row returned is re-read through the ORCH-11 loader
975            // afterwards, and `ran` narrates exactly what was performed.
976            "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            // ORCH-9 observed tier: the approval requests waiting for an
983            // answer. At the pinned harness versions the only uniform source
984            // is a LIVE request held by an open runtime connection, plus
985            // supercode's own queued subagent approvals — neither Hermes
986            // 0.21.0 nor OpenClaw 2026.7.1-2 has an approvals door to read
987            // (see `crate::approvals`). A harness whose runtime cannot carry
988            // a protocol request at all is refused by name.
989            "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                // Claude Code is the one harness that publishes its RUNNING
1005                // sessions. The registry is read once per discovery and joined
1006                // by session id; every record in it has already survived a
1007                // `kill(pid, 0)` liveness check inside `read_registry`.
1008                let doors = crate::mail_route::LiveSessions::read(&query.homes);
1009                // A running session's name (`twin-76`) lives only in its
1010                // harness's registry, not in the transcript discovery searches:
1011                // a query naming one also finds that session, by its id.
1012                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                        // Only that session's own folder is scanned, not the whole store again.
1040                        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                            // A Codex prompt (an approval, a choice) writes nothing to its rollout, which
1078                            // reads that turn as working; the live row reads its pane, and says so.
1079                            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                // Preserve the metadata-only wire shape, but carry the catalog's
1101                // proof/counts when the caller explicitly requests preview search.
1102                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) = &params.options {
1112                    options.validate()?;
1113                    if let Some(result) = indexed_claude_window(&params.read.locator, options)? {
1114                        return Ok(result);
1115                    }
1116                    return load_session(&params.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                            &params.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(&params.read.locator, params.read.read_fidelity())
1130                } else {
1131                    self.catalog
1132                        .load_parent_with_fidelity(
1133                            &params.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                        &params.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(&params.subscription);
1175                Ok(json!({
1176                    "removed": self.followers.remove(&params.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(&params.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(&params.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(&params.subscription).is_some()
1242                }))
1243            }
1244            "harness.v1.sessions.import" => {
1245                let params = decode::<ImportSessionParams>(params)?;
1246                let session = Session::load_str(&params.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(&params.locator).map_err(operation)?;
1253                let artifact = session_artifact(&params.locator, &session, params.target_harness)?;
1254                if method == "harness.v1.sessions.export"
1255                    && params.target_harness == TransferFormat::Hermes
1256                {
1257                    // UNI-18: write through Hermes's own door, never into its store
1258                    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(&params.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(&params.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(&params.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(&params.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                // An artifact from another machine carries its whole source as a recovery file;
1313                // keeping its segments here lets a later write back to that format restore it
1314                // byte for byte on this machine too (docs/plans/portable-residue.md).
1315                for file in &params.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                    &params.cwd,
1327                    &params.homes,
1328                )
1329                .map_err(ServiceError::Operation)?;
1330                Ok(json!({"locator": locator}))
1331            }
1332            // ORCH-7 observed tier. Read-only: the handlers open the harness's
1333            // own job store (Claude Code's session JSONL, Hermes's and
1334            // OpenClaw's `cron/jobs.json`) and never write, fire, or schedule.
1335            "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(&params.harness, "jobs.get")?;
1347                match crate::jobs::get_job(&params.harness, &params.id, &params.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            // ORCH-18 controlled tier. Each verb translates to the HARNESS'S
1358            // OWN cron verb and runs it (`crate::jobs_control`); supercode
1359            // schedules nothing. The row returned is re-read from the
1360            // harness's store afterwards, and `ran` narrates the exact command
1361            // with any credential redacted.
1362            "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            // ORCH-8 observed tier. Read-only: the handlers open the harness's
1383            // own run store (Hermes's `cron/executions.db`, OpenClaw's
1384            // `cron_run_logs`) and never claim, retry, or prune a fire.
1385            "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(&params.harness, "runs.get")?;
1397                match crate::runs::get_run(&params.harness, &params.id, &params.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(&params.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                    &params.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(&params.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        // Durability order is intentional: the full source of truth lands
1477        // before either object that can refer to it. A crash may leave an
1478        // unused sidecar, but can never leave a reduced view whose originals
1479        // were not durably written first.
1480        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        // The receipt is issued only after a real disk reload. This proves
1507        // the exact files another process will consume, not the convenient
1508        // in-memory values that produced them.
1509        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        // `sc.reduction` is deliberately in-memory-only metadata: it must
1521        // never leak onto a provider-facing transcript. Reapplying the
1522        // durable log to the durable sidecar restores those ids. Comparing
1523        // its wire form with the transcript reloaded above proves that the
1524        // persisted view is exactly the deterministic projection before we
1525        // use the restamped form for inversion.
1526        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            &params.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    /// Recognize the one request family whose waiting happens entirely
1583    /// outside this service's state, and hand a transport the half it can run
1584    /// off the task that owns the service.
1585    ///
1586    /// Opening a runtime is the only door here that waits on a foreign
1587    /// program: it spawns the harness's own binary and completes that
1588    /// program's protocol handshake, which takes as long as the program takes
1589    /// to answer. A transport that awaited the whole request inline would
1590    /// stop reading its own input for that whole time, so ONE slow launch
1591    /// would queue every later request on the same server — including reads
1592    /// like `sessions.discover` that touch no runtime at all. Splitting the
1593    /// request lets the transport spawn [`RuntimeOpen::open`] and keep
1594    /// reading, then pay only the short bookkeeping half
1595    /// ([`Self::register_open_runtime`]) when the runtime is up.
1596    ///
1597    /// `None` for every other method: those are answered by
1598    /// [`Self::handle_async`] as before.
1599    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    /// Recognize a [`DETACHED_METHODS`] request and hand a transport the
1615    /// whole of it: the service-state half is read here and now, and what
1616    /// remains waits on a foreign program with nothing of this service's in
1617    /// hand.
1618    ///
1619    /// Same reason as [`Self::runtime_open`], different doors. Probing a
1620    /// harness starts it and completes its handshake; messaging a live
1621    /// session waits on its Claude relay's send; a conversation verb runs the
1622    /// harness's own CLI or calls its HTTP API. A transport that awaited any
1623    /// of those inline would stop reading its own input for that whole time,
1624    /// so one probe of an unhealthy harness would queue every later request
1625    /// on the same server.
1626    ///
1627    /// Unlike an opening runtime there is no bookkeeping half: the answer
1628    /// [`DetachedCall::run`] produces is the caller's complete response, so a
1629    /// transport writes it without coming back here.
1630    ///
1631    /// `None` for every other method — including the LIVE `sessions.new` /
1632    /// `sessions.reset` door and `runtimes.close`, which wait on a runtime
1633    /// connection this service owns and so are split off by
1634    /// [`Self::detach_runtime`] instead.
1635    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                            // The live door needs the open runtime connection
1663                            // this service owns; it stays inline.
1664                            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    /// Recognize the two doors that wait on a runtime THIS SERVICE OWNS, and
1681    /// hand a transport the whole of each by lending the connection out.
1682    ///
1683    /// `runtimes.close` surrenders its runtime for good; the LIVE
1684    /// `sessions.new` / `sessions.reset` door borrows one for the length of
1685    /// the slash command and gives it back through
1686    /// [`Self::finish_detached`]. Both are bounded by
1687    /// [`RUNTIME_CONTROL_DEADLINE`], and a wedged runtime spends all of it —
1688    /// which is exactly as long as a transport that awaited them inline would
1689    /// stop reading its own input.
1690    ///
1691    /// `None` for every other method, and for the `sessions.new` /
1692    /// `sessions.reset` doors that are not live: [`Self::detach`] owns those.
1693    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(&params))
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                // Everything but the live door — including a refusal and a
1717                // request naming no connection — is `detach`'s or
1718                // `handle_async`'s to answer.
1719                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    /// Take back whatever a detached call borrowed and hand over the caller's
1750    /// response. Every answer from [`DetachedCall::run`] comes through here,
1751    /// so a lent-out connection is back in the service before the response
1752    /// that used it is written.
1753    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    /// Answer a request split out by [`Self::runtime_open`] and already
1767    /// awaited by [`RuntimeOpen::open`]: register the runtime this service now
1768    /// owns and build its JSON-RPC response.
1769    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    /// Take ownership of an opened runtime.
1779    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                // A live endpoint is attached through the hosted-runtime backend whatever harness it
1806                // hosts (runtimes.attach_existing routes it there, and verifies its receipt), so it is
1807                // answered with that backend's capabilities, not the harness adapter's own.
1808                #[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": &params.harness,
1816                        "capabilities": SupercodeHttpRuntimeBackend::live_capabilities(),
1817                    }));
1818                }
1819                let backend = runtime_backend(&params)?;
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(&params.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(&params.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(&params.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(&params.connection)?
1879                        .respond(params.request_id, params.response),
1880                )
1881                .await?
1882                .map_err(operation)?;
1883                // ORCH-9: an answered request is no longer waiting for one.
1884                self.approvals.answered(&params.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(&params.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(&params.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(&params.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(&params.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(&params)?;
1933                close_runtime(runtime, process_group).await
1934            }
1935            _ => Err(ServiceError::MethodNotFound),
1936        }
1937    }
1938
1939    /// Deliver one message into a session that is running right now.
1940    #[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(&params).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, &params.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                    &params.harness,
1966                    &params.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            // Pin attachment to the executable hosting this runtime. A bare
2026            // `supercode` could resolve to an older global install whose CLI
2027            // does not understand the receipt it is being asked to open.
2028            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    /// What a caller is told about a connection that is out on a detached
2081    /// call. It is not gone and it is not free: it is mid-call, which is the
2082    /// same answer the runtime itself gives a second turn.
2083    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    /// Take a runtime OUT of the service for the duration of one detached
2090    /// call, leaving its name marked as lent out.
2091    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    /// Surrender a runtime for good: the connection and everything the
2106    /// service hung off it are gone before its teardown is even attempted.
2107    ///
2108    /// `close` is what a caller reaches for when a runtime has stopped
2109    /// answering, and a runtime that has stopped answering is exactly the one
2110    /// whose graceful close cannot complete: a hosted runtime's own loop
2111    /// parks on the call the runtime never answered, so it never dequeues the
2112    /// shutdown either. Keeping the entry until teardown succeeded made a
2113    /// wedged runtime permanent — every later call on that connection, and
2114    /// every new turn, answered "a harness turn is already in progress" with
2115    /// no way to take the connection back.
2116    /// [`Self::surrender_runtime`] for a `runtimes.close` call, refusing a connection that no longer holds
2117    /// the runtime the caller named.
2118    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(&params.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(&params.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    /// SIGKILL the process group of every runtime this service owns, without
2156    /// waiting on any of them.
2157    ///
2158    /// A host leaving for good calls this BEFORE dropping the service. The
2159    /// handle this service holds is not the runtime's connection: a hosted
2160    /// runtime's real transport lives in the task driving it, so neither
2161    /// exiting the process nor dropping these handles reaches the harness
2162    /// process — while dropping them does remove each runtime's live-runtime
2163    /// receipt. Signalling first is what keeps a removed receipt from
2164    /// advertising a harness that is still running.
2165    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    /// ORCH-19: run one conversation-lifecycle verb through the harness's own
2173    /// door.
2174    ///
2175    /// Two doors, one shape. A CLI / HTTP / own-store door is self-contained
2176    /// in [`crate::sessions_control`]. A LIVE door (Hermes's and OpenClaw's
2177    /// `/new` and `/reset`, which are slash commands their gateway interprets
2178    /// INSIDE a session) is performed here, because only the service owns the
2179    /// open runtime connection — the command is typed through the very same
2180    /// `send_input` path a human's message takes, so supercode invents no
2181    /// private channel.
2182    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            // The live door types the slash command through an open hosted
2192            // runtime, which only exists with the `adapter-api` feature; the
2193            // CLI / HTTP / own-store doors below need nothing extra.
2194            #[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                // Typing into a live session is a control call on an open
2221                // runtime, and a wedged runtime never accepts one, so it is
2222                // bounded exactly like the other control verbs. A transport
2223                // with a loop of its own lends the connection out instead of
2224                // waiting here: see [`Self::detach_runtime`].
2225                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    /// Answer an inventory request whole, for callers that have nowhere to
2234    /// put the waiting half. A transport with a loop of its own splits it
2235    /// instead: see [`Self::detach`].
2236    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    /// The half of an inventory request that reads this service's state:
2245    /// resolve the selection and count the persisted sessions each row
2246    /// reports. What remains — finding executables, asking them their
2247    /// version, and (at `probe: handshake`) starting each harness and
2248    /// completing its protocol handshake — touches no service state at all.
2249    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, &params.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), &params.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(&params.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                    &params.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            // ORCH-20: answering needs the live connection, so it takes the
2439            // async door and ends in `harness.v1.runtimes.respond`.
2440            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            // ORCH-19 controlled tier. Every verb goes through the HARNESS'S
2449            // OWN door — its CLI, its HTTP API, or its slash command typed
2450            // into a live driven session — and returns the row re-read from
2451            // the harness's store afterwards.
2452            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    // the lease must forward every verb its capabilities advertise; without
2511    // this, steer fell to the trait default and refused a turn it claimed
2512    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
2525/// One inventory request's waiting half, already separated from the service
2526/// state it reads. See [`HarnessSessionService::inventory_work`].
2527struct 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
2534/// Perform one conversation-lifecycle verb through a door that is
2535/// self-contained in [`crate::sessions_control`]: the harness's own CLI, its
2536/// HTTP API, the orchestrator daemon's socket, or supercode's own store.
2537/// Touches no service state, so this runs on any task. The LIVE door is not
2538/// here — it types its slash command through a runtime connection the service
2539/// owns, and is performed by [`HarnessSessionService::mutate_session`].
2540async fn run_session_mutation(
2541    verb: crate::SessionVerb,
2542    mutation: &crate::SessionMutation,
2543) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2544    // Only the HTTP door actually awaits anything. The CLI, store and daemon
2545    // doors run the harness's own program, or its store, with calls that
2546    // block the calling THREAD from start to finish — a future that never
2547    // yields, which no timeout around it can interrupt and which would hold a
2548    // runtime worker for as long as the harness takes. They go to a blocking
2549    // task, where blocking is what the thread is for.
2550    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
2566/// Probe every selected harness and assemble the report. Touches no service
2567/// state, so this runs on any task.
2568async 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, &params, 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    // ORC-7: the orchestrator publishes no runtime launch — it is not an
2601    // adapter supercode connects a turn to. What "installed" means for it
2602    // is that its Node daemon entry is present, so the row answers from
2603    // that instead of from a PATH lookup it could never satisfy.
2604    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        // The orchestrator's "executable" is a Node module, not a CLI
2614        // with a `--version` flag; running it to ask would start a daemon.
2615        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        // These two adapters have explicit native status/login contracts
2630        // and complete local evidence coverage (including Claude's macOS
2631        // Keychain-backed oauthAccount marker). Treating absent evidence
2632        // as unknown advertises a start that will only fail interactively.
2633        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                // Keep the isolated home alive through process teardown.
2805                // Otherwise the compiler may release the last meaningful
2806                // use after cloning `launch`, and a still-starting CLI can
2807                // recreate its state directory after Drop removed it.
2808                // Some Node-based launchers finish a short asynchronous
2809                // installation-id write just after their parent process
2810                // is reaped. Remove once immediately, allow that bounded
2811                // writer to settle, then perform the authoritative pass.
2812                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        // Stderr is retained as a runtime event, but is not transport health.
2915        // Grok, for example, can log an AuthorizationRequired error from an
2916        // optional background worker while its ACP session continues to send
2917        // updates and complete prompts normally.
2918        _ => 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    // Exact parent-only window: recursive/full-artifact requests retain the
2928    // existing owner. This is not a bounded display-history substitution.
2929    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                // Codex narrates while it works (`phase: commentary`); only its `final_answer` ends a turn
3028                && 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        // The reported window describes the top-level transcript. Applying it
3076        // recursively would silently truncate subagents without returning a
3077        // window for each child. Keep their histories complete while carrying
3078        // the caller's media policy through the tree.
3079        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    /// Optional fidelity for the READ surfaces (`sessions.load`,
3177    /// `sessions.follow`).
3178    ///
3179    /// Omitted means [`Fidelity::Semantic`]: these two methods only ever
3180    /// produce a read-only view, and a compacted or resumed-across-files
3181    /// transcript — the everyday shape of a long Claude Code session — has no
3182    /// losslessly reconstructable record graph, so refusing to render it made
3183    /// the mirror unusable rather than accurate. A caller that intends to
3184    /// CONTINUE from what it reads asks for a lossless level explicitly and
3185    /// gets the strict refusal back. Every other method (export, translate,
3186    /// branch, handoff, resume_instructions) is lossless-only and has no
3187    /// such knob.
3188    #[serde(default)]
3189    fidelity: Option<Fidelity>,
3190    /// Optional bounded frontend projection. Absent preserves the historical
3191    /// complete-session read contract.
3192    #[serde(default)]
3193    view: Option<SessionReadView>,
3194}
3195
3196#[derive(Deserialize)]
3197struct SessionReadView {
3198    /// Number of trailing normalized messages to return. Zero is treated as
3199    /// one so a caller cannot accidentally request an unbounded empty mode.
3200    #[serde(default)]
3201    tail_messages: Option<usize>,
3202    /// Whether Claude Code child transcripts belong in this view. The
3203    /// frontend default is false; the legacy no-view path remains true.
3204    #[serde(default)]
3205    include_subagents: bool,
3206    /// Preserve human-visible native history across model-context compaction.
3207    #[serde(default)]
3208    display_history: bool,
3209    /// Bound each individual text field so a single tool result cannot turn a
3210    /// small message window into a hundred-megabyte RPC response.
3211    #[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    /// A stable caller key. Retrying delivery after a lost receipt keeps the same envelope.
3317    #[serde(default)]
3318    idempotency_key: Option<String>,
3319    /// Channel observations retain their trust boundary instead of becoming peer instructions.
3320    #[serde(default)]
3321    channel: bool,
3322    /// Name the sender is known by (`fleet-board`, `aaron`). Replies to the
3323    /// message are filed in this sender's mailbox, read with
3324    /// `sessions.inbox`. Defaults to `supercode`.
3325    #[serde(default)]
3326    from_name: Option<String>,
3327    /// A voice front for the receiving session, never a separate agent.
3328    #[serde(default)]
3329    voice_for: Option<crate::mailbox::MailAddress>,
3330    /// Human-readable sender label, independent of its return mailbox.
3331    #[serde(default)]
3332    sender_name: Option<String>,
3333    /// Id of the message this one answers.
3334    #[serde(default)]
3335    in_reply_to: Option<String>,
3336    /// File one notice in the sender's mailbox when the receiver next goes idle.
3337    #[serde(default)]
3338    notify_when_idle: bool,
3339    /// The text is the session's own user speaking (a voice bridge, a board
3340    /// the owner types in): user mail, delivered as the user's own turn.
3341    /// `harness.v1` is served only to the machine's owner (its daemon admits
3342    /// operators only), which is the authority this carries.
3343    #[serde(default)]
3344    as_user: bool,
3345    /// Same storage roots discovery accepts, so a caller (and a test) can
3346    /// point the live-session registry somewhere other than `$HOME`.
3347    #[serde(default)]
3348    homes: crate::HarnessHomes,
3349}
3350
3351#[derive(Deserialize)]
3352#[serde(deny_unknown_fields)]
3353struct InboxParams {
3354    /// Sender name used with `sessions.message` (its mailbox), or
3355    #[serde(default)]
3356    from_name: Option<String>,
3357    /// an explicit session address `sc:<machine>:<harness>:<id>`.
3358    #[serde(default)]
3359    address: Option<String>,
3360    /// Include messages already read.
3361    #[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
3396/// Deliver `text` into a session that is running right now, or say why not.
3397///
3398/// A refusal is a RESULT, not a JSON-RPC error: "that session is persisted
3399/// only" is an answer about the session, which a mirror renders next to the
3400/// transcript, and this service's error envelope carries no structured data
3401/// field a machine-readable reason could survive in.
3402///
3403/// The door is chosen by the one router every sender uses
3404/// ([`crate::mail_route`]): a session supercode controls gets it through its
3405/// runtime (the default tier); a Claude session it does not control through a
3406/// Claude relay, so the reply comes back; a Codex session through its hook;
3407/// anything else is stored. Replies are filed in the sender's mailbox under
3408/// `reply_to`, read with `sessions.inbox`.
3409async 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(&params.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        &params.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    // A Claude session's registry name must still name only that session.
3450    if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3451        && crate::runtime_mail::controlled_runtime(
3452            HarnessId::CLAUDE_CODE,
3453            &params.locator.session_id,
3454        )
3455        .is_none()
3456    {
3457        if let Err(refusal) =
3458            crate::claude_peer::resolve_live_session(&params.homes, &params.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(&params.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) = &params.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
3609/// `sessions.message` with `as_user`: user mail, through the one door that
3610/// carries the user's authority (a hosted runtime's input, or the session's
3611/// pane when its composer is empty; otherwise it waits in the mailbox).
3612async 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(&params.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    /// Root processes (a terminal pane's shell, say).
3667    pids: Vec<u32>,
3668    #[serde(default)]
3669    homes: crate::HarnessHomes,
3670}
3671
3672/// `harness.v1.sessions.activity_under`: the harness session running beneath
3673/// each root process, with its activity, so a terminal's attention takes a
3674/// harness pane's liveness from the harness's own lifecycle.
3675async 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(&params.pids, &params.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
3693/// Mailbox address of an operator sender (a board, a person's shell) that is
3694/// not itself a harness session: `sc:<machine>:operator:<name>`.
3695fn 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
3717/// `harness.v1.sessions.inbox`: a sender's mailbox, each message with the
3718/// text its reader sees. Unread messages are claimed and acknowledged by this
3719/// read, so a second read does not return them again.
3720fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
3721    let address = match (&params.address, &params.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/// Source identity of one follow subscription, plus the last lifecycle state
3753/// already reported on it. The follower itself stays purely persistence-facing.
3754// Only the adapter-api poll reads these; the subscription bookkeeping itself is
3755// shared by both builds.
3756#[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
3770/// Add what makes an indexed row behaviorally equivalent to a discovered row: the attach
3771/// endpoint of a session supercode hosts, and the door a message reaches it by.
3772///
3773/// The durable index owns only persistence metadata. Both are projections: every message/attach
3774/// operation revalidates its authority, so publishing one here never trusts a stale browser-held
3775/// handle. The doors are read once per batch ([`crate::mail_route::LiveSessions`]).
3776fn 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    // A running session with no title of its own is known by its live name.
3783    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    // How a message reaches this session right now, chosen by the same
3807    // router every sender uses: `runtime` (supercode controls it), `native`,
3808    // `hook` or `stored`. Absent when no process is running it.
3809    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    // What a session waiting on its user is asking, so every surface can show the question and not only the state.
3816    // Read as `message list` reads it: the transcript's open call, else the prompt on the session's pane.
3817    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        // Waiting on its user (a question, an approval): the word every other surface uses.
3861        (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
3862        // The normalized activity object can honestly report a live owner even
3863        // when the stock harness never published a turn status. Preserve the
3864        // older field's stricter contract instead of guessing `running`.
3865        (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    /// UNI-18: a Hermes target. Its artifact is the Codex rollout that
3886    /// `hermes sessions import --from codex` reads; `sessions.export` performs
3887    /// that import into the Hermes home.
3888    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            // a Hermes artifact is the Codex rollout Hermes imports
3917            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    /// Where the continuation is written; unset roots are the environment's own, as discovery reads them.
3966    #[serde(default)]
3967    homes: HarnessHomes,
3968}
3969
3970/// A session supercode starts or resumes runs without approval prompts
3971/// (`yolo`) unless the caller asks for the harness's own (`default`).
3972#[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/// `harness.v1.workflow.load` parameters: which harness's board, and its home.
3990#[derive(Deserialize)]
3991struct WorkflowLoadParams {
3992    from: crate::workflow_doors::WorkflowHarness,
3993    home: PathBuf,
3994}
3995
3996/// ONT-4 `harness.v1.orchestration.load` parameters. `flavor` says which layout the
3997/// folder is read as; our own is the default.
3998#[derive(Deserialize)]
3999struct OrchestrationLoadParams {
4000    root: PathBuf,
4001    #[serde(default)]
4002    flavor: crate::orchestration_doors::HomeFlavor,
4003}
4004
4005/// ONT-4 `harness.v1.orchestration.save` parameters. `vault` is merged into the
4006/// home's own secrets; a caller that sends none keeps what is on disk.
4007#[derive(Deserialize)]
4008struct OrchestrationSaveParams {
4009    root: PathBuf,
4010    orchestration: crate::orchestration::Orchestration,
4011    #[serde(default)]
4012    vault: BTreeMap<String, String>,
4013}
4014
4015/// ONT-4 `harness.v1.orchestration.compile` parameters.
4016#[derive(Deserialize)]
4017struct OrchestrationCompileParams {
4018    from: crate::orchestration_doors::OrchestrationHarness,
4019    home: PathBuf,
4020}
4021
4022/// ONT-4 `harness.v1.orchestration.decompile` parameters. `source` is the home the
4023/// orchestration was compiled from: it is re-compiled to recover the io bookkeeping
4024/// that byte reuse and the live-store refusal (UNI-18) are decided from.
4025#[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/// `harness.v1.orchestration.import` parameters: another harness's home, and the
4038/// folder of ours it becomes.
4039#[derive(Deserialize)]
4040struct OrchestrationImportParams {
4041    from: crate::orchestration_doors::OrchestrationHarness,
4042    home: PathBuf,
4043    into: PathBuf,
4044}
4045
4046/// `harness.v1.orchestration.export` parameters: a folder of ours, and the home of
4047/// another harness it becomes.
4048#[derive(Deserialize)]
4049struct OrchestrationExportParams {
4050    to: crate::orchestration_doors::OrchestrationHarness,
4051    root: PathBuf,
4052    dest: PathBuf,
4053}
4054
4055/// `harness.v1.jobs.get` parameters.
4056#[derive(Deserialize)]
4057struct JobsGetParams {
4058    harness: String,
4059    id: String,
4060    #[serde(default)]
4061    homes: crate::HarnessHomes,
4062}
4063
4064/// ORCH-18: run one mutating job verb through the harness's own CLI.
4065///
4066/// The refusal ladder is deliberate: a harness with no scheduled-job concept
4067/// at all answers with the SAME sentence `jobs.list` gives it, and a harness
4068/// that has jobs but publishes no client-callable verb (Claude Code, whose
4069/// jobs are created by the model inside a session) answers with its own
4070/// reason. Neither is ever a silent no-op.
4071fn 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
4081/// ORCH-22: run one mutating skills verb through the harness's own door.
4082///
4083/// The refusal ladder mirrors `jobs.*`: a harness with no skills root at all
4084/// answers with the same sentence `skills.list` gives it, and a harness whose
4085/// door does not publish this verb (OpenClaw has no `skills remove` at the
4086/// pin) answers with its own reason. Neither is ever a silent no-op.
4087fn 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
4105/// The skills twin of [`job_control_error`], with the same mapping rule.
4106fn 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
4120/// ORCH-21: run one mutating profile verb through the harness's own CLI.
4121///
4122/// The refusal ladder mirrors `mutate_job`'s: a harness with no profile
4123/// concept at all answers with the SAME sentence `profiles.list` gives it, and
4124/// a harness that HAS profiles but publishes no client-callable verb (Codex's
4125/// file-authored `[profiles.<name>]` tables, supercode's compiled-in presets)
4126/// answers with its own reason. Neither is ever a silent no-op.
4127fn 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
4137/// The same mapping `job_control_error` applies, for the profile noun.
4138fn 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
4152/// Map a controlled-tier failure onto the service's error vocabulary. A verb
4153/// the harness lacks is `UnsupportedAction`; a harness verb that RAN and
4154/// failed carries its own stderr through as the operation error.
4155fn 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
4167/// Map an ORCH-19 controlled-tier failure onto the service's error
4168/// vocabulary. A verb the harness has no door for is `UnsupportedAction`; a
4169/// door that RAN and failed carries the harness's own stderr / HTTP body
4170/// through as the operation error.
4171fn 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
4181/// A harness without a scheduled-job concept refuses the verb rather than
4182/// answering with an empty list — an absent capability and an empty inventory
4183/// are different answers (the same rule `runtimes.capabilities` applies to
4184/// `steer`).
4185fn 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/// `harness.v1.runs.get` parameters.
4196#[derive(Deserialize)]
4197struct RunsGetParams {
4198    harness: String,
4199    id: String,
4200    #[serde(default)]
4201    homes: crate::HarnessHomes,
4202}
4203
4204/// A harness with no run store refuses the verb rather than answering with an
4205/// empty history — the same rule `jobs.list` applies. Claude Code lands here
4206/// on purpose: its cron fires are ordinary turns inside the session that
4207/// created the job, so there is no fire record to list.
4208fn 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    /// Omit subprocess-based `--version` calls when a latency-sensitive UI only needs readiness.
4278    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/// Receipt-backed evidence that a harness has a RUNNING instance right now,
4332/// distinct from being merely installed (UNI-7). Detection is passive and
4333/// default-on: a gateway liveness connect for daemon harnesses, a fresh
4334/// SQLite WAL stamp for store-writer harnesses (precedent: the opencode
4335/// follower's -wal/-shm freshness). Control stays behind per-connection
4336/// grants — this reports observations only.
4337/// ORCH-17: the gateway-health noun on an inventory row. Derived from the
4338/// UNI-7 running-instance probe (Hermes: `state.db-wal` freshness; OpenClaw:
4339/// a TCP connect to the gateway endpoint resolved from its OWN config) plus
4340/// the executable version — never by starting anything.
4341#[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/// ORCH-17: `gateway` on a `harness.v1.harnesses.list` row.
4350#[derive(Debug, Clone, Serialize)]
4351pub struct GatewayHealth {
4352    pub state: GatewayState,
4353    /// The endpoint supercode would connect to (OpenClaw: the gateway
4354    /// WebSocket resolved from `openclaw.json`; core harnesses: their
4355    /// declared connect address when one exists). `None` when the harness
4356    /// has no single endpoint (Hermes multiplexes platforms).
4357    #[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    /// What the verdict rests on, or why it is `unknown`.
4362    pub evidence: String,
4363    pub checked_at_ms: u64,
4364}
4365
4366/// OpenClaw's gateway WebSocket endpoint, resolved from its own config the
4367/// way the registry's connect descriptor prescribes (`gateway.url`, else
4368/// `gateway.port`, else the documented default).
4369fn 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
4390/// Ask Hermes itself (`hermes gateway status`, read-only, ~1 s) whether its
4391/// gateway is up. The command is per-host launchd/systemd text without a JSON
4392/// form at 0.19–0.21; the verdict is read from the lines it prints:
4393/// "supervised by launchd (PID …)" / "is running" → up, "not running" /
4394/// "not installed" → down, anything else → no verdict. `SUPERCODE_HERMES_BIN`
4395/// overrides the executable so a fake can stand in under test.
4396fn 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                    // The harness's own door outranks the WAL heuristic: an idle
4448                    // gateway writes nothing for minutes yet is up.
4449                    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        // ORC-7: the orchestrator's gateway IS its daemon, and the daemon's
4472        // own lease file is the record of it. A lease naming a live pid is
4473        // up; a lease whose process is gone is down and says so as a STALE
4474        // lease, never as "no lease"; no lease at all is down. Nothing is
4475        // started, and no port is guessed — the daemon multiplexes adapters
4476        // the way Hermes does, so it has no single endpoint either.
4477        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    /// How the instance was detected.
4526    method: RunningInstanceMethod,
4527    /// The evidence the verdict rests on (endpoint reached / WAL path+age).
4528    evidence: String,
4529    /// Epoch-ms instant the probe executed.
4530    checked_at_ms: u64,
4531}
4532
4533#[derive(Debug, Clone, Copy, Serialize)]
4534#[serde(rename_all = "snake_case")]
4535enum RunningInstanceMethod {
4536    /// A TCP connect to the harness's own configured gateway endpoint
4537    /// succeeded.
4538    GatewayConnect,
4539    /// The harness's session store has an active SQLite WAL (a live writer
4540    /// holds the store open and stamped it recently).
4541    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
4551/// OpenClaw: the gateway endpoint comes from the harness's OWN config
4552/// (`<home>/.openclaw/openclaw.json` — `gateway.url` or `gateway.port`,
4553/// default port 18789); a successful TCP connect is the running signal.
4554fn 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
4597/// Hermes: `<home>/.hermes/state.db-wal` freshly modified means a live writer
4598/// holds the store open (SQLite WAL exists only while a connection is open;
4599/// a recent stamp distinguishes an active instance from a stale crash
4600/// leftover).
4601fn 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
4618/// Default-on running-instance detection for the harnesses that have one.
4619fn 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    /// Receipt-backed running-instance detection (None = not detected or the
4645    /// harness has no running-instance concept). Distinct from `installed`.
4646    #[serde(skip_serializing_if = "Option::is_none")]
4647    running: Option<RunningInstance>,
4648    /// ORCH-17: gateway health derived from `running` + the harness's own config.
4649    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/// A session supercode starts or resumes runs without approval prompts
4668/// (`yolo`) unless the caller asks for the harness's own (`default`).
4669#[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    /// MCP servers to mount into the new session through the harness's own
4683    /// start door (ORC-6). Backends without such a door ignore them.
4684    #[serde(default)]
4685    mcp_servers: Vec<crate::McpServerLaunch>,
4686    /// The session's approval policy, where the harness's start door takes one (Codex).
4687    #[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    /// MCP servers to mount into the resumed session (the start door's own
4699    /// field, carried again because a session's tools die with its process).
4700    #[serde(default)]
4701    mcp_servers: Vec<crate::McpServerLaunch>,
4702    /// The session's approval policy, carried again on resume as on start (Codex).
4703    #[serde(default)]
4704    approval_policy: Option<String>,
4705}
4706
4707#[derive(Deserialize)]
4708struct RuntimeConnectionParams {
4709    connection: String,
4710}
4711
4712/// `runtimes.close`: a caller that recorded `runtime_id` with the connection closes it only while the
4713/// connection still holds that runtime (a restarted daemon numbers its connections afresh).
4714#[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        // A session that came from `format` before returns its source records verbatim for the
4873        // prefix the residue store holds (docs/plans/portable-residue.md).
4874        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
5305/// Stage the resolved gateway credential in a private (0600) file so the
5306/// bridge can read it via `--token-file` — the delivery the real `openclaw
5307/// acp` accepts. One stable file per endpoint (keyed by an address digest,
5308/// no secret material in the name), overwritten on every connect so files
5309/// never accumulate and a rotated token never goes stale on disk.
5310fn 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
5333/// Open a connect-mode descriptor: resolve the endpoint address and
5334/// credential from the harness's own config file and build the backend that
5335/// joins the already-running endpoint. Fails closed with a specific
5336/// diagnostic when the config cannot be resolved or the declared protocol has
5337/// no connect-capable client yet.
5338fn 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            // OpenClaw's own `openclaw acp` binary is the gateway client: a
5361            // stdio ACP bridge that joins the RUNNING gateway at the resolved
5362            // endpoint. Blind-walk finding 2026-08-31: the real bridge does
5363            // NOT honor OPENCLAW_GATEWAY_TOKEN from the environment — the
5364            // credential must arrive via `--token-file` (never bare `--token`
5365            // on argv, where process listings could read it). The env var is
5366            // still set for older bridges that did read it. Requires openclaw
5367            // >= 2026.7: the 2026.2 bridge drops its gateway socket
5368            // mid-prompt and advertises no session resume (executed finding,
5369            // docs/interop/research/openclaw-acp-dialect-2026-08-30.json).
5370            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            // The bridge program comes from the descriptor's own default
5384            // launch (the compiled registry pins `openclaw`), so tests can
5385            // substitute an absolute mock-bridge path without touching
5386            // process-global state.
5387            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
5411/// The registry's connect-mode launch for this harness, honored only when the
5412/// caller supplied neither an explicit launch nor a base URL.
5413fn 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
5437/// The doors that open a runtime: each spawns or joins a program and waits on
5438/// that program's protocol handshake before it can answer.
5439pub 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
5446/// How long a runtime gets to finish opening before its caller is answered an
5447/// error instead. A program that never speaks the protocol at all — the wrong
5448/// binary, a shim that prints usage and waits — never answers the handshake,
5449/// so the wait is unbounded without this.
5450pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5451
5452/// How long a control call on an ALREADY-open runtime — send input, interrupt,
5453/// steer, respond, close — gets before its caller is answered an error
5454/// instead. A live runtime answers these in milliseconds; a wedged one never
5455/// answers at all, and `close` is exactly what a caller reaches for when it
5456/// suspects that.
5457pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5458
5459/// The doors whose work happens entirely OUTSIDE this service's state once
5460/// its state has been read: probing harnesses, relaying a message into a
5461/// live session, and performing a conversation verb through a harness's own
5462/// CLI / HTTP / store door. Every one of them waits on a child process or a
5463/// network peer. See [`HarnessSessionService::detach`].
5464pub 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
5474/// How long a request moved off a transport's loop gets before its caller is
5475/// answered an error instead. Each of these already bounds its own inner
5476/// waits (a probe's handshake, the relay's send); this is the backstop for
5477/// the ones that do not — a harness CLI that never exits — so no caller waits
5478/// forever on a detached task no one is watching.
5479pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5480
5481/// How long `sessions.discover` gets before its caller is answered an error
5482/// instead. Discovery reads each harness's own store, and a store on a cold
5483/// or unavailable mount answers at the filesystem's pace rather than its own.
5484///
5485/// Deliberately shorter than the clients' own request deadline (30s): the
5486/// server's answer names the store that did not answer, and it is only read
5487/// if it lands before the client stops listening.
5488pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5489
5490/// Bound one control call on an open runtime by [`RUNTIME_CONTROL_DEADLINE`],
5491/// naming the method and the bound when it blows.
5492async 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
5506/// One [`RUNTIME_OPEN_METHODS`] request, parsed but not yet started. See
5507/// [`HarnessSessionService::runtime_open`] for why it exists apart from
5508/// [`HarnessSessionService::handle_async`].
5509pub struct RuntimeOpen {
5510    id: Value,
5511    method: String,
5512    params: Value,
5513}
5514
5515impl RuntimeOpen {
5516    /// Do the waiting: spawn or join the program and complete its handshake,
5517    /// bounded by [`RUNTIME_OPEN_DEADLINE`]. Touches no service state, so this
5518    /// runs on any task.
5519    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
5526/// The result of [`RuntimeOpen::open`], ready for
5527/// [`HarnessSessionService::finish_runtime_open`].
5528pub struct OpenedRuntime {
5529    id: Value,
5530    outcome: std::result::Result<OpenRuntime, ServiceError>,
5531}
5532
5533/// One detached request: the half that reads this service's state already
5534/// done, and the half that waits not yet started. See
5535/// [`HarnessSessionService::detach`] and
5536/// [`HarnessSessionService::detach_runtime`].
5537pub struct DetachedCall {
5538    id: Value,
5539    method: String,
5540    work: std::result::Result<Work, ServiceError>,
5541}
5542
5543impl DetachedCall {
5544    /// Do the waiting and answer. Runs on any task: whatever this call needed
5545    /// from the service was taken before it left.
5546    pub async fn run(self) -> DetachedAnswer {
5547        let Self { id, method, work } = self;
5548        match work {
5549            // A call holding a runtime is already bounded by
5550            // RUNTIME_CONTROL_DEADLINE, and its future OWNS that connection:
5551            // a second timeout around it would drop the connection mid-call
5552            // and take down a runtime its caller still has.
5553            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
5581/// One detached call's complete answer, plus whatever it must hand back to
5582/// the service before that answer is written. See
5583/// [`HarnessSessionService::finish_detached`].
5584pub struct DetachedAnswer {
5585    response: Value,
5586    returned: Option<ReturnedRuntime>,
5587}
5588
5589impl DetachedAnswer {
5590    /// The caller's JSON-RPC response, for a transport that owns no service
5591    /// to give a borrowed connection back to.
5592    pub fn into_response(self) -> Value {
5593        self.response
5594    }
5595}
5596
5597/// A connection lent to a detached call, on its way back to the service that
5598/// owns it.
5599pub struct ReturnedRuntime {
5600    connection: String,
5601    runtime: Box<dyn RuntimeConnection>,
5602}
5603
5604/// The waiting half of one detached request: with nothing of the service's
5605/// in hand, or holding a connection the service lent out for the call.
5606enum Work {
5607    Free(DetachedWork),
5608    Runtime(RuntimeWork),
5609}
5610
5611/// The waiting half of one detached request that holds nothing of the
5612/// service's.
5613enum DetachedWork {
5614    /// Probe the selected harnesses: find their executables, ask each its
5615    /// version, and at `probe: handshake` start each one and complete its
5616    /// protocol handshake.
5617    Inventory(InventoryWork),
5618    /// Relay one message into a live session.
5619    Message(MessageSessionParams),
5620    /// Perform one conversation verb through the harness's own CLI, HTTP API,
5621    /// daemon socket, or supercode's own store.
5622    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(&params).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
5642/// One detached call that holds a runtime connection for its whole run.
5643enum RuntimeWork {
5644    /// Tear down a runtime the service has already surrendered.
5645    Close {
5646        runtime: Box<dyn RuntimeConnection>,
5647        process_group: Option<u32>,
5648    },
5649    /// Type one live slash command through a borrowed connection, then give
5650    /// the connection back.
5651    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
5661/// What one [`RuntimeWork`] answers with: the caller's result, and the
5662/// connection to give back when the call only borrowed one.
5663type 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
5697/// Tear down a runtime already out of the service, within
5698/// [`RUNTIME_CONTROL_DEADLINE`].
5699async 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            // Dropping the handle is not enough: the process that stopped
5710            // answering is held by a task parked on it, so nothing here runs
5711            // its Drop. Signal the group the graceful path would have
5712            // signalled, then say so.
5713            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
5724/// The conversation a live `sessions.new` / `sessions.reset` acts on: the one
5725/// the request named, or the runtime's own session.
5726fn 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
5734/// Type one harness slash command into a live session through the very same
5735/// `send_input` path a human's message takes, within
5736/// [`RUNTIME_CONTROL_DEADLINE`].
5737async 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
5758/// A runtime that is up and whose handshake completed, with what the service
5759/// needs to take ownership of it.
5760enum OpenRuntime {
5761    /// supercode spawned this process, so it also hosts it: a frontend server,
5762    /// a live-runtime registration and a terminal launch of its own.
5763    Hosted {
5764        runtime: Box<dyn RuntimeConnection>,
5765        capabilities: crate::RuntimeCapabilities,
5766        workspace: PathBuf,
5767        /// A new session (`start`): it has no transcript yet, so there is no history to look for.
5768        fresh: bool,
5769    },
5770    /// `attach_existing` joined a process supercode does not own. It is
5771    /// registered as a bare connection and hosts nothing.
5772    Joined { runtime: Box<dyn RuntimeConnection> },
5773}
5774
5775/// Open the runtime one [`RUNTIME_OPEN_METHODS`] request asks for, within
5776/// [`RUNTIME_OPEN_DEADLINE`]. The error a blown deadline answers names the
5777/// method and the bound, so a caller reads why it was cut loose instead of
5778/// waiting on a handshake that is never coming.
5779async 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(&params.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(&params.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(&params.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(&params.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(&params.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(&params.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
5906/// Wrap one service outcome in its JSON-RPC 2.0 envelope.
5907fn 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 &params.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) = &params.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
6070/// Disposable harness state for a no-prompt readiness probe. Merely opening
6071/// several stock CLIs writes a session header or migrates configuration, so a
6072/// handshake must never point at the user's real home. Authentication files
6073/// are copied into the private temporary home; all writes disappear with the
6074/// guard after the connection closes.
6075struct 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        // supercode reads its own config home ($SUPERCODE_HOME, else
6098        // $XDG_CONFIG_HOME/supercode, else ~/.config/supercode), not a fixed
6099        // place under HOME: a login kept under XDG_CONFIG_HOME probed as
6100        // "no API key found" while `supercode run` answered.
6101        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        // The gateway endpoint + token live in openclaw's own config; without
6166        // it the isolated probe dials the default endpoint unauthenticated
6167        // (PARITY-24 finding 2026-08-31).
6168        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        // Hermes keeps its provider selection in config.yaml, its OAuth
6182        // credential pool in auth.json, and API keys in .env; without them
6183        // the isolated probe sees "No LLM provider configured" for a
6184        // hermes that answers fine from the user's real home.
6185        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    // macOS keeps Claude Code's OAuth login in the Keychain, so
6343    // `.claude/.credentials.json` never exists there and the file probe above
6344    // reports a signed-in install as unauthenticated forever. A completed
6345    // login also writes an `oauthAccount` record into `~/.claude.json` on
6346    // every platform — file-based, prompt-free evidence (querying the
6347    // Keychain itself from an unsigned daemon can raise a UI prompt).
6348    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
6395/// The process group a runtime's own handle names, when it names one.
6396///
6397/// Every adapter that spawns a local process spawns it as its own group
6398/// leader (`Command::process_group(0)`), so the endpoint's pid IS the group
6399/// id. A runtime reached over HTTP, or one supercode joined rather than
6400/// spawned, names no group here and is left alone.
6401fn 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
6408/// SIGKILL a wedged runtime's whole process group, reporting whether there
6409/// was one to signal. This is the same group teardown a graceful `close`
6410/// performs; it runs here only when the graceful path blew its deadline,
6411/// because the task parked on the unanswered call still owns the process
6412/// handle and so no `Drop` of ours can reach it.
6413fn 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/// ORCH-12 `harness.v1.memory.show|search` params. `homes` is the same
6508/// storage-root override every read-only method accepts, so a caller can
6509/// point the read at a fixture home without touching the real ones.
6510#[derive(Debug, Clone, Deserialize, Default)]
6511#[serde(default)]
6512struct MemoryRequest {
6513    /// Harness whose store is read. Required.
6514    harness: Option<String>,
6515    /// The needle, required by `search`.
6516    query: Option<String>,
6517    /// Hermes profile, OpenClaw agent, or Claude Code project.
6518    profile: Option<String>,
6519    /// Claude Code session id selecting a project store (`show` only).
6520    session: Option<String>,
6521    /// Include each document's whole text (`show` only).
6522    full: bool,
6523    /// Treat `query` as a regular expression (`search` only).
6524    regex: bool,
6525    /// Working tree whose project store is read.
6526    cwd: Option<std::path::PathBuf>,
6527    /// Storage roots to read.
6528    homes: crate::HarnessHomes,
6529}
6530
6531/// Read the memory noun. A harness with no memory store fails with
6532/// `UnsupportedAction` (RPC `-32020`), never an empty list.
6533fn 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/// ORCH-10 `harness.v1.profiles.list|get` params. `homes` is the same
6585/// storage-root override every read-only method accepts, so a caller can
6586/// point the read at a fixture home without touching the real ones.
6587#[derive(Debug, Clone, Deserialize)]
6588#[serde(default)]
6589struct ProfilesQuery {
6590    /// Restrict the listing to one harness. `get` requires it.
6591    harness: Option<String>,
6592    /// Profile name, required by `get`.
6593    name: Option<String>,
6594    /// Storage roots to read.
6595    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
6608/// Read the profile noun. A harness with no profile concept fails with
6609/// `UnsupportedAction` (RPC `-32020`), never an empty list.
6610fn 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/// ORCH-14 `harness.v1.channels.list|status` params, the same storage-root
6648/// override every read-only method accepts so a caller can point the read at
6649/// a fixture home without touching the real ones.
6650#[derive(Debug, Clone, Deserialize)]
6651#[serde(default)]
6652struct ChannelsQuery {
6653    /// Restrict the listing to one harness. `status` requires it.
6654    harness: Option<String>,
6655    /// Channel name, required by `status`.
6656    name: Option<String>,
6657    /// Storage roots to read.
6658    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/// Read the channel noun. A harness with no channel concept fails with
6672/// `UnsupportedAction` (RPC `-32020`), never an empty list. No row carries a
6673/// token, key or secret — see `crate::channels` "Secrecy".
6674#[derive(Debug, Clone, Deserialize)]
6675#[serde(default)]
6676struct RoutesQuery {
6677    harness: Option<String>,
6678    /// Restrict to routes targeting one profile / agent.
6679    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}