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