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_with_metadata, resolve_live_runtime};
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.usage",
52    "harness.v1.sessions.load",
53    "harness.v1.sessions.recorded_config",
54    "harness.v1.sessions.follow",
55    "harness.v1.sessions.unfollow",
56    "harness.v1.sessions.activity.subscribe",
57    "harness.v1.sessions.activity.unsubscribe",
58    "harness.v1.sessions.activity_under",
59    "harness.v1.sessions.index.subscribe",
60    "harness.v1.sessions.index.resize",
61    "harness.v1.sessions.index.unsubscribe",
62    "harness.v1.sessions.message",
63    "harness.v1.sessions.inbox",
64    "harness.v1.sessions.import",
65    "harness.v1.sessions.export",
66    "harness.v1.sessions.translate",
67    "harness.v1.sessions.reduce",
68    "harness.v1.sessions.branch",
69    "harness.v1.sessions.handoff",
70    "harness.v1.sessions.materialize",
71    "harness.v1.sessions.resume_instructions",
72    "harness.v1.skills.list",
73    "harness.v1.skills.install",
74    "harness.v1.skills.remove",
75    "harness.v1.memory.show",
76    "harness.v1.memory.search",
77    "harness.v1.jobs.list",
78    "harness.v1.jobs.get",
79    "harness.v1.jobs.create",
80    "harness.v1.jobs.update",
81    "harness.v1.jobs.pause",
82    "harness.v1.jobs.resume",
83    "harness.v1.jobs.run",
84    "harness.v1.jobs.delete",
85    "harness.v1.jobs.notepad",
86    "harness.v1.jobs.notepad_set",
87    "harness.v1.jobs.notepad_delete",
88    "harness.v1.sessions.new",
89    "harness.v1.sessions.reset",
90    "harness.v1.sessions.archive",
91    "harness.v1.sessions.delete",
92    "harness.v1.runs.list",
93    "harness.v1.runs.get",
94    "harness.v1.approvals.list",
95    "harness.v1.approvals.resolve",
96    "harness.v1.runtimes.capabilities",
97    "harness.v1.runtimes.start",
98    "harness.v1.runtimes.resume",
99    "harness.v1.runtimes.attach_existing",
100    "harness.v1.runtimes.attach",
101    "harness.v1.runtimes.send_input",
102    "harness.v1.runtimes.interrupt",
103    "harness.v1.runtimes.steer",
104    "harness.v1.runtimes.respond",
105    "harness.v1.runtimes.terminal_instructions",
106    "harness.v1.runtimes.acquire_control",
107    "harness.v1.runtimes.heartbeat",
108    "harness.v1.runtimes.detach",
109    "harness.v1.runtimes.close",
110    "harness.v1.profiles.list",
111    "harness.v1.profiles.get",
112    "harness.v1.profiles.create",
113    "harness.v1.profiles.delete",
114    "harness.v1.channels.list",
115    "harness.v1.routes.list",
116    "harness.v1.triggers.list",
117    "harness.v1.channels.status",
118    "harness.v1.orchestration.load",
119    "harness.v1.orchestration.save",
120    "harness.v1.orchestration.compile",
121    "harness.v1.orchestration.decompile",
122    "harness.v1.orchestration.import",
123    "harness.v1.orchestration.export",
124    "harness.v1.workflow.load",
125];
126
127/// Protocol namespace implemented by this service.
128pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
129/// Notification method emitted for followed-session changes.
130pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
131/// Notification method emitted for normalized session-activity transitions.
132pub const SESSION_ACTIVITY_EVENT_METHOD: &str = "harness.v1.sessions.activity_event";
133/// Notification method emitted for revisioned session-list changes.
134pub const SESSION_INDEX_EVENT_METHOD: &str = "harness.v1.sessions.index_event";
135/// Notification method emitted for live runtime events.
136pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
137
138/// Stateful persisted-session service. Each instance owns its follow
139/// subscriptions; discovery and loading remain read-only.
140pub struct HarnessSessionService {
141    catalog: HarnessCatalog,
142    followers: BTreeMap<String, SessionFollower>,
143    followed_sources: BTreeMap<String, FollowedSource>,
144    activity_subscriptions: BTreeMap<String, ActivitySubscription>,
145    index_subscriptions: BTreeMap<String, crate::session_index::SessionIndexSubscription>,
146    index_notifier: Arc<Notify>,
147    #[cfg(feature = "adapter-api")]
148    activity_monitor: crate::session_activity::SessionActivityMonitor,
149    next_subscription: u64,
150    runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
151    /// Connections lent to a detached call that is running right now. The
152    /// runtime itself is OUT of `runtimes` for that whole call, and these
153    /// names are how a second caller is told the connection is busy rather
154    /// than unknown.
155    runtimes_in_flight: BTreeSet<String>,
156    terminal_launches: BTreeMap<String, StructuredLaunch>,
157    runtime_sequences: BTreeMap<String, u64>,
158    next_runtime: u64,
159    reduction_store_root: Option<PathBuf>,
160    /// ORCH-9: live permission/approval requests outstanding on the open
161    /// runtime connections above, fed by the same event pump that publishes
162    /// `harness.v1.runtimes.event`.
163    approvals: crate::approvals::ApprovalRegistry,
164    /// ORCH-9: supercode's own queued subagent approvals, when the host that
165    /// owns this service publishes its parent queue here.
166    subagent_approvals: Option<Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>>,
167}
168
169impl Default for HarnessSessionService {
170    fn default() -> Self {
171        Self::new()
172    }
173}
174
175impl HarnessSessionService {
176    /// Create an empty service instance.
177    pub fn new() -> Self {
178        Self {
179            catalog: HarnessCatalog::new(),
180            followers: BTreeMap::new(),
181            followed_sources: BTreeMap::new(),
182            activity_subscriptions: BTreeMap::new(),
183            index_subscriptions: BTreeMap::new(),
184            index_notifier: Arc::new(Notify::new()),
185            #[cfg(feature = "adapter-api")]
186            activity_monitor: Default::default(),
187            next_subscription: 1,
188            runtimes: BTreeMap::new(),
189            runtimes_in_flight: BTreeSet::new(),
190            terminal_launches: BTreeMap::new(),
191            runtime_sequences: BTreeMap::new(),
192            next_runtime: 1,
193            reduction_store_root: None,
194            approvals: crate::approvals::ApprovalRegistry::new(),
195            subagent_approvals: None,
196        }
197    }
198
199    /// Override the trusted, service-owned store used for durable reduction
200    /// bundles. Embedders and tests use this to keep all writes inside an
201    /// explicitly selected root; the CLI otherwise uses the normal
202    /// `$SUPERCODE_HOME/sessions` location.
203    pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
204        self.reduction_store_root = Some(root.into());
205        self
206    }
207
208    /// ORCH-9: publish the parent's own subagent-approval queue into
209    /// `harness.v1.approvals.list`.
210    ///
211    /// This is the SAME `Arc` an [`crate::Agent`] pushes into
212    /// (`Agent::pending_child_approvals`), so a host that runs supercode's own
213    /// loop beside this service surfaces those requests through the uniform
214    /// door without copying them anywhere.
215    pub fn observe_subagent_approvals(
216        &mut self,
217        queue: Arc<std::sync::Mutex<Vec<crate::subagents::QueuedApproval>>>,
218    ) {
219        self.subagent_approvals = Some(queue);
220    }
221
222    /// ORCH-9: every approval request this service can see, newest last.
223    ///
224    /// Two sources, both live: the requests outstanding on the open runtime
225    /// connections, and supercode's own queued subagent approvals. There is
226    /// no file or database source at the pinned harness versions (see
227    /// [`crate::approvals`]), so a stored or proposal row is never produced.
228    pub fn approvals(&self, query: &crate::approvals::ApprovalsQuery) -> Vec<crate::ApprovalRow> {
229        let now = crate::approvals::now_ms();
230        let mut rows = self.approvals.rows(now);
231        if let Some(queue) = self.subagent_approvals.as_ref() {
232            let queued = queue
233                .lock()
234                .unwrap_or_else(std::sync::PoisonError::into_inner)
235                .clone();
236            rows.extend(crate::approvals::subagent_rows(&queued, now));
237        }
238        rows.retain(|row| query.matches(row));
239        rows.sort_by(|left, right| {
240            left.requested_at_ms
241                .cmp(&right.requested_at_ms)
242                .then_with(|| left.id.cmp(&right.id))
243        });
244        rows
245    }
246
247    /// ORCH-20 (controlled tier): answer one listed approval request by its
248    /// row id and one uniform decision.
249    ///
250    /// The decision is translated into the option token and reply envelope
251    /// the door that raised the request already accepts
252    /// ([`crate::approvals::plan_reply`]), and the answer is then sent by
253    /// calling `harness.v1.runtimes.respond` itself — the same code path, the
254    /// same adapter, the same bookkeeping that drops the row. This verb adds
255    /// a translation and nothing else.
256    async fn approvals_resolve(
257        &mut self,
258        params: Value,
259    ) -> std::result::Result<Value, ServiceError> {
260        let params = decode::<crate::approvals::ApprovalsResolveParams>(params)?;
261        if params.id.trim().is_empty() {
262            return Err(ServiceError::InvalidParams(
263                "approvals resolve requires the `id` of a listed approval row".into(),
264            ));
265        }
266        let choice = match (params.decision, params.option_id.as_deref()) {
267            (Some(_), Some(_)) => {
268                return Err(ServiceError::InvalidParams(
269                    "approvals resolve takes either `decision` or `option_id`, not both".into(),
270                ))
271            }
272            (Some(decision), None) => crate::approvals::ApprovalChoice::Decision(decision),
273            (None, Some(option)) => crate::approvals::ApprovalChoice::Option(option.to_string()),
274            (None, None) => {
275                return Err(ServiceError::InvalidParams(format!(
276                    "approvals resolve requires `decision` ({}) or an explicit `option_id`",
277                    crate::approvals::ApprovalDecision::ALL
278                        .map(|decision| decision.as_str())
279                        .join(" | "),
280                )))
281            }
282        };
283        let resolution = self
284            .approvals
285            .resolution(&params.id, &choice)
286            .map_err(|error| ServiceError::InvalidParams(error.to_string()))?;
287        // The harness's own door, unchanged: this is the identical call
288        // `harness.v1.runtimes.respond` performs for a caller who built the
289        // envelope by hand, including dropping the answered row.
290        self.runtime_call(
291            "harness.v1.runtimes.respond",
292            json!({
293                "connection": resolution.connection,
294                "request_id": resolution.request_id,
295                "response": resolution.response,
296            }),
297        )
298        .await?;
299        Ok(json!({
300            "id": params.id,
301            "decision": params.decision.map(|decision| decision.as_str()),
302            "option_id": resolution.option_id,
303            "resolved": true,
304        }))
305    }
306
307    /// Return the edge-triggered wakeup used by session-index filesystem
308    /// subscriptions. Transports can await this instead of polling indexes.
309    #[cfg(feature = "adapter-api")]
310    pub fn session_index_notifier(&self) -> Arc<Notify> {
311        Arc::clone(&self.index_notifier)
312    }
313
314    /// Handle one JSON-RPC 2.0 request and return one JSON-RPC response.
315    #[cfg(feature = "adapter-api")]
316    pub fn handle(&mut self, request: Value) -> Value {
317        let id = request.get("id").cloned().unwrap_or(Value::Null);
318        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
319            return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
320        }
321        let Some(method) = request.get("method").and_then(Value::as_str) else {
322            return rpc_error(id, -32600, "request is missing `method`");
323        };
324        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
325        match self.call(method, params) {
326            Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
327            Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
328            Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
329            Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
330            Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
331            Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
332        }
333    }
334
335    /// Handle either a persisted-session request or an asynchronous live
336    /// runtime request.
337    #[cfg(feature = "adapter-api")]
338    pub async fn handle_async(&mut self, request: Value) -> Value {
339        let method = request
340            .get("method")
341            .and_then(Value::as_str)
342            .unwrap_or_default();
343        if matches!(
344            method,
345            "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
346        ) {
347            let id = request.get("id").cloned().unwrap_or(Value::Null);
348            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
349                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
350            }
351            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
352            return match self.inventory_call(method, params).await {
353                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
354                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
355                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
356                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
357                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
358                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
359            };
360        }
361        if matches!(
362            method,
363            "harness.v1.harnesses.auth.methods"
364                | "harness.v1.harnesses.auth.begin"
365                | "harness.v1.harnesses.auth.verify"
366        ) {
367            let id = request.get("id").cloned().unwrap_or(Value::Null);
368            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
369                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
370            }
371            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
372            return match self.harness_authentication_call(method, params).await {
373                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
374                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
375                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
376                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
377                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
378                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
379            };
380        }
381        // ORCH-19 controlled tier. Answered here rather than through the SDK
382        // operation dispatch below so the harness's OWN refusal reaches the
383        // caller: `sdk_error` collapses every `UnsupportedAction` to one
384        // generic sentence, and the whole point of this tier is that a
385        // refusal names which door the harness does have.
386        if matches!(
387            method,
388            "harness.v1.sessions.new"
389                | "harness.v1.sessions.reset"
390                | "harness.v1.sessions.archive"
391                | "harness.v1.sessions.delete"
392        ) {
393            let id = request.get("id").cloned().unwrap_or(Value::Null);
394            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
395                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
396            }
397            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
398            let verb = match method {
399                "harness.v1.sessions.new" => crate::SessionVerb::New,
400                "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
401                "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
402                _ => crate::SessionVerb::Delete,
403            };
404            return match self.mutate_session(verb, params).await {
405                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
406                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
407                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
408                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
409                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
410                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
411            };
412        }
413        if method == "harness.v1.sessions.message" {
414            let id = request.get("id").cloned().unwrap_or(Value::Null);
415            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
416                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
417            }
418            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
419            return match self.message_call(params).await {
420                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
421                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
422                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
423                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
424                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
425                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
426            };
427        }
428        if matches!(
429            method,
430            "harness.v1.harnesses.settings" | "harness.v1.harnesses.configure"
431        ) {
432            let id = request.get("id").cloned().unwrap_or(Value::Null);
433            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
434                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
435            }
436            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
437            return match self.harness_settings_call(method, params) {
438                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
439                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
440                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
441                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
442                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
443                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
444            };
445        }
446        if method == "harness.v1.sessions.activity_under" {
447            let id = request.get("id").cloned().unwrap_or(Value::Null);
448            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
449            return match activity_under_call(params).await {
450                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
451                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
452                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
453                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
454                Err(_) => rpc_error(id, -32000, "sessions.activity_under failed"),
455            };
456        }
457        if method == "harness.v1.sessions.activity.subscribe" {
458            let id = request.get("id").cloned().unwrap_or(Value::Null);
459            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
460                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
461            }
462            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
463            return match self.subscribe_session_activity(params).await {
464                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
465                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
466                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
467                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
468                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
469                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
470            };
471        }
472        if let Some(operation) = SdkOperation::from_method(method) {
473            let id = request.get("id").cloned().unwrap_or(Value::Null);
474            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
475                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
476            }
477            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
478            return match self.execute(SdkRequest { operation, params }).await {
479                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
480                Err(error) => sdk_rpc_error(id, &error),
481            };
482        }
483        if !method.starts_with("harness.v1.runtimes.") {
484            return self.handle(request);
485        }
486        let id = request.get("id").cloned().unwrap_or(Value::Null);
487        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
488            return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
489        }
490        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
491        match self.runtime_call(method, params).await {
492            Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
493            Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
494            Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
495            Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
496            Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
497            Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
498        }
499    }
500
501    /// Poll all active subscriptions once and return zero or more JSON-RPC
502    /// notifications. Recoverable follower errors are delivered as events.
503    #[cfg(feature = "adapter-api")]
504    pub fn poll(&mut self) -> Vec<Value> {
505        let mut notifications = Vec::new();
506        for (subscription, follower) in &mut self.followers {
507            match follower.poll() {
508                Ok(Some(event)) => notifications.push(json!({
509                    "jsonrpc": "2.0",
510                    "method": SESSION_EVENT_METHOD,
511                    "params": {
512                        "subscription": subscription,
513                        "event": event.to_json(),
514                    }
515                })),
516                Ok(None) => {}
517                Err(error) => notifications.push(json!({
518                    "jsonrpc": "2.0",
519                    "method": SESSION_EVENT_METHOD,
520                    "params": {
521                        "subscription": subscription,
522                        "event": {
523                            "type": "watch_error",
524                            "recoverable": true,
525                            "message": error.to_string(),
526                        },
527                    }
528                })),
529            }
530        }
531        notifications
532    }
533
534    /// Report each followed session's live-runtime lifecycle state on that
535    /// session's own subscription, emitting only when the state changes.
536    ///
537    /// A growing transcript is not evidence that an agent is working, so the
538    /// state comes from the live-runtime registry and nowhere else. A followed
539    /// session with no registered Supercode runtime — a harness running outside
540    /// Supercode — reports `persisted`, which says plainly that its activity is
541    /// unknown rather than guessing at it. These events carry no sequence
542    /// number and no transcript content; they never interleave with the
543    /// content follower's sequenced stream.
544    #[cfg(feature = "adapter-api")]
545    pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
546        let registry = crate::LocalRuntimeRegistry::new();
547        let authorization = crate::RuntimeAuthorization::observer();
548        let mut notifications = Vec::new();
549        for (subscription, source) in &mut self.followed_sources {
550            let state = match registry
551                .source_state(&source.harness, &source.session_id, &authorization)
552                .await
553            {
554                Ok(Some(state)) => state,
555                Ok(None) => crate::RuntimeRegistryState::Persisted,
556                // A failed registry read is not evidence of a state change.
557                Err(_) => continue,
558            };
559            if source.reported.as_deref() == Some(state.as_str()) {
560                continue;
561            }
562            source.reported = Some(state.as_str().to_string());
563            notifications.push(json!({
564                "jsonrpc": "2.0",
565                "method": SESSION_EVENT_METHOD,
566                "params": {
567                    "subscription": subscription,
568                    "event": {"type": "runtime_state", "state": state.as_str()},
569                },
570            }));
571        }
572        notifications
573    }
574
575    /// Poll normalized activity subscriptions, emitting only proven state
576    /// transitions. Every subscription is bulk-sampled so stock-harness
577    /// process and registry discovery happens once per UI, not once per row.
578    #[cfg(feature = "adapter-api")]
579    pub async fn poll_session_activities(&mut self) -> Vec<Value> {
580        let subscriptions = self
581            .activity_subscriptions
582            .iter()
583            .map(|(id, subscription)| {
584                (
585                    id.clone(),
586                    subscription.locators.clone(),
587                    subscription.homes.clone(),
588                )
589            })
590            .collect::<Vec<_>>();
591        let mut notifications = Vec::new();
592        for (subscription_id, locators, homes) in subscriptions {
593            let Ok(mut activities) = self.activity_monitor.resolve(&locators, &homes).await else {
594                // A failed evidence read proves no transition. Retain the last
595                // good state instead of flashing every row to persisted.
596                continue;
597            };
598            pane_prompt_override(&mut activities, &homes);
599            let Some(subscription) = self.activity_subscriptions.get_mut(&subscription_id) else {
600                continue;
601            };
602            let mut changed = Vec::new();
603            for activity in activities {
604                let key = activity.key();
605                if subscription
606                    .reported
607                    .get(&key)
608                    .is_some_and(|previous| previous.same_state(&activity))
609                {
610                    continue;
611                }
612                subscription.reported.insert(key, activity.clone());
613                changed.push(activity);
614            }
615            if !changed.is_empty() {
616                notifications.push(json!({
617                    "jsonrpc": "2.0",
618                    "method": SESSION_ACTIVITY_EVENT_METHOD,
619                    "params": {
620                        "subscription": subscription_id,
621                        "activities": changed,
622                    },
623                }));
624            }
625        }
626        notifications
627    }
628
629    /// Drain native-store invalidations and emit revisioned descriptor deltas.
630    /// An idle subscription performs no catalog or transcript reads between
631    /// its minute-scale recovery reconciliations.
632    #[cfg(feature = "adapter-api")]
633    pub fn poll_session_indexes(&mut self) -> Vec<Value> {
634        let mut notifications = Vec::new();
635        for (subscription, index) in &mut self.index_subscriptions {
636            let homes = index.homes().clone();
637            match index.poll() {
638                Ok(Some(delta)) => match live_index_changes(delta.changes, &homes) {
639                    Ok(changes) => notifications.push(json!({
640                        "jsonrpc": "2.0",
641                        "method": SESSION_INDEX_EVENT_METHOD,
642                        "params": {
643                            "subscription": subscription,
644                            "revision": delta.revision,
645                            "changes": changes,
646                        },
647                    })),
648                    Err(error) => notifications.push(json!({
649                        "jsonrpc": "2.0",
650                        "method": SESSION_INDEX_EVENT_METHOD,
651                        "params": {
652                            "subscription": subscription,
653                            "error": {"recoverable": true, "message": error_message(error)},
654                        },
655                    })),
656                },
657                Ok(None) => {}
658                Err(error) => notifications.push(json!({
659                    "jsonrpc": "2.0",
660                    "method": SESSION_INDEX_EVENT_METHOD,
661                    "params": {
662                        "subscription": subscription,
663                        "error": {"recoverable": true, "message": error},
664                    },
665                })),
666            }
667        }
668        notifications
669    }
670
671    #[cfg(feature = "adapter-api")]
672    async fn subscribe_session_activity(
673        &mut self,
674        params: Value,
675    ) -> std::result::Result<Value, ServiceError> {
676        let params = decode::<ActivitySubscribeParams>(params)?;
677        if params.locators.is_empty() {
678            return Err(ServiceError::InvalidParams(
679                "sessions.activity.subscribe requires at least one locator".into(),
680            ));
681        }
682        if params.locators.len() > 2_048 {
683            return Err(ServiceError::InvalidParams(
684                "sessions.activity.subscribe accepts at most 2048 locators".into(),
685            ));
686        }
687        let mut initial = self
688            .activity_monitor
689            .resolve(&params.locators, &params.homes)
690            .await
691            .map_err(ServiceError::Sdk)?;
692        pane_prompt_override(&mut initial, &params.homes);
693        let subscription = format!("activity-sub-{}", self.next_subscription);
694        self.next_subscription += 1;
695        let reported = initial
696            .iter()
697            .cloned()
698            .map(|activity| (activity.key(), activity))
699            .collect();
700        self.activity_subscriptions.insert(
701            subscription.clone(),
702            ActivitySubscription {
703                locators: params.locators,
704                homes: params.homes,
705                reported,
706            },
707        );
708        Ok(json!({"subscription": subscription, "initial": initial}))
709    }
710
711    /// Non-blockingly sample one event from every connected live runtime.
712    #[cfg(feature = "adapter-api")]
713    pub async fn poll_runtimes(&mut self) -> Vec<Value> {
714        self.poll_sdk_events()
715            .await
716            .into_iter()
717            .map(|(connection, runtime_event)| {
718                json!({
719                    "jsonrpc": "2.0",
720                    "method": RUNTIME_EVENT_METHOD,
721                    "params": {
722                        "connection": connection,
723                        "session_id": runtime_event.session_id,
724                        "sequence": runtime_event.event.sequence,
725                        "event": {
726                            "kind": runtime_event.event.kind,
727                            "payload": runtime_event.event.payload,
728                        },
729                    },
730                })
731            })
732            .collect()
733    }
734
735    async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
736        let mut events = Vec::new();
737        let mut closed = Vec::new();
738        let now_ms = crate::approvals::now_ms();
739        for (connection, runtime) in &mut self.runtimes {
740            let session_id = runtime.handle().runtime_id.clone();
741            let harness = runtime.handle().harness.clone();
742            // Drain what the runtime already has: a turn is several events
743            // (updates, then the protocol's completion), and delivering one
744            // per poll would cost a poll interval each. A zero timeout takes
745            // only what is ready — an idle runtime costs nothing.
746            for _ in 0..256 {
747                match tokio::time::timeout(Duration::ZERO, runtime.next_event()).await {
748                    Ok(Ok(Some(event))) => {
749                        let terminal = event.kind == "transport_closed";
750                        // ORCH-9: a permission/approval request arrives as an
751                        // ordinary event; it becomes listable here and stops
752                        // being listable when `runtimes.respond` answers it.
753                        self.approvals
754                            .observe(connection, &harness, &session_id, &event, now_ms);
755                        let next_sequence = self
756                            .runtime_sequences
757                            .entry(session_id.clone())
758                            .or_insert(0);
759                        let sequence = event.sequence.unwrap_or_else(|| {
760                            *next_sequence = next_sequence.saturating_add(1);
761                            *next_sequence
762                        });
763                        *next_sequence = (*next_sequence).max(sequence);
764                        events.push((
765                            connection.clone(),
766                            SdkRuntimeEvent {
767                                session_id: session_id.clone(),
768                                event: SdkEvent {
769                                    sequence,
770                                    kind: event.kind,
771                                    payload: event.payload,
772                                },
773                            },
774                        ));
775                        if terminal {
776                            closed.push(connection.clone());
777                            break;
778                        }
779                    }
780                    Ok(Ok(None)) => {
781                        let sequence = self
782                            .runtime_sequences
783                            .entry(session_id.clone())
784                            .or_insert(0);
785                        *sequence = sequence.saturating_add(1);
786                        events.push((
787                        connection.clone(),
788                        SdkRuntimeEvent {
789                            session_id,
790                            event: SdkEvent {
791                                sequence: *sequence,
792                                kind: "transport_closed".into(),
793                                payload: json!({"message": "Harness runtime transport closed."}),
794                            },
795                        },
796                    ));
797                        closed.push(connection.clone());
798                        break;
799                    }
800                    Err(_) => break,
801                    Ok(Err(error)) => {
802                        let sequence = self
803                            .runtime_sequences
804                            .entry(session_id.clone())
805                            .or_insert(0);
806                        *sequence = sequence.saturating_add(1);
807                        events.push((
808                        connection.clone(),
809                        SdkRuntimeEvent {
810                            session_id,
811                            event: SdkEvent {
812                                sequence: *sequence,
813                                kind: "transport_error".into(),
814                                payload: json!({"message": error.to_string(), "terminal": true}),
815                            },
816                        },
817                    ));
818                        closed.push(connection.clone());
819                        break;
820                    }
821                }
822            }
823        }
824        for connection in closed {
825            if let Some(runtime) = self.runtimes.remove(&connection) {
826                self.runtime_sequences.remove(&runtime.handle().runtime_id);
827            }
828            self.terminal_launches.remove(&connection);
829            // A connection that is gone cannot answer anything it was
830            // holding; those requests stop being listable with it.
831            self.approvals.forget(&connection);
832        }
833        events
834    }
835
836    fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
837        match method {
838            "harness.v1.capabilities" => Ok(json!({
839                "version": HARNESS_SERVICE_VERSION,
840                "sdk": self.capabilities(),
841                "methods": HARNESS_SERVICE_METHODS,
842                "notifications": [
843                    SESSION_EVENT_METHOD,
844                    SESSION_ACTIVITY_EVENT_METHOD,
845                    SESSION_INDEX_EVENT_METHOD,
846                    RUNTIME_EVENT_METHOD
847                ],
848                "harnesses": harness_support_registry()
849                    .harnesses
850                    .into_iter()
851                    .map(|harness| harness.id)
852                    .collect::<Vec<_>>(),
853            })),
854            "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
855                .map_err(|error| ServiceError::Operation(error.to_string())),
856            "harness.v1.profiles.list" | "harness.v1.profiles.get" => profiles_call(method, params),
857            // ORCH-21 controlled tier. Each verb translates to the HARNESS'S
858            // OWN profile verb and runs it (`crate::profiles_control`);
859            // supercode makes and removes nothing itself. The row returned is
860            // re-read through the ORCH-10 loader afterwards, and `ran`
861            // narrates the exact command.
862            "harness.v1.profiles.create" => {
863                mutate_profile(crate::profiles_control::ProfileVerb::Create, params)
864            }
865            "harness.v1.profiles.delete" => {
866                mutate_profile(crate::profiles_control::ProfileVerb::Delete, params)
867            }
868            "harness.v1.channels.list" | "harness.v1.channels.status" => {
869                channels_call(method, params)
870            }
871            // ORCH-15 observed tier: which profile / agent a surface tuple
872            // resolves to, read from each gateway harness's own config.
873            "harness.v1.routes.list" => routes_call(params),
874            // ORCH-16 observed tier: inbound webhook routes / hook mappings.
875            "harness.v1.triggers.list" => triggers_call(params),
876            // ONT-4: the orchestration doors. One home folder in, one typed orchestration
877            // value out (and back). Every one of the four is
878            // `crate::orchestration_doors`, which the `supercode orchestration` verbs call
879            // too — the RPC adds nothing but the envelope. A vault VALUE
880            // never crosses this wire: a load or a compile answers with the
881            // `.env` KEY NAMES, and a caller that needs a value reads the
882            // home's own `.env`.
883            // the workflow layer's read door: a harness's board as one typed value,
884            // the same code the `supercode workflow load` verb calls
885            "harness.v1.workflow.load" => {
886                let params = decode::<WorkflowLoadParams>(params)?;
887                let read = crate::workflow_doors::load(
888                    params.from.unwrap_or_else(|| {
889                        crate::workflow_doors::WorkflowHarness::of_home(&params.home)
890                    }),
891                    &params.home,
892                )
893                .map_err(operation)?;
894                serde_json::to_value(read)
895                    .map_err(|error| ServiceError::Operation(error.to_string()))
896            }
897            "harness.v1.orchestration.load" => {
898                let params = decode::<OrchestrationLoadParams>(params)?;
899                let read = crate::orchestration_doors::load(&params.root, params.flavor)
900                    .map_err(operation)?;
901                serde_json::to_value(read)
902                    .map_err(|error| ServiceError::Operation(error.to_string()))
903            }
904            "harness.v1.orchestration.save" => {
905                let params = decode::<OrchestrationSaveParams>(params)?;
906                let saved = crate::orchestration_doors::save(
907                    &params.root,
908                    params.orchestration,
909                    params.vault,
910                )
911                .map_err(operation)?;
912                serde_json::to_value(saved)
913                    .map_err(|error| ServiceError::Operation(error.to_string()))
914            }
915            "harness.v1.orchestration.compile" => {
916                let params = decode::<OrchestrationCompileParams>(params)?;
917                let read = crate::orchestration_doors::compile(params.from, &params.home)
918                    .map_err(operation)?;
919                serde_json::to_value(read)
920                    .map_err(|error| ServiceError::Operation(error.to_string()))
921            }
922            "harness.v1.orchestration.decompile" => {
923                let params = decode::<OrchestrationDecompileParams>(params)?;
924                let report = crate::orchestration_doors::decompile(
925                    params.to,
926                    params.orchestration,
927                    &params.source,
928                    params.source_flavor,
929                    &params.dest,
930                    params.vault,
931                )
932                .map_err(operation)?;
933                serde_json::to_value(report)
934                    .map_err(|error| ServiceError::Operation(error.to_string()))
935            }
936            // a migration keeps the credential in this process: a compile and
937            // a save (import), a load and a decompile (export), composed here
938            // because composed by a client the secret would have to cross
939            // the wire
940            "harness.v1.orchestration.import" => {
941                let params = decode::<OrchestrationImportParams>(params)?;
942                let imported =
943                    crate::orchestration_doors::import(params.from, &params.home, &params.into)
944                        .map_err(operation)?;
945                serde_json::to_value(imported)
946                    .map_err(|error| ServiceError::Operation(error.to_string()))
947            }
948            "harness.v1.orchestration.export" => {
949                let params = decode::<OrchestrationExportParams>(params)?;
950                let report =
951                    crate::orchestration_doors::export(params.to, &params.root, &params.dest)
952                        .map_err(operation)?;
953                serde_json::to_value(report)
954                    .map_err(|error| ServiceError::Operation(error.to_string()))
955            }
956            // ORCH-12 observed tier: read and search the persistent memory
957            // documents a harness keeps on disk. Read-only — every write
958            // (`hermes memory off`, `openclaw memory forget|reset`, Claude
959            // Code's `/memory`) stays the harness's own verb. A harness with
960            // no memory store is refused with UnsupportedAction.
961            "harness.v1.memory.show" | "harness.v1.memory.search" => memory_call(method, params),
962            // ORCH-11 observed tier: read-only enumeration of every harness's
963            // installed skill packages. An unknown harness id is refused with
964            // UnsupportedAction — every harness supports skills, so a filter
965            // that matches nothing is a caller error, never an empty listing.
966            "harness.v1.skills.list" => {
967                let query = decode::<crate::skills::SkillsQuery>(params)?;
968                if let Some(harness) = query.harness.as_deref() {
969                    if !crate::skills::SKILL_HARNESSES.contains(&harness) {
970                        return Err(ServiceError::UnsupportedAction(format!(
971                            "`{harness}` has no skills root Volter Harness reads"
972                        )));
973                    }
974                }
975                serde_json::to_value(crate::skills::list_skills(&query))
976                    .map_err(|error| ServiceError::Operation(error.to_string()))
977            }
978            // ORCH-22 controlled tier: each verb goes through the door the
979            // HARNESS publishes — `hermes skills install|uninstall`,
980            // `openclaw skills install`, and for the core four the loader's
981            // own directory, which is the only skills door those harnesses
982            // have. supercode resolves no registry and unpacks no archive.
983            // The row returned is re-read through the ORCH-11 loader
984            // afterwards, and `ran` narrates exactly what was performed.
985            "harness.v1.skills.install" => {
986                mutate_skill(crate::skills_control::SkillVerb::Install, params)
987            }
988            "harness.v1.skills.remove" => {
989                mutate_skill(crate::skills_control::SkillVerb::Remove, params)
990            }
991            // ORCH-9 observed tier: the approval requests waiting for an
992            // answer. At the pinned harness versions the only uniform source
993            // is a LIVE request held by an open runtime connection, plus
994            // supercode's own queued subagent approvals — neither Hermes
995            // 0.21.0 nor OpenClaw 2026.7.1-2 has an approvals door to read
996            // (see `crate::approvals`). A harness whose runtime cannot carry
997            // a protocol request at all is refused by name.
998            "harness.v1.approvals.list" => {
999                let query = decode::<crate::approvals::ApprovalsQuery>(params)?;
1000                if let Some(harness) = query.harness.as_deref() {
1001                    if !crate::approvals::lists_approvals(harness) {
1002                        return Err(ServiceError::UnsupportedAction(format!(
1003                            "`{harness}` has no runtime door that carries an approval request"
1004                        )));
1005                    }
1006                }
1007                serde_json::to_value(self.approvals(&query))
1008                    .map_err(|error| ServiceError::Operation(error.to_string()))
1009            }
1010            "harness.v1.sessions.discover" => {
1011                let query = decode::<DiscoveryQuery>(params)?;
1012                let mut page = discover_session_page(&query).map_err(operation)?;
1013                // Claude Code is the one harness that publishes its RUNNING
1014                // sessions. The registry is read once per discovery and joined
1015                // by session id; every record in it has already survived a
1016                // `kill(pid, 0)` liveness check inside `read_registry`.
1017                let doors = crate::mail_route::LiveSessions::read(&query.homes);
1018                // A running session's name (`twin-76`) lives only in its
1019                // harness's registry, not in the transcript discovery searches:
1020                // a query naming one also finds that session, by its id.
1021                if let Some(text) = query
1022                    .query
1023                    .as_deref()
1024                    .map(str::trim)
1025                    .filter(|text| !text.is_empty())
1026                {
1027                    let needle = text.to_lowercase();
1028                    for live in doors.all() {
1029                        let name = live
1030                            .name
1031                            .split('@')
1032                            .next()
1033                            .unwrap_or(&live.name)
1034                            .to_lowercase();
1035                        if !name.contains(&needle)
1036                            || page.sessions.iter().any(|session| {
1037                                session.locator.session_id == live.address.session_id
1038                            })
1039                        {
1040                            continue;
1041                        }
1042                        if query
1043                            .limit
1044                            .is_some_and(|limit| page.sessions.len() >= limit)
1045                        {
1046                            break;
1047                        }
1048                        // Only that session's own folder is scanned, not the whole store again.
1049                        let mut by_id = query.clone();
1050                        by_id.query = Some(live.address.session_id.clone());
1051                        by_id.cursor = None;
1052                        if let Some(cwd) = &live.cwd {
1053                            by_id.workspace = Some(cwd.clone());
1054                            by_id.workspace_subtree = false;
1055                            by_id.workspace_family = None;
1056                        }
1057                        if let Ok(found) = discover_session_page(&by_id) {
1058                            page.sessions
1059                                .extend(found.sessions.into_iter().filter(|session| {
1060                                    session.locator.session_id == live.address.session_id
1061                                }));
1062                        }
1063                    }
1064                }
1065                let activities = crate::session_activity::resolve_stock_session_activities(
1066                    &page
1067                        .sessions
1068                        .iter()
1069                        .map(|session| session.locator.clone())
1070                        .collect::<Vec<_>>(),
1071                    &query.homes,
1072                )
1073                .into_iter()
1074                .map(|activity| (activity.key(), activity))
1075                .collect::<BTreeMap<_, _>>();
1076                let sessions = page
1077                    .sessions
1078                    .into_iter()
1079                    .map(|session| {
1080                        let mut value = live_descriptor_value(&session, &doors)?;
1081                        let activity_key = (
1082                            session.locator.harness.as_str().to_string(),
1083                            session.locator.session_id.clone(),
1084                        );
1085                        if let Some(activity) = activities.get(&activity_key) {
1086                            // A Codex prompt (an approval, a choice) writes nothing to its rollout, which
1087                            // reads that turn as working; the live row reads its pane, and says so.
1088                            let mut activity = activity.clone();
1089                            let waiting = doors.all().iter().any(|live| {
1090                                live.address.harness == activity_key.0
1091                                    && live.address.session_id == activity_key.1
1092                                    && live.status == "waiting"
1093                            });
1094                            if waiting && activity.turn != crate::SessionTurnState::NeedsInput {
1095                                activity.turn = crate::SessionTurnState::NeedsInput;
1096                                activity.evidence.source = "pane_prompt".into();
1097                            }
1098                            let activity = &activity;
1099                            value["activity"] = serde_json::to_value(activity)
1100                                .map_err(|error| ServiceError::Operation(error.to_string()))?;
1101                            if let Some(status) = legacy_live_status(activity) {
1102                                value["live_status"] = json!(status);
1103                            }
1104                        }
1105                        Ok(value)
1106                    })
1107                    .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1108                let mut result = json!({"sessions": sessions, "next_cursor": page.next_cursor});
1109                // Preserve the metadata-only wire shape, but carry the catalog's
1110                // proof/counts when the caller explicitly requests preview search.
1111                if query.search_previews {
1112                    result["receipt"] = serde_json::to_value(page.receipt)
1113                        .map_err(|error| ServiceError::Operation(error.to_string()))?;
1114                }
1115                Ok(result)
1116            }
1117            "harness.v1.sessions.inbox" => inbox_call(decode::<InboxParams>(params)?),
1118            "harness.v1.sessions.usage" => {
1119                let params = decode::<LocatorParams>(params)?;
1120                let session = load_session(&params.locator).map_err(operation)?;
1121                Ok(
1122                    json!({"harness": params.locator.harness, "session_id": session.meta.session_id,
1123                    "usage": recorded_session_usage(&session)}),
1124                )
1125            }
1126            "harness.v1.sessions.recorded_config" => {
1127                let params = decode::<LocatorParams>(params)?;
1128                recorded_config(&params.locator).map_err(operation)
1129            }
1130            "harness.v1.sessions.load" => {
1131                let params = decode::<LoadSessionParams>(params)?;
1132                load_session_door(params)
1133            }
1134            "harness.v1.sessions.follow" => {
1135                let params = decode::<LocatorParams>(params)?;
1136                let mut follower = self
1137                    .catalog
1138                    .follow_read_view(
1139                        &params.locator,
1140                        params.read_fidelity(),
1141                        params.include_subagents(),
1142                        params.tail_messages(),
1143                        params.max_message_chars(),
1144                        params.display_history(),
1145                    )
1146                    .map_err(operation)?;
1147                let initial = follower
1148                    .poll()
1149                    .map_err(operation)?
1150                    .map(|event| event.to_json());
1151                let subscription = format!("sub-{}", self.next_subscription);
1152                self.next_subscription += 1;
1153                self.followers.insert(subscription.clone(), follower);
1154                self.followed_sources.insert(
1155                    subscription.clone(),
1156                    FollowedSource {
1157                        harness: params.locator.harness.as_str().to_string(),
1158                        session_id: params.locator.session_id.clone(),
1159                        reported: None,
1160                    },
1161                );
1162                Ok(json!({"subscription": subscription, "initial": initial}))
1163            }
1164            "harness.v1.sessions.unfollow" => {
1165                let params = decode::<UnfollowParams>(params)?;
1166                self.followed_sources.remove(&params.subscription);
1167                Ok(json!({
1168                    "removed": self.followers.remove(&params.subscription).is_some()
1169                }))
1170            }
1171            "harness.v1.sessions.activity.unsubscribe" => {
1172                let params = decode::<UnfollowParams>(params)?;
1173                Ok(json!({
1174                    "removed": self.activity_subscriptions.remove(&params.subscription).is_some()
1175                }))
1176            }
1177            "harness.v1.sessions.index.subscribe" => {
1178                let query = decode::<DiscoveryQuery>(params)?;
1179                crate::session_index::validate_query(&query)
1180                    .map_err(ServiceError::InvalidParams)?;
1181                let homes = query.homes.clone();
1182                let (index, initial) = crate::session_index::SessionIndexSubscription::open(
1183                    query,
1184                    Arc::clone(&self.index_notifier),
1185                )
1186                .map_err(ServiceError::Operation)?;
1187                let doors = crate::mail_route::LiveSessions::read(&homes);
1188                let initial = initial
1189                    .iter()
1190                    .map(|descriptor| live_descriptor_value(descriptor, &doors))
1191                    .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1192                let subscription = format!("index-sub-{}", self.next_subscription);
1193                self.next_subscription += 1;
1194                self.index_subscriptions.insert(subscription.clone(), index);
1195                Ok(json!({
1196                    "subscription": subscription,
1197                    "revision": 1,
1198                    "initial": initial,
1199                }))
1200            }
1201            "harness.v1.sessions.index.resize" => {
1202                let params = decode::<IndexResizeParams>(params)?;
1203                crate::session_index::validate_limit(params.limit)
1204                    .map_err(ServiceError::InvalidParams)?;
1205                let index = self
1206                    .index_subscriptions
1207                    .get_mut(&params.subscription)
1208                    .ok_or_else(|| {
1209                        ServiceError::InvalidParams("unknown session index subscription".into())
1210                    })?;
1211                let prepared = index
1212                    .prepare_resize(params.limit)
1213                    .map_err(ServiceError::Operation)?;
1214                let doors = crate::mail_route::LiveSessions::read(index.homes());
1215                let initial = prepared
1216                    .page
1217                    .sessions
1218                    .iter()
1219                    .map(|descriptor| live_descriptor_value(descriptor, &doors))
1220                    .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
1221                let response = json!({
1222                    "subscription": params.subscription,
1223                    "revision": prepared.revision,
1224                    "initial": initial,
1225                    "receipt": prepared.page.receipt,
1226                });
1227                index.commit_resize(prepared);
1228                Ok(response)
1229            }
1230            "harness.v1.sessions.index.unsubscribe" => {
1231                let params = decode::<UnfollowParams>(params)?;
1232                Ok(json!({
1233                    "removed": self.index_subscriptions.remove(&params.subscription).is_some()
1234                }))
1235            }
1236            "harness.v1.sessions.import" => {
1237                let params = decode::<ImportSessionParams>(params)?;
1238                let session = Session::load_str(&params.content, params.source_harness.into())
1239                    .map_err(operation)?;
1240                Ok(json!({"session": normalized_session_json(&session)}))
1241            }
1242            "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
1243                let params = decode::<ExportSessionParams>(params)?;
1244                let session = load_session(&params.locator).map_err(operation)?;
1245                let artifact = session_artifact(&params.locator, &session, params.target_harness)?;
1246                if method == "harness.v1.sessions.export"
1247                    && params.target_harness == TransferFormat::Hermes
1248                {
1249                    // UNI-18: write through Hermes's own door, never into its store
1250                    let imported = crate::hermes_import::import_into_hermes(&session, None)
1251                        .map_err(operation)?;
1252                    return Ok(json!({"artifact": artifact, "imported": imported}));
1253                }
1254                Ok(json!({"artifact": artifact}))
1255            }
1256            "harness.v1.sessions.reduce" => {
1257                let params = decode::<ReduceSessionParams>(params)?;
1258                self.reduce_session(params)
1259            }
1260            "harness.v1.sessions.branch" => {
1261                let params = decode::<BranchSessionParams>(params)?;
1262                let session = load_session(&params.locator).map_err(operation)?;
1263                let storage = params.locator.storage.path().display().to_string();
1264                let bootstrap_prompt = format!(
1265                    "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.",
1266                    params.locator.harness.as_str(), params.locator.session_id, storage
1267                );
1268                let artifact = params
1269                    .target_harness
1270                    .map(|target| session_artifact(&params.locator, &session, target))
1271                    .transpose()?;
1272                Ok(json!({
1273                    "parent": params.locator,
1274                    "session": normalized_session_json(&session),
1275                    "bootstrap_prompt": bootstrap_prompt,
1276                    "artifact": artifact,
1277                }))
1278            }
1279            "harness.v1.sessions.handoff" => {
1280                let params = decode::<HandoffSessionParams>(params)?;
1281                let session = load_session(&params.locator).map_err(operation)?;
1282                let cwd = params
1283                    .cwd
1284                    .or_else(|| session.meta.cwd.clone())
1285                    .unwrap_or_else(|| PathBuf::from("."));
1286                let artifact = handoff_artifact(&params.locator, &session, params.target_harness)?;
1287                let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
1288                    ServiceError::Operation(
1289                        "handoff artifact omitted target session identity".into(),
1290                    )
1291                })?;
1292                let instructions =
1293                    handoff_instructions(params.target_harness, target_session_id, &cwd);
1294                Ok(json!({
1295                    "artifact": artifact,
1296                    "launch": instructions.launch,
1297                    "materialize": instructions.materialize,
1298                    "requires_materialization": instructions.requires_materialization,
1299                    "note": instructions.note,
1300                }))
1301            }
1302            "harness.v1.sessions.materialize" => {
1303                let params = decode::<MaterializeSessionParams>(params)?;
1304                // An artifact from another machine carries its whole source as a recovery file;
1305                // keeping its segments here lets a later write back to that format restore it
1306                // byte for byte on this machine too (docs/plans/portable-residue.md).
1307                for file in &params.artifact.files {
1308                    if file.role == "source_recovery"
1309                        && file.path == "recovery/source.supercode.jsonl"
1310                    {
1311                        if let Ok(source) = Session::from_native_str(&file.content) {
1312                            crate::residue_store::store_segments(&source);
1313                        }
1314                    }
1315                }
1316                let locator = crate::native_materialize::materialize_native_artifact(
1317                    params.artifact,
1318                    &params.cwd,
1319                    &params.homes,
1320                )
1321                .map_err(ServiceError::Operation)?;
1322                Ok(json!({"locator": locator}))
1323            }
1324            // ORCH-7 observed tier. Read-only: the handlers open the harness's
1325            // own job store (Claude Code's session JSONL, Hermes's and
1326            // OpenClaw's `cron/jobs.json`) and never write, fire, or schedule.
1327            "harness.v1.jobs.list" => {
1328                let query = decode::<crate::jobs::JobsQuery>(params)?;
1329                if let Some(harness) = query.harness.as_deref() {
1330                    refuse_harness_without_jobs(harness, "jobs.list")?;
1331                }
1332                let listing = crate::jobs::list_jobs(&query).map_err(operation)?;
1333                serde_json::to_value(listing)
1334                    .map_err(|error| ServiceError::Operation(error.to_string()))
1335            }
1336            "harness.v1.jobs.get" => {
1337                let params = decode::<JobsGetParams>(params)?;
1338                refuse_harness_without_jobs(&params.harness, "jobs.get")?;
1339                match crate::jobs::get_job(&params.harness, &params.id, &params.homes)
1340                    .map_err(operation)?
1341                {
1342                    Some((job, source)) => Ok(json!({"job": job, "source": source})),
1343                    None => Err(ServiceError::Operation(format!(
1344                        "`{}` has no scheduled job `{}`",
1345                        params.harness, params.id
1346                    ))),
1347                }
1348            }
1349            // ORCH-18 controlled tier. Each verb translates to the HARNESS'S
1350            // OWN cron verb and runs it (`crate::jobs_control`); supercode
1351            // schedules nothing. The row returned is re-read from the
1352            // harness's store afterwards, and `ran` narrates the exact command
1353            // with any credential redacted.
1354            "harness.v1.jobs.create" => mutate_job(crate::jobs_control::JobVerb::Create, params),
1355            "harness.v1.jobs.update" => mutate_job(crate::jobs_control::JobVerb::Update, params),
1356            "harness.v1.jobs.pause" => mutate_job(crate::jobs_control::JobVerb::Pause, params),
1357            "harness.v1.jobs.resume" => mutate_job(crate::jobs_control::JobVerb::Resume, params),
1358            "harness.v1.jobs.run" => mutate_job(crate::jobs_control::JobVerb::Run, params),
1359            "harness.v1.jobs.delete" => mutate_job(crate::jobs_control::JobVerb::Delete, params),
1360            "harness.v1.jobs.notepad"
1361            | "harness.v1.jobs.notepad_set"
1362            | "harness.v1.jobs.notepad_delete" => {
1363                let request = decode::<crate::jobs_notepad::JobNotepadRequest>(params)?;
1364                refuse_harness_without_jobs(&request.harness, "jobs.notepad")?;
1365                let answer = match method {
1366                    "harness.v1.jobs.notepad_set" => crate::jobs_notepad::set(&request),
1367                    "harness.v1.jobs.notepad_delete" => crate::jobs_notepad::delete(&request),
1368                    _ => crate::jobs_notepad::read(&request),
1369                }
1370                .map_err(job_control_error)?;
1371                serde_json::to_value(answer)
1372                    .map_err(|error| ServiceError::Operation(error.to_string()))
1373            }
1374            // ORCH-8 observed tier. Read-only: the handlers open the harness's
1375            // own run store (Hermes's `cron/executions.db`, OpenClaw's
1376            // `cron_run_logs`) and never claim, retry, or prune a fire.
1377            "harness.v1.runs.list" => {
1378                let query = decode::<crate::runs::RunsQuery>(params)?;
1379                if let Some(harness) = query.harness.as_deref() {
1380                    refuse_harness_without_runs(harness, "runs.list")?;
1381                }
1382                let listing = crate::runs::list_runs(&query).map_err(operation)?;
1383                serde_json::to_value(listing)
1384                    .map_err(|error| ServiceError::Operation(error.to_string()))
1385            }
1386            "harness.v1.runs.get" => {
1387                let params = decode::<RunsGetParams>(params)?;
1388                refuse_harness_without_runs(&params.harness, "runs.get")?;
1389                match crate::runs::get_run(&params.harness, &params.id, &params.homes)
1390                    .map_err(operation)?
1391                {
1392                    Some((run, source)) => Ok(json!({"run": run, "source": source})),
1393                    None => Err(ServiceError::Operation(format!(
1394                        "`{}` has no run `{}`",
1395                        params.harness, params.id
1396                    ))),
1397                }
1398            }
1399            "harness.v1.sessions.resume_instructions" => {
1400                let params = decode::<ResumeInstructionsParams>(params)?;
1401                let session = load_session(&params.locator).map_err(operation)?;
1402                let cwd = params
1403                    .cwd
1404                    .or(session.meta.cwd)
1405                    .unwrap_or_else(|| PathBuf::from("."));
1406                let launch = resume_launch(
1407                    params.locator.harness.as_str(),
1408                    &params.locator.session_id,
1409                    &cwd,
1410                    params.policy,
1411                )?;
1412                Ok(json!({"launch": launch}))
1413            }
1414            _ => Err(ServiceError::MethodNotFound),
1415        }
1416    }
1417
1418    fn reduce_session(
1419        &self,
1420        params: ReduceSessionParams,
1421    ) -> std::result::Result<Value, ServiceError> {
1422        let session = load_session(&params.locator).map_err(operation)?;
1423        if session.messages.is_empty() {
1424            return Err(ServiceError::InvalidParams(
1425                "cannot reduce an empty session".into(),
1426            ));
1427        }
1428        let keep_last = params.keep_last.clamp(1, 128);
1429        let policy = reduce::ReductionPolicy {
1430            clear_turns_older_than: Some(keep_last),
1431            ..Default::default()
1432        };
1433        let (view, log) =
1434            reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
1435        if log.reductions.is_empty() {
1436            return Err(ServiceError::UnsupportedAction(format!(
1437                "session `{}` is already too small for a meaningful reversible reduction",
1438                params.locator.session_id
1439            )));
1440        }
1441        let source_tokens = supercode_runtime::estimate_view_tokens(&session.messages);
1442        let reduced_tokens = supercode_runtime::estimate_view_tokens(&view);
1443        if reduced_tokens >= source_tokens {
1444            return Err(ServiceError::UnsupportedAction(format!(
1445                "session `{}` has no token-reducing reversible projection",
1446                params.locator.session_id
1447            )));
1448        }
1449
1450        let store_root = self
1451            .reduction_store_root
1452            .clone()
1453            .unwrap_or_else(default_reduction_store_root);
1454        let store = crate::SessionStore::open(&store_root).map_err(operation)?;
1455        let rescue_id = format!("rescue-{}", generated_session_id());
1456        let imported = session
1457            .imported_message_count
1458            .unwrap_or(session.messages.len())
1459            .min(session.messages.len());
1460        let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
1461        let view_jsonl = messages_jsonl(&view)?;
1462        let title = format!(
1463            "Reduced {} continuation from {}",
1464            params.target_harness.id(),
1465            params.locator.session_id
1466        );
1467
1468        // Durability order is intentional: the full source of truth lands
1469        // before either object that can refer to it. A crash may leave an
1470        // unused sidecar, but can never leave a reduced view whose originals
1471        // were not durably written first.
1472        store
1473            .save_sidecar(&rescue_id, &sidecar_jsonl)
1474            .map_err(operation)?;
1475        store
1476            .save_reduction_log(&rescue_id, &log)
1477            .map_err(operation)?;
1478        store
1479            .save(&rescue_id, &title, &view_jsonl)
1480            .map_err(operation)?;
1481
1482        let source_bytes = serde_json::to_vec(&session.messages)
1483            .map_err(|error| ServiceError::Operation(error.to_string()))?
1484            .len() as u64;
1485        let reduced_bytes = serde_json::to_vec(&view)
1486            .map_err(|error| ServiceError::Operation(error.to_string()))?
1487            .len() as u64;
1488        store
1489            .set_reduction_stats(
1490                &rescue_id,
1491                &title,
1492                source_bytes,
1493                reduced_bytes,
1494                log.reductions.len() as u32,
1495            )
1496            .map_err(operation)?;
1497
1498        // The receipt is issued only after a real disk reload. This proves
1499        // the exact files another process will consume, not the convenient
1500        // in-memory values that produced them.
1501        let reloaded_sidecar = store
1502            .load_sidecar(&rescue_id)
1503            .map_err(operation)?
1504            .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
1505        let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
1506        let reloaded_log = store
1507            .load_reduction_log(&rescue_id)
1508            .map_err(operation)?
1509            .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
1510        let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
1511        reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
1512        // `sc.reduction` is deliberately in-memory-only metadata: it must
1513        // never leak onto a provider-facing transcript. Reapplying the
1514        // durable log to the durable sidecar restores those ids. Comparing
1515        // its wire form with the transcript reloaded above proves that the
1516        // persisted view is exactly the deterministic projection before we
1517        // use the restamped form for inversion.
1518        let (restamped_view, restamped_log) =
1519            reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
1520        if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
1521            return Err(ServiceError::Operation(
1522                "persisted reduction view does not match its durable log and sidecar".into(),
1523            ));
1524        }
1525        if restamped_log != reloaded_log {
1526            return Err(ServiceError::Operation(
1527                "reapplying the durable reduction log changed its identity".into(),
1528            ));
1529        }
1530        let inverted =
1531            reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
1532        if inverted != session.messages {
1533            return Err(ServiceError::Operation(
1534                "reduction inversion did not restore the source messages byte-exactly".into(),
1535            ));
1536        }
1537
1538        let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
1539        let sidecar_path = store.sidecar_path(&rescue_id);
1540        let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
1541        let bootstrap_prompt = reduced_bootstrap_prompt(
1542            &params.locator,
1543            params.target_harness,
1544            &view_jsonl,
1545            &sidecar_path,
1546            &reduction_log_path,
1547        );
1548        let mut reduced_session = session.clone();
1549        reduced_session.meta.session_id = Some(rescue_id.clone());
1550        reduced_session.messages = view;
1551
1552        Ok(json!({
1553            "session": normalized_session_json(&reduced_session),
1554            "bootstrap_prompt": bootstrap_prompt,
1555            "receipt": {
1556                "id": rescue_id,
1557                "sidecar_id": rescue_id,
1558                "source_harness": params.locator.harness,
1559                "target_harness": params.target_harness.id(),
1560                "source_tokens": source_tokens,
1561                "reduced_tokens": reduced_tokens,
1562                "ratio": ratio,
1563                "source_bytes": source_bytes,
1564                "reduced_bytes": reduced_bytes,
1565                "reductions": reloaded_log.reductions.len(),
1566                "sidecar_path": sidecar_path,
1567                "reduction_log_path": reduction_log_path,
1568                "verified": true,
1569                "reversible": true,
1570            }
1571        }))
1572    }
1573
1574    /// Recognize the one request family whose waiting happens entirely
1575    /// outside this service's state, and hand a transport the half it can run
1576    /// off the task that owns the service.
1577    ///
1578    /// Opening a runtime is the only door here that waits on a foreign
1579    /// program: it spawns the harness's own binary and completes that
1580    /// program's protocol handshake, which takes as long as the program takes
1581    /// to answer. A transport that awaited the whole request inline would
1582    /// stop reading its own input for that whole time, so ONE slow launch
1583    /// would queue every later request on the same server — including reads
1584    /// like `sessions.discover` that touch no runtime at all. Splitting the
1585    /// request lets the transport spawn [`RuntimeOpen::open`] and keep
1586    /// reading, then pay only the short bookkeeping half
1587    /// ([`Self::register_open_runtime`]) when the runtime is up.
1588    ///
1589    /// `None` for every other method: those are answered by
1590    /// [`Self::handle_async`] as before.
1591    pub fn runtime_open(request: &Value) -> Option<RuntimeOpen> {
1592        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1593            return None;
1594        }
1595        let method = request.get("method").and_then(Value::as_str)?;
1596        if !RUNTIME_OPEN_METHODS.contains(&method) {
1597            return None;
1598        }
1599        Some(RuntimeOpen {
1600            id: request.get("id").cloned().unwrap_or(Value::Null),
1601            method: method.to_string(),
1602            params: request.get("params").cloned().unwrap_or_else(|| json!({})),
1603        })
1604    }
1605
1606    /// Recognize a [`DETACHED_METHODS`] request and hand a transport the
1607    /// whole of it: the service-state half is read here and now, and what
1608    /// remains waits on a foreign program with nothing of this service's in
1609    /// hand.
1610    ///
1611    /// Same reason as [`Self::runtime_open`], different doors. Probing a
1612    /// harness starts it and completes its handshake; messaging a live
1613    /// session waits on its Claude relay's send; a conversation verb runs the
1614    /// harness's own CLI or calls its HTTP API. A transport that awaited any
1615    /// of those inline would stop reading its own input for that whole time,
1616    /// so one probe of an unhealthy harness would queue every later request
1617    /// on the same server.
1618    ///
1619    /// Unlike an opening runtime there is no bookkeeping half: the answer
1620    /// [`DetachedCall::run`] produces is the caller's complete response, so a
1621    /// transport writes it without coming back here.
1622    ///
1623    /// `None` for every other method — including the LIVE `sessions.new` /
1624    /// `sessions.reset` door and `runtimes.close`, which wait on a runtime
1625    /// connection this service owns and so are split off by
1626    /// [`Self::detach_runtime`] instead.
1627    pub fn detach(&self, request: &Value) -> Option<DetachedCall> {
1628        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1629            return None;
1630        }
1631        let method = request.get("method").and_then(Value::as_str)?;
1632        if !DETACHED_METHODS.contains(&method) {
1633            return None;
1634        }
1635        let id = request.get("id").cloned().unwrap_or(Value::Null);
1636        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1637        let work = match method {
1638            "harness.v1.harnesses.list" | "harness.v1.harnesses.probe" => self
1639                .inventory_work(method, params)
1640                .map(DetachedWork::Inventory),
1641            "harness.v1.sessions.message" => {
1642                decode::<MessageSessionParams>(params).map(DetachedWork::Message)
1643            }
1644            "harness.v1.sessions.load" => {
1645                decode::<LoadSessionParams>(params).map(DetachedWork::Load)
1646            }
1647            _ => {
1648                let verb = match method {
1649                    "harness.v1.sessions.new" => crate::SessionVerb::New,
1650                    "harness.v1.sessions.reset" => crate::SessionVerb::Reset,
1651                    "harness.v1.sessions.archive" => crate::SessionVerb::Archive,
1652                    _ => crate::SessionVerb::Delete,
1653                };
1654                match decode::<crate::SessionMutation>(params) {
1655                    Ok(mutation) => {
1656                        match crate::sessions_control::door(&mutation.harness, verb) {
1657                            // The live door needs the open runtime connection
1658                            // this service owns; it stays inline.
1659                            Ok(crate::SessionDoor::Live(_)) => return None,
1660                            Ok(_) => Ok(DetachedWork::SessionMutation { verb, mutation }),
1661                            Err(error) => Err(session_control_error(error)),
1662                        }
1663                    }
1664                    Err(error) => Err(error),
1665                }
1666            }
1667        };
1668        Some(DetachedCall {
1669            id,
1670            method: method.to_string(),
1671            work: work.map(Work::Free),
1672        })
1673    }
1674
1675    /// Recognize the two doors that wait on a runtime THIS SERVICE OWNS, and
1676    /// hand a transport the whole of each by lending the connection out.
1677    ///
1678    /// `runtimes.close` surrenders its runtime for good; the LIVE
1679    /// `sessions.new` / `sessions.reset` door borrows one for the length of
1680    /// the slash command and gives it back through
1681    /// [`Self::finish_detached`]. Both are bounded by
1682    /// [`RUNTIME_CONTROL_DEADLINE`], and a wedged runtime spends all of it —
1683    /// which is exactly as long as a transport that awaited them inline would
1684    /// stop reading its own input.
1685    ///
1686    /// `None` for every other method, and for the `sessions.new` /
1687    /// `sessions.reset` doors that are not live: [`Self::detach`] owns those.
1688    pub fn detach_runtime(&mut self, request: &Value) -> Option<DetachedCall> {
1689        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
1690            return None;
1691        }
1692        let method = request.get("method").and_then(Value::as_str)?;
1693        let id = request.get("id").cloned().unwrap_or(Value::Null);
1694        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1695        let work = match method {
1696            "harness.v1.runtimes.close" => decode::<RuntimeCloseParams>(params)
1697                .and_then(|params| self.surrender_runtime_of(&params))
1698                .map(|(runtime, process_group)| {
1699                    Work::Runtime(RuntimeWork::Close {
1700                        runtime,
1701                        process_group,
1702                    })
1703                }),
1704            "harness.v1.sessions.new" | "harness.v1.sessions.reset" => {
1705                let verb = if method == "harness.v1.sessions.new" {
1706                    crate::SessionVerb::New
1707                } else {
1708                    crate::SessionVerb::Reset
1709                };
1710                let mutation = decode::<crate::SessionMutation>(params).ok()?;
1711                // Everything but the live door — including a refusal and a
1712                // request naming no connection — is `detach`'s or
1713                // `handle_async`'s to answer.
1714                let Ok(crate::SessionDoor::Live(command)) =
1715                    crate::sessions_control::door(&mutation.harness, verb)
1716                else {
1717                    return None;
1718                };
1719                let connection = mutation
1720                    .connection
1721                    .clone()
1722                    .filter(|value| !value.trim().is_empty())?;
1723                self.lend_runtime(&connection).map(|runtime| {
1724                    let session = live_session_name(runtime.as_ref(), &mutation);
1725                    Work::Runtime(RuntimeWork::LiveCommand {
1726                        connection,
1727                        runtime,
1728                        verb,
1729                        mutation,
1730                        command,
1731                        session,
1732                    })
1733                })
1734            }
1735            _ => return None,
1736        };
1737        Some(DetachedCall {
1738            id,
1739            method: method.to_string(),
1740            work,
1741        })
1742    }
1743
1744    /// Take back whatever a detached call borrowed and hand over the caller's
1745    /// response. Every answer from [`DetachedCall::run`] comes through here,
1746    /// so a lent-out connection is back in the service before the response
1747    /// that used it is written.
1748    pub fn finish_detached(&mut self, answer: DetachedAnswer) -> Value {
1749        let DetachedAnswer { response, returned } = answer;
1750        if let Some(ReturnedRuntime {
1751            connection,
1752            runtime,
1753        }) = returned
1754        {
1755            self.runtimes_in_flight.remove(&connection);
1756            self.runtimes.insert(connection, runtime);
1757        }
1758        response
1759    }
1760
1761    /// Answer a request split out by [`Self::runtime_open`] and already
1762    /// awaited by [`RuntimeOpen::open`]: register the runtime this service now
1763    /// owns and build its JSON-RPC response.
1764    pub async fn finish_runtime_open(&mut self, opened: OpenedRuntime) -> Value {
1765        let OpenedRuntime { id, outcome } = opened;
1766        let result = match outcome {
1767            Ok(open) => self.register_open_runtime(open).await,
1768            Err(error) => Err(error),
1769        };
1770        service_response(id, result)
1771    }
1772
1773    /// Take ownership of an opened runtime.
1774    async fn register_open_runtime(
1775        &mut self,
1776        open: OpenRuntime,
1777    ) -> std::result::Result<Value, ServiceError> {
1778        match open {
1779            OpenRuntime::Hosted {
1780                runtime,
1781                capabilities,
1782                workspace,
1783                fresh,
1784            } => {
1785                self.insert_hosted_runtime(runtime, capabilities, workspace, fresh)
1786                    .await
1787            }
1788            OpenRuntime::Joined { runtime } => self.insert_runtime(runtime),
1789        }
1790    }
1791
1792    async fn runtime_call(
1793        &mut self,
1794        method: &str,
1795        params: Value,
1796    ) -> std::result::Result<Value, ServiceError> {
1797        match method {
1798            "harness.v1.runtimes.capabilities" => {
1799                let params = decode::<RuntimeBackendParams>(params)?;
1800                // A live endpoint is attached through the hosted-runtime backend whatever harness it
1801                // hosts (runtimes.attach_existing routes it there, and verifies its receipt), so it is
1802                // answered with that backend's capabilities, not the harness adapter's own.
1803                #[cfg(feature = "adapter-api")]
1804                if params
1805                    .base_url
1806                    .as_deref()
1807                    .is_some_and(|value| LiveRuntimeEndpoint::parse(value).is_ok())
1808                {
1809                    return Ok(json!({
1810                        "harness": &params.harness,
1811                        "capabilities": SupercodeHttpRuntimeBackend::live_capabilities(),
1812                    }));
1813                }
1814                let backend = runtime_backend(&params)?;
1815                Ok(json!({
1816                    "harness": backend.harness(),
1817                    "capabilities": backend.capabilities(),
1818                }))
1819            }
1820            method if RUNTIME_OPEN_METHODS.contains(&method) => {
1821                self.register_open_runtime(open_runtime(method, params).await?)
1822                    .await
1823            }
1824            "harness.v1.runtimes.send_input" => {
1825                let params = decode::<RuntimeInputParams>(params)?;
1826                let image_urls = validate_runtime_image_urls(params.image_urls)?;
1827                let runtime = self.runtime_mut(&params.connection)?;
1828                let turn_id = within_control_deadline(
1829                    method,
1830                    runtime.send_input(RuntimeInput {
1831                        text: params.text,
1832                        image_urls,
1833                    }),
1834                )
1835                .await?
1836                .map_err(operation)?;
1837                Ok(json!({"turn_id": turn_id}))
1838            }
1839            "harness.v1.runtimes.interrupt" => {
1840                let params = decode::<RuntimeConnectionParams>(params)?;
1841                within_control_deadline(method, self.runtime_mut(&params.connection)?.interrupt())
1842                    .await?
1843                    .map_err(operation)?;
1844                Ok(json!({}))
1845            }
1846            "harness.v1.runtimes.steer" => {
1847                let params = decode::<RuntimeInputParams>(params)?;
1848                if !params.image_urls.is_empty() {
1849                    return Err(ServiceError::InvalidParams(
1850                        "runtime steering accepts text only".into(),
1851                    ));
1852                }
1853                let text = params.text.trim();
1854                if text.is_empty() || text.chars().count() > 50_000 {
1855                    return Err(ServiceError::InvalidParams(
1856                        "runtime steering requires 1 to 50,000 text characters".into(),
1857                    ));
1858                }
1859                within_control_deadline(
1860                    method,
1861                    self.runtime_mut(&params.connection)?
1862                        .steer(text.to_string()),
1863                )
1864                .await?
1865                .map_err(operation)?;
1866                Ok(json!({}))
1867            }
1868            "harness.v1.runtimes.respond" => {
1869                let params = decode::<RuntimeRespondParams>(params)?;
1870                let request_id = params.request_id.clone();
1871                within_control_deadline(
1872                    method,
1873                    self.runtime_mut(&params.connection)?
1874                        .respond(params.request_id, params.response),
1875                )
1876                .await?
1877                .map_err(operation)?;
1878                // ORCH-9: an answered request is no longer waiting for one.
1879                self.approvals.answered(&params.connection, &request_id);
1880                Ok(json!({}))
1881            }
1882            "harness.v1.runtimes.acquire_control" => {
1883                let params = decode::<RuntimeConnectionParams>(params)?;
1884                let snapshot = within_control_deadline(
1885                    method,
1886                    self.runtime_mut(&params.connection)?.acquire_control(),
1887                )
1888                .await?
1889                .map_err(operation)?;
1890                serde_json::to_value(snapshot)
1891                    .map_err(|error| ServiceError::Operation(error.to_string()))
1892            }
1893            "harness.v1.runtimes.heartbeat" => {
1894                let params = decode::<RuntimeConnectionParams>(params)?;
1895                let snapshot = within_control_deadline(
1896                    method,
1897                    self.runtime_mut(&params.connection)?.heartbeat(),
1898                )
1899                .await?
1900                .map_err(operation)?;
1901                serde_json::to_value(snapshot)
1902                    .map_err(|error| ServiceError::Operation(error.to_string()))
1903            }
1904            "harness.v1.runtimes.detach" => {
1905                let params = decode::<RuntimeConnectionParams>(params)?;
1906                let snapshot =
1907                    within_control_deadline(method, self.runtime_mut(&params.connection)?.detach())
1908                        .await?
1909                        .map_err(operation)?;
1910                serde_json::to_value(snapshot)
1911                    .map_err(|error| ServiceError::Operation(error.to_string()))
1912            }
1913            "harness.v1.runtimes.terminal_instructions" => {
1914                let params = decode::<RuntimeConnectionParams>(params)?;
1915                let launch = self
1916                    .terminal_launches
1917                    .get(&params.connection)
1918                    .ok_or_else(|| {
1919                        ServiceError::Operation(
1920                            "this runtime is not hosted for terminal attachment".into(),
1921                        )
1922                    })?;
1923                Ok(json!({"launch":launch}))
1924            }
1925            "harness.v1.runtimes.close" => {
1926                let params = decode::<RuntimeCloseParams>(params)?;
1927                let (runtime, process_group) = self.surrender_runtime_of(&params)?;
1928                close_runtime(runtime, process_group).await
1929            }
1930            _ => Err(ServiceError::MethodNotFound),
1931        }
1932    }
1933
1934    /// Deliver one message into a session that is running right now.
1935    #[cfg(feature = "adapter-api")]
1936    async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1937        let params = decode::<MessageSessionParams>(params)?;
1938        Ok(message_live_session(&params).await)
1939    }
1940
1941    #[cfg(feature = "adapter-api")]
1942    fn harness_settings_call(
1943        &self,
1944        method: &str,
1945        params: Value,
1946    ) -> std::result::Result<Value, ServiceError> {
1947        let homes = crate::HarnessHomes::default();
1948        match method {
1949            "harness.v1.harnesses.settings" => {
1950                let params = decode::<HarnessSettingsParams>(params)?;
1951                let report = crate::inspect_harness_interop_settings(&homes, &params.harness)
1952                    .map_err(|error| ServiceError::Operation(error.to_string()))?;
1953                serde_json::to_value(report)
1954                    .map_err(|error| ServiceError::Operation(error.to_string()))
1955            }
1956            "harness.v1.harnesses.configure" => {
1957                let params = decode::<ConfigureHarnessParams>(params)?;
1958                let report = crate::configure_harness_interop_settings(
1959                    &homes,
1960                    &params.harness,
1961                    &params.changes,
1962                    params.expected_revision.as_deref(),
1963                )
1964                .map_err(|error| ServiceError::Operation(error.to_string()))?;
1965                serde_json::to_value(report)
1966                    .map_err(|error| ServiceError::Operation(error.to_string()))
1967            }
1968            _ => Err(ServiceError::MethodNotFound),
1969        }
1970    }
1971
1972    fn insert_runtime(
1973        &mut self,
1974        runtime: Box<dyn RuntimeConnection>,
1975    ) -> std::result::Result<Value, ServiceError> {
1976        let connection = format!("runtime-{}", self.next_runtime);
1977        self.next_runtime += 1;
1978        let handle = runtime.handle().clone();
1979        self.runtime_sequences
1980            .entry(handle.runtime_id.clone())
1981            .or_insert(0);
1982        self.runtimes.insert(connection.clone(), runtime);
1983        Ok(json!({"connection": connection, "handle": handle}))
1984    }
1985
1986    #[cfg(feature = "adapter-api")]
1987    async fn insert_hosted_runtime(
1988        &mut self,
1989        runtime: Box<dyn RuntimeConnection>,
1990        capabilities: crate::RuntimeCapabilities,
1991        workspace: PathBuf,
1992        fresh: bool,
1993    ) -> std::result::Result<Value, ServiceError> {
1994        let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities, fresh);
1995        let token: std::sync::Arc<str> = crate::server::generate_token().into();
1996        let server = crate::server::run_frontend_http(
1997            host.clone(),
1998            host.frontend_sender(),
1999            "127.0.0.1:0",
2000            token.clone(),
2001            connection.handle().runtime_id.clone(),
2002        )
2003        .await
2004        .map_err(|error| ServiceError::Operation(error.to_string()))?;
2005        let source = LiveRuntimeSource {
2006            harness: connection.handle().harness.as_str().to_string(),
2007            session_id: connection.handle().runtime_id.clone(),
2008            workspace: workspace.clone(),
2009        };
2010        let registration = register_live_runtime_with_metadata(
2011            connection.handle().runtime_id.clone(),
2012            source.clone(),
2013            format!("http://{}", server.address()),
2014            token.to_string(),
2015            crate::LiveRuntimeMetadata {
2016                child_pid: match &connection.handle().endpoint {
2017                    crate::RuntimeEndpoint::LocalProcess { pid, command, .. }
2018                        if connection.handle().harness.as_str() != "claude-code"
2019                            || command.windows(2).any(|args| {
2020                                args[0] == "--name"
2021                                    && args[1].starts_with(crate::claude_relay::RELAY_NAME_PREFIX)
2022                            }) =>
2023                    {
2024                        *pid
2025                    }
2026                    _ => None,
2027                },
2028                endpoint_capabilities: vec!["http".into(), "acp".into()],
2029                ..Default::default()
2030            },
2031        )
2032        .map_err(|error| ServiceError::Operation(error.to_string()))?;
2033        let endpoint = registration.endpoint().to_string();
2034        let launch = StructuredLaunch {
2035            cwd: workspace,
2036            // Pin attachment to the executable hosting this runtime. A bare
2037            // `supercode` could resolve to an older global install whose CLI
2038            // does not understand the receipt it is being asked to open.
2039            program: std::env::current_exe()
2040                .ok()
2041                .map(|path| path.to_string_lossy().into_owned())
2042                .unwrap_or_else(|| "supercode".into()),
2043            arguments: vec![
2044                "open".into(),
2045                endpoint,
2046                "--harness".into(),
2047                source.harness,
2048                "--session".into(),
2049                source.session_id,
2050            ],
2051            env: BTreeMap::new(),
2052        };
2053        host.retain_registration(registration);
2054        let lease = HostedRuntimeLease {
2055            connection,
2056            _host: host,
2057            _server: server,
2058        };
2059        let opened = self.insert_runtime(Box::new(lease))?;
2060        let connection_id = opened["connection"]
2061            .as_str()
2062            .expect("insert_runtime returns a connection id")
2063            .to_string();
2064        self.terminal_launches.insert(connection_id, launch);
2065        Ok(opened)
2066    }
2067
2068    #[cfg(not(feature = "adapter-api"))]
2069    async fn insert_hosted_runtime(
2070        &mut self,
2071        runtime: Box<dyn RuntimeConnection>,
2072        _capabilities: crate::RuntimeCapabilities,
2073        _workspace: PathBuf,
2074        _fresh: bool,
2075    ) -> std::result::Result<Value, ServiceError> {
2076        self.insert_runtime(runtime)
2077    }
2078
2079    fn runtime_mut(
2080        &mut self,
2081        connection: &str,
2082    ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
2083        if self.runtimes_in_flight.contains(connection) {
2084            return Err(self.lent_out(connection));
2085        }
2086        self.runtimes.get_mut(connection).ok_or_else(|| {
2087            ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2088        })
2089    }
2090
2091    /// What a caller is told about a connection that is out on a detached
2092    /// call. It is not gone and it is not free: it is mid-call, which is the
2093    /// same answer the runtime itself gives a second turn.
2094    fn lent_out(&self, connection: &str) -> ServiceError {
2095        ServiceError::Operation(format!(
2096            "runtime connection `{connection}`: a harness turn is already in progress"
2097        ))
2098    }
2099
2100    /// Take a runtime OUT of the service for the duration of one detached
2101    /// call, leaving its name marked as lent out.
2102    fn lend_runtime(
2103        &mut self,
2104        connection: &str,
2105    ) -> std::result::Result<Box<dyn RuntimeConnection>, ServiceError> {
2106        if self.runtimes_in_flight.contains(connection) {
2107            return Err(self.lent_out(connection));
2108        }
2109        let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2110            ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2111        })?;
2112        self.runtimes_in_flight.insert(connection.to_string());
2113        Ok(runtime)
2114    }
2115
2116    /// Surrender a runtime for good: the connection and everything the
2117    /// service hung off it are gone before its teardown is even attempted.
2118    ///
2119    /// `close` is what a caller reaches for when a runtime has stopped
2120    /// answering, and a runtime that has stopped answering is exactly the one
2121    /// whose graceful close cannot complete: a hosted runtime's own loop
2122    /// parks on the call the runtime never answered, so it never dequeues the
2123    /// shutdown either. Keeping the entry until teardown succeeded made a
2124    /// wedged runtime permanent — every later call on that connection, and
2125    /// every new turn, answered "a harness turn is already in progress" with
2126    /// no way to take the connection back.
2127    /// [`Self::surrender_runtime`] for a `runtimes.close` call, refusing a connection that no longer holds
2128    /// the runtime the caller named.
2129    fn surrender_runtime_of(
2130        &mut self,
2131        params: &RuntimeCloseParams,
2132    ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2133        if params.connection.is_empty() {
2134            let expected = params.runtime_id.as_deref().ok_or_else(|| {
2135                ServiceError::InvalidParams("close requires connection or runtime_id".into())
2136            })?;
2137            let connection = self
2138                .runtimes
2139                .iter()
2140                .find(|(_, runtime)| runtime.handle().runtime_id == expected)
2141                .map(|(connection, _)| connection.clone())
2142                .ok_or_else(|| {
2143                    ServiceError::InvalidParams(format!("unknown runtime `{expected}`"))
2144                })?;
2145            return self.surrender_runtime(&connection);
2146        }
2147        if let Some(expected) = params.runtime_id.as_deref() {
2148            let held = self
2149                .runtimes
2150                .get(&params.connection)
2151                .map(|runtime| runtime.handle().runtime_id.clone());
2152            if held.as_deref() != Some(expected) {
2153                return Err(ServiceError::InvalidParams(format!(
2154                    "runtime connection `{}` does not hold runtime `{expected}`",
2155                    params.connection
2156                )));
2157            }
2158        }
2159        self.surrender_runtime(&params.connection)
2160    }
2161
2162    fn surrender_runtime(
2163        &mut self,
2164        connection: &str,
2165    ) -> std::result::Result<(Box<dyn RuntimeConnection>, Option<u32>), ServiceError> {
2166        if self.runtimes_in_flight.contains(connection) {
2167            return Err(self.lent_out(connection));
2168        }
2169        let runtime = self.runtimes.remove(connection).ok_or_else(|| {
2170            ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
2171        })?;
2172        let process_group = runtime_process_group(runtime.handle());
2173        let runtime_id = runtime.handle().runtime_id.clone();
2174        self.terminal_launches.remove(connection);
2175        self.runtime_sequences.remove(&runtime_id);
2176        self.approvals.forget(connection);
2177        Ok((runtime, process_group))
2178    }
2179
2180    /// SIGKILL the process group of every runtime this service owns, without
2181    /// waiting on any of them.
2182    ///
2183    /// A host leaving for good calls this BEFORE dropping the service. The
2184    /// handle this service holds is not the runtime's connection: a hosted
2185    /// runtime's real transport lives in the task driving it, so neither
2186    /// exiting the process nor dropping these handles reaches the harness
2187    /// process — while dropping them does remove each runtime's live-runtime
2188    /// receipt. Signalling first is what keeps a removed receipt from
2189    /// advertising a harness that is still running.
2190    pub fn kill_all_runtime_groups(&self) -> usize {
2191        self.runtimes
2192            .values()
2193            .filter(|runtime| kill_runtime_process_group(runtime_process_group(runtime.handle())))
2194            .count()
2195    }
2196
2197    /// ORCH-19: run one conversation-lifecycle verb through the harness's own
2198    /// door.
2199    ///
2200    /// Two doors, one shape. A CLI / HTTP / own-store door is self-contained
2201    /// in [`crate::sessions_control`]. A LIVE door (Hermes's and OpenClaw's
2202    /// `/new` and `/reset`, which are slash commands their gateway interprets
2203    /// INSIDE a session) is performed here, because only the service owns the
2204    /// open runtime connection — the command is typed through the very same
2205    /// `send_input` path a human's message takes, so supercode invents no
2206    /// private channel.
2207    async fn mutate_session(
2208        &mut self,
2209        verb: crate::SessionVerb,
2210        params: Value,
2211    ) -> std::result::Result<Value, ServiceError> {
2212        let mutation = decode::<crate::SessionMutation>(params)?;
2213        let door = crate::sessions_control::door(&mutation.harness, verb)
2214            .map_err(session_control_error)?;
2215        let outcome = match door {
2216            // The live door types the slash command through an open hosted
2217            // runtime, which only exists with the `adapter-api` feature; the
2218            // CLI / HTTP / own-store doors below need nothing extra.
2219            #[cfg(not(feature = "adapter-api"))]
2220            crate::SessionDoor::Live(command) => {
2221                return Err(ServiceError::Operation(format!(
2222                    "`{}` performs `sessions.{}` by typing `{command}` into a live driven \
2223                     session, which needs this build's `adapter-api` feature",
2224                    mutation.harness,
2225                    verb.as_str()
2226                )));
2227            }
2228            #[cfg(feature = "adapter-api")]
2229            crate::SessionDoor::Live(command) => {
2230                let connection = mutation
2231                    .connection
2232                    .clone()
2233                    .filter(|value| !value.trim().is_empty())
2234                    .ok_or_else(|| {
2235                        ServiceError::InvalidParams(format!(
2236                            "`{}` performs `sessions.{}` by typing `{command}` into a live \
2237                             driven session: pass the `connection` of an open runtime \
2238                             (`harness.v1.runtimes.start`)",
2239                            mutation.harness,
2240                            verb.as_str()
2241                        ))
2242                    })?;
2243                let runtime = self.runtime_mut(&connection)?;
2244                let session = live_session_name(runtime.as_ref(), &mutation);
2245                // Typing into a live session is a control call on an open
2246                // runtime, and a wedged runtime never accepts one, so it is
2247                // bounded exactly like the other control verbs. A transport
2248                // with a loop of its own lends the connection out instead of
2249                // waiting here: see [`Self::detach_runtime`].
2250                return type_live_command(runtime.as_mut(), verb, &mutation, command, session)
2251                    .await;
2252            }
2253            _ => run_session_mutation(verb, &mutation).await?,
2254        };
2255        serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
2256    }
2257
2258    /// Answer an inventory request whole, for callers that have nowhere to
2259    /// put the waiting half. A transport with a loop of its own splits it
2260    /// instead: see [`Self::detach`].
2261    async fn inventory_call(
2262        &self,
2263        method: &str,
2264        params: Value,
2265    ) -> std::result::Result<Value, ServiceError> {
2266        run_inventory(self.inventory_work(method, params)?).await
2267    }
2268
2269    /// The half of an inventory request that reads this service's state:
2270    /// resolve the selection and count the persisted sessions each row
2271    /// reports. What remains — finding executables, asking them their
2272    /// version, and (at `probe: handshake`) starting each harness and
2273    /// completing its protocol handshake — touches no service state at all.
2274    fn inventory_work(
2275        &self,
2276        method: &str,
2277        params: Value,
2278    ) -> std::result::Result<InventoryWork, ServiceError> {
2279        let mut params = decode::<HarnessInventoryParams>(params)?;
2280        if method == "harness.v1.harnesses.probe" {
2281            let harness = params.harness.take().ok_or_else(|| {
2282                ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
2283            })?;
2284            params.harnesses = vec![harness];
2285        }
2286        let selected = params
2287            .harnesses
2288            .iter()
2289            .map(HarnessId::as_str)
2290            .collect::<std::collections::BTreeSet<_>>();
2291        let supported = harness_support_registry()
2292            .harnesses
2293            .into_iter()
2294            .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
2295            .collect::<Vec<_>>();
2296        if !params.harnesses.is_empty() && supported.len() != selected.len() {
2297            let known = supported
2298                .iter()
2299                .map(|harness| harness.id.as_str())
2300                .collect::<std::collections::BTreeSet<_>>();
2301            let missing = params
2302                .harnesses
2303                .iter()
2304                .filter(|id| !known.contains(id.as_str()))
2305                .map(HarnessId::as_str)
2306                .collect::<Vec<_>>();
2307            return Err(ServiceError::InvalidParams(format!(
2308                "unknown harness(es): {}",
2309                missing.join(", ")
2310            )));
2311        }
2312        let global_counts = params
2313            .include_sessions
2314            .then(|| self.session_counts(None, &params.harnesses));
2315        let workspace_counts = params
2316            .include_sessions
2317            .then(|| {
2318                params
2319                    .workspace
2320                    .as_deref()
2321                    .map(|workspace| self.session_counts(Some(workspace), &params.harnesses))
2322            })
2323            .flatten();
2324        Ok(InventoryWork {
2325            params,
2326            supported,
2327            global_counts,
2328            workspace_counts,
2329        })
2330    }
2331
2332    #[cfg(feature = "adapter-api")]
2333    async fn harness_authentication_call(
2334        &self,
2335        method: &str,
2336        params: Value,
2337    ) -> std::result::Result<Value, ServiceError> {
2338        match method {
2339            "harness.v1.harnesses.auth.methods" | "harness.v1.harnesses.auth.verify" => {
2340                let params = decode::<HarnessAuthenticationParams>(params)?;
2341                serde_json::to_value(crate::inspect_harness_authentication(&params.harness).await)
2342                    .map_err(|error| ServiceError::Operation(error.to_string()))
2343            }
2344            "harness.v1.harnesses.auth.begin" => {
2345                let params = decode::<BeginHarnessAuthenticationParams>(params)?;
2346                let cwd = params
2347                    .cwd
2348                    .or_else(|| std::env::current_dir().ok())
2349                    .unwrap_or_else(|| PathBuf::from("."));
2350                let plan = crate::harness_authentication_plan(
2351                    &params.harness,
2352                    params.environment,
2353                    params.method,
2354                    &cwd,
2355                )
2356                .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
2357                serde_json::to_value(plan)
2358                    .map_err(|error| ServiceError::Operation(error.to_string()))
2359            }
2360            _ => Err(ServiceError::MethodNotFound),
2361        }
2362    }
2363
2364    fn session_counts(
2365        &self,
2366        workspace: Option<&Path>,
2367        harnesses: &[HarnessId],
2368    ) -> BTreeMap<String, usize> {
2369        let mut counts = BTreeMap::new();
2370        for session in self
2371            .catalog
2372            .discover(&DiscoveryQuery {
2373                workspace: workspace.map(Path::to_path_buf),
2374                harnesses: harnesses.to_vec(),
2375                ..DiscoveryQuery::default()
2376            })
2377            .unwrap_or_default()
2378        {
2379            *counts
2380                .entry(session.locator.harness.as_str().to_string())
2381                .or_insert(0) += 1;
2382        }
2383        counts
2384    }
2385}
2386
2387#[async_trait::async_trait]
2388impl SdkService for HarnessSessionService {
2389    fn capabilities(&self) -> SdkCapabilities {
2390        SdkCapabilities::default()
2391    }
2392
2393    async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
2394        if request.operation == SdkOperation::Events {
2395            let events = self
2396                .poll_sdk_events()
2397                .await
2398                .into_iter()
2399                .map(|(_, event)| event)
2400                .collect::<Vec<_>>();
2401            return serde_json::to_value(events).map_err(|error| {
2402                SdkError::new(
2403                    SdkErrorCode::Execution,
2404                    request.operation,
2405                    error.to_string(),
2406                )
2407            });
2408        }
2409        if self.runtimes.is_empty()
2410            && matches!(
2411                request.operation,
2412                SdkOperation::Input
2413                    | SdkOperation::Interrupt
2414                    | SdkOperation::Steer
2415                    | SdkOperation::Respond
2416                    | SdkOperation::Close
2417            )
2418        {
2419            return Err(SdkError::unsupported(request.operation));
2420        }
2421        let method = request
2422            .operation
2423            .method()
2424            .ok_or_else(|| SdkError::unsupported(request.operation))?;
2425        let result = match request.operation {
2426            SdkOperation::Discover
2427            | SdkOperation::Load
2428            | SdkOperation::Export
2429            | SdkOperation::ProfilesList
2430            | SdkOperation::ProfilesGet
2431            | SdkOperation::ProfilesCreate
2432            | SdkOperation::ProfilesDelete
2433            | SdkOperation::SkillsList
2434            | SdkOperation::SkillsInstall
2435            | SdkOperation::SkillsRemove
2436            | SdkOperation::ChannelsList
2437            | SdkOperation::RoutesList
2438            | SdkOperation::TriggersList
2439            | SdkOperation::ChannelsStatus
2440            | SdkOperation::MemoryShow
2441            | SdkOperation::MemorySearch
2442            | SdkOperation::JobsList
2443            | SdkOperation::JobsGet
2444            | SdkOperation::JobsCreate
2445            | SdkOperation::JobsUpdate
2446            | SdkOperation::JobsPause
2447            | SdkOperation::JobsResume
2448            | SdkOperation::JobsRun
2449            | SdkOperation::JobsDelete
2450            | SdkOperation::JobsNotepad
2451            | SdkOperation::JobsNotepadSet
2452            | SdkOperation::JobsNotepadDelete
2453            | SdkOperation::RunsList
2454            | SdkOperation::RunsGet
2455            | SdkOperation::ApprovalsList
2456            | SdkOperation::OrchestrationLoad
2457            | SdkOperation::OrchestrationSave
2458            | SdkOperation::OrchestrationCompile
2459            | SdkOperation::OrchestrationDecompile
2460            | SdkOperation::OrchestrationImport
2461            | SdkOperation::OrchestrationExport
2462            | SdkOperation::WorkflowLoad => self.call(method, request.params),
2463            // ORCH-20: answering needs the live connection, so it takes the
2464            // async door and ends in `harness.v1.runtimes.respond`.
2465            SdkOperation::ApprovalsResolve => self.approvals_resolve(request.params).await,
2466            SdkOperation::Start
2467            | SdkOperation::Resume
2468            | SdkOperation::Input
2469            | SdkOperation::Interrupt
2470            | SdkOperation::Steer
2471            | SdkOperation::Respond
2472            | SdkOperation::Close => self.runtime_call(method, request.params).await,
2473            // ORCH-19 controlled tier. Every verb goes through the HARNESS'S
2474            // OWN door — its CLI, its HTTP API, or its slash command typed
2475            // into a live driven session — and returns the row re-read from
2476            // the harness's store afterwards.
2477            SdkOperation::SessionsNew => {
2478                self.mutate_session(crate::SessionVerb::New, request.params)
2479                    .await
2480            }
2481            SdkOperation::SessionsReset => {
2482                self.mutate_session(crate::SessionVerb::Reset, request.params)
2483                    .await
2484            }
2485            SdkOperation::SessionsArchive => {
2486                self.mutate_session(crate::SessionVerb::Archive, request.params)
2487                    .await
2488            }
2489            SdkOperation::SessionsDelete => {
2490                self.mutate_session(crate::SessionVerb::Delete, request.params)
2491                    .await
2492            }
2493            SdkOperation::Events => unreachable!("handled before method dispatch"),
2494        };
2495        result.map_err(|error| sdk_error(request.operation, error))
2496    }
2497
2498    async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
2499        Ok(self
2500            .poll_sdk_events()
2501            .await
2502            .into_iter()
2503            .map(|(_, event)| event)
2504            .collect())
2505    }
2506}
2507
2508#[cfg(feature = "adapter-api")]
2509struct HostedRuntimeLease {
2510    connection: HostedHarnessConnection,
2511    _host: std::sync::Arc<HostedHarnessRuntime>,
2512    _server: crate::server::FrontendHttpServer,
2513}
2514
2515#[async_trait::async_trait]
2516#[cfg(feature = "adapter-api")]
2517impl RuntimeConnection for HostedRuntimeLease {
2518    fn handle(&self) -> &crate::RuntimeHandle {
2519        self.connection.handle()
2520    }
2521
2522    async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
2523        self.connection.send_input(input).await
2524    }
2525
2526    async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
2527        self.connection.next_event().await
2528    }
2529
2530    async fn interrupt(&mut self) -> crate::Result<()> {
2531        self.connection.interrupt().await
2532    }
2533
2534    // the lease must forward every verb its capabilities advertise; without
2535    // this, steer fell to the trait default and refused a turn it claimed
2536    async fn steer(&mut self, text: String) -> crate::Result<()> {
2537        self.connection.steer(text).await
2538    }
2539
2540    async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
2541        self.connection.respond(request_id, response).await
2542    }
2543
2544    async fn close(&mut self) -> crate::Result<()> {
2545        self.connection.close().await
2546    }
2547}
2548
2549/// One inventory request's waiting half, already separated from the service
2550/// state it reads. See [`HarnessSessionService::inventory_work`].
2551struct InventoryWork {
2552    params: HarnessInventoryParams,
2553    supported: Vec<crate::HarnessSupportDescriptor>,
2554    global_counts: Option<BTreeMap<String, usize>>,
2555    workspace_counts: Option<BTreeMap<String, usize>>,
2556}
2557
2558/// Perform one conversation-lifecycle verb through a door that is
2559/// self-contained in [`crate::sessions_control`]: the harness's own CLI, its
2560/// HTTP API, the orchestrator daemon's socket, or supercode's own store.
2561/// Touches no service state, so this runs on any task. The LIVE door is not
2562/// here — it types its slash command through a runtime connection the service
2563/// owns, and is performed by [`HarnessSessionService::mutate_session`].
2564async fn run_session_mutation(
2565    verb: crate::SessionVerb,
2566    mutation: &crate::SessionMutation,
2567) -> std::result::Result<crate::SessionMutationOutcome, ServiceError> {
2568    // Only the HTTP door actually awaits anything. The CLI, store and daemon
2569    // doors run the harness's own program, or its store, with calls that
2570    // block the calling THREAD from start to finish — a future that never
2571    // yields, which no timeout around it can interrupt and which would hold a
2572    // runtime worker for as long as the harness takes. They go to a blocking
2573    // task, where blocking is what the thread is for.
2574    let door =
2575        crate::sessions_control::door(&mutation.harness, verb).map_err(session_control_error)?;
2576    if let crate::SessionDoor::Http = door {
2577        return crate::sessions_control::mutate(verb, mutation)
2578            .await
2579            .map_err(session_control_error);
2580    }
2581    let mutation = mutation.clone();
2582    tokio::task::spawn_blocking(move || crate::sessions_control::mutate_blocking(verb, &mutation))
2583        .await
2584        .map_err(|error| {
2585            ServiceError::Operation(format!("the conversation verb could not be run: {error}"))
2586        })?
2587        .map_err(session_control_error)
2588}
2589
2590/// Probe every selected harness and assemble the report. Touches no service
2591/// state, so this runs on any task.
2592async fn run_inventory(work: InventoryWork) -> std::result::Result<Value, ServiceError> {
2593    let InventoryWork {
2594        params,
2595        supported,
2596        global_counts,
2597        workspace_counts,
2598    } = work;
2599    let probes = supported.into_iter().map(|descriptor| {
2600        let global = global_counts
2601            .as_ref()
2602            .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2603        let workspace = workspace_counts
2604            .as_ref()
2605            .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
2606        probe_harness(descriptor, &params, global, workspace)
2607    });
2608    let harnesses = futures::future::join_all(probes).await;
2609    serde_json::to_value(HarnessInventoryReport {
2610        probe: params.probe,
2611        workspace: params.workspace,
2612        harnesses,
2613    })
2614    .map_err(|error| ServiceError::Operation(error.to_string()))
2615}
2616
2617async fn probe_harness(
2618    descriptor: crate::HarnessSupportDescriptor,
2619    params: &HarnessInventoryParams,
2620    global: Option<usize>,
2621    workspace: Option<usize>,
2622) -> LocalHarness {
2623    let launch = descriptor.runtime.default_launch.as_ref();
2624    // ORC-7: the orchestrator publishes no runtime launch — it is not an
2625    // adapter supercode connects a turn to. What "installed" means for it
2626    // is that its Node daemon entry is present, so the row answers from
2627    // that instead of from a PATH lookup it could never satisfy.
2628    let orchestrator_entry = (descriptor.id.as_str() == HarnessId::ORCHESTRATOR)
2629        .then(crate::orchestrator::daemon_entry)
2630        .and_then(Result::ok);
2631    let executable = match &orchestrator_entry {
2632        Some(entry) => Some(entry.clone()),
2633        None => launch.and_then(|launch| find_executable(&launch.program)),
2634    };
2635    let installed = executable.is_some();
2636    let version = if params.skip_versions || orchestrator_entry.is_some() {
2637        // The orchestrator's "executable" is a Node module, not a CLI
2638        // with a `--version` flag; running it to ask would start a daemon.
2639        None
2640    } else {
2641        match executable.as_deref() {
2642            Some(path) => executable_version(path).await,
2643            None => None,
2644        }
2645    };
2646    let configured = auth_evidence(descriptor.id.as_str());
2647    let mut auth = if configured {
2648        HarnessAuthState::Configured
2649    } else if matches!(
2650        descriptor.id.as_str(),
2651        HarnessId::CLAUDE_CODE | HarnessId::CODEX
2652    ) {
2653        // These two adapters have explicit native status/login contracts
2654        // and complete local evidence coverage (including Claude's macOS
2655        // Keychain-backed oauthAccount marker). Treating absent evidence
2656        // as unknown advertises a start that will only fail interactively.
2657        HarnessAuthState::Required
2658    } else {
2659        HarnessAuthState::Unknown
2660    };
2661    let mut runtime = if installed {
2662        HarnessRuntimeState::Degraded
2663    } else {
2664        HarnessRuntimeState::Unavailable
2665    };
2666    let is_orchestrator = descriptor.id.as_str() == HarnessId::ORCHESTRATOR;
2667    let mut reason = (!installed).then(|| {
2668        if is_orchestrator {
2669            format!(
2670                "{} is supported but its daemon entry `{}` was not found",
2671                descriptor.display_name,
2672                crate::orchestrator::DAEMON_ENTRY
2673            )
2674        } else {
2675            format!(
2676                "{} is supported but `{}` was not found on PATH",
2677                descriptor.display_name,
2678                launch
2679                    .map(|launch| launch.program.as_str())
2680                    .unwrap_or("executable")
2681            )
2682        }
2683    });
2684    let mut repair = (!installed).then(|| {
2685        if is_orchestrator {
2686            format!(
2687                "Install the `supercode-orchestrator` package so `{}` resolves.",
2688                crate::orchestrator::DAEMON_ENTRY
2689            )
2690        } else {
2691            format!(
2692                "Install {} and ensure `{}` is on PATH.",
2693                descriptor.display_name,
2694                launch
2695                    .map(|launch| launch.program.as_str())
2696                    .unwrap_or("its executable")
2697            )
2698        }
2699    });
2700
2701    if installed && params.probe == HarnessProbeLevel::Handshake {
2702        let backend_params = RuntimeBackendParams {
2703            harness: descriptor.id.clone(),
2704            protocol: None,
2705            launch: None,
2706            base_url: None,
2707            policy: RuntimePolicy::Default,
2708        };
2709        match runtime_backend(&backend_params) {
2710            Ok(backend) => {
2711                let cwd = params
2712                    .workspace
2713                    .clone()
2714                    .or_else(|| std::env::current_dir().ok())
2715                    .unwrap_or_else(|| PathBuf::from("."));
2716                let isolated = descriptor
2717                    .runtime
2718                    .default_launch
2719                    .clone()
2720                    .and_then(|launch| IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok());
2721                let Some(isolated) = isolated else {
2722                    reason = Some(
2723                        "No-prompt runtime handshake could not create its isolated harness home."
2724                            .into(),
2725                    );
2726                    repair = Some(
2727                        "Check temporary-directory permissions, then run the handshake probe again."
2728                            .into(),
2729                    );
2730                    let running = probe_running_instance(descriptor.id.as_str());
2731                    return LocalHarness {
2732                        gateway: gateway_health(
2733                            descriptor.id.as_str(),
2734                            installed,
2735                            running.as_ref(),
2736                            version.as_deref(),
2737                        ),
2738                        id: descriptor.id,
2739                        display_name: descriptor.display_name,
2740                        supported: true,
2741                        installed,
2742                        executable: executable.map(|path| path.to_string_lossy().into_owned()),
2743                        version,
2744                        auth,
2745                        runtime,
2746                        protocol: descriptor.runtime.protocol,
2747                        capabilities: descriptor.runtime.capabilities.clone(),
2748                        effective_capabilities: descriptor.runtime.capabilities,
2749                        sessions: HarnessSessionCounts { global, workspace },
2750                        running,
2751                        reason,
2752                        repair,
2753                    };
2754                };
2755                match tokio::time::timeout(
2756                    Duration::from_secs(30),
2757                    backend.start(RuntimeStartRequest {
2758                        cwd,
2759                        launch: Some(isolated.launch.clone()),
2760                        mcp_servers: Vec::new(),
2761                        approval_policy: None,
2762                    }),
2763                )
2764                .await
2765                {
2766                    Ok(Ok(mut connection)) => {
2767                        match stabilize_handshake(connection.as_mut()).await {
2768                            Ok(()) => {
2769                                auth = HarnessAuthState::Ready;
2770                                runtime = HarnessRuntimeState::Ready;
2771                                reason = Some(
2772                                    "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
2773                                        .into(),
2774                                );
2775                                repair = None;
2776                            }
2777                            Err(message) => {
2778                                auth = if looks_like_auth_error(&message) {
2779                                    HarnessAuthState::Required
2780                                } else if configured {
2781                                    HarnessAuthState::Configured
2782                                } else {
2783                                    HarnessAuthState::Unknown
2784                                };
2785                                reason = Some(format!(
2786                                    "No-prompt runtime handshake became unhealthy during startup: {message}"
2787                                ));
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                                    "Run the harness directly to inspect its startup failure, then probe again."
2795                                        .into()
2796                                });
2797                            }
2798                        }
2799                        let _ =
2800                            tokio::time::timeout(Duration::from_secs(3), connection.close()).await;
2801                    }
2802                    Ok(Err(error)) => {
2803                        let message = truncate_text(&error.to_string(), 500);
2804                        auth = if looks_like_auth_error(&message) {
2805                            HarnessAuthState::Required
2806                        } else if configured {
2807                            HarnessAuthState::Configured
2808                        } else {
2809                            HarnessAuthState::Unknown
2810                        };
2811                        reason = Some(format!("No-prompt runtime handshake failed: {message}"));
2812                        repair = Some(if auth == HarnessAuthState::Required {
2813                            format!(
2814                                "Run `{}` interactively once and complete sign-in, then probe again.",
2815                                launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
2816                            )
2817                        } else {
2818                            "Check the harness installation and run the handshake probe again."
2819                                .into()
2820                        });
2821                    }
2822                    Err(_) => {
2823                        reason =
2824                            Some("No-prompt runtime handshake timed out after 30 seconds.".into());
2825                        repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
2826                    }
2827                }
2828                // Keep the isolated home alive through process teardown.
2829                // Otherwise the compiler may release the last meaningful
2830                // use after cloning `launch`, and a still-starting CLI can
2831                // recreate its state directory after Drop removed it.
2832                // Some Node-based launchers finish a short asynchronous
2833                // installation-id write just after their parent process
2834                // is reaped. Remove once immediately, allow that bounded
2835                // writer to settle, then perform the authoritative pass.
2836                let _ = isolated.cleanup();
2837                tokio::time::sleep(Duration::from_millis(250)).await;
2838                if let Err(error) = isolated.cleanup() {
2839                    auth = if configured {
2840                        HarnessAuthState::Configured
2841                    } else {
2842                        HarnessAuthState::Unknown
2843                    };
2844                    runtime = HarnessRuntimeState::Degraded;
2845                    reason = Some(format!(
2846                        "No-prompt runtime handshake could not remove its isolated harness home: {error}"
2847                    ));
2848                    repair = Some(
2849                        "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
2850                            .into(),
2851                    );
2852                }
2853            }
2854            Err(error) => {
2855                reason = Some(error_message(error));
2856            }
2857        }
2858    } else if installed && configured {
2859        reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
2860    } else if installed && auth == HarnessAuthState::Required {
2861        reason = Some("Executable found, but no native authentication evidence is present.".into());
2862        repair = Some(format!(
2863            "Run `supercode harness login {}` to use the harness-owned sign-in flow.",
2864            descriptor.id.as_str()
2865        ));
2866    } else if installed {
2867        reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
2868        repair = Some(format!(
2869            "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
2870            launch
2871                .map(|launch| launch.program.as_str())
2872                .unwrap_or("the harness")
2873        ));
2874    }
2875
2876    let effective_capabilities = if installed {
2877        descriptor.runtime.capabilities.clone()
2878    } else {
2879        unavailable_capabilities()
2880    };
2881    let running = probe_running_instance(descriptor.id.as_str());
2882    LocalHarness {
2883        gateway: gateway_health(
2884            descriptor.id.as_str(),
2885            installed,
2886            running.as_ref(),
2887            version.as_deref(),
2888        ),
2889        id: descriptor.id,
2890        display_name: descriptor.display_name,
2891        supported: true,
2892        installed,
2893        executable: executable.map(|path| path.to_string_lossy().into_owned()),
2894        version,
2895        auth,
2896        runtime,
2897        protocol: descriptor.runtime.protocol,
2898        capabilities: descriptor.runtime.capabilities,
2899        effective_capabilities,
2900        sessions: HarnessSessionCounts { global, workspace },
2901        running,
2902        reason,
2903        repair,
2904    }
2905}
2906
2907async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
2908    let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
2909    loop {
2910        let now = tokio::time::Instant::now();
2911        if now >= deadline {
2912            return Ok(());
2913        }
2914        match tokio::time::timeout(deadline - now, connection.next_event()).await {
2915            Err(_) => return Ok(()),
2916            Ok(Ok(Some(event))) => {
2917                if let Some(message) = handshake_event_failure(&event) {
2918                    return Err(truncate_text(&message, 500));
2919                }
2920            }
2921            Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
2922            Ok(Err(error)) => return Err(error.to_string()),
2923        }
2924    }
2925}
2926
2927fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
2928    let detail = event
2929        .payload
2930        .get("message")
2931        .or_else(|| event.payload.get("line"))
2932        .and_then(Value::as_str)
2933        .unwrap_or(event.kind.as_str());
2934    match event.kind.as_str() {
2935        "transport_closed" => Some("runtime transport closed during startup".into()),
2936        "transport_error" => Some(format!("runtime transport error: {detail}")),
2937        "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
2938        // Stderr is retained as a runtime event, but is not transport health.
2939        // Grok, for example, can log an AuthorizationRequired error from an
2940        // optional background worker while its ACP session continues to send
2941        // updates and complete prompts normally.
2942        _ => None,
2943    }
2944}
2945
2946fn indexed_claude_window(
2947    locator: &SessionLocator,
2948    options: &SessionLoadOptions,
2949) -> std::result::Result<Option<Value>, ServiceError> {
2950    use supercode_interchange::session::ClaudeReadIndex;
2951    // Exact parent-only window: recursive/full-artifact requests retain the
2952    // existing owner. This is not a bounded display-history substitution.
2953    if locator.harness.as_str() != HarnessId::CLAUDE_CODE
2954        || options.include_subagents != Some(false)
2955    {
2956        return Ok(None);
2957    }
2958    let crate::StorageLocator::File { path } = &locator.storage else {
2959        return Ok(None);
2960    };
2961    if !ClaudeReadIndex::supports(path)
2962        .map_err(|error| ServiceError::Operation(error.to_string()))?
2963    {
2964        return Ok(None);
2965    }
2966    let mut index = ClaudeReadIndex::open(path, Fidelity::ByteLossless)
2967        .map_err(|error| ServiceError::Operation(error.to_string()))?;
2968    let total = index.len();
2969    let (offset, end) = projected_message_window(total, options);
2970    let session = index
2971        .read_messages(offset..end)
2972        .map_err(|error| ServiceError::Operation(error.to_string()))?;
2973    let summary = index
2974        .read_summary()
2975        .map_err(|error| ServiceError::Operation(error.to_string()))?;
2976    let selected_options = SessionLoadOptions {
2977        message_offset: None,
2978        message_limit: None,
2979        message_tail: None,
2980        ..options.clone()
2981    };
2982    let mut selected = projected_session_json(&session, &selected_options);
2983    selected["raw_record_count"] = json!(index.raw_record_count());
2984    Ok(Some(json!({
2985        "session": selected,
2986        "summary": projected_session_summary(&summary, options),
2987        "window": {
2988            "has_more": offset > 0 || end < total, "has_newer": end < total,
2989            "has_older": offset > 0, "newer_items": index.item_count(end..total),
2990            "offset": offset, "older_items": index.item_count(0..offset),
2991            "returned": end - offset, "total_messages": total,
2992        }
2993    })))
2994}
2995
2996fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
2997    let total_messages = session.messages.len();
2998    let (offset, end) = projected_message_window(total_messages, options);
2999    json!({
3000        "session": projected_session_json(session, options),
3001        "summary": projected_session_summary(session, options),
3002        "window": {
3003            "has_more": offset > 0 || end < total_messages,
3004            "has_newer": end < total_messages,
3005            "has_older": offset > 0,
3006            "newer_items": normalized_item_count(&session.messages[end..]),
3007            "offset": offset,
3008            "older_items": normalized_item_count(&session.messages[..offset]),
3009            "returned": end.saturating_sub(offset),
3010            "total_messages": total_messages,
3011        }
3012    })
3013}
3014
3015fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
3016    messages
3017        .iter()
3018        .map(|message| {
3019            let conversation = usize::from(
3020                matches!(message.role, Role::Assistant | Role::User)
3021                    && message_has_content(message),
3022            );
3023            let tool_result =
3024                usize::from(message.role == Role::Tool && message_has_content(message));
3025            conversation + tool_result + message.tool_calls().len()
3026        })
3027        .sum()
3028}
3029
3030fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
3031    let mut conversational = session.messages.iter().filter(|message| {
3032        matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
3033    });
3034    let first_message = conversational.clone().next();
3035    let last_message = conversational.next_back();
3036    let mut assistant = session
3037        .messages
3038        .iter()
3039        .filter(|message| message.role == Role::Assistant && message_has_content(message));
3040    let first_assistant_message = assistant.clone().next();
3041    let last_assistant_message = assistant.next_back();
3042    let end_of_turn = session
3043        .messages
3044        .iter()
3045        .rev()
3046        .find(|message| message.role != Role::System)
3047        .is_some_and(|message| {
3048            message.role == Role::Assistant
3049                && message_has_content(message)
3050                && message.tool_calls().is_empty()
3051                // Codex narrates while it works (`phase: commentary`); only its `final_answer` ends a turn
3052                && message.metadata.get("phase").map(String::as_str) != Some("commentary")
3053        });
3054    let project = |message: Option<&crate::ChatMessage>| {
3055        message.map(|message| project_inline_media(message_json(message), options))
3056    };
3057    json!({
3058        "end_of_turn": end_of_turn,
3059        "first_assistant_message": project(first_assistant_message),
3060        "first_message": project(first_message),
3061        "last_assistant_message": project(last_assistant_message),
3062        "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
3063        "last_message": project(last_message),
3064    })
3065}
3066
3067fn message_has_content(message: &crate::ChatMessage) -> bool {
3068    message
3069        .content
3070        .as_deref()
3071        .is_some_and(|content| !content.trim().is_empty())
3072        || message
3073            .content_parts
3074            .as_ref()
3075            .is_some_and(|parts| !parts.is_empty())
3076}
3077
3078fn message_text(message: &crate::ChatMessage) -> String {
3079    if let Some(content) = &message.content {
3080        return content.clone();
3081    }
3082    message
3083        .content_parts
3084        .as_ref()
3085        .into_iter()
3086        .flatten()
3087        .filter_map(|part| part.get("text").and_then(Value::as_str))
3088        .collect::<Vec<_>>()
3089        .join("\n")
3090}
3091
3092fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
3093    let (offset, end) = projected_message_window(session.messages.len(), options);
3094    let messages = session.messages[offset..end]
3095        .iter()
3096        .map(|message| project_inline_media(message_json(message), options))
3097        .collect::<Vec<_>>();
3098    let subagents = if options.include_subagents.unwrap_or(true) {
3099        // The reported window describes the top-level transcript. Applying it
3100        // recursively would silently truncate subagents without returning a
3101        // window for each child. Keep their histories complete while carrying
3102        // the caller's media policy through the tree.
3103        let subagent_options = SessionLoadOptions {
3104            message_limit: None,
3105            message_offset: None,
3106            message_tail: None,
3107            ..options.clone()
3108        };
3109        session
3110            .subagents
3111            .iter()
3112            .map(|subagent| projected_session_json(subagent, &subagent_options))
3113            .collect::<Vec<_>>()
3114    } else {
3115        Vec::new()
3116    };
3117    json!({
3118        "source": match session.meta.source {
3119            SessionSource::ClaudeCode => "claude_code",
3120            SessionSource::Codex => "codex",
3121            SessionSource::Gemini => "gemini",
3122            SessionSource::Goose => "goose",
3123            SessionSource::Grok => "grok",
3124            SessionSource::Native => "native",
3125            SessionSource::OpenClaw => "openclaw",
3126            SessionSource::Hermes => "hermes",
3127            SessionSource::OpenCode => "opencode",
3128            SessionSource::Pi => "pi",
3129        },
3130        "session_id": session.meta.session_id,
3131        "ended_at": session.meta.ended_at,
3132        "end_reason": session.meta.end_reason,
3133        "model": session.meta.model,
3134        "resolved_profile": recorded_session_profile(session),
3135        "resolved_profile_history": recorded_session_profile_history(session),
3136        "cwd": session.meta.cwd,
3137        "system_prompt": session.meta.system_prompt,
3138        "agent_id": session.meta.agent_id,
3139        "parent_tool_use_id": session.meta.parent_tool_use_id,
3140        "lineage": session.meta.lineage,
3141        "messages": messages,
3142        "subagents": subagents,
3143        "raw_record_count": session.raw.len(),
3144        "parse_error_lines": session.parse_error_lines,
3145    })
3146}
3147
3148/// G3: the model use a native session recorded, from its own records, per UTC day and model. Claude Code
3149/// writes `message.usage` on every line of a response (deduplicated by message id); Codex writes cumulative
3150/// `token_count` totals (each change is the delta). Recomputed whole from the transcript, so a reader can
3151/// replace what it held; nothing is estimated where the harness recorded nothing.
3152/// A session's recorded configuration (`harness.v1.sessions.recorded_config`, ADR 0012): what its harness recorded of
3153/// its directory, model and permissions, without transcript content. Native records are observations; permissions are
3154/// never reconstructed from current defaults.
3155pub fn recorded_config(locator: &SessionLocator) -> crate::Result<Value> {
3156    let session = load_session(locator)?;
3157    let mut config = json!({"harness": locator.harness, "session_id": session.meta.session_id,
3158        "cwd": session.meta.cwd, "model": session.meta.model,
3159        "approval_policy": null, "sandbox_policy": null, "permission_mode": null});
3160    for header in &session.meta.codex_headers {
3161        let payload = header.get("payload").unwrap_or(header);
3162        for key in ["cwd", "model", "approval_policy", "sandbox_policy"] {
3163            if let Some(value) = payload.get(key) {
3164                config[key] = value.clone();
3165            }
3166        }
3167    }
3168    // A Claude Code session's latest recorded permission mode (its transcript's permission-mode records).
3169    if let Ok(manifest) = crate::claude_runtime_state::ClaudeRuntimeManifest::from_session(&session)
3170    {
3171        config["permission_mode"] = json!(manifest.posture.permission_mode);
3172    }
3173    Ok(config)
3174}
3175
3176fn recorded_session_usage(session: &Session) -> Vec<Value> {
3177    let mut days: std::collections::BTreeMap<(String, String), [u64; 4]> =
3178        std::collections::BTreeMap::new();
3179    let mut add =
3180        |at: Option<&str>, model: &str, input: u64, output: u64, read: u64, write: u64| {
3181            let day = at
3182                .and_then(|t| t.get(..10))
3183                .unwrap_or("unknown")
3184                .to_string();
3185            let entry = days.entry((day, model.to_string())).or_insert([0; 4]);
3186            entry[0] += input;
3187            entry[1] += output;
3188            entry[2] += read;
3189            entry[3] += write;
3190        };
3191    let records = session
3192        .raw
3193        .iter()
3194        .filter_map(|line| serde_json::from_str::<Value>(line).ok());
3195    let n = |v: Option<&Value>| v.and_then(Value::as_u64).unwrap_or(0);
3196    match session.meta.source {
3197        SessionSource::ClaudeCode => {
3198            let mut responses: std::collections::BTreeMap<String, Value> =
3199                std::collections::BTreeMap::new();
3200            for record in records {
3201                let Some(message) = record.get("message") else {
3202                    continue;
3203                };
3204                let (Some(id), Some(_)) = (
3205                    message.get("id").and_then(Value::as_str),
3206                    message.get("usage"),
3207                ) else {
3208                    continue;
3209                };
3210                if message.get("model").and_then(Value::as_str) == Some("<synthetic>") {
3211                    continue;
3212                }
3213                responses.insert(id.to_string(), record.clone());
3214            }
3215            for record in responses.values() {
3216                let message = &record["message"];
3217                let usage = &message["usage"];
3218                add(
3219                    record.get("timestamp").and_then(Value::as_str),
3220                    message
3221                        .get("model")
3222                        .and_then(Value::as_str)
3223                        .unwrap_or("unknown"),
3224                    n(usage.get("input_tokens")),
3225                    n(usage.get("output_tokens")),
3226                    n(usage.get("cache_read_input_tokens")),
3227                    n(usage.get("cache_creation_input_tokens")),
3228                );
3229            }
3230        }
3231        SessionSource::Codex => {
3232            let mut model = session
3233                .meta
3234                .model
3235                .clone()
3236                .unwrap_or_else(|| "unknown".into());
3237            let mut previous = [0u64; 3];
3238            for record in records {
3239                let payload = record.get("payload").unwrap_or(&Value::Null);
3240                if record.get("type").and_then(Value::as_str) == Some("turn_context") {
3241                    if let Some(m) = payload.get("model").and_then(Value::as_str) {
3242                        model = m.to_string();
3243                    }
3244                    continue;
3245                }
3246                if payload.get("type").and_then(Value::as_str) != Some("token_count") {
3247                    continue;
3248                }
3249                let Some(total) = payload.get("info").and_then(|i| i.get("total_token_usage"))
3250                else {
3251                    continue;
3252                };
3253                let now = [
3254                    n(total.get("input_tokens")),
3255                    n(total.get("output_tokens")),
3256                    n(total.get("cached_input_tokens")),
3257                ];
3258                if now == previous {
3259                    continue;
3260                }
3261                add(
3262                    record.get("timestamp").and_then(Value::as_str),
3263                    &model,
3264                    now[0].saturating_sub(previous[0]),
3265                    now[1].saturating_sub(previous[1]),
3266                    now[2].saturating_sub(previous[2]),
3267                    0,
3268                );
3269                previous = now;
3270            }
3271        }
3272        _ => {}
3273    }
3274    days.into_iter()
3275        .map(|((day, model), [input, output, read, write])| {
3276            let mut usage = json!({"at": format!("{day}T00:00:00Z"), "model": model,
3277            "input_tokens": input, "output_tokens": output,
3278            "cache_read_tokens": read, "cache_write_tokens": write});
3279            // G3: the harness's records carry no price, so a model the built-in rate table prices gets its list cost
3280            // (cached prompt billed at the input rate, as `pricing` bills it); an unpriced model gets none.
3281            if let Some(price) = crate::pricing::built_in(&model) {
3282                usage["cost"] = json!({"amount": price.cost_usd(input + read + write, output),
3283                    "currency": "USD", "source": "rate_table"});
3284            }
3285            usage
3286        })
3287        .collect()
3288}
3289
3290/// Configuration recorded by the native session, independent of the caller's intended profile.
3291/// Unknown values remain null; neither the reader's binary nor today's folder config proves a past run.
3292fn recorded_session_profile(session: &Session) -> Value {
3293    let harness = match session.meta.source {
3294        SessionSource::ClaudeCode => "claude-code",
3295        SessionSource::Codex => "codex",
3296        SessionSource::Gemini => "gemini",
3297        SessionSource::Goose => "goose",
3298        SessionSource::Grok => "grok",
3299        SessionSource::Native => "native",
3300        SessionSource::OpenClaw => "openclaw",
3301        SessionSource::Hermes => "hermes",
3302        SessionSource::OpenCode => "opencode",
3303        SessionSource::Pi => "pi",
3304    };
3305    let mut model = session
3306        .messages
3307        .iter()
3308        .rev()
3309        .find_map(|message| message.metadata.get("model").cloned())
3310        .or_else(|| session.meta.model.clone());
3311    let mut cwd = session
3312        .meta
3313        .cwd
3314        .as_ref()
3315        .map(|path| path.to_string_lossy().into_owned());
3316    let mut version: Option<String> = None;
3317    let mut approval = Value::Null;
3318    let mut sandbox = Value::Null;
3319    for header in &session.meta.codex_headers {
3320        let payload = header.get("payload").unwrap_or(header);
3321        if let Some(value) = payload.get("cli_version").and_then(Value::as_str) {
3322            version = Some(value.to_owned());
3323        }
3324        if let Some(value) = payload.get("model").and_then(Value::as_str) {
3325            model = Some(value.to_owned());
3326        }
3327        if let Some(value) = payload.get("cwd").and_then(Value::as_str) {
3328            cwd = Some(value.to_owned());
3329        }
3330        if let Some(value) = payload.get("approval_policy") {
3331            approval = value.clone();
3332        }
3333        if let Some(value) = payload.get("sandbox_policy") {
3334            sandbox = value.clone();
3335        }
3336    }
3337    if session.meta.source == SessionSource::ClaudeCode {
3338        version = session
3339            .raw
3340            .iter()
3341            .rev()
3342            .filter_map(|line| serde_json::from_str::<Value>(line).ok())
3343            .find_map(|record| {
3344                record
3345                    .get("version")
3346                    .and_then(Value::as_str)
3347                    .map(str::to_owned)
3348            });
3349    }
3350    json!({"worker":{"harness":harness,"model":model,"cwd":cwd,"binary":{"version":version}},
3351        "permissions":{"native":{"approval":approval,"sandbox":sandbox}},"source":"native-session"})
3352}
3353
3354/// Each native configuration change is retained even when several happened between connector captures.
3355/// No current model, executable or intended configuration is substituted for an earlier observation.
3356fn recorded_session_profile_history(session: &Session) -> Vec<Value> {
3357    let mut profile = json!({"worker":{"harness":if session.meta.source == SessionSource::Codex {"codex"} else {"claude-code"},
3358        "model":null,"cwd":null,"binary":{"version":null}},
3359        "permissions":{"native":{"approval":null,"sandbox":null}},"source":"native-session"});
3360    let records: Box<dyn Iterator<Item = (usize, Value)> + '_> =
3361        if session.meta.source == SessionSource::Codex {
3362            Box::new(session.meta.codex_headers.iter().cloned().enumerate())
3363        } else if session.meta.source == SessionSource::ClaudeCode {
3364            Box::new(session.raw.iter().enumerate().filter_map(|(index, line)| {
3365                serde_json::from_str::<Value>(line)
3366                    .ok()
3367                    .map(|record| (index, record))
3368            }))
3369        } else {
3370            return Vec::new();
3371        };
3372    let mut history = Vec::new();
3373    for (index, record) in records {
3374        let before = profile.clone();
3375        let payload = record.get("payload").unwrap_or(&record);
3376        for (native, target) in [("cwd", "cwd"), ("model", "model")] {
3377            if let Some(value) = payload.get(native).and_then(Value::as_str) {
3378                profile["worker"][target] = json!(value);
3379            }
3380        }
3381        let version = payload
3382            .get("cli_version")
3383            .or_else(|| record.get("version"))
3384            .and_then(Value::as_str);
3385        if let Some(value) = version {
3386            profile["worker"]["binary"]["version"] = json!(value);
3387        }
3388        if let Some(value) = record
3389            .get("message")
3390            .and_then(|m| m.get("model"))
3391            .and_then(Value::as_str)
3392        {
3393            profile["worker"]["model"] = json!(value);
3394        }
3395        for (native, target) in [
3396            ("approval_policy", "approval"),
3397            ("sandbox_policy", "sandbox"),
3398        ] {
3399            if let Some(value) = payload.get(native) {
3400                profile["permissions"]["native"][target] = value.clone();
3401            }
3402        }
3403        if profile != before {
3404            history.push(json!({"native_key":format!("profile-record:{index}"),"at":record.get("timestamp").and_then(Value::as_str),"profile":profile}));
3405        }
3406    }
3407    history
3408}
3409
3410fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
3411    if let Some(tail) = options.message_tail {
3412        return (total.saturating_sub(tail), total);
3413    }
3414    let offset = options.message_offset.unwrap_or(0).min(total);
3415    let end = options
3416        .message_limit
3417        .map(|limit| offset.saturating_add(limit).min(total))
3418        .unwrap_or(total);
3419    (offset, end)
3420}
3421
3422fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
3423    let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
3424        return message;
3425    };
3426    for part in parts {
3427        let Some(url) = part
3428            .get("image_url")
3429            .and_then(|image| image.get("url"))
3430            .and_then(Value::as_str)
3431        else {
3432            continue;
3433        };
3434        let Some(rest) = url.strip_prefix("data:") else {
3435            continue;
3436        };
3437        let Some((media_type, encoded)) = rest.split_once(";base64,") else {
3438            continue;
3439        };
3440        let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
3441        let decoded_bytes = encoded.len().saturating_mul(3) / 4;
3442        let decoded_bytes = decoded_bytes.saturating_sub(padding);
3443        let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
3444            || options
3445                .max_inline_media_bytes
3446                .is_some_and(|limit| decoded_bytes > limit);
3447        if should_elide {
3448            *part = json!({
3449                "type": "media_reference",
3450                "media_type": media_type,
3451                "encoding": "base64",
3452                "encoded_bytes": encoded.len(),
3453                "decoded_bytes": decoded_bytes,
3454                "omitted": true,
3455            });
3456        }
3457    }
3458    message
3459}
3460
3461#[derive(Deserialize)]
3462struct LocatorParams {
3463    locator: SessionLocator,
3464    /// Optional fidelity for the READ surfaces (`sessions.load`,
3465    /// `sessions.follow`).
3466    ///
3467    /// Omitted means [`Fidelity::Semantic`]: these two methods only ever
3468    /// produce a read-only view, and a compacted or resumed-across-files
3469    /// transcript — the everyday shape of a long Claude Code session — has no
3470    /// losslessly reconstructable record graph, so refusing to render it made
3471    /// the mirror unusable rather than accurate. A caller that intends to
3472    /// CONTINUE from what it reads asks for a lossless level explicitly and
3473    /// gets the strict refusal back. Every other method (export, translate,
3474    /// branch, handoff, resume_instructions) is lossless-only and has no
3475    /// such knob.
3476    #[serde(default)]
3477    fidelity: Option<Fidelity>,
3478    /// Optional bounded frontend projection. Absent preserves the historical
3479    /// complete-session read contract.
3480    #[serde(default)]
3481    view: Option<SessionReadView>,
3482}
3483
3484#[derive(Deserialize)]
3485struct SessionReadView {
3486    /// Number of trailing normalized messages to return. Zero is treated as
3487    /// one so a caller cannot accidentally request an unbounded empty mode.
3488    #[serde(default)]
3489    tail_messages: Option<usize>,
3490    /// Whether Claude Code child transcripts belong in this view. The
3491    /// frontend default is false; the legacy no-view path remains true.
3492    #[serde(default)]
3493    include_subagents: bool,
3494    /// Preserve human-visible native history across model-context compaction.
3495    #[serde(default)]
3496    display_history: bool,
3497    /// Bound each individual text field so a single tool result cannot turn a
3498    /// small message window into a hundred-megabyte RPC response.
3499    #[serde(default)]
3500    max_message_chars: Option<usize>,
3501}
3502
3503impl LocatorParams {
3504    fn read_fidelity(&self) -> Fidelity {
3505        self.fidelity.unwrap_or(Fidelity::Semantic)
3506    }
3507
3508    fn include_subagents(&self) -> bool {
3509        self.view
3510            .as_ref()
3511            .map(|view| view.include_subagents)
3512            .unwrap_or(true)
3513    }
3514
3515    fn tail_messages(&self) -> Option<usize> {
3516        self.view
3517            .as_ref()
3518            .and_then(|view| view.tail_messages)
3519            .map(|limit| limit.clamp(1, 5_000))
3520    }
3521
3522    fn display_history(&self) -> bool {
3523        self.view.as_ref().is_some_and(|view| view.display_history)
3524    }
3525
3526    fn max_message_chars(&self) -> Option<usize> {
3527        self.view
3528            .as_ref()
3529            .and_then(|view| view.max_message_chars)
3530            .map(|limit| limit.clamp(256, 64_000))
3531    }
3532
3533    fn bound_session(&self, session: &mut Session) {
3534        bound_session_view(session, self.tail_messages(), self.max_message_chars());
3535    }
3536}
3537
3538#[derive(Debug, Clone, Copy, Default, Deserialize)]
3539#[serde(rename_all = "snake_case")]
3540enum InlineMediaMode {
3541    #[default]
3542    Full,
3543    Metadata,
3544}
3545
3546#[derive(Debug, Clone, Default, Deserialize)]
3547#[serde(default)]
3548struct SessionLoadOptions {
3549    include_subagents: Option<bool>,
3550    inline_media: InlineMediaMode,
3551    max_inline_media_bytes: Option<usize>,
3552    message_limit: Option<usize>,
3553    message_offset: Option<usize>,
3554    message_tail: Option<usize>,
3555}
3556
3557impl SessionLoadOptions {
3558    fn validate(&self) -> std::result::Result<(), ServiceError> {
3559        if self.message_tail.is_some()
3560            && (self.message_limit.is_some() || self.message_offset.is_some())
3561        {
3562            return Err(ServiceError::InvalidParams(
3563                "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
3564                    .into(),
3565            ));
3566        }
3567        Ok(())
3568    }
3569}
3570
3571/// `sessions.load`: a read of one session's own store, sharing no service state (the catalog is stateless), so a
3572/// transport runs it beside its loop: a large transcript's load (seconds) held every later request on the server.
3573fn load_session_door(params: LoadSessionParams) -> std::result::Result<Value, ServiceError> {
3574    if let Some(options) = &params.options {
3575        options.validate()?;
3576        if let Some(result) = indexed_claude_window(&params.read.locator, options)? {
3577            return Ok(result);
3578        }
3579        return load_session(&params.read.locator)
3580            .map(|session| projected_session_result(&session, options))
3581            .map_err(operation);
3582    }
3583    let mut session = if params.read.display_history() {
3584        HarnessCatalog::new()
3585            .load_display_view(
3586                &params.read.locator,
3587                params.read.read_fidelity(),
3588                params.read.tail_messages().unwrap_or(500),
3589            )
3590            .map_err(crate::Error::from)
3591    } else if params.read.include_subagents() {
3592        load_session_with_fidelity(&params.read.locator, params.read.read_fidelity())
3593    } else {
3594        HarnessCatalog::new()
3595            .load_parent_with_fidelity(&params.read.locator, params.read.read_fidelity())
3596            .map_err(crate::Error::from)
3597    }
3598    .map_err(operation)?;
3599    params.read.bound_session(&mut session);
3600    // The session's own record of how it ran (model, version, cwd, policies): a capture reports it to its
3601    // agent's record, as the projected read carries it.
3602    let mut value = normalized_session_json(&session);
3603    value["resolved_profile"] = recorded_session_profile(&session);
3604    value["resolved_profile_history"] = Value::Array(recorded_session_profile_history(&session));
3605    Ok(json!({"session": value}))
3606}
3607
3608#[derive(Deserialize)]
3609struct LoadSessionParams {
3610    #[serde(flatten)]
3611    read: LocatorParams,
3612    #[serde(default)]
3613    options: Option<SessionLoadOptions>,
3614}
3615
3616#[derive(Deserialize)]
3617struct UnfollowParams {
3618    subscription: String,
3619}
3620
3621#[derive(Debug, Deserialize)]
3622#[serde(deny_unknown_fields)]
3623struct IndexResizeParams {
3624    subscription: String,
3625    limit: usize,
3626}
3627
3628/// The Codex sessions whose daemon pane shows a prompt, as last read, and when; and whether a read is
3629/// running. Read on a thread of its own, never on serve's loop.
3630#[cfg(feature = "adapter-api")]
3631#[derive(Default)]
3632struct PanePrompts {
3633    waiting: std::collections::HashSet<String>,
3634    read_at: Option<std::time::Instant>,
3635    reading: bool,
3636}
3637
3638#[cfg(feature = "adapter-api")]
3639static PANE_PROMPTS: std::sync::Mutex<Option<PanePrompts>> = std::sync::Mutex::new(None);
3640
3641/// How old the pane-prompt reading may be before a new one is started.
3642#[cfg(feature = "adapter-api")]
3643const PANE_PROMPT_REFRESH: std::time::Duration = std::time::Duration::from_secs(3);
3644
3645/// A Codex prompt (an approval, a choice) writes nothing to its rollout, which reads that turn as
3646/// working (or unknown); its daemon pane shows the prompt. A running Codex session whose live row
3647/// is waiting needs input, the reading discovery makes. The live rows (ps, the daemon's panes, their
3648/// screens) take about a second to read, and serve's loop answers every request: they are read on a
3649/// thread of their own at most every few seconds while some Codex session runs, and this applies
3650/// the last reading without waiting for it. A prompt reaches a subscription within one refresh.
3651#[cfg(feature = "adapter-api")]
3652fn pane_prompt_override(activities: &mut [crate::SessionActivity], homes: &crate::HarnessHomes) {
3653    let running_codex = |activity: &crate::SessionActivity| {
3654        activity.harness.as_str() == "codex"
3655            && activity.presence == crate::session_activity::SessionPresence::Running
3656            // An approval or a choice comes in a turn: an idle session has none to show, so idle Codex TUIs
3657            // never start the read.
3658            && matches!(
3659                activity.turn,
3660                crate::SessionTurnState::Working | crate::SessionTurnState::Unknown
3661            )
3662    };
3663    if !activities.iter().any(running_codex) {
3664        return;
3665    }
3666    let waiting = {
3667        let Ok(mut guard) = PANE_PROMPTS.lock() else {
3668            return;
3669        };
3670        let prompts = guard.get_or_insert_with(PanePrompts::default);
3671        let stale = prompts
3672            .read_at
3673            .is_none_or(|at| at.elapsed() >= PANE_PROMPT_REFRESH);
3674        if stale && !prompts.reading {
3675            prompts.reading = true;
3676            let homes = homes.clone();
3677            std::thread::spawn(move || {
3678                let doors = crate::mail_route::LiveSessions::read(&homes);
3679                let waiting = doors
3680                    .all()
3681                    .iter()
3682                    .filter(|live| live.address.harness == "codex" && live.status == "waiting")
3683                    .map(|live| live.address.session_id.clone())
3684                    .collect();
3685                if let Ok(mut guard) = PANE_PROMPTS.lock() {
3686                    let prompts = guard.get_or_insert_with(PanePrompts::default);
3687                    prompts.waiting = waiting;
3688                    prompts.read_at = Some(std::time::Instant::now());
3689                    prompts.reading = false;
3690                }
3691            });
3692        }
3693        prompts.waiting.clone()
3694    };
3695    for activity in activities
3696        .iter_mut()
3697        .filter(|activity| running_codex(activity))
3698    {
3699        if waiting.contains(&activity.session_id) {
3700            activity.turn = crate::SessionTurnState::NeedsInput;
3701            activity.evidence.source = "pane_prompt".into();
3702        }
3703    }
3704}
3705
3706#[derive(Deserialize)]
3707struct ActivitySubscribeParams {
3708    locators: Vec<SessionLocator>,
3709    #[serde(default)]
3710    homes: crate::HarnessHomes,
3711}
3712
3713/// Whom `sessions.message` writes to: a session's harness and id, as its mail address names it. A session's storage is
3714/// not needed to deliver mail, and an agent (`{harness: "agent", session_id: <name>}`, docs/adr/0008-agent-mailbox.md)
3715/// has none; a full catalog locator is accepted as well.
3716#[derive(Debug, Clone, Deserialize)]
3717struct MessageTarget {
3718    harness: HarnessId,
3719    session_id: String,
3720    #[serde(default)]
3721    #[allow(dead_code)]
3722    storage: Option<supercode_interchange::catalog::StorageLocator>,
3723}
3724
3725#[derive(Deserialize)]
3726struct MessageSessionParams {
3727    locator: MessageTarget,
3728    text: String,
3729    #[serde(default)]
3730    subject: Option<String>,
3731    /// A stable caller key. Retrying delivery after a lost receipt keeps the same envelope.
3732    #[serde(default)]
3733    idempotency_key: Option<String>,
3734    /// Channel observations retain their trust boundary instead of becoming peer instructions.
3735    #[serde(default)]
3736    channel: bool,
3737    /// Name the sender is known by (`fleet-board`, `aaron`). Replies to the
3738    /// message are filed in this sender's mailbox, read with
3739    /// `sessions.inbox`. Defaults to `supercode`.
3740    #[serde(default)]
3741    from_name: Option<String>,
3742    /// A voice front for the receiving session, never a separate agent.
3743    #[serde(default)]
3744    voice_for: Option<crate::mailbox::MailAddress>,
3745    /// Human-readable sender label, independent of its return mailbox.
3746    #[serde(default)]
3747    sender_name: Option<String>,
3748    /// Id of the message this one answers.
3749    #[serde(default)]
3750    in_reply_to: Option<String>,
3751    /// A channel's own id for this line (`rh2:<message id>`), so a reply naming it
3752    /// by `reply_marker` joins its thread (docs/adr/0008-agent-mailbox.md).
3753    #[serde(default)]
3754    marker: Option<String>,
3755    /// The channel's id for the line this one answers, when the platform marks
3756    /// it a reply (a Slack `thread_ts`, an RH2 `parentMessageId`).
3757    #[serde(default)]
3758    reply_marker: Option<String>,
3759    /// Where the channel shows the thread (`rh2:<room>/<conversation>`).
3760    #[serde(default)]
3761    surface: Option<String>,
3762    /// File one notice in the sender's mailbox when the receiver next goes idle.
3763    #[serde(default)]
3764    notify_when_idle: bool,
3765    /// The text is the session's own user speaking (a voice bridge, a board
3766    /// the owner types in): user mail, delivered as the user's own turn.
3767    /// `harness.v1` is served only to the machine's owner (its daemon admits
3768    /// operators only), which is the authority this carries.
3769    #[serde(default)]
3770    as_user: bool,
3771    /// Same storage roots discovery accepts, so a caller (and a test) can
3772    /// point the live-session registry somewhere other than `$HOME`.
3773    #[serde(default)]
3774    homes: crate::HarnessHomes,
3775}
3776
3777#[derive(Deserialize)]
3778#[serde(deny_unknown_fields)]
3779struct InboxParams {
3780    /// Sender name used with `sessions.message` (its mailbox), or
3781    #[serde(default)]
3782    from_name: Option<String>,
3783    /// an explicit session address `sc:<machine>:<harness>:<id>`.
3784    #[serde(default)]
3785    address: Option<String>,
3786    /// Include messages already read.
3787    #[serde(default)]
3788    all: bool,
3789}
3790
3791#[derive(Deserialize)]
3792#[serde(deny_unknown_fields)]
3793struct HarnessSettingsParams {
3794    harness: String,
3795}
3796
3797#[derive(Deserialize)]
3798#[serde(deny_unknown_fields)]
3799struct ConfigureHarnessParams {
3800    harness: String,
3801    #[serde(default)]
3802    changes: Vec<crate::HarnessSettingChange>,
3803    #[serde(default)]
3804    expected_revision: Option<String>,
3805}
3806
3807fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3808    match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3809        Ok(report) => (
3810            serde_json::to_value(report).unwrap_or(Value::Null),
3811            Value::Null,
3812        ),
3813        Err(error) => (
3814            Value::Null,
3815            Value::String(format!(
3816                "Volter Harness could not inspect Claude Code inbound controls: {error}"
3817            )),
3818        ),
3819    }
3820}
3821
3822/// Deliver `text` into a session that is running right now, or say why not.
3823///
3824/// A refusal is a RESULT, not a JSON-RPC error: "that session is persisted
3825/// only" is an answer about the session, which a mirror renders next to the
3826/// transcript, and this service's error envelope carries no structured data
3827/// field a machine-readable reason could survive in.
3828///
3829/// The door is chosen by the one router every sender uses
3830/// ([`crate::mail_route`]): a session supercode controls gets it through its
3831/// runtime (the default tier); a Claude session it does not control through a
3832/// Claude relay, so the reply comes back; a Codex session through its hook;
3833/// anything else is stored. Replies are filed in the sender's mailbox under
3834/// `reply_to`, read with `sessions.inbox`.
3835async fn message_live_session(params: &MessageSessionParams) -> Value {
3836    use crate::mail_route::{Delivered, NoDoor, Refused};
3837    let (inbound_controls, inbound_controls_error) =
3838        claude_inbound_controls_or_error(&params.homes);
3839    let refused = |reason: &str, message: String| {
3840        json!({
3841            "delivered_to_bus": false,
3842            "refusal": {"reason": reason, "message": message},
3843            "inbound_controls": inbound_controls,
3844            "inbound_controls_error": inbound_controls_error,
3845        })
3846    };
3847    if params.text.trim().is_empty() {
3848        return refused(
3849            crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3850            "refusing to deliver an empty message".into(),
3851        );
3852    }
3853    let sender = match operator_address(params.from_name.as_deref()) {
3854        Ok(sender) => sender,
3855        Err(message) => return refused("invalid_sender", message),
3856    };
3857    let receiver = match crate::mailbox::MailAddress::new(
3858        crate::mailbox::local_machine_name(),
3859        params.locator.harness.as_str(),
3860        &params.locator.session_id,
3861    ) {
3862        Ok(receiver) => receiver,
3863        Err(error) => return refused("delivery_failed", error.to_string()),
3864    };
3865    if params
3866        .voice_for
3867        .as_ref()
3868        .is_some_and(|voice_for| voice_for != &receiver)
3869    {
3870        return refused(
3871            "invalid_sender",
3872            "a voice front must represent the receiving session".into(),
3873        );
3874    }
3875    // A Claude session's registry name must still name only that session. A thread's holder the mailbox launched may
3876    // have stopped (decision 9): its mail is the mailbox's to route, which resumes it (decision 5).
3877    if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3878        && !crate::mail_agent::resumable(&receiver)
3879        && crate::runtime_mail::controlled_runtime(
3880            HarnessId::CLAUDE_CODE,
3881            &params.locator.session_id,
3882        )
3883        .is_none()
3884    {
3885        if let Err(refusal) =
3886            crate::claude_peer::resolve_live_session(&params.homes, &params.locator.session_id)
3887        {
3888            return refused(refusal.reason.as_str(), refusal.message);
3889        }
3890    }
3891    if params.as_user {
3892        // The owner writing to an agent types into its main session (docs/adr/0008-agent-mailbox.md, decisions 3 and
3893        // 12): a typed line there is the owner's, and main answers or delegates it.
3894        let agent_address =
3895            (receiver.harness == crate::mail_agent::AGENT_HARNESS).then(|| receiver.clone());
3896        let receiver = if receiver.harness == crate::mail_agent::AGENT_HARNESS {
3897            match crate::mail_agent::load(&receiver.session_id) {
3898                Some(agent) => agent.main_session,
3899                None => {
3900                    return refused(
3901                        "delivery_failed",
3902                        format!(
3903                            "no agent named {} is declared on this machine",
3904                            receiver.session_id
3905                        ),
3906                    )
3907                }
3908            }
3909        } else {
3910            receiver
3911        };
3912        return message_as_user(
3913            params,
3914            sender,
3915            receiver,
3916            agent_address,
3917            inbound_controls,
3918            inbound_controls_error,
3919        )
3920        .await;
3921    }
3922    // An agent's mail (a root to it, a reply in one of its threads) is routed by its
3923    // mailbox, which resumes a stopped receiver; anything else needs its receiver running.
3924    let door = crate::mail_route::door_for(&params.homes, &receiver);
3925    let mut envelope = match crate::mailbox::Envelope::new(
3926        sender.clone(),
3927        params
3928            .sender_name
3929            .clone()
3930            .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3931        if params.channel {
3932            crate::mailbox::MailKind::Channel
3933        } else {
3934            crate::mailbox::MailKind::Peer
3935        },
3936        crate::mailbox::ReplyVia::Command,
3937        params.text.clone(),
3938    ) {
3939        Ok(envelope) => envelope,
3940        Err(error) => return refused("delivery_failed", error.to_string()),
3941    };
3942    if let Some(key) = &params.idempotency_key {
3943        if key.is_empty() || key.len() > 256 {
3944            return refused("invalid_key", "idempotency_key needs 1–256 bytes".into());
3945        }
3946        let identity = serde_json::to_vec(&(sender.to_string(), receiver.to_string(), key))
3947            .expect("string tuple serializes");
3948        envelope.id = format!("m-{}", &blake3::hash(&identity).to_hex()[..24]);
3949        let mailbox = match crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &receiver) {
3950            Ok(mailbox) => mailbox,
3951            Err(error) => return refused("delivery_failed", error.to_string()),
3952        };
3953        match mailbox.find(&envelope.id) {
3954            Ok(Some(previous)) => {
3955                let saved = previous.envelope;
3956                let subject = params.subject.as_deref().map(|value| {
3957                    value
3958                        .lines()
3959                        .next()
3960                        .unwrap_or("")
3961                        .chars()
3962                        .take(200)
3963                        .collect::<String>()
3964                });
3965                if saved.body != envelope.body
3966                    || saved.subject != subject
3967                    || saved.kind != envelope.kind
3968                    || saved.from_name != envelope.from_name
3969                    || saved.in_reply_to != params.in_reply_to
3970                    || saved.voice_for != params.voice_for
3971                {
3972                    return refused(
3973                        "idempotency_conflict",
3974                        "this key already names different mail".into(),
3975                    );
3976                }
3977            }
3978            Ok(None) => {}
3979            Err(error) => return refused("delivery_failed", error.to_string()),
3980        }
3981    }
3982    envelope.in_reply_to = params.in_reply_to.clone();
3983    envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
3984    envelope.voice_for = params.voice_for.clone();
3985    envelope.subject = params.subject.as_deref().map(|value| {
3986        value
3987            .lines()
3988            .next()
3989            .unwrap_or("")
3990            .chars()
3991            .take(200)
3992            .collect()
3993    });
3994    let channel = crate::mail_agent::Channel {
3995        marker: params.marker.clone(),
3996        reply_marker: params.reply_marker.clone(),
3997        surface: params.surface.clone(),
3998    };
3999    match crate::mail_agent::plan(&mut envelope, &receiver, &channel) {
4000        Err(error) => return refused("delivery_failed", error.to_string()),
4001        Ok(Some(plan)) => {
4002            let caller = crate::mail_route::Caller {
4003                address: sender.clone(),
4004                name: envelope.from_name.clone(),
4005            };
4006            let outcome = crate::mail_send::deliver_planned(
4007                &params.homes,
4008                &caller,
4009                &envelope,
4010                &plan,
4011                false,
4012                params.notify_when_idle,
4013            )
4014            .await;
4015            return json!({
4016                "delivered_to_bus": outcome.code == 0,
4017                "message_id": envelope.id,
4018                "reply_to": sender.to_string(),
4019                "thread": plan.thread.as_ref().map(|thread| thread.id.clone()),
4020                "target": {
4021                    "harness": params.locator.harness.as_str(),
4022                    "session_id": params.locator.session_id,
4023                },
4024                "delivery": {"door": "agent", "how": outcome.text},
4025                "inbound_controls": inbound_controls,
4026                "inbound_controls_error": inbound_controls_error,
4027            });
4028        }
4029        Ok(None) => {}
4030    }
4031    let door = match door {
4032        Ok(door) => door,
4033        Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
4034            return refused(
4035                crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
4036                format!(
4037                    "no running `{}` session `{}` is reachable; its transcript is persisted only",
4038                    params.locator.harness.as_str(),
4039                    params.locator.session_id
4040                ),
4041            )
4042        }
4043    };
4044    let how = match crate::mail_route::deliver(
4045        &envelope,
4046        &receiver,
4047        &door,
4048        true,
4049        params.notify_when_idle,
4050    )
4051    .await
4052    {
4053        Err(detail) => {
4054            return refused(
4055                crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
4056                detail,
4057            )
4058        }
4059        Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
4060        Ok(Err(Refused::TooLong(bytes))) => {
4061            return refused(
4062                "too_long",
4063                format!(
4064                    "the message is {bytes} bytes; the limit is {}",
4065                    crate::mail_route::MAX_RELAYED_BYTES
4066                ),
4067            )
4068        }
4069        Ok(Ok(delivered)) => match delivered {
4070            Delivered::Steered => "steered",
4071            Delivered::Started => "started",
4072            Delivered::Native { busy: true } => "next_tool_call",
4073            Delivered::Native { busy: false } => "started",
4074            Delivered::Hooked => "hook",
4075            Delivered::HookWoken => "woken",
4076            Delivered::Queued => "queued",
4077            Delivered::Stored => "stored",
4078            Delivered::Operator => "filed",
4079        },
4080    };
4081    json!({
4082        "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
4083        "message_id": envelope.id,
4084        "reply_to": sender.to_string(),
4085        "target": {
4086            "harness": params.locator.harness.as_str(),
4087            "session_id": params.locator.session_id,
4088            "name": match &door {
4089                crate::mail_route::Door::Native(session) => Some(session.name.clone()),
4090                _ => None,
4091            },
4092        },
4093        "delivery": {"door": door.name(), "how": how},
4094        "inbound_controls": inbound_controls,
4095        "inbound_controls_error": inbound_controls_error,
4096    })
4097}
4098
4099/// `sessions.message` with `as_user`: user mail, through the one door that
4100/// carries the user's authority (a hosted runtime's input, or the session's
4101/// pane when its composer is empty; otherwise it waits in the mailbox).
4102async fn message_as_user(
4103    params: &MessageSessionParams,
4104    sender: crate::mailbox::MailAddress,
4105    receiver: crate::mailbox::MailAddress,
4106    agent: Option<crate::mailbox::MailAddress>,
4107    inbound_controls: Value,
4108    inbound_controls_error: Value,
4109) -> Value {
4110    let mut envelope = match crate::mailbox::Envelope::new(
4111        sender.clone(),
4112        format!("{}@{}", sender.session_id, sender.machine),
4113        crate::mailbox::MailKind::User,
4114        crate::mailbox::ReplyVia::None,
4115        params.text.clone(),
4116    ) {
4117        Ok(envelope) => envelope,
4118        Err(error) => {
4119            return json!({
4120                "delivered_to_bus": false,
4121                "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4122                "inbound_controls": inbound_controls,
4123                "inbound_controls_error": inbound_controls_error,
4124            })
4125        }
4126    };
4127    // A line to an agent is a root (or a reply in its thread): the agent's mailbox files its thread, with the owner a
4128    // participant, before it is typed into main.
4129    // The session the plan wakes takes the typed line (main for a root, a thread's holder for a reply in it); every
4130    // other participant has it filed unread (main as CC, the owner's own mailbox).
4131    let mut receiver = receiver;
4132    if let Some(agent) = &agent {
4133        envelope.in_reply_to = params.in_reply_to.clone();
4134        envelope.thread = crate::mailbox::thread_of_reply(params.in_reply_to.as_deref());
4135        let channel = crate::mail_agent::Channel::default();
4136        let planned = match crate::mail_agent::plan(&mut envelope, agent, &channel) {
4137            Ok(planned) => planned,
4138            Err(error) => {
4139                return json!({
4140                    "delivered_to_bus": false,
4141                    "refusal": {"reason": "delivery_failed", "message": error.to_string()},
4142                    "inbound_controls": inbound_controls,
4143                    "inbound_controls_error": inbound_controls_error,
4144                })
4145            }
4146        };
4147        if let Some(planned) = planned {
4148            let typed_to = planned
4149                .recipients
4150                .iter()
4151                .find(|r| r.wake && r.address.harness != "operator")
4152                .map(|r| r.address.clone());
4153            if let Some(typed_to) = typed_to {
4154                receiver = typed_to;
4155            }
4156            if let Some(agent) = &planned.agent {
4157                // The agent's own mailbox keeps the record of what was addressed to it.
4158                let _ = crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), agent)
4159                    .and_then(|mailbox| mailbox.deliver_read(&envelope));
4160            }
4161            // A copy of the owner's line is a record of a line typed elsewhere (as decision 12's copies are), never a
4162            // turn of the copy's session: pending user mail would be typed into it.
4163            let mut copy = envelope.clone();
4164            copy.kind = crate::mailbox::MailKind::Typed;
4165            for other in planned
4166                .recipients
4167                .iter()
4168                .filter(|r| r.address != receiver && r.address != sender)
4169            {
4170                let _ = crate::mail_agent::file_unread(&other.address, &copy);
4171            }
4172            // Every owner line to any agent is the owner's account manager's as CC (decision 17), unless it is that
4173            // agent's own.
4174            if let Some(account_manager) = crate::mail_agent::owners_account_manager() {
4175                let main = account_manager.main_session.clone();
4176                let already =
4177                    main == receiver || planned.recipients.iter().any(|r| r.address == main);
4178                if account_manager.name != agent.session_id && !already {
4179                    let _ = crate::mail_agent::file_unread(&main, &copy);
4180                }
4181            }
4182        }
4183    }
4184    // A thread's holder the mailbox launched and that has since closed (decision 9) is resumed the way it was
4185    // launched, and takes the line once its door is back (decision 5); the line is never a new root.
4186    if agent.is_some()
4187        && crate::mail_route::door_for(&params.homes, &receiver).is_err()
4188        && crate::mail_agent::resumable(&receiver)
4189    {
4190        if let Err(error) = crate::mail_agent::resume(&receiver) {
4191            return json!({
4192                "delivered_to_bus": false,
4193                "refusal": {"reason": "no_user_door", "message": format!("{receiver} was not resumed: {error}")},
4194                "inbound_controls": inbound_controls,
4195                "inbound_controls_error": inbound_controls_error,
4196            });
4197        }
4198        let started = std::time::Instant::now();
4199        while crate::mail_route::door_for(&params.homes, &receiver).is_err()
4200            && started.elapsed() < crate::mail_agent::RESUME_WAIT
4201        {
4202            tokio::time::sleep(std::time::Duration::from_secs(1)).await;
4203        }
4204    }
4205    match crate::mail_route::deliver_user_turn(&params.homes, &envelope, &receiver).await {
4206        Ok(turn) => json!({
4207            "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
4208            "message_id": envelope.id,
4209            "target": {
4210                "harness": params.locator.harness.as_str(),
4211                "session_id": params.locator.session_id,
4212            },
4213            "delivery": {
4214                "door": match turn {
4215                    crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
4216                    _ => "pane",
4217                },
4218                "how": turn.as_str(),
4219            },
4220            "inbound_controls": inbound_controls,
4221            "inbound_controls_error": inbound_controls_error,
4222        }),
4223        Err(message) => json!({
4224            "delivered_to_bus": false,
4225            "refusal": {"reason": "no_user_door", "message": message},
4226            "inbound_controls": inbound_controls,
4227            "inbound_controls_error": inbound_controls_error,
4228        }),
4229    }
4230}
4231
4232#[derive(Deserialize)]
4233#[serde(deny_unknown_fields)]
4234struct ActivityUnderParams {
4235    /// Root processes (a terminal pane's shell, say).
4236    pids: Vec<u32>,
4237    #[serde(default)]
4238    homes: crate::HarnessHomes,
4239}
4240
4241/// `harness.v1.sessions.activity_under`: the harness session running beneath
4242/// each root process, with its activity, so a terminal's attention takes a
4243/// harness pane's liveness from the harness's own lifecycle.
4244async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
4245    let params = decode::<ActivityUnderParams>(params)?;
4246    if params.pids.len() > 1024 {
4247        return Err(ServiceError::InvalidParams(
4248            "sessions.activity_under accepts at most 1024 pids".into(),
4249        ));
4250    }
4251    let found = crate::session_activity::activity_under(&params.pids, &params.homes)
4252        .await
4253        .map_err(ServiceError::Sdk)?;
4254    Ok(json!({
4255        "activities": found
4256            .into_iter()
4257            .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
4258            .collect::<Vec<_>>(),
4259    }))
4260}
4261
4262/// Mailbox address of an operator sender (a board, a person's shell) that is
4263/// not itself a harness session: `sc:<machine>:operator:<name>`.
4264fn operator_address(
4265    name: Option<&str>,
4266) -> std::result::Result<crate::mailbox::MailAddress, String> {
4267    let name: String = name
4268        .unwrap_or("supercode")
4269        .trim()
4270        .chars()
4271        .map(|character| {
4272            if character.is_whitespace() || character == '@' {
4273                '-'
4274            } else {
4275                character
4276            }
4277        })
4278        .collect();
4279    if name.is_empty() {
4280        return Err("from_name must not be empty".into());
4281    }
4282    crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
4283        .map_err(|error| error.to_string())
4284}
4285
4286/// `harness.v1.sessions.inbox`: a sender's mailbox, each message with the
4287/// text its reader sees. Unread messages are claimed and acknowledged by this
4288/// read, so a second read does not return them again.
4289fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
4290    let address = match (&params.address, &params.from_name) {
4291        (Some(address), None) => crate::mailbox::MailAddress::parse(address)
4292            .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
4293        (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
4294        (Some(_), Some(_)) => {
4295            return Err(ServiceError::InvalidParams(
4296                "sessions.inbox takes from_name or address, not both".into(),
4297            ))
4298        }
4299    };
4300    let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
4301    let mailbox =
4302        crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
4303    let claimed = mailbox.claim_unread().map_err(operation)?;
4304    let mut messages: Vec<Value> = Vec::new();
4305    if params.all {
4306        for stored in mailbox.list().map_err(operation)? {
4307            if stored.state == crate::mailbox::MailState::Read {
4308                messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4309            }
4310        }
4311    }
4312    for stored in &claimed {
4313        messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
4314    }
4315    for stored in &claimed {
4316        mailbox.acknowledge(stored).map_err(operation)?;
4317    }
4318    Ok(json!({"address": address.to_string(), "messages": messages}))
4319}
4320
4321/// Source identity of one follow subscription, plus the last lifecycle state
4322/// already reported on it. The follower itself stays purely persistence-facing.
4323// Only the adapter-api poll reads these; the subscription bookkeeping itself is
4324// shared by both builds.
4325#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4326struct FollowedSource {
4327    harness: String,
4328    session_id: String,
4329    reported: Option<String>,
4330}
4331
4332#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
4333struct ActivitySubscription {
4334    locators: Vec<SessionLocator>,
4335    homes: crate::HarnessHomes,
4336    reported: BTreeMap<(String, String), crate::SessionActivity>,
4337}
4338
4339/// Add what makes an indexed row behaviorally equivalent to a discovered row: the attach
4340/// endpoint of a session supercode hosts, and the door a message reaches it by.
4341///
4342/// The durable index owns only persistence metadata. Both are projections: every message/attach
4343/// operation revalidates its authority, so publishing one here never trusts a stale browser-held
4344/// handle. The doors are read once per batch ([`crate::mail_route::LiveSessions`]).
4345fn live_descriptor_value(
4346    session: &SessionDescriptor,
4347    doors: &crate::mail_route::LiveSessions,
4348) -> std::result::Result<Value, ServiceError> {
4349    let mut value = serde_json::to_value(session)
4350        .map_err(|error| ServiceError::Operation(error.to_string()))?;
4351    // A running session with no title of its own is known by its live name.
4352    if value.get("title").is_none_or(Value::is_null) {
4353        if let Some(live) = doors.all().iter().find(|live| {
4354            live.address.harness == session.locator.harness.as_str()
4355                && live.address.session_id == session.locator.session_id
4356        }) {
4357            let name = live.name.split('@').next().unwrap_or(&live.name);
4358            if !name.is_empty() {
4359                value["title"] = json!(name);
4360            }
4361        }
4362    }
4363    if let Some(workspace) = &session.cwd {
4364        let source = LiveRuntimeSource {
4365            harness: session.locator.harness.as_str().to_string(),
4366            session_id: session.locator.session_id.clone(),
4367            workspace: workspace.clone(),
4368        };
4369        if let Some(endpoint) = discover_live_runtime(&source)
4370            .map_err(|error| ServiceError::Operation(error.to_string()))?
4371        {
4372            value["live_endpoint"] = json!(endpoint.as_str());
4373        }
4374    }
4375    // How a message reaches this session right now, chosen by the same
4376    // router every sender uses: `runtime` (supercode controls it), `native`,
4377    // `hook` or `stored`. Absent when no process is running it.
4378    if let Some(door) = doors.door(
4379        session.locator.harness.as_str(),
4380        &session.locator.session_id,
4381    ) {
4382        value["delivery"] = json!(door);
4383    }
4384    // What a session waiting on its user is asking, so every surface can show the question and not only the state.
4385    // Read as `message list` reads it: the transcript's open call, else the prompt on the session's pane.
4386    if let Some(live) = doors.all().iter().find(|live| {
4387        live.address.harness == session.locator.harness.as_str()
4388            && live.address.session_id == session.locator.session_id
4389            && live.status == "waiting"
4390    }) {
4391        if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
4392            value["pending_request"] = request;
4393        }
4394    }
4395    Ok(value)
4396}
4397
4398fn live_index_changes(
4399    changes: Vec<crate::session_index::SessionIndexChange>,
4400    homes: &HarnessHomes,
4401) -> std::result::Result<Vec<Value>, ServiceError> {
4402    use crate::session_index::SessionIndexChange;
4403    let doors = crate::mail_route::LiveSessions::read(homes);
4404    changes
4405        .into_iter()
4406        .map(|change| match change {
4407            SessionIndexChange::Added { descriptor } => Ok(json!({
4408                "kind": "added",
4409                "descriptor": live_descriptor_value(&descriptor, &doors)?,
4410            })),
4411            SessionIndexChange::Updated { descriptor } => Ok(json!({
4412                "kind": "updated",
4413                "descriptor": live_descriptor_value(&descriptor, &doors)?,
4414            })),
4415            SessionIndexChange::Removed { key } => Ok(json!({
4416                "kind": "removed",
4417                "key": key,
4418            })),
4419        })
4420        .collect()
4421}
4422
4423fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
4424    use crate::{SessionPresence, SessionTurnState};
4425    match (activity.presence, activity.turn) {
4426        (SessionPresence::Persisted, _) => None,
4427        (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
4428        (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
4429        // Waiting on its user (a question, an approval): the word every other surface uses.
4430        (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
4431        // The normalized activity object can honestly report a live owner even
4432        // when the stock harness never published a turn status. Preserve the
4433        // older field's stricter contract instead of guessing `running`.
4434        (SessionPresence::Running, SessionTurnState::Unknown)
4435            if activity.evidence.native_state.is_none() =>
4436        {
4437            None
4438        }
4439        (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
4440    }
4441}
4442
4443#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
4444#[serde(rename_all = "kebab-case")]
4445enum TransferFormat {
4446    ClaudeCode,
4447    Codex,
4448    #[serde(rename = "opencode", alias = "open-code")]
4449    OpenCode,
4450    Pi,
4451    Grok,
4452    Gemini,
4453    Goose,
4454    /// UNI-18: a Hermes target. Its artifact is the Codex rollout that
4455    /// `hermes sessions import --from codex` reads; `sessions.export` performs
4456    /// that import into the Hermes home.
4457    Hermes,
4458}
4459
4460impl TransferFormat {
4461    fn id(self) -> &'static str {
4462        match self {
4463            Self::ClaudeCode => HarnessId::CLAUDE_CODE,
4464            Self::Codex => HarnessId::CODEX,
4465            Self::OpenCode => HarnessId::OPENCODE,
4466            Self::Pi => HarnessId::PI,
4467            Self::Grok => HarnessId::GROK,
4468            Self::Gemini => HarnessId::GEMINI,
4469            Self::Goose => HarnessId::GOOSE,
4470            Self::Hermes => HarnessId::HERMES,
4471        }
4472    }
4473}
4474
4475impl From<TransferFormat> for SessionFormat {
4476    fn from(value: TransferFormat) -> Self {
4477        match value {
4478            TransferFormat::ClaudeCode => Self::ClaudeCode,
4479            TransferFormat::Codex => Self::Codex,
4480            TransferFormat::OpenCode => Self::OpenCode,
4481            TransferFormat::Pi => Self::Pi,
4482            TransferFormat::Grok => Self::Grok,
4483            TransferFormat::Gemini => Self::Gemini,
4484            TransferFormat::Goose => Self::Goose,
4485            // a Hermes artifact is the Codex rollout Hermes imports
4486            TransferFormat::Hermes => Self::Codex,
4487        }
4488    }
4489}
4490
4491#[derive(Deserialize)]
4492struct ImportSessionParams {
4493    source_harness: TransferFormat,
4494    content: String,
4495}
4496
4497#[derive(Deserialize)]
4498struct ExportSessionParams {
4499    locator: SessionLocator,
4500    target_harness: TransferFormat,
4501}
4502
4503#[derive(Deserialize)]
4504struct ReduceSessionParams {
4505    locator: SessionLocator,
4506    target_harness: TransferFormat,
4507    #[serde(default = "default_keep_last")]
4508    keep_last: usize,
4509}
4510
4511fn default_keep_last() -> usize {
4512    6
4513}
4514
4515#[derive(Deserialize)]
4516struct BranchSessionParams {
4517    locator: SessionLocator,
4518    #[serde(default)]
4519    target_harness: Option<TransferFormat>,
4520}
4521
4522#[derive(Deserialize)]
4523struct HandoffSessionParams {
4524    locator: SessionLocator,
4525    target_harness: TransferFormat,
4526    #[serde(default)]
4527    cwd: Option<PathBuf>,
4528}
4529
4530#[derive(Deserialize)]
4531struct MaterializeSessionParams {
4532    artifact: crate::native_materialize::MaterializeArtifact,
4533    cwd: PathBuf,
4534    /// Where the continuation is written; unset roots are the environment's own, as discovery reads them.
4535    #[serde(default)]
4536    homes: HarnessHomes,
4537}
4538
4539/// A session supercode starts or resumes runs without approval prompts
4540/// (`yolo`) unless the caller asks for the harness's own (`default`).
4541#[derive(Debug, Clone, Copy, Default, Deserialize)]
4542#[serde(rename_all = "snake_case")]
4543enum ResumePolicy {
4544    Default,
4545    #[default]
4546    Yolo,
4547}
4548
4549#[derive(Deserialize)]
4550struct ResumeInstructionsParams {
4551    locator: SessionLocator,
4552    #[serde(default)]
4553    cwd: Option<PathBuf>,
4554    #[serde(default)]
4555    policy: ResumePolicy,
4556}
4557
4558/// `harness.v1.workflow.load` parameters: which harness's board, and its home.
4559#[derive(Deserialize)]
4560struct WorkflowLoadParams {
4561    /// The store, when the caller names it; else the one the home declares.
4562    #[serde(default)]
4563    from: Option<crate::workflow_doors::WorkflowHarness>,
4564    home: PathBuf,
4565}
4566
4567/// ONT-4 `harness.v1.orchestration.load` parameters. `flavor` says which layout the
4568/// folder is read as; our own is the default.
4569#[derive(Deserialize)]
4570struct OrchestrationLoadParams {
4571    root: PathBuf,
4572    #[serde(default)]
4573    flavor: crate::orchestration_doors::HomeFlavor,
4574}
4575
4576/// ONT-4 `harness.v1.orchestration.save` parameters. `vault` is merged into the
4577/// home's own secrets; a caller that sends none keeps what is on disk.
4578#[derive(Deserialize)]
4579struct OrchestrationSaveParams {
4580    root: PathBuf,
4581    orchestration: crate::orchestration::Orchestration,
4582    #[serde(default)]
4583    vault: BTreeMap<String, String>,
4584}
4585
4586/// ONT-4 `harness.v1.orchestration.compile` parameters.
4587#[derive(Deserialize)]
4588struct OrchestrationCompileParams {
4589    from: crate::orchestration_doors::OrchestrationHarness,
4590    home: PathBuf,
4591}
4592
4593/// ONT-4 `harness.v1.orchestration.decompile` parameters. `source` is the home the
4594/// orchestration was compiled from: it is re-compiled to recover the io bookkeeping
4595/// that byte reuse and the live-store refusal (UNI-18) are decided from.
4596#[derive(Deserialize)]
4597struct OrchestrationDecompileParams {
4598    to: crate::orchestration_doors::OrchestrationHarness,
4599    orchestration: crate::orchestration::Orchestration,
4600    source: PathBuf,
4601    #[serde(default)]
4602    source_flavor: crate::orchestration_doors::SourceFlavor,
4603    dest: PathBuf,
4604    #[serde(default)]
4605    vault: BTreeMap<String, String>,
4606}
4607
4608/// `harness.v1.orchestration.import` parameters: another harness's home, and the
4609/// folder of ours it becomes.
4610#[derive(Deserialize)]
4611struct OrchestrationImportParams {
4612    from: crate::orchestration_doors::OrchestrationHarness,
4613    home: PathBuf,
4614    into: PathBuf,
4615}
4616
4617/// `harness.v1.orchestration.export` parameters: a folder of ours, and the home of
4618/// another harness it becomes.
4619#[derive(Deserialize)]
4620struct OrchestrationExportParams {
4621    to: crate::orchestration_doors::OrchestrationHarness,
4622    root: PathBuf,
4623    dest: PathBuf,
4624}
4625
4626/// `harness.v1.jobs.get` parameters.
4627#[derive(Deserialize)]
4628struct JobsGetParams {
4629    harness: String,
4630    id: String,
4631    #[serde(default)]
4632    homes: crate::HarnessHomes,
4633}
4634
4635/// ORCH-18: run one mutating job verb through the harness's own CLI.
4636///
4637/// The refusal ladder is deliberate: a harness with no scheduled-job concept
4638/// at all answers with the SAME sentence `jobs.list` gives it, and a harness
4639/// that has jobs but publishes no client-callable verb (Claude Code, whose
4640/// jobs are created by the model inside a session) answers with its own
4641/// reason. Neither is ever a silent no-op.
4642fn mutate_job(
4643    verb: crate::jobs_control::JobVerb,
4644    params: Value,
4645) -> std::result::Result<Value, ServiceError> {
4646    let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
4647    refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
4648    let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
4649    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4650}
4651
4652/// ORCH-22: run one mutating skills verb through the harness's own door.
4653///
4654/// The refusal ladder mirrors `jobs.*`: a harness with no skills root at all
4655/// answers with the same sentence `skills.list` gives it, and a harness whose
4656/// door does not publish this verb (OpenClaw has no `skills remove` at the
4657/// pin) answers with its own reason. Neither is ever a silent no-op.
4658fn mutate_skill(
4659    verb: crate::skills_control::SkillVerb,
4660    params: Value,
4661) -> std::result::Result<Value, ServiceError> {
4662    let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4663    if !crate::skills_control::supports_skill_control(&mutation.harness) {
4664        return Err(ServiceError::UnsupportedAction(format!(
4665            "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4666            mutation.harness,
4667            verb.as_str(),
4668            crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4669        )));
4670    }
4671    let outcome =
4672        crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4673    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4674}
4675
4676/// The skills twin of [`job_control_error`], with the same mapping rule.
4677fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4678    match error {
4679        crate::skills_control::SkillControlError::Unsupported(message) => {
4680            ServiceError::UnsupportedAction(message)
4681        }
4682        crate::skills_control::SkillControlError::Invalid(message) => {
4683            ServiceError::InvalidParams(message)
4684        }
4685        crate::skills_control::SkillControlError::Failed(message) => {
4686            ServiceError::Operation(message)
4687        }
4688    }
4689}
4690
4691/// ORCH-21: run one mutating profile verb through the harness's own CLI.
4692///
4693/// The refusal ladder mirrors `mutate_job`'s: a harness with no profile
4694/// concept at all answers with the SAME sentence `profiles.list` gives it, and
4695/// a harness that HAS profiles but publishes no client-callable verb (Codex's
4696/// file-authored `[profiles.<name>]` tables, supercode's compiled-in presets)
4697/// answers with its own reason. Neither is ever a silent no-op.
4698fn mutate_profile(
4699    verb: crate::profiles_control::ProfileVerb,
4700    params: Value,
4701) -> std::result::Result<Value, ServiceError> {
4702    let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4703    let outcome =
4704        crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4705    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4706}
4707
4708/// The same mapping `job_control_error` applies, for the profile noun.
4709fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4710    match error {
4711        crate::profiles_control::ProfileControlError::Unsupported(message) => {
4712            ServiceError::UnsupportedAction(message)
4713        }
4714        crate::profiles_control::ProfileControlError::Invalid(message) => {
4715            ServiceError::InvalidParams(message)
4716        }
4717        crate::profiles_control::ProfileControlError::Failed(message) => {
4718            ServiceError::Operation(message)
4719        }
4720    }
4721}
4722
4723/// Map a controlled-tier failure onto the service's error vocabulary. A verb
4724/// the harness lacks is `UnsupportedAction`; a harness verb that RAN and
4725/// failed carries its own stderr through as the operation error.
4726fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4727    match error {
4728        crate::jobs_control::JobControlError::Unsupported(message) => {
4729            ServiceError::UnsupportedAction(message)
4730        }
4731        crate::jobs_control::JobControlError::Invalid(message) => {
4732            ServiceError::InvalidParams(message)
4733        }
4734        crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4735    }
4736}
4737
4738/// Map an ORCH-19 controlled-tier failure onto the service's error
4739/// vocabulary. A verb the harness has no door for is `UnsupportedAction`; a
4740/// door that RAN and failed carries the harness's own stderr / HTTP body
4741/// through as the operation error.
4742fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4743    match error {
4744        crate::SessionControlError::Unsupported(message) => {
4745            ServiceError::UnsupportedAction(message)
4746        }
4747        crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4748        crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4749    }
4750}
4751
4752/// A harness without a scheduled-job concept refuses the verb rather than
4753/// answering with an empty list — an absent capability and an empty inventory
4754/// are different answers (the same rule `runtimes.capabilities` applies to
4755/// `steer`).
4756fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4757    if crate::jobs::supports_jobs(harness) {
4758        return Ok(());
4759    }
4760    Err(ServiceError::UnsupportedAction(format!(
4761        "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4762        crate::jobs::JOB_HARNESSES.join(", ")
4763    )))
4764}
4765
4766/// `harness.v1.runs.get` parameters.
4767#[derive(Deserialize)]
4768struct RunsGetParams {
4769    harness: String,
4770    id: String,
4771    #[serde(default)]
4772    homes: crate::HarnessHomes,
4773}
4774
4775/// A harness with no run store refuses the verb rather than answering with an
4776/// empty history — the same rule `jobs.list` applies. Claude Code lands here
4777/// on purpose: its cron fires are ordinary turns inside the session that
4778/// created the job, so there is no fire record to list.
4779fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4780    if crate::runs::supports_runs(harness) {
4781        return Ok(());
4782    }
4783    Err(ServiceError::UnsupportedAction(format!(
4784        "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4785        crate::runs::RUN_HARNESSES.join(", ")
4786    )))
4787}
4788
4789#[derive(Serialize)]
4790struct SessionArtifact {
4791    source_harness: HarnessId,
4792    target_harness: &'static str,
4793    session_id: Option<String>,
4794    content: String,
4795    suggested_filename: String,
4796    files: Vec<SessionArtifactFile>,
4797    fidelity: Fidelity,
4798    residue: Vec<String>,
4799}
4800
4801#[derive(Serialize)]
4802struct SessionArtifactFile {
4803    path: String,
4804    content: String,
4805    role: ArtifactFileRole,
4806}
4807
4808#[derive(Serialize)]
4809#[serde(rename_all = "snake_case")]
4810enum ArtifactFileRole {
4811    Primary,
4812    Subagent,
4813    Bundle,
4814    SourceRecovery,
4815}
4816
4817#[derive(Serialize)]
4818struct StructuredLaunch {
4819    cwd: PathBuf,
4820    program: String,
4821    arguments: Vec<String>,
4822    env: BTreeMap<String, String>,
4823}
4824
4825struct HandoffInstructions {
4826    launch: StructuredLaunch,
4827    materialize: Option<StructuredLaunch>,
4828    requires_materialization: bool,
4829    note: String,
4830}
4831
4832#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4833#[serde(rename_all = "snake_case")]
4834enum HarnessProbeLevel {
4835    #[default]
4836    Passive,
4837    Handshake,
4838}
4839
4840#[derive(Default, Deserialize)]
4841#[serde(default)]
4842struct HarnessInventoryParams {
4843    harness: Option<HarnessId>,
4844    harnesses: Vec<HarnessId>,
4845    workspace: Option<PathBuf>,
4846    probe: HarnessProbeLevel,
4847    include_sessions: bool,
4848    /// Omit subprocess-based `--version` calls when a latency-sensitive UI only needs readiness.
4849    skip_versions: bool,
4850}
4851
4852#[derive(Deserialize)]
4853struct HarnessAuthenticationParams {
4854    harness: HarnessId,
4855}
4856
4857#[derive(Deserialize)]
4858struct BeginHarnessAuthenticationParams {
4859    harness: HarnessId,
4860    #[serde(default = "local_browser_authentication_environment")]
4861    environment: crate::HarnessAuthenticationEnvironment,
4862    #[serde(default)]
4863    method: Option<crate::HarnessAuthenticationMethodId>,
4864    #[serde(default)]
4865    cwd: Option<PathBuf>,
4866}
4867
4868fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4869    crate::HarnessAuthenticationEnvironment::LocalBrowser
4870}
4871
4872#[derive(Serialize)]
4873struct HarnessInventoryReport {
4874    probe: HarnessProbeLevel,
4875    workspace: Option<PathBuf>,
4876    harnesses: Vec<LocalHarness>,
4877}
4878
4879#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4880#[serde(rename_all = "snake_case")]
4881enum HarnessAuthState {
4882    Ready,
4883    Configured,
4884    Required,
4885    Unknown,
4886}
4887
4888#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4889#[serde(rename_all = "snake_case")]
4890enum HarnessRuntimeState {
4891    Ready,
4892    Degraded,
4893    Unavailable,
4894}
4895
4896#[derive(Serialize)]
4897struct HarnessSessionCounts {
4898    global: Option<usize>,
4899    workspace: Option<usize>,
4900}
4901
4902/// Receipt-backed evidence that a harness has a RUNNING instance right now,
4903/// distinct from being merely installed (UNI-7). Detection is passive and
4904/// default-on: a gateway liveness connect for daemon harnesses, a fresh
4905/// SQLite WAL stamp for store-writer harnesses (precedent: the opencode
4906/// follower's -wal/-shm freshness). Control stays behind per-connection
4907/// grants — this reports observations only.
4908/// ORCH-17: the gateway-health noun on an inventory row. Derived from the
4909/// UNI-7 running-instance probe (Hermes: `state.db-wal` freshness; OpenClaw:
4910/// a TCP connect to the gateway endpoint resolved from its OWN config) plus
4911/// the executable version — never by starting anything.
4912#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4913#[serde(rename_all = "snake_case")]
4914pub enum GatewayState {
4915    Up,
4916    Down,
4917    Unknown,
4918}
4919
4920/// ORCH-17: `gateway` on a `harness.v1.harnesses.list` row.
4921#[derive(Debug, Clone, Serialize)]
4922pub struct GatewayHealth {
4923    pub state: GatewayState,
4924    /// The endpoint supercode would connect to (OpenClaw: the gateway
4925    /// WebSocket resolved from `openclaw.json`; core harnesses: their
4926    /// declared connect address when one exists). `None` when the harness
4927    /// has no single endpoint (Hermes multiplexes platforms).
4928    #[serde(skip_serializing_if = "Option::is_none")]
4929    pub endpoint: Option<String>,
4930    #[serde(skip_serializing_if = "Option::is_none")]
4931    pub version: Option<String>,
4932    /// What the verdict rests on, or why it is `unknown`.
4933    pub evidence: String,
4934    pub checked_at_ms: u64,
4935}
4936
4937/// OpenClaw's gateway WebSocket endpoint, resolved from its own config the
4938/// way the registry's connect descriptor prescribes (`gateway.url`, else
4939/// `gateway.port`, else the documented default).
4940fn openclaw_gateway_endpoint(home: &Path) -> String {
4941    let config_path = home.join(".openclaw/openclaw.json");
4942    let gateway = std::fs::read_to_string(&config_path)
4943        .ok()
4944        .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4945        .and_then(|config| config.get("gateway").cloned());
4946    if let Some(url) = gateway
4947        .as_ref()
4948        .and_then(|gateway| gateway.get("url"))
4949        .and_then(serde_json::Value::as_str)
4950    {
4951        return url.to_string();
4952    }
4953    let port = gateway
4954        .as_ref()
4955        .and_then(|gateway| gateway.get("port"))
4956        .and_then(serde_json::Value::as_u64)
4957        .unwrap_or(18789);
4958    format!("ws://127.0.0.1:{port}")
4959}
4960
4961/// Ask Hermes itself (`hermes gateway status`, read-only, ~1 s) whether its
4962/// gateway is up. The command is per-host launchd/systemd text without a JSON
4963/// form at 0.19–0.21; the verdict is read from the lines it prints:
4964/// "supervised by launchd (PID …)" / "is running" → up, "not running" /
4965/// "not installed" → down, anything else → no verdict. `SUPERCODE_HERMES_BIN`
4966/// overrides the executable so a fake can stand in under test.
4967fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4968    let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4969    let output = std::process::Command::new(&program)
4970        .args(["gateway", "status"])
4971        .stdin(std::process::Stdio::null())
4972        .output()
4973        .ok()?;
4974    let text = format!(
4975        "{}{}",
4976        String::from_utf8_lossy(&output.stdout),
4977        String::from_utf8_lossy(&output.stderr)
4978    );
4979    let verdict = text.lines().find_map(|line| {
4980        let l = line.trim();
4981        if l.contains("supervised by launchd (PID")
4982            || l.contains("supervised by systemd (PID")
4983            || l.contains("Gateway is running")
4984            || l.contains("process is running")
4985        {
4986            Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4987        } else if l.contains("not running") || l.contains("not installed") {
4988            Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4989        } else {
4990            None
4991        }
4992    });
4993    verdict
4994}
4995
4996fn gateway_health(
4997    id: &str,
4998    installed: bool,
4999    running: Option<&RunningInstance>,
5000    version: Option<&str>,
5001) -> GatewayHealth {
5002    let checked_at_ms = now_epoch_ms();
5003    let home = supercode_interchange::user_home()
5004        .map(std::path::PathBuf::into_os_string)
5005        .map(PathBuf::from);
5006    match id {
5007        HarnessId::HERMES | HarnessId::OPENCLAW => {
5008            let endpoint = (id == HarnessId::OPENCLAW)
5009                .then(|| home.as_deref().map(openclaw_gateway_endpoint))
5010                .flatten();
5011            let (state, evidence) = match running {
5012                Some(instance) => (GatewayState::Up, instance.evidence.clone()),
5013                None if !installed => (
5014                    GatewayState::Unknown,
5015                    format!("`{id}` is not installed; no gateway to probe"),
5016                ),
5017                None if id == HarnessId::HERMES => match hermes_gateway_status() {
5018                    // The harness's own door outranks the WAL heuristic: an idle
5019                    // gateway writes nothing for minutes yet is up.
5020                    Some((state, evidence)) => (state, evidence),
5021                    None => (
5022                        GatewayState::Down,
5023                        "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
5024                    ),
5025                },
5026                None => (
5027                    GatewayState::Down,
5028                    format!(
5029                        "no TCP listener at {}",
5030                        endpoint.as_deref().unwrap_or("the gateway endpoint")
5031                    ),
5032                ),
5033            };
5034            GatewayHealth {
5035                state,
5036                endpoint,
5037                version: version.map(str::to_string),
5038                evidence,
5039                checked_at_ms,
5040            }
5041        }
5042        // ORC-7: the orchestrator's gateway IS its daemon, and the daemon's
5043        // own lease file is the record of it. A lease naming a live pid is
5044        // up; a lease whose process is gone is down and says so as a STALE
5045        // lease, never as "no lease"; no lease at all is down. Nothing is
5046        // started, and no port is guessed — the daemon multiplexes adapters
5047        // the way Hermes does, so it has no single endpoint either.
5048        HarnessId::ORCHESTRATOR => {
5049            let root = crate::HarnessHomes::default().orchestrator;
5050            let (state, evidence) = match crate::orchestrator::read_lease(&root) {
5051                Some(lease) if lease.is_live() => (
5052                    GatewayState::Up,
5053                    format!(
5054                        "`{}` names pid {} (started {}), which is live",
5055                        crate::orchestrator::lock_path(&root).display(),
5056                        lease.pid,
5057                        lease.started_at
5058                    ),
5059                ),
5060                Some(lease) => (
5061                    GatewayState::Down,
5062                    format!(
5063                        "stale lease `{}`: pid {} is gone",
5064                        crate::orchestrator::lock_path(&root).display(),
5065                        lease.pid
5066                    ),
5067                ),
5068                None => (
5069                    GatewayState::Down,
5070                    format!(
5071                        "no lease at `{}`; `supercode orchestrator start` writes one",
5072                        crate::orchestrator::lock_path(&root).display()
5073                    ),
5074                ),
5075            };
5076            GatewayHealth {
5077                state,
5078                endpoint: None,
5079                version: version.map(str::to_string),
5080                evidence,
5081                checked_at_ms,
5082            }
5083        }
5084        _ => GatewayHealth {
5085            state: GatewayState::Unknown,
5086            endpoint: None,
5087            version: version.map(str::to_string),
5088            evidence: format!("`{id}` runs per session, not as a gateway"),
5089            checked_at_ms,
5090        },
5091    }
5092}
5093
5094#[derive(Debug, Clone, Serialize)]
5095struct RunningInstance {
5096    /// How the instance was detected.
5097    method: RunningInstanceMethod,
5098    /// The evidence the verdict rests on (endpoint reached / WAL path+age).
5099    evidence: String,
5100    /// Epoch-ms instant the probe executed.
5101    checked_at_ms: u64,
5102}
5103
5104#[derive(Debug, Clone, Copy, Serialize)]
5105#[serde(rename_all = "snake_case")]
5106enum RunningInstanceMethod {
5107    /// A TCP connect to the harness's own configured gateway endpoint
5108    /// succeeded.
5109    GatewayConnect,
5110    /// The harness's session store has an active SQLite WAL (a live writer
5111    /// holds the store open and stamped it recently).
5112    StoreWalActivity,
5113}
5114
5115fn now_epoch_ms() -> u64 {
5116    std::time::SystemTime::now()
5117        .duration_since(std::time::UNIX_EPOCH)
5118        .map(|elapsed| elapsed.as_millis() as u64)
5119        .unwrap_or(0)
5120}
5121
5122/// OpenClaw: the gateway endpoint comes from the harness's OWN config
5123/// (`<home>/.openclaw/openclaw.json` — `gateway.url` or `gateway.port`,
5124/// default port 18789); a successful TCP connect is the running signal.
5125fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
5126    let config_path = home.join(".openclaw/openclaw.json");
5127    let text = std::fs::read_to_string(&config_path).ok();
5128    let gateway = text
5129        .as_deref()
5130        .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
5131        .and_then(|config| config.get("gateway").cloned());
5132    let address = gateway
5133        .as_ref()
5134        .and_then(|gateway| gateway.get("url"))
5135        .and_then(serde_json::Value::as_str)
5136        .and_then(|url| {
5137            url.split("://").nth(1).map(|rest| {
5138                rest.trim_end_matches('/')
5139                    .split('/')
5140                    .next()
5141                    .unwrap_or(rest)
5142                    .to_string()
5143            })
5144        })
5145        .unwrap_or_else(|| {
5146            let port = gateway
5147                .as_ref()
5148                .and_then(|gateway| gateway.get("port"))
5149                .and_then(serde_json::Value::as_u64)
5150                .unwrap_or(18789);
5151            format!("127.0.0.1:{port}")
5152        });
5153    let reachable = std::net::TcpStream::connect_timeout(
5154        &address.parse().ok()?,
5155        std::time::Duration::from_millis(400),
5156    )
5157    .is_ok();
5158    reachable.then(|| RunningInstance {
5159        method: RunningInstanceMethod::GatewayConnect,
5160        evidence: format!(
5161            "gateway endpoint {address} accepted a TCP connect (from {})",
5162            config_path.display()
5163        ),
5164        checked_at_ms: now_epoch_ms(),
5165    })
5166}
5167
5168/// Hermes: `<home>/.hermes/state.db-wal` freshly modified means a live writer
5169/// holds the store open (SQLite WAL exists only while a connection is open;
5170/// a recent stamp distinguishes an active instance from a stale crash
5171/// leftover).
5172fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
5173    let wal = home.join(".hermes/state.db-wal");
5174    let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
5175    let age_ms = std::time::SystemTime::now()
5176        .duration_since(modified)
5177        .map(|age| age.as_millis() as u64)
5178        .unwrap_or(u64::MAX);
5179    (age_ms <= max_wal_age_ms).then(|| RunningInstance {
5180        method: RunningInstanceMethod::StoreWalActivity,
5181        evidence: format!(
5182            "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
5183            wal.display()
5184        ),
5185        checked_at_ms: now_epoch_ms(),
5186    })
5187}
5188
5189/// Default-on running-instance detection for the harnesses that have one.
5190fn probe_running_instance(id: &str) -> Option<RunningInstance> {
5191    let home = supercode_interchange::user_home()
5192        .map(std::path::PathBuf::into_os_string)
5193        .map(PathBuf::from)?;
5194    match id {
5195        HarnessId::OPENCLAW => probe_openclaw_running(&home),
5196        HarnessId::HERMES => probe_hermes_running(&home, 300_000),
5197        _ => None,
5198    }
5199}
5200
5201#[derive(Serialize)]
5202struct LocalHarness {
5203    id: HarnessId,
5204    display_name: String,
5205    supported: bool,
5206    installed: bool,
5207    executable: Option<String>,
5208    version: Option<String>,
5209    auth: HarnessAuthState,
5210    runtime: HarnessRuntimeState,
5211    protocol: String,
5212    capabilities: crate::RuntimeCapabilities,
5213    effective_capabilities: crate::RuntimeCapabilities,
5214    sessions: HarnessSessionCounts,
5215    /// Receipt-backed running-instance detection (None = not detected or the
5216    /// harness has no running-instance concept). Distinct from `installed`.
5217    #[serde(skip_serializing_if = "Option::is_none")]
5218    running: Option<RunningInstance>,
5219    /// ORCH-17: gateway health derived from `running` + the harness's own config.
5220    gateway: GatewayHealth,
5221    reason: Option<String>,
5222    repair: Option<String>,
5223}
5224
5225#[derive(Clone, Deserialize)]
5226struct RuntimeBackendParams {
5227    harness: HarnessId,
5228    #[serde(default)]
5229    protocol: Option<String>,
5230    #[serde(default)]
5231    launch: Option<RuntimeLaunch>,
5232    #[serde(default)]
5233    base_url: Option<String>,
5234    #[serde(default)]
5235    policy: RuntimePolicy,
5236}
5237
5238/// A session supercode starts or resumes runs without approval prompts
5239/// (`yolo`) unless the caller asks for the harness's own (`default`).
5240#[derive(Debug, Clone, Copy, Default, Deserialize)]
5241#[serde(rename_all = "snake_case")]
5242enum RuntimePolicy {
5243    Default,
5244    #[default]
5245    Yolo,
5246}
5247
5248#[derive(Deserialize)]
5249struct RuntimeStartParams {
5250    #[serde(flatten)]
5251    backend: RuntimeBackendParams,
5252    cwd: PathBuf,
5253    /// MCP servers to mount into the new session through the harness's own
5254    /// start door (ORC-6). Backends without such a door ignore them.
5255    #[serde(default)]
5256    mcp_servers: Vec<crate::McpServerLaunch>,
5257    /// The session's approval policy, where the harness's start door takes one (Codex).
5258    #[serde(default)]
5259    approval_policy: Option<String>,
5260}
5261
5262#[derive(Deserialize)]
5263struct RuntimeAttachParams {
5264    #[serde(flatten)]
5265    backend: RuntimeBackendParams,
5266    runtime_id: String,
5267    #[serde(default)]
5268    cwd: Option<PathBuf>,
5269    /// MCP servers to mount into the resumed session (the start door's own
5270    /// field, carried again because a session's tools die with its process).
5271    #[serde(default)]
5272    mcp_servers: Vec<crate::McpServerLaunch>,
5273    /// The session's approval policy, carried again on resume as on start (Codex).
5274    #[serde(default)]
5275    approval_policy: Option<String>,
5276}
5277
5278#[derive(Deserialize)]
5279struct RuntimeConnectionParams {
5280    connection: String,
5281}
5282
5283/// `runtimes.close`: a caller that recorded `runtime_id` with the connection closes it only while the
5284/// connection still holds that runtime (a restarted daemon numbers its connections afresh).
5285#[derive(Deserialize)]
5286struct RuntimeCloseParams {
5287    #[serde(default)]
5288    connection: String,
5289    #[serde(default)]
5290    runtime_id: Option<String>,
5291}
5292
5293#[derive(Deserialize)]
5294struct RuntimeInputParams {
5295    connection: String,
5296    text: String,
5297    #[serde(default)]
5298    image_urls: Vec<String>,
5299}
5300
5301const MAX_RUNTIME_IMAGES: usize = 4;
5302const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
5303const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
5304
5305fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
5306    if image_urls.len() > MAX_RUNTIME_IMAGES {
5307        return Err(ServiceError::InvalidParams(format!(
5308            "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
5309        )));
5310    }
5311    let mut total = 0usize;
5312    for url in &image_urls {
5313        if !(url.starts_with("data:image/")
5314            || url.starts_with("https://")
5315            || url.starts_with("http://"))
5316        {
5317            return Err(ServiceError::InvalidParams(
5318                "runtime images must be image data URLs or HTTP(S) URLs".into(),
5319            ));
5320        }
5321        if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
5322            return Err(ServiceError::InvalidParams(format!(
5323                "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
5324            )));
5325        }
5326        total = total.saturating_add(url.len());
5327    }
5328    if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
5329        return Err(ServiceError::InvalidParams(format!(
5330            "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
5331        )));
5332    }
5333    Ok(image_urls)
5334}
5335
5336#[derive(Deserialize)]
5337struct RuntimeRespondParams {
5338    connection: String,
5339    request_id: Value,
5340    response: Value,
5341}
5342
5343fn default_reduction_store_root() -> PathBuf {
5344    if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
5345        return PathBuf::from(root).join("sessions");
5346    }
5347    if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
5348        return PathBuf::from(home).join(".supercode").join("sessions");
5349    }
5350    PathBuf::from(".supercode").join("sessions")
5351}
5352
5353fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
5354    let mut output = String::new();
5355    for message in messages {
5356        output.push_str(
5357            &serde_json::to_string(message)
5358                .map_err(|error| ServiceError::Operation(error.to_string()))?,
5359        );
5360        output.push('\n');
5361    }
5362    Ok(output)
5363}
5364
5365fn parse_messages_jsonl(
5366    content: &str,
5367) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
5368    content
5369        .lines()
5370        .enumerate()
5371        .filter(|(_, line)| !line.trim().is_empty())
5372        .map(|(index, line)| {
5373            serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
5374                ServiceError::Operation(format!(
5375                    "reduced transcript line {} is invalid: {error}",
5376                    index + 1
5377                ))
5378            })
5379        })
5380        .collect()
5381}
5382
5383fn reduced_bootstrap_prompt(
5384    source: &SessionLocator,
5385    target: TransferFormat,
5386    view_jsonl: &str,
5387    sidecar_path: &Path,
5388    reduction_log_path: &Path,
5389) -> String {
5390    format!(
5391        "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
5392         \n\
5393         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\
5394         \n\
5395         <supercode-reduced-session source-session=\"{source_id}\">\n\
5396         {view_jsonl}\
5397         </supercode-reduced-session>\n\
5398         \n\
5399         Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
5400        source_harness = source.harness.as_str(),
5401        target_harness = target.id(),
5402        sidecar = sidecar_path.display(),
5403        log = reduction_log_path.display(),
5404        source_id = source.session_id,
5405    )
5406}
5407
5408fn session_artifact(
5409    locator: &SessionLocator,
5410    session: &Session,
5411    target: TransferFormat,
5412) -> std::result::Result<SessionArtifact, ServiceError> {
5413    session_artifact_with_id(locator, session, target, None)
5414}
5415
5416fn session_artifact_with_id(
5417    locator: &SessionLocator,
5418    session: &Session,
5419    target: TransferFormat,
5420    target_session_id: Option<&str>,
5421) -> std::result::Result<SessionArtifact, ServiceError> {
5422    let format: SessionFormat = target.into();
5423    let diagonal = format.source() == session.meta.source;
5424    crate::residue_store::store_segments(session);
5425    let has_appended_turns = session
5426        .imported_message_count
5427        .is_some_and(|imported| imported < session.messages.len());
5428    let mut restoration = None;
5429    let content = if let Some(id) = target_session_id {
5430        if diagonal && format != SessionFormat::OpenCode {
5431            session
5432                .to_jsonl_spliced(format, Some(id))
5433                .map_err(operation)?
5434        } else {
5435            let mut rewritten = session.clone();
5436            rewritten.meta.session_id = Some(id.to_string());
5437            rewritten.to_jsonl(format).map_err(operation)?
5438        }
5439    } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
5440        session.raw_verbatim()
5441    } else if diagonal {
5442        session.to_jsonl_spliced(format, None).map_err(operation)?
5443    } else {
5444        // A session that came from `format` before returns its source records verbatim for the
5445        // prefix the residue store holds (docs/plans/portable-residue.md).
5446        match session
5447            .restore_residue(format, crate::residue_store::lookup)
5448            .map_err(operation)?
5449        {
5450            Some((content, report)) => {
5451                restoration = Some(report);
5452                content
5453            }
5454            None => session.to_jsonl(format).map_err(operation)?,
5455        }
5456    };
5457    let stem = sanitize_filename(
5458        target_session_id
5459            .or(session.meta.session_id.as_deref())
5460            .unwrap_or(&locator.session_id),
5461    );
5462    let suggested_filename = if diagonal && target == TransferFormat::Grok {
5463        "chat_history.jsonl".to_string()
5464    } else if target == TransferFormat::Goose {
5465        format!("{stem}.goose.json")
5466    } else {
5467        format!("{stem}.{}.jsonl", target.id())
5468    };
5469    let mut files = vec![SessionArtifactFile {
5470        path: suggested_filename.clone(),
5471        content: content.clone(),
5472        role: ArtifactFileRole::Primary,
5473    }];
5474    if target == TransferFormat::ClaudeCode {
5475        let bundle_stem = Path::new(&suggested_filename)
5476            .file_stem()
5477            .and_then(|stem| stem.to_str())
5478            .unwrap_or(&stem);
5479        let mut child_paths = BTreeSet::new();
5480        for (index, subagent) in session.subagents.iter().enumerate() {
5481            let agent_id = subagent
5482                .meta
5483                .agent_id
5484                .as_deref()
5485                .map(|id| id.strip_prefix("agent-").unwrap_or(id))
5486                .map(sanitize_filename)
5487                .filter(|id| !id.is_empty())
5488                .unwrap_or_else(|| format!("subagent-{}", index + 1));
5489            let child_has_appended_turns = subagent
5490                .imported_message_count
5491                .is_some_and(|imported| imported < subagent.messages.len());
5492            let child_content = if target_session_id.is_none()
5493                && subagent.meta.source == SessionSource::ClaudeCode
5494                && subagent.raw_is_verbatim
5495                && !child_has_appended_turns
5496            {
5497                subagent.raw_verbatim()
5498            } else if subagent.meta.source == SessionSource::ClaudeCode {
5499                subagent
5500                    .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
5501                    .map_err(operation)?
5502            } else {
5503                let mut child = subagent.clone();
5504                if let Some(id) = target_session_id {
5505                    child.meta.session_id = Some(id.to_string());
5506                }
5507                child
5508                    .to_jsonl(SessionFormat::ClaudeCode)
5509                    .map_err(operation)?
5510            };
5511            let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
5512            if !child_paths.insert(path.clone()) {
5513                return Err(ServiceError::Operation(format!(
5514                    "Claude subagent ids collide at artifact path `{path}`"
5515                )));
5516            }
5517            files.push(SessionArtifactFile {
5518                path,
5519                content: child_content,
5520                role: ArtifactFileRole::Subagent,
5521            });
5522        }
5523    }
5524    if diagonal && target == TransferFormat::Grok {
5525        append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
5526    }
5527    if !diagonal || !session.raw_is_verbatim {
5528        files.push(SessionArtifactFile {
5529            path: "recovery/source.supercode.jsonl".into(),
5530            content: session.to_native_jsonl(),
5531            role: ArtifactFileRole::SourceRecovery,
5532        });
5533        for (index, subagent) in session.subagents.iter().enumerate() {
5534            let id = subagent
5535                .meta
5536                .agent_id
5537                .as_deref()
5538                .map(sanitize_filename)
5539                .unwrap_or_else(|| format!("subagent-{}", index + 1));
5540            files.push(SessionArtifactFile {
5541                path: format!("recovery/subagents/{id}.supercode.jsonl"),
5542                content: subagent.to_native_jsonl(),
5543                role: ArtifactFileRole::SourceRecovery,
5544            });
5545        }
5546    }
5547    if !diagonal && session.meta.source == SessionSource::Grok {
5548        append_grok_bundle_files(
5549            locator,
5550            "recovery/grok/",
5551            ArtifactFileRole::SourceRecovery,
5552            &mut files,
5553        )?;
5554    }
5555    let (fidelity, residue) = if diagonal
5556        && target_session_id.is_none()
5557        && session.raw_is_verbatim
5558        && !has_appended_turns
5559    {
5560        (Fidelity::ByteLossless, Vec::new())
5561    } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
5562        (
5563            Fidelity::ValueLossless,
5564            vec![if target_session_id.is_some() {
5565                "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
5566            } else {
5567                "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
5568            }],
5569        )
5570    } else {
5571        match restoration {
5572            Some(report) if report.rendered_messages == 0 => (
5573                Fidelity::ByteLossless,
5574                vec![format!(
5575                    "restored verbatim from this conversation's {} source records in the residue store",
5576                    target.id()
5577                )],
5578            ),
5579            Some(report) => (
5580                Fidelity::Semantic,
5581                vec![format!(
5582                    "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
5583                    report.restored_messages,
5584                    report.restored_messages + report.rendered_messages,
5585                    report.rendered_messages,
5586                    target.id()
5587                )],
5588            ),
5589            None => (
5590                Fidelity::Semantic,
5591                vec!["target schema has no portable slot for every source-native record and metadata field".into()],
5592            ),
5593        }
5594    };
5595    Ok(SessionArtifact {
5596        source_harness: locator.harness.clone(),
5597        target_harness: target.id(),
5598        session_id: target_session_id
5599            .map(str::to_string)
5600            .or_else(|| session.meta.session_id.clone()),
5601        content,
5602        suggested_filename,
5603        files,
5604        fidelity,
5605        residue,
5606    })
5607}
5608
5609fn append_grok_bundle_files(
5610    locator: &SessionLocator,
5611    prefix: &str,
5612    role: ArtifactFileRole,
5613    files: &mut Vec<SessionArtifactFile>,
5614) -> std::result::Result<(), ServiceError> {
5615    let primary = locator.storage.path();
5616    if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
5617        return Err(ServiceError::Operation(format!(
5618            "Grok bundle locator must name chat_history.jsonl, got {}",
5619            primary.display()
5620        )));
5621    }
5622    let parent = primary.parent().ok_or_else(|| {
5623        ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
5624    })?;
5625    for name in ["summary.json", "updates.jsonl"] {
5626        let path = parent.join(name);
5627        let metadata = match std::fs::symlink_metadata(&path) {
5628            Ok(metadata) => metadata,
5629            Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
5630            Err(error) => return Err(ServiceError::Operation(error.to_string())),
5631        };
5632        if metadata.file_type().is_symlink() || !metadata.is_file() {
5633            return Err(ServiceError::Operation(format!(
5634                "refusing non-regular Grok bundle member {}",
5635                path.display()
5636            )));
5637        }
5638        let content = std::fs::read_to_string(&path).map_err(|error| {
5639            ServiceError::Operation(format!(
5640                "Grok bundle member {} is not representable as UTF-8: {error}",
5641                path.display()
5642            ))
5643        })?;
5644        files.push(SessionArtifactFile {
5645            path: format!("{prefix}{name}"),
5646            content,
5647            role: match role {
5648                ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
5649                _ => ArtifactFileRole::SourceRecovery,
5650            },
5651        });
5652    }
5653    Ok(())
5654}
5655
5656fn handoff_artifact(
5657    locator: &SessionLocator,
5658    session: &Session,
5659    target: TransferFormat,
5660) -> std::result::Result<SessionArtifact, ServiceError> {
5661    let target_session_id = target_session_id(target);
5662    session_artifact_with_id(locator, session, target, Some(&target_session_id))
5663}
5664
5665fn target_session_id(target: TransferFormat) -> String {
5666    let uuid = generated_session_id();
5667    match target {
5668        TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5669        TransferFormat::ClaudeCode
5670        | TransferFormat::Codex
5671        | TransferFormat::Pi
5672        | TransferFormat::Grok
5673        | TransferFormat::Gemini
5674        | TransferFormat::Goose
5675        | TransferFormat::Hermes => uuid,
5676    }
5677}
5678
5679fn sanitize_filename(value: &str) -> String {
5680    let value = value
5681        .chars()
5682        .map(|character| {
5683            if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5684                character
5685            } else {
5686                '-'
5687            }
5688        })
5689        .collect::<String>();
5690    let value = value.trim_matches('-');
5691    if value.is_empty() {
5692        "session".into()
5693    } else {
5694        value.chars().take(100).collect()
5695    }
5696}
5697
5698fn handoff_instructions(
5699    target: TransferFormat,
5700    session_id: &str,
5701    cwd: &Path,
5702) -> HandoffInstructions {
5703    let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5704        cwd: cwd.to_path_buf(),
5705        program: program.into(),
5706        arguments,
5707        env: BTreeMap::new(),
5708    };
5709    match target {
5710        TransferFormat::ClaudeCode => HandoffInstructions {
5711            launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5712            materialize: None,
5713            requires_materialization: true,
5714            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(),
5715        },
5716        TransferFormat::Hermes => HandoffInstructions {
5717            launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5718            materialize: None,
5719            requires_materialization: true,
5720            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(),
5721        },
5722        TransferFormat::Codex => HandoffInstructions {
5723            launch: launch("codex", vec!["resume".into(), session_id.into()]),
5724            materialize: None,
5725            requires_materialization: true,
5726            note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5727        },
5728        TransferFormat::OpenCode => HandoffInstructions {
5729            launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5730            materialize: Some(launch(
5731                "opencode",
5732                vec!["import".into(), "{artifact_path}".into()],
5733            )),
5734            requires_materialization: true,
5735            note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5736        },
5737        TransferFormat::Pi => HandoffInstructions {
5738            launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5739            materialize: None,
5740            requires_materialization: true,
5741            note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5742        },
5743        TransferFormat::Grok => HandoffInstructions {
5744            launch: launch(
5745                "grok",
5746                vec!["--resume".into(), "{materialized_session_id}".into()],
5747            ),
5748            materialize: None,
5749            requires_materialization: true,
5750            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(),
5751        },
5752        TransferFormat::Gemini => HandoffInstructions {
5753            launch: launch(
5754                "gemini",
5755                vec!["--session-file".into(), "{artifact_path}".into()],
5756            ),
5757            materialize: None,
5758            requires_materialization: true,
5759            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(),
5760        },
5761        TransferFormat::Goose => HandoffInstructions {
5762            launch: launch(
5763                "goose",
5764                vec![
5765                    "session".into(),
5766                    "--resume".into(),
5767                    "--session-id".into(),
5768                    "{imported_session_id}".into(),
5769                ],
5770            ),
5771            materialize: Some(launch(
5772                "goose",
5773                vec!["session".into(), "import".into(), "{artifact_path}".into()],
5774            )),
5775            requires_materialization: true,
5776            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(),
5777        },
5778    }
5779}
5780
5781fn resume_launch(
5782    harness: &str,
5783    session_id: &str,
5784    cwd: &Path,
5785    policy: ResumePolicy,
5786) -> std::result::Result<StructuredLaunch, ServiceError> {
5787    let mut arguments = Vec::new();
5788    let program = match harness {
5789        HarnessId::GROK => {
5790            if matches!(policy, ResumePolicy::Yolo) {
5791                if crate::support::self_sandbox_supported() {
5792                    arguments.extend(["--sandbox".into(), "workspace".into()]);
5793                }
5794                arguments.push("--always-approve".into());
5795            }
5796            arguments.extend(["--resume".into(), session_id.into()]);
5797            "grok"
5798        }
5799        HarnessId::CODEX => {
5800            arguments.extend(crate::startup_prompts::startup_arguments(
5801                harness,
5802                Some(cwd),
5803                &[],
5804                matches!(policy, ResumePolicy::Yolo),
5805            ));
5806            arguments.extend(["resume".into(), session_id.into()]);
5807            "codex"
5808        }
5809        HarnessId::CLAUDE_CODE => {
5810            arguments.extend(crate::startup_prompts::startup_arguments(
5811                harness,
5812                Some(cwd),
5813                &[],
5814                matches!(policy, ResumePolicy::Yolo),
5815            ));
5816            arguments.extend(["--resume".into(), session_id.into()]);
5817            "claude"
5818        }
5819        HarnessId::GEMINI => {
5820            arguments.extend(crate::startup_prompts::startup_arguments(
5821                harness,
5822                Some(cwd),
5823                &[],
5824                matches!(policy, ResumePolicy::Yolo),
5825            ));
5826            arguments.extend(["--resume".into(), session_id.into()]);
5827            "gemini"
5828        }
5829        HarnessId::GOOSE => {
5830            arguments.extend([
5831                "session".into(),
5832                "--resume".into(),
5833                "--session-id".into(),
5834                session_id.into(),
5835            ]);
5836            "goose"
5837        }
5838        HarnessId::PI => {
5839            arguments.extend(crate::startup_prompts::startup_arguments(
5840                harness,
5841                Some(cwd),
5842                &[],
5843                matches!(policy, ResumePolicy::Yolo),
5844            ));
5845            arguments.extend(["--session".into(), session_id.into()]);
5846            "pi"
5847        }
5848        HarnessId::OPENCODE => {
5849            arguments.extend(["--session".into(), session_id.into()]);
5850            "opencode"
5851        }
5852        HarnessId::SUPERCODE => {
5853            if matches!(policy, ResumePolicy::Yolo) {
5854                arguments.push("--dangerous".into());
5855            }
5856            arguments.extend(["resume".into(), session_id.into()]);
5857            "supercode"
5858        }
5859        other => {
5860            return Err(ServiceError::InvalidParams(format!(
5861                "no structured resume launch is registered for harness `{other}`"
5862            )))
5863        }
5864    };
5865    Ok(StructuredLaunch {
5866        cwd: cwd.to_path_buf(),
5867        env: if program == "grok" {
5868            crate::support::grok_home_env()
5869        } else {
5870            BTreeMap::new()
5871        },
5872        program: program.into(),
5873        arguments,
5874    })
5875}
5876
5877/// Stage the resolved gateway credential in a private (0600) file so the
5878/// bridge can read it via `--token-file` — the delivery the real `openclaw
5879/// acp` accepts. One stable file per endpoint (keyed by an address digest,
5880/// no secret material in the name), overwritten on every connect so files
5881/// never accumulate and a rotated token never goes stale on disk.
5882fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5883    let digest = blake3::hash(address.as_bytes()).to_hex();
5884    let path = std::env::temp_dir().join(format!(
5885        "supercode-openclaw-gateway-token-{}",
5886        &digest.as_str()[..16]
5887    ));
5888    #[cfg(unix)]
5889    {
5890        use std::io::Write;
5891        use std::os::unix::fs::OpenOptionsExt;
5892        let mut file = std::fs::OpenOptions::new()
5893            .write(true)
5894            .create(true)
5895            .truncate(true)
5896            .mode(0o600)
5897            .open(&path)?;
5898        file.write_all(secret.as_bytes())?;
5899    }
5900    #[cfg(not(unix))]
5901    std::fs::write(&path, secret)?;
5902    Ok(path)
5903}
5904
5905/// Open a connect-mode descriptor: resolve the endpoint address and
5906/// credential from the harness's own config file and build the backend that
5907/// joins the already-running endpoint. Fails closed with a specific
5908/// diagnostic when the config cannot be resolved or the declared protocol has
5909/// no connect-capable client yet.
5910fn open_connect_descriptor(
5911    descriptor: &crate::HarnessSupportDescriptor,
5912    home: &Path,
5913) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5914    let Some(connect) = &descriptor.runtime.connect_launch else {
5915        return Err(ServiceError::InvalidParams(format!(
5916            "harness `{}` has no registered connect-mode launch",
5917            descriptor.id.as_str()
5918        )));
5919    };
5920    let resolved = connect
5921        .resolve(home)
5922        .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5923    match (descriptor.id.as_str(), connect.protocol.as_str()) {
5924        (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5925            let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5926            if let Some(token) = resolved.auth {
5927                backend = backend.with_bearer(token);
5928            }
5929            Ok(Box::new(backend))
5930        }
5931        (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5932            // OpenClaw's own `openclaw acp` binary is the gateway client: a
5933            // stdio ACP bridge that joins the RUNNING gateway at the resolved
5934            // endpoint. Blind-walk finding 2026-08-31: the real bridge does
5935            // NOT honor OPENCLAW_GATEWAY_TOKEN from the environment — the
5936            // credential must arrive via `--token-file` (never bare `--token`
5937            // on argv, where process listings could read it). The env var is
5938            // still set for older bridges that did read it. Requires openclaw
5939            // >= 2026.7: the 2026.2 bridge drops its gateway socket
5940            // mid-prompt and advertises no session resume (executed finding,
5941            // docs/interop/research/openclaw-acp-dialect-2026-08-30.json).
5942            let mut env = BTreeMap::new();
5943            let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5944            if let Some(token) = resolved.auth {
5945                let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5946                    .map_err(|error| {
5947                        ServiceError::UnsupportedAction(format!(
5948                            "could not stage the gateway credential for the bridge: {error}"
5949                        ))
5950                    })?;
5951                arguments.push("--token-file".into());
5952                arguments.push(token_path.to_string_lossy().into_owned());
5953                env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5954            }
5955            // The bridge program comes from the descriptor's own default
5956            // launch (the compiled registry pins `openclaw`), so tests can
5957            // substitute an absolute mock-bridge path without touching
5958            // process-global state.
5959            let program = descriptor
5960                .runtime
5961                .default_launch
5962                .as_ref()
5963                .map(|launch| launch.program.clone())
5964                .unwrap_or_else(|| "openclaw".into());
5965            let launch = RuntimeLaunch {
5966                program,
5967                arguments,
5968                env,
5969            };
5970            Ok(Box::new(
5971                crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5972                    .with_resume_support(descriptor.runtime.capabilities.resume_session),
5973            ))
5974        }
5975        _ => Err(ServiceError::UnsupportedAction(format!(
5976            "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5977            descriptor.id.as_str(),
5978            connect.protocol
5979        ))),
5980    }
5981}
5982
5983/// The registry's connect-mode launch for this harness, honored only when the
5984/// caller supplied neither an explicit launch nor a base URL.
5985fn registry_connect_descriptor(
5986    params: &RuntimeBackendParams,
5987) -> Option<crate::HarnessSupportDescriptor> {
5988    if params.launch.is_some() || params.base_url.is_some() {
5989        return None;
5990    }
5991    harness_support_registry()
5992        .harnesses
5993        .into_iter()
5994        .find(|descriptor| descriptor.id == params.harness)
5995        .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5996}
5997
5998fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5999    supercode_interchange::user_home()
6000        .map(std::path::PathBuf::into_os_string)
6001        .map(PathBuf::from)
6002        .ok_or_else(|| {
6003            ServiceError::UnsupportedAction(
6004                "connect-mode launches need HOME to locate the harness config".into(),
6005            )
6006        })
6007}
6008
6009/// The doors that open a runtime: each spawns or joins a program and waits on
6010/// that program's protocol handshake before it can answer.
6011pub const RUNTIME_OPEN_METHODS: &[&str] = &[
6012    "harness.v1.runtimes.start",
6013    "harness.v1.runtimes.resume",
6014    "harness.v1.runtimes.attach",
6015    "harness.v1.runtimes.attach_existing",
6016];
6017
6018/// How long a runtime gets to finish opening before its caller is answered an
6019/// error instead. A program that never speaks the protocol at all — the wrong
6020/// binary, a shim that prints usage and waits — never answers the handshake,
6021/// so the wait is unbounded without this.
6022pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
6023
6024/// How long a control call on an ALREADY-open runtime — send input, interrupt,
6025/// steer, respond, close — gets before its caller is answered an error
6026/// instead. A live runtime answers these in milliseconds; a wedged one never
6027/// answers at all, and `close` is exactly what a caller reaches for when it
6028/// suspects that.
6029pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
6030
6031/// The doors whose work happens entirely OUTSIDE this service's state once
6032/// its state has been read: probing harnesses, relaying a message into a
6033/// live session, and performing a conversation verb through a harness's own
6034/// CLI / HTTP / store door. Every one of them waits on a child process or a
6035/// network peer. See [`HarnessSessionService::detach`].
6036pub const DETACHED_METHODS: &[&str] = &[
6037    "harness.v1.harnesses.list",
6038    "harness.v1.sessions.load",
6039    "harness.v1.harnesses.probe",
6040    "harness.v1.sessions.message",
6041    "harness.v1.sessions.new",
6042    "harness.v1.sessions.reset",
6043    "harness.v1.sessions.archive",
6044    "harness.v1.sessions.delete",
6045];
6046
6047/// How long a request moved off a transport's loop gets before its caller is
6048/// answered an error instead. Each of these already bounds its own inner
6049/// waits (a probe's handshake, the relay's send); this is the backstop for
6050/// the ones that do not — a harness CLI that never exits — so no caller waits
6051/// forever on a detached task no one is watching.
6052pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
6053
6054/// How long `sessions.discover` gets before its caller is answered an error
6055/// instead. Discovery reads each harness's own store, and a store on a cold
6056/// or unavailable mount answers at the filesystem's pace rather than its own.
6057///
6058/// Deliberately shorter than the clients' own request deadline (30s): the
6059/// server's answer names the store that did not answer, and it is only read
6060/// if it lands before the client stops listening.
6061pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
6062
6063/// Bound one control call on an open runtime by [`RUNTIME_CONTROL_DEADLINE`],
6064/// naming the method and the bound when it blows.
6065async fn within_control_deadline<F: std::future::Future>(
6066    method: &str,
6067    call: F,
6068) -> std::result::Result<F::Output, ServiceError> {
6069    tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
6070        .await
6071        .map_err(|_| {
6072            ServiceError::Operation(format!(
6073                "`{method}` gave up after {}s: the runtime did not answer",
6074                RUNTIME_CONTROL_DEADLINE.as_secs()
6075            ))
6076        })
6077}
6078
6079/// One [`RUNTIME_OPEN_METHODS`] request, parsed but not yet started. See
6080/// [`HarnessSessionService::runtime_open`] for why it exists apart from
6081/// [`HarnessSessionService::handle_async`].
6082pub struct RuntimeOpen {
6083    id: Value,
6084    method: String,
6085    params: Value,
6086}
6087
6088impl RuntimeOpen {
6089    /// Do the waiting: spawn or join the program and complete its handshake,
6090    /// bounded by [`RUNTIME_OPEN_DEADLINE`]. Touches no service state, so this
6091    /// runs on any task.
6092    pub async fn open(self) -> OpenedRuntime {
6093        let Self { id, method, params } = self;
6094        let outcome = open_runtime(&method, params).await;
6095        OpenedRuntime { id, outcome }
6096    }
6097}
6098
6099/// The result of [`RuntimeOpen::open`], ready for
6100/// [`HarnessSessionService::finish_runtime_open`].
6101pub struct OpenedRuntime {
6102    id: Value,
6103    outcome: std::result::Result<OpenRuntime, ServiceError>,
6104}
6105
6106/// One detached request: the half that reads this service's state already
6107/// done, and the half that waits not yet started. See
6108/// [`HarnessSessionService::detach`] and
6109/// [`HarnessSessionService::detach_runtime`].
6110pub struct DetachedCall {
6111    id: Value,
6112    method: String,
6113    work: std::result::Result<Work, ServiceError>,
6114}
6115
6116impl DetachedCall {
6117    /// Do the waiting and answer. Runs on any task: whatever this call needed
6118    /// from the service was taken before it left.
6119    pub async fn run(self) -> DetachedAnswer {
6120        let Self { id, method, work } = self;
6121        match work {
6122            // A call holding a runtime is already bounded by
6123            // RUNTIME_CONTROL_DEADLINE, and its future OWNS that connection:
6124            // a second timeout around it would drop the connection mid-call
6125            // and take down a runtime its caller still has.
6126            Ok(Work::Runtime(work)) => {
6127                let (result, returned) = work.run().await;
6128                DetachedAnswer {
6129                    response: service_response(id, result),
6130                    returned,
6131                }
6132            }
6133            Ok(Work::Free(work)) => {
6134                let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
6135                    Ok(result) => result,
6136                    Err(_) => Err(ServiceError::Operation(format!(
6137                        "`{method}` gave up after {}s: the harness it waits on did not answer",
6138                        DETACHED_CALL_DEADLINE.as_secs()
6139                    ))),
6140                };
6141                DetachedAnswer {
6142                    response: service_response(id, result),
6143                    returned: None,
6144                }
6145            }
6146            Err(error) => DetachedAnswer {
6147                response: service_response(id, Err(error)),
6148                returned: None,
6149            },
6150        }
6151    }
6152}
6153
6154/// One detached call's complete answer, plus whatever it must hand back to
6155/// the service before that answer is written. See
6156/// [`HarnessSessionService::finish_detached`].
6157pub struct DetachedAnswer {
6158    response: Value,
6159    returned: Option<ReturnedRuntime>,
6160}
6161
6162impl DetachedAnswer {
6163    /// The caller's JSON-RPC response, for a transport that owns no service
6164    /// to give a borrowed connection back to.
6165    pub fn into_response(self) -> Value {
6166        self.response
6167    }
6168}
6169
6170/// A connection lent to a detached call, on its way back to the service that
6171/// owns it.
6172pub struct ReturnedRuntime {
6173    connection: String,
6174    runtime: Box<dyn RuntimeConnection>,
6175}
6176
6177/// The waiting half of one detached request: with nothing of the service's
6178/// in hand, or holding a connection the service lent out for the call.
6179enum Work {
6180    Free(DetachedWork),
6181    Runtime(RuntimeWork),
6182}
6183
6184/// The waiting half of one detached request that holds nothing of the
6185/// service's.
6186enum DetachedWork {
6187    /// Probe the selected harnesses: find their executables, ask each its
6188    /// version, and at `probe: handshake` start each one and complete its
6189    /// protocol handshake.
6190    Inventory(InventoryWork),
6191    /// Relay one message into a live session.
6192    Message(MessageSessionParams),
6193    /// Read one session from its own store, on the blocking pool.
6194    Load(LoadSessionParams),
6195    /// Perform one conversation verb through the harness's own CLI, HTTP API,
6196    /// daemon socket, or supercode's own store.
6197    SessionMutation {
6198        verb: crate::SessionVerb,
6199        mutation: crate::SessionMutation,
6200    },
6201}
6202
6203impl DetachedWork {
6204    async fn run(self) -> std::result::Result<Value, ServiceError> {
6205        match self {
6206            Self::Inventory(work) => run_inventory(work).await,
6207            Self::Message(params) => Ok(message_live_session(&params).await),
6208            Self::Load(params) => tokio::task::spawn_blocking(move || load_session_door(params))
6209                .await
6210                .map_err(|error| {
6211                    ServiceError::Operation(format!("the session read stopped: {error}"))
6212                })?,
6213            Self::SessionMutation { verb, mutation } => {
6214                let outcome = run_session_mutation(verb, &mutation).await?;
6215                serde_json::to_value(outcome)
6216                    .map_err(|error| ServiceError::Operation(error.to_string()))
6217            }
6218        }
6219    }
6220}
6221
6222/// One detached call that holds a runtime connection for its whole run.
6223enum RuntimeWork {
6224    /// Tear down a runtime the service has already surrendered.
6225    Close {
6226        runtime: Box<dyn RuntimeConnection>,
6227        process_group: Option<u32>,
6228    },
6229    /// Type one live slash command through a borrowed connection, then give
6230    /// the connection back.
6231    LiveCommand {
6232        connection: String,
6233        runtime: Box<dyn RuntimeConnection>,
6234        verb: crate::SessionVerb,
6235        mutation: crate::SessionMutation,
6236        command: &'static str,
6237        session: String,
6238    },
6239}
6240
6241/// What one [`RuntimeWork`] answers with: the caller's result, and the
6242/// connection to give back when the call only borrowed one.
6243type RuntimeWorkAnswer = (
6244    std::result::Result<Value, ServiceError>,
6245    Option<ReturnedRuntime>,
6246);
6247
6248impl RuntimeWork {
6249    async fn run(self) -> RuntimeWorkAnswer {
6250        match self {
6251            Self::Close {
6252                runtime,
6253                process_group,
6254            } => (close_runtime(runtime, process_group).await, None),
6255            Self::LiveCommand {
6256                connection,
6257                mut runtime,
6258                verb,
6259                mutation,
6260                command,
6261                session,
6262            } => {
6263                let result =
6264                    type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
6265                (
6266                    result,
6267                    Some(ReturnedRuntime {
6268                        connection,
6269                        runtime,
6270                    }),
6271                )
6272            }
6273        }
6274    }
6275}
6276
6277/// Tear down a runtime already out of the service, within
6278/// [`RUNTIME_CONTROL_DEADLINE`].
6279async fn close_runtime(
6280    mut runtime: Box<dyn RuntimeConnection>,
6281    process_group: Option<u32>,
6282) -> std::result::Result<Value, ServiceError> {
6283    match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
6284        Ok(result) => {
6285            result.map_err(operation)?;
6286            Ok(json!({"closed": true}))
6287        }
6288        Err(deadline) => {
6289            // Dropping the handle is not enough: the process that stopped
6290            // answering is held by a task parked on it, so nothing here runs
6291            // its Drop. Signal the group the graceful path would have
6292            // signalled, then say so.
6293            let killed = kill_runtime_process_group(process_group);
6294            drop(runtime);
6295            Ok(json!({
6296                "closed": true,
6297                "killed": killed,
6298                "detail": error_message(deadline),
6299            }))
6300        }
6301    }
6302}
6303
6304/// The conversation a live `sessions.new` / `sessions.reset` acts on: the one
6305/// the request named, or the runtime's own session.
6306fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
6307    mutation
6308        .session
6309        .clone()
6310        .filter(|value| !value.trim().is_empty())
6311        .unwrap_or_else(|| runtime.handle().runtime_id.clone())
6312}
6313
6314/// Type one harness slash command into a live session through the very same
6315/// `send_input` path a human's message takes, within
6316/// [`RUNTIME_CONTROL_DEADLINE`].
6317async fn type_live_command(
6318    runtime: &mut dyn RuntimeConnection,
6319    verb: crate::SessionVerb,
6320    mutation: &crate::SessionMutation,
6321    command: &str,
6322    session: String,
6323) -> std::result::Result<Value, ServiceError> {
6324    within_control_deadline(
6325        &format!("sessions.{}", verb.as_str()),
6326        runtime.send_input(RuntimeInput {
6327            text: command.to_string(),
6328            image_urls: Vec::new(),
6329        }),
6330    )
6331    .await?
6332    .map_err(operation)?;
6333    let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
6334        .map_err(session_control_error)?;
6335    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
6336}
6337
6338/// A runtime that is up and whose handshake completed, with what the service
6339/// needs to take ownership of it.
6340enum OpenRuntime {
6341    /// supercode spawned this process, so it also hosts it: a frontend server,
6342    /// a live-runtime registration and a terminal launch of its own.
6343    Hosted {
6344        runtime: Box<dyn RuntimeConnection>,
6345        capabilities: crate::RuntimeCapabilities,
6346        workspace: PathBuf,
6347        /// A new session (`start`): it has no transcript yet, so there is no history to look for.
6348        fresh: bool,
6349    },
6350    /// `attach_existing` joined a process supercode does not own. It is
6351    /// registered as a bare connection and hosts nothing.
6352    Joined { runtime: Box<dyn RuntimeConnection> },
6353}
6354
6355/// Open the runtime one [`RUNTIME_OPEN_METHODS`] request asks for, within
6356/// [`RUNTIME_OPEN_DEADLINE`]. The error a blown deadline answers names the
6357/// method and the bound, so a caller reads why it was cut loose instead of
6358/// waiting on a handshake that is never coming.
6359async fn open_runtime(
6360    method: &str,
6361    params: Value,
6362) -> std::result::Result<OpenRuntime, ServiceError> {
6363    match tokio::time::timeout(
6364        RUNTIME_OPEN_DEADLINE,
6365        open_runtime_unbounded(method, params),
6366    )
6367    .await
6368    {
6369        Ok(result) => result,
6370        Err(_) => Err(ServiceError::Operation(format!(
6371            "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
6372            RUNTIME_OPEN_DEADLINE.as_secs()
6373        ))),
6374    }
6375}
6376
6377async fn open_runtime_unbounded(
6378    method: &str,
6379    params: Value,
6380) -> std::result::Result<OpenRuntime, ServiceError> {
6381    match method {
6382        "harness.v1.runtimes.start" => {
6383            let params = decode::<RuntimeStartParams>(params)?;
6384            let backend = runtime_backend(&params.backend)?;
6385            let capabilities = backend.capabilities();
6386            let workspace = params.cwd.clone();
6387            let runtime = backend
6388                .start(RuntimeStartRequest {
6389                    cwd: params.cwd,
6390                    launch: runtime_launch(&params.backend),
6391                    mcp_servers: params.mcp_servers,
6392                    approval_policy: params.approval_policy,
6393                })
6394                .await
6395                .map_err(operation)?;
6396            Ok(OpenRuntime::Hosted {
6397                runtime,
6398                capabilities,
6399                workspace,
6400                fresh: true,
6401            })
6402        }
6403        "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
6404            let params = decode::<RuntimeAttachParams>(params)?;
6405            let backend = runtime_backend(&params.backend)?;
6406            let capabilities = backend.capabilities();
6407            let workspace = params
6408                .cwd
6409                .clone()
6410                .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
6411            let runtime = backend
6412                .attach(RuntimeAttachRequest {
6413                    runtime_id: params.runtime_id,
6414                    cwd: params.cwd,
6415                    launch: runtime_launch(&params.backend),
6416                    mcp_servers: params.mcp_servers,
6417                    approval_policy: params.approval_policy,
6418                })
6419                .await
6420                .map_err(operation)?;
6421            Ok(OpenRuntime::Hosted {
6422                runtime,
6423                capabilities,
6424                workspace,
6425                fresh: false,
6426            })
6427        }
6428        "harness.v1.runtimes.attach_existing" => {
6429            let params = decode::<RuntimeAttachParams>(params)?;
6430            let backend: Box<dyn RuntimeBackend> = match params
6431                .backend
6432                .base_url
6433                .as_deref()
6434                .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
6435            {
6436                Some(endpoint) => {
6437                    #[cfg(not(feature = "adapter-api"))]
6438                    {
6439                        let _ = endpoint;
6440                        return Err(ServiceError::UnsupportedAction(
6441                            "live HTTP attachment adapter is not compiled".into(),
6442                        ));
6443                    }
6444                    #[cfg(feature = "adapter-api")]
6445                    {
6446                        let workspace = params.cwd.clone().ok_or_else(|| {
6447                            ServiceError::InvalidParams(
6448                                "Volter Harness live attach requires the project cwd".into(),
6449                            )
6450                        })?;
6451                        let source = LiveRuntimeSource {
6452                            harness: params.backend.harness.as_str().to_string(),
6453                            session_id: params.runtime_id.clone(),
6454                            workspace,
6455                        };
6456                        let receipt = resolve_live_runtime(&endpoint, &source)
6457                            .map_err(|error| ServiceError::Operation(error.to_string()))?;
6458                        Box::new(SupercodeHttpRuntimeBackend::new(receipt))
6459                    }
6460                }
6461                None => runtime_backend(&params.backend)?,
6462            };
6463            let capabilities = backend.capabilities();
6464            if !capabilities.attach_existing_process {
6465                return Err(ServiceError::Operation(format!(
6466                    "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
6467                    backend.harness().as_str()
6468                )));
6469            }
6470            let runtime = backend
6471                .attach_existing(RuntimeAttachRequest {
6472                    runtime_id: params.runtime_id,
6473                    cwd: params.cwd,
6474                    launch: runtime_launch(&params.backend),
6475                    mcp_servers: params.mcp_servers,
6476                    approval_policy: params.approval_policy,
6477                })
6478                .await
6479                .map_err(operation)?;
6480            Ok(OpenRuntime::Joined { runtime })
6481        }
6482        _ => Err(ServiceError::MethodNotFound),
6483    }
6484}
6485
6486/// Wrap one service outcome in its JSON-RPC 2.0 envelope.
6487fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
6488    match result {
6489        Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
6490        Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
6491        Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
6492        Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
6493        Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
6494        Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
6495    }
6496}
6497
6498fn runtime_backend(
6499    params: &RuntimeBackendParams,
6500) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
6501    if let Some(descriptor) = registry_connect_descriptor(params) {
6502        return open_connect_descriptor(&descriptor, &service_home()?);
6503    }
6504    if params.protocol.as_deref() == Some("acp") {
6505        let launch = params
6506            .launch
6507            .clone()
6508            .or_else(|| {
6509                harness_support_registry()
6510                    .harnesses
6511                    .into_iter()
6512                    .find(|harness| harness.id == params.harness)
6513                    .filter(|harness| {
6514                        harness.runtime.implementation == ImplementationKind::GenericProtocol
6515                            && harness.runtime.protocol.starts_with("acp")
6516                    })
6517                    .and_then(|harness| harness.runtime.default_launch)
6518            })
6519            .ok_or_else(|| {
6520                ServiceError::InvalidParams(
6521                    "an ACP runtime requires `launch` unless the harness has a registered default"
6522                        .into(),
6523                )
6524            })?;
6525        let resume_session = harness_support_registry()
6526            .harnesses
6527            .into_iter()
6528            .find(|harness| harness.id == params.harness)
6529            .is_some_and(|harness| harness.runtime.capabilities.resume_session);
6530        return Ok(Box::new(
6531            AcpRuntimeBackend::new(params.harness.clone(), launch)
6532                .with_resume_support(resume_session),
6533        ));
6534    }
6535    let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
6536        HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
6537        HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
6538        HarnessId::PI => Box::new(PiRuntimeBackend::new()),
6539        HarnessId::OPENCODE => match &params.base_url {
6540            Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
6541            None => Box::new(OpenCodeRuntimeBackend::new()),
6542        },
6543        harness => {
6544            let descriptor = harness_support_registry()
6545                .harnesses
6546                .into_iter()
6547                .find(|descriptor| descriptor.id.as_str() == harness)
6548                .filter(|descriptor| {
6549                    descriptor.runtime.implementation == ImplementationKind::GenericProtocol
6550                        && descriptor.runtime.protocol.starts_with("acp")
6551                });
6552            let Some(descriptor) = descriptor else {
6553                return Err(ServiceError::InvalidParams(format!(
6554                    "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
6555                )));
6556            };
6557            let resume = descriptor.runtime.capabilities.resume_session;
6558            Box::new(
6559                AcpRuntimeBackend::new(
6560                    descriptor.id,
6561                    descriptor
6562                        .runtime
6563                        .default_launch
6564                        .expect("generic ACP registry entry includes its launch"),
6565                )
6566                .with_resume_support(resume),
6567            )
6568        }
6569    };
6570    Ok(backend)
6571}
6572
6573fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
6574    if let Some(launch) = &params.launch {
6575        return Some(launch.clone());
6576    }
6577    if !matches!(params.policy, RuntimePolicy::Yolo) {
6578        return None;
6579    }
6580    let launch = match params.harness.as_str() {
6581        HarnessId::GROK => RuntimeLaunch {
6582            program: "grok".into(),
6583            arguments: {
6584                let mut arguments: Vec<String> = Vec::new();
6585                if crate::support::self_sandbox_supported() {
6586                    arguments.extend(["--sandbox".into(), "workspace".into()]);
6587                }
6588                arguments.extend([
6589                    "--always-approve".into(),
6590                    "agent".into(),
6591                    "--no-leader".into(),
6592                    "stdio".into(),
6593                ]);
6594                arguments
6595            },
6596            env: crate::support::grok_env(),
6597        },
6598        HarnessId::CODEX => RuntimeLaunch {
6599            program: "codex".into(),
6600            arguments: vec![
6601                "--dangerously-bypass-approvals-and-sandbox".into(),
6602                "--dangerously-bypass-hook-trust".into(),
6603                "app-server".into(),
6604            ],
6605            env: BTreeMap::new(),
6606        },
6607        HarnessId::CLAUDE_CODE => RuntimeLaunch {
6608            program: "claude".into(),
6609            arguments: vec![
6610                "--dangerously-skip-permissions".into(),
6611                "--print".into(),
6612                "--input-format".into(),
6613                "stream-json".into(),
6614                "--output-format".into(),
6615                "stream-json".into(),
6616                "--verbose".into(),
6617            ],
6618            env: BTreeMap::new(),
6619        },
6620        HarnessId::PI => RuntimeLaunch {
6621            program: "pi".into(),
6622            arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
6623            env: BTreeMap::new(),
6624        },
6625        HarnessId::OPENCODE => RuntimeLaunch {
6626            program: "opencode".into(),
6627            arguments: vec!["serve".into()],
6628            env: BTreeMap::new(),
6629        },
6630        HarnessId::GEMINI => RuntimeLaunch {
6631            program: "gemini".into(),
6632            arguments: vec!["--acp".into(), "--yolo".into()],
6633            env: BTreeMap::new(),
6634        },
6635        HarnessId::GOOSE => RuntimeLaunch {
6636            program: "goose".into(),
6637            arguments: vec!["acp".into()],
6638            env: BTreeMap::new(),
6639        },
6640        HarnessId::SUPERCODE => RuntimeLaunch {
6641            program: "supercode".into(),
6642            arguments: vec!["acp".into(), "--dangerous".into()],
6643            env: BTreeMap::new(),
6644        },
6645        _ => return None,
6646    };
6647    Some(launch)
6648}
6649
6650/// Disposable harness state for a no-prompt readiness probe. Merely opening
6651/// several stock CLIs writes a session header or migrates configuration, so a
6652/// handshake must never point at the user's real home. Authentication files
6653/// are copied into the private temporary home; all writes disappear with the
6654/// guard after the connection closes.
6655struct IsolatedProbeHome {
6656    launch: RuntimeLaunch,
6657    root: PathBuf,
6658}
6659
6660impl IsolatedProbeHome {
6661    fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
6662        let root = std::env::temp_dir().join(format!(
6663            "supercode-harness-probe-{harness}-{}",
6664            generated_session_id()
6665        ));
6666        std::fs::create_dir_all(&root)?;
6667        set_private_dir_permissions(&root)?;
6668
6669        if let Some(source_home) = supercode_interchange::user_home()
6670            .map(std::path::PathBuf::into_os_string)
6671            .map(PathBuf::from)
6672        {
6673            for relative in probe_auth_files(harness) {
6674                copy_probe_file(&source_home, &root, relative)?;
6675            }
6676        }
6677        // supercode reads its own config home ($SUPERCODE_HOME, else
6678        // $XDG_CONFIG_HOME/supercode, else ~/.config/supercode), not a fixed
6679        // place under HOME: a login kept under XDG_CONFIG_HOME probed as
6680        // "no API key found" while `supercode run` answered.
6681        if harness == HarnessId::SUPERCODE {
6682            let config_home = crate::agent::global_instructions_dir();
6683            for file in ["config.toml", "credentials.toml"] {
6684                copy_probe_path(
6685                    &config_home.join(file),
6686                    &root.join(".config/supercode").join(file),
6687                )?;
6688            }
6689        }
6690        configure_isolated_probe_auth(harness, &root)?;
6691
6692        let root_text = root.to_string_lossy().into_owned();
6693        for (key, value) in [
6694            ("HOME", root_text.clone()),
6695            (
6696                "XDG_CACHE_HOME",
6697                root.join(".cache").to_string_lossy().into_owned(),
6698            ),
6699            (
6700                "XDG_CONFIG_HOME",
6701                root.join(".config").to_string_lossy().into_owned(),
6702            ),
6703            (
6704                "XDG_DATA_HOME",
6705                root.join(".local/share").to_string_lossy().into_owned(),
6706            ),
6707        ] {
6708            launch.env.insert(key.into(), value);
6709        }
6710        let scoped = match harness {
6711            HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6712            HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6713            HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6714            HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6715            HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6716            HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6717            _ => None,
6718        };
6719        if let Some((key, value)) = scoped {
6720            launch
6721                .env
6722                .insert(key.into(), value.to_string_lossy().into_owned());
6723        }
6724        Ok(Self { launch, root })
6725    }
6726
6727    fn cleanup(&self) -> std::io::Result<()> {
6728        match std::fs::remove_dir_all(&self.root) {
6729            Ok(()) => Ok(()),
6730            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6731            Err(error) => Err(error),
6732        }
6733    }
6734}
6735
6736impl Drop for IsolatedProbeHome {
6737    fn drop(&mut self) {
6738        let _ = self.cleanup();
6739    }
6740}
6741
6742fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6743    match harness {
6744        HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6745        // The gateway endpoint + token live in openclaw's own config; without
6746        // it the isolated probe dials the default endpoint unauthenticated
6747        // (PARITY-24 finding 2026-08-31).
6748        HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6749        HarnessId::CODEX => &[".codex/auth.json"],
6750        HarnessId::GEMINI => &[
6751            ".gemini/google_accounts.json",
6752            ".gemini/oauth_creds.json",
6753            ".gemini/settings.json",
6754        ],
6755        HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6756        HarnessId::OPENCODE => &[
6757            ".config/opencode/auth.json",
6758            ".local/share/opencode/auth.json",
6759        ],
6760        HarnessId::PI => &[".pi/agent/auth.json"],
6761        // Hermes keeps its provider selection in config.yaml, its OAuth
6762        // credential pool in auth.json, and API keys in .env; without them
6763        // the isolated probe sees "No LLM provider configured" for a
6764        // hermes that answers fine from the user's real home.
6765        HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6766        _ => &[],
6767    }
6768}
6769
6770fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6771    copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6772}
6773
6774fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6775    if !source.is_file() {
6776        return Ok(());
6777    }
6778    if let Some(parent) = destination.parent() {
6779        std::fs::create_dir_all(parent)?;
6780        set_private_dir_permissions(parent)?;
6781    }
6782    std::fs::copy(source, destination)?;
6783    set_private_file_permissions(destination)
6784}
6785
6786fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6787    if harness != HarnessId::GEMINI {
6788        return Ok(());
6789    }
6790    let oauth = probe_home.join(".gemini/oauth_creds.json");
6791    if !oauth.is_file() {
6792        return Ok(());
6793    }
6794    let settings_path = probe_home.join(".gemini/settings.json");
6795    let mut settings = std::fs::read_to_string(&settings_path)
6796        .ok()
6797        .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6798        .unwrap_or_else(|| json!({}));
6799    settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6800    std::fs::write(
6801        &settings_path,
6802        serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6803    )?;
6804    set_private_file_permissions(&settings_path)
6805}
6806
6807#[cfg(unix)]
6808fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6809    use std::os::unix::fs::PermissionsExt;
6810    std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6811}
6812
6813#[cfg(not(unix))]
6814fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6815    Ok(())
6816}
6817
6818#[cfg(unix)]
6819fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6820    use std::os::unix::fs::PermissionsExt;
6821    std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6822}
6823
6824#[cfg(not(unix))]
6825fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6826    Ok(())
6827}
6828
6829fn find_executable(program: &str) -> Option<PathBuf> {
6830    let candidate = PathBuf::from(program);
6831    if candidate.components().count() > 1 {
6832        return candidate.is_file().then_some(candidate);
6833    }
6834    let path = std::env::var_os("PATH")?;
6835    for directory in std::env::split_paths(&path) {
6836        let candidate = directory.join(program);
6837        if candidate.is_file() {
6838            return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6839        }
6840        #[cfg(windows)]
6841        {
6842            for extension in ["exe", "cmd", "bat"] {
6843                let candidate = directory.join(format!("{program}.{extension}"));
6844                if candidate.is_file() {
6845                    return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6846                }
6847            }
6848        }
6849    }
6850    None
6851}
6852
6853async fn executable_version(executable: &Path) -> Option<String> {
6854    let mut command = tokio::process::Command::new(executable);
6855    command
6856        .arg("--version")
6857        .stdin(std::process::Stdio::null())
6858        .stdout(std::process::Stdio::piped())
6859        .stderr(std::process::Stdio::piped())
6860        .kill_on_drop(true);
6861    let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6862        .await
6863        .ok()?
6864        .ok()?;
6865    let stdout = String::from_utf8_lossy(&output.stdout);
6866    let stderr = String::from_utf8_lossy(&output.stderr);
6867    stdout
6868        .lines()
6869        .chain(stderr.lines())
6870        .map(str::trim)
6871        .find(|line| !line.is_empty())
6872        .map(|line| truncate_text(line, 200))
6873}
6874
6875pub(crate) fn auth_evidence(harness: &str) -> bool {
6876    let env_names: &[&str] = match harness {
6877        HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6878        HarnessId::CODEX => &["OPENAI_API_KEY"],
6879        HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6880        HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6881        HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6882        HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6883        HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6884        _ => &[],
6885    };
6886    if env_names
6887        .iter()
6888        .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6889    {
6890        return true;
6891    }
6892    let Some(home) = supercode_interchange::user_home()
6893        .map(std::path::PathBuf::into_os_string)
6894        .map(PathBuf::from)
6895    else {
6896        return false;
6897    };
6898    let files: Vec<PathBuf> = match harness {
6899        HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6900        HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6901        HarnessId::OPENCODE => vec![
6902            home.join(".local/share/opencode/auth.json"),
6903            home.join(".config/opencode/auth.json"),
6904        ],
6905        HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6906        HarnessId::GROK => vec![home.join(".grok/auth.json")],
6907        HarnessId::GEMINI => vec![
6908            home.join(".gemini/oauth_creds.json"),
6909            home.join(".gemini/google_accounts.json"),
6910        ],
6911        HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6912        HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6913        _ => Vec::new(),
6914    };
6915    if files.into_iter().any(|path| {
6916        std::fs::metadata(path)
6917            .map(|metadata| metadata.is_file() && metadata.len() > 2)
6918            .unwrap_or(false)
6919    }) {
6920        return true;
6921    }
6922    // macOS keeps Claude Code's OAuth login in the Keychain, so
6923    // `.claude/.credentials.json` never exists there and the file probe above
6924    // reports a signed-in install as unauthenticated forever. A completed
6925    // login also writes an `oauthAccount` record into `~/.claude.json` on
6926    // every platform — file-based, prompt-free evidence (querying the
6927    // Keychain itself from an unsigned daemon can raise a UI prompt).
6928    if harness == HarnessId::CLAUDE_CODE {
6929        return std::fs::read_to_string(home.join(".claude.json"))
6930            .map(|text| text.contains("\"oauthAccount\""))
6931            .unwrap_or(false);
6932    }
6933    false
6934}
6935
6936fn looks_like_auth_error(message: &str) -> bool {
6937    let message = message.to_ascii_lowercase();
6938    [
6939        "auth",
6940        "login",
6941        "sign in",
6942        "sign-in",
6943        "credential",
6944        "unauthorized",
6945        "forbidden",
6946        "token",
6947    ]
6948    .iter()
6949    .any(|needle| message.contains(needle))
6950}
6951
6952fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6953    crate::RuntimeCapabilities {
6954        start_session: false,
6955        resume_session: false,
6956        attach_existing_process: false,
6957        send_input: false,
6958        stream_events: false,
6959        interrupt: false,
6960        steer: false,
6961        respond_to_requests: false,
6962    }
6963}
6964
6965fn truncate_text(text: &str, max_chars: usize) -> String {
6966    let mut chars = text.chars();
6967    let truncated = chars.by_ref().take(max_chars).collect::<String>();
6968    if chars.next().is_some() {
6969        format!("{truncated}…")
6970    } else {
6971        truncated
6972    }
6973}
6974
6975/// The process group a runtime's own handle names, when it names one.
6976///
6977/// Every adapter that spawns a local process spawns it as its own group
6978/// leader (`Command::process_group(0)`), so the endpoint's pid IS the group
6979/// id. A runtime reached over HTTP, or one supercode joined rather than
6980/// spawned, names no group here and is left alone.
6981fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6982    match &handle.endpoint {
6983        crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6984        crate::RuntimeEndpoint::Http { .. } => None,
6985    }
6986}
6987
6988/// SIGKILL a wedged runtime's whole process group, reporting whether there
6989/// was one to signal. This is the same group teardown a graceful `close`
6990/// performs; it runs here only when the graceful path blew its deadline,
6991/// because the task parked on the unanswered call still owns the process
6992/// handle and so no `Drop` of ours can reach it.
6993fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6994    match process_group {
6995        #[cfg(unix)]
6996        Some(pid) => {
6997            crate::lsp::kill_process_group(pid);
6998            true
6999        }
7000        #[cfg(not(unix))]
7001        Some(_) => false,
7002        None => false,
7003    }
7004}
7005
7006fn error_message(error: ServiceError) -> String {
7007    match error {
7008        ServiceError::InvalidParams(message)
7009        | ServiceError::Operation(message)
7010        | ServiceError::UnsupportedAction(message) => message,
7011        ServiceError::MethodNotFound => "runtime adapter is not available".into(),
7012        ServiceError::Sdk(error) => error.to_string(),
7013    }
7014}
7015
7016#[derive(Debug)]
7017enum ServiceError {
7018    InvalidParams(String),
7019    MethodNotFound,
7020    UnsupportedAction(String),
7021    Operation(String),
7022    Sdk(SdkError),
7023}
7024
7025fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
7026    match error {
7027        ServiceError::InvalidParams(message) => {
7028            SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
7029        }
7030        ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
7031            SdkError::unsupported(operation)
7032        }
7033        ServiceError::Operation(message) => {
7034            let code = if message.contains("already in progress") {
7035                SdkErrorCode::Busy
7036            } else if message.contains("not supported by this runtime") {
7037                SdkErrorCode::UnsupportedAction
7038            } else if message.contains("unknown runtime connection") {
7039                SdkErrorCode::NotFound
7040            } else {
7041                SdkErrorCode::Execution
7042            };
7043            SdkError::new(code, operation, message)
7044        }
7045        ServiceError::Sdk(error) => error,
7046    }
7047}
7048
7049fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
7050    let error_code = error.code();
7051    let code = match error_code {
7052        SdkErrorCode::Unauthenticated => -32030,
7053        SdkErrorCode::Unauthorized => -32031,
7054        SdkErrorCode::ControllerRequired => -32032,
7055        SdkErrorCode::LeaseExpired => -32033,
7056        SdkErrorCode::InvalidArgument => -32602,
7057        SdkErrorCode::NotFound => -32004,
7058        SdkErrorCode::Busy => -32000,
7059        SdkErrorCode::UnsupportedAction => -32020,
7060        SdkErrorCode::Execution => -32002,
7061        SdkErrorCode::Transport => -32003,
7062    };
7063    json!({
7064        "jsonrpc": "2.0",
7065        "id": id,
7066        "error": {
7067            "code": code,
7068            "name": error_code,
7069            "operation": error.operation(),
7070            "message": error.to_string(),
7071        },
7072    })
7073}
7074
7075fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
7076    serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
7077}
7078
7079fn operation(error: impl Into<crate::Error>) -> ServiceError {
7080    let error = error.into();
7081    match error {
7082        crate::Error::Sdk(error) => ServiceError::Sdk(error),
7083        error => ServiceError::Operation(error.to_string()),
7084    }
7085}
7086
7087/// ORCH-12 `harness.v1.memory.show|search` params. `homes` is the same
7088/// storage-root override every read-only method accepts, so a caller can
7089/// point the read at a fixture home without touching the real ones.
7090#[derive(Debug, Clone, Deserialize, Default)]
7091#[serde(default)]
7092struct MemoryRequest {
7093    /// Harness whose store is read. Required.
7094    harness: Option<String>,
7095    /// The needle, required by `search`.
7096    query: Option<String>,
7097    /// Hermes profile, OpenClaw agent, or Claude Code project.
7098    profile: Option<String>,
7099    /// Claude Code session id selecting a project store (`show` only).
7100    session: Option<String>,
7101    /// Include each document's whole text (`show` only).
7102    full: bool,
7103    /// Treat `query` as a regular expression (`search` only).
7104    regex: bool,
7105    /// Working tree whose project store is read.
7106    cwd: Option<std::path::PathBuf>,
7107    /// Storage roots to read.
7108    homes: crate::HarnessHomes,
7109}
7110
7111/// Read the memory noun. A harness with no memory store fails with
7112/// `UnsupportedAction` (RPC `-32020`), never an empty list.
7113fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7114    let request = decode::<MemoryRequest>(params)?;
7115    let harness = request
7116        .harness
7117        .clone()
7118        .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7119    let to_service = |error: crate::memory::MemoryError| match error {
7120        crate::memory::MemoryError::UnsupportedHarness { .. }
7121        | crate::memory::MemoryError::SessionNotScoped { .. } => {
7122            ServiceError::UnsupportedAction(error.to_string())
7123        }
7124        other => ServiceError::InvalidParams(other.to_string()),
7125    };
7126    match method {
7127        "harness.v1.memory.show" => {
7128            let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
7129                harness,
7130                profile: request.profile,
7131                session: request.session,
7132                full: request.full,
7133                cwd: request.cwd,
7134                homes: request.homes,
7135            })
7136            .map_err(to_service)?;
7137            Ok(json!({
7138                "schema": crate::memory::MEMORY_SCHEMA,
7139                "documents": documents,
7140            }))
7141        }
7142        "harness.v1.memory.search" => {
7143            let query = request
7144                .query
7145                .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
7146            let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
7147                harness,
7148                query,
7149                profile: request.profile,
7150                regex: request.regex,
7151                cwd: request.cwd,
7152                homes: request.homes,
7153            })
7154            .map_err(to_service)?;
7155            Ok(json!({
7156                "schema": crate::memory::MEMORY_SCHEMA,
7157                "matches": matches,
7158            }))
7159        }
7160        _ => Err(ServiceError::MethodNotFound),
7161    }
7162}
7163
7164/// ORCH-10 `harness.v1.profiles.list|get` params. `homes` is the same
7165/// storage-root override every read-only method accepts, so a caller can
7166/// point the read at a fixture home without touching the real ones.
7167#[derive(Debug, Clone, Deserialize)]
7168#[serde(default)]
7169struct ProfilesQuery {
7170    /// Restrict the listing to one harness. `get` requires it.
7171    harness: Option<String>,
7172    /// Profile name, required by `get`.
7173    name: Option<String>,
7174    /// Storage roots to read.
7175    homes: crate::HarnessHomes,
7176}
7177
7178impl Default for ProfilesQuery {
7179    fn default() -> Self {
7180        Self {
7181            harness: None,
7182            name: None,
7183            homes: crate::HarnessHomes::default(),
7184        }
7185    }
7186}
7187
7188/// Read the profile noun. A harness with no profile concept fails with
7189/// `UnsupportedAction` (RPC `-32020`), never an empty list.
7190fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7191    let query = decode::<ProfilesQuery>(params)?;
7192    let to_service = |error: crate::profiles::ProfileError| match error {
7193        crate::profiles::ProfileError::UnsupportedHarness { .. } => {
7194            ServiceError::UnsupportedAction(error.to_string())
7195        }
7196        crate::profiles::ProfileError::NotFound { .. } => {
7197            ServiceError::InvalidParams(error.to_string())
7198        }
7199    };
7200    match method {
7201        "harness.v1.profiles.list" => {
7202            let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
7203                .map_err(to_service)?;
7204            Ok(json!({
7205                "schema": crate::profiles::PROFILES_SCHEMA,
7206                "profiles": profiles,
7207            }))
7208        }
7209        "harness.v1.profiles.get" => {
7210            let harness = query
7211                .harness
7212                .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7213            let name = query
7214                .name
7215                .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7216            let profile =
7217                crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
7218            Ok(json!({
7219                "schema": crate::profiles::PROFILES_SCHEMA,
7220                "profile": profile,
7221            }))
7222        }
7223        _ => Err(ServiceError::MethodNotFound),
7224    }
7225}
7226
7227/// ORCH-14 `harness.v1.channels.list|status` params, the same storage-root
7228/// override every read-only method accepts so a caller can point the read at
7229/// a fixture home without touching the real ones.
7230#[derive(Debug, Clone, Deserialize)]
7231#[serde(default)]
7232struct ChannelsQuery {
7233    /// Restrict the listing to one harness. `status` requires it.
7234    harness: Option<String>,
7235    /// Channel name, required by `status`.
7236    name: Option<String>,
7237    /// Storage roots to read.
7238    homes: crate::HarnessHomes,
7239}
7240
7241impl Default for ChannelsQuery {
7242    fn default() -> Self {
7243        Self {
7244            harness: None,
7245            name: None,
7246            homes: crate::HarnessHomes::default(),
7247        }
7248    }
7249}
7250
7251/// Read the channel noun. A harness with no channel concept fails with
7252/// `UnsupportedAction` (RPC `-32020`), never an empty list. No row carries a
7253/// token, key or secret — see `crate::channels` "Secrecy".
7254#[derive(Debug, Clone, Deserialize)]
7255#[serde(default)]
7256struct RoutesQuery {
7257    harness: Option<String>,
7258    /// Restrict to routes targeting one profile / agent.
7259    profile: Option<String>,
7260    homes: crate::HarnessHomes,
7261}
7262
7263impl Default for RoutesQuery {
7264    fn default() -> Self {
7265        Self {
7266            harness: None,
7267            profile: None,
7268            homes: crate::HarnessHomes::default(),
7269        }
7270    }
7271}
7272
7273#[derive(Debug, Clone, Deserialize)]
7274#[serde(default)]
7275struct TriggersQuery {
7276    harness: Option<String>,
7277    homes: crate::HarnessHomes,
7278}
7279
7280impl Default for TriggersQuery {
7281    fn default() -> Self {
7282        Self {
7283            harness: None,
7284            homes: crate::HarnessHomes::default(),
7285        }
7286    }
7287}
7288
7289fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
7290    let query = decode::<TriggersQuery>(params)?;
7291    let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
7292        .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7293    Ok(json!({
7294        "schema": crate::triggers::TRIGGERS_SCHEMA,
7295        "triggers": triggers,
7296    }))
7297}
7298
7299fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
7300    let query = decode::<RoutesQuery>(params)?;
7301    let routes = crate::routes::list_routes(
7302        &query.homes,
7303        query.harness.as_deref(),
7304        query.profile.as_deref(),
7305    )
7306    .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
7307    Ok(json!({
7308        "schema": crate::routes::ROUTES_SCHEMA,
7309        "routes": routes,
7310    }))
7311}
7312
7313fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
7314    let query = decode::<ChannelsQuery>(params)?;
7315    let to_service = |error: crate::channels::ChannelError| match error {
7316        crate::channels::ChannelError::UnsupportedHarness { .. } => {
7317            ServiceError::UnsupportedAction(error.to_string())
7318        }
7319        crate::channels::ChannelError::NotFound { .. } => {
7320            ServiceError::InvalidParams(error.to_string())
7321        }
7322    };
7323    match method {
7324        "harness.v1.channels.list" => {
7325            let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
7326                .map_err(to_service)?;
7327            Ok(json!({
7328                "schema": crate::channels::CHANNELS_SCHEMA,
7329                "channels": channels,
7330            }))
7331        }
7332        "harness.v1.channels.status" => {
7333            let harness = query
7334                .harness
7335                .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
7336            let name = query
7337                .name
7338                .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
7339            let channel = crate::channels::channel_status(&query.homes, &harness, &name)
7340                .map_err(to_service)?;
7341            Ok(json!({
7342                "schema": crate::channels::CHANNELS_SCHEMA,
7343                "channel": channel,
7344            }))
7345        }
7346        _ => Err(ServiceError::MethodNotFound),
7347    }
7348}
7349
7350fn rpc_error(id: Value, code: i64, message: &str) -> Value {
7351    json!({
7352        "jsonrpc": "2.0",
7353        "id": id,
7354        "error": {"code": code, "message": message},
7355    })
7356}