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