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