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