Skip to main content

supercode_harness/
harness_service.rs

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