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    #[serde(default)]
3280    subject: Option<String>,
3281    /// Name the sender is known by (`fleet-board`, `aaron`). Replies to the
3282    /// message are filed in this sender's mailbox, read with
3283    /// `sessions.inbox`. Defaults to `supercode`.
3284    #[serde(default)]
3285    from_name: Option<String>,
3286    /// A voice front for the receiving session, never a separate agent.
3287    #[serde(default)]
3288    voice_for: Option<crate::mailbox::MailAddress>,
3289    /// Human-readable sender label, independent of its return mailbox.
3290    #[serde(default)]
3291    sender_name: Option<String>,
3292    /// Id of the message this one answers.
3293    #[serde(default)]
3294    in_reply_to: Option<String>,
3295    /// File one notice in the sender's mailbox when the receiver next goes idle.
3296    #[serde(default)]
3297    notify_when_idle: bool,
3298    /// The text is the session's own user speaking (a voice bridge, a board
3299    /// the owner types in): user mail, delivered as the user's own turn.
3300    /// `harness.v1` is served only to the machine's owner (its daemon admits
3301    /// operators only), which is the authority this carries.
3302    #[serde(default)]
3303    as_user: bool,
3304    /// Same storage roots discovery accepts, so a caller (and a test) can
3305    /// point the live-session registry somewhere other than `$HOME`.
3306    #[serde(default)]
3307    homes: crate::HarnessHomes,
3308}
3309
3310#[derive(Deserialize)]
3311#[serde(deny_unknown_fields)]
3312struct InboxParams {
3313    /// Sender name used with `sessions.message` (its mailbox), or
3314    #[serde(default)]
3315    from_name: Option<String>,
3316    /// an explicit session address `sc:<machine>:<harness>:<id>`.
3317    #[serde(default)]
3318    address: Option<String>,
3319    /// Include messages already read.
3320    #[serde(default)]
3321    all: bool,
3322}
3323
3324#[derive(Deserialize)]
3325#[serde(deny_unknown_fields)]
3326struct HarnessSettingsParams {
3327    harness: String,
3328}
3329
3330#[derive(Deserialize)]
3331#[serde(deny_unknown_fields)]
3332struct ConfigureHarnessParams {
3333    harness: String,
3334    #[serde(default)]
3335    changes: Vec<crate::HarnessSettingChange>,
3336    #[serde(default)]
3337    expected_revision: Option<String>,
3338}
3339
3340fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
3341    match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
3342        Ok(report) => (
3343            serde_json::to_value(report).unwrap_or(Value::Null),
3344            Value::Null,
3345        ),
3346        Err(error) => (
3347            Value::Null,
3348            Value::String(format!(
3349                "Volter Harness could not inspect Claude Code inbound controls: {error}"
3350            )),
3351        ),
3352    }
3353}
3354
3355/// Deliver `text` into a session that is running right now, or say why not.
3356///
3357/// A refusal is a RESULT, not a JSON-RPC error: "that session is persisted
3358/// only" is an answer about the session, which a mirror renders next to the
3359/// transcript, and this service's error envelope carries no structured data
3360/// field a machine-readable reason could survive in.
3361///
3362/// The door is chosen by the one router every sender uses
3363/// ([`crate::mail_route`]): a session supercode controls gets it through its
3364/// runtime (the default tier); a Claude session it does not control through a
3365/// Claude relay, so the reply comes back; a Codex session through its hook;
3366/// anything else is stored. Replies are filed in the sender's mailbox under
3367/// `reply_to`, read with `sessions.inbox`.
3368async fn message_live_session(params: &MessageSessionParams) -> Value {
3369    use crate::mail_route::{Delivered, NoDoor, Refused};
3370    let (inbound_controls, inbound_controls_error) =
3371        claude_inbound_controls_or_error(&params.homes);
3372    let refused = |reason: &str, message: String| {
3373        json!({
3374            "delivered_to_bus": false,
3375            "refusal": {"reason": reason, "message": message},
3376            "inbound_controls": inbound_controls,
3377            "inbound_controls_error": inbound_controls_error,
3378        })
3379    };
3380    if params.text.trim().is_empty() {
3381        return refused(
3382            crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3383            "refusing to deliver an empty message".into(),
3384        );
3385    }
3386    let sender = match operator_address(params.from_name.as_deref()) {
3387        Ok(sender) => sender,
3388        Err(message) => return refused("invalid_sender", message),
3389    };
3390    let receiver = match crate::mailbox::MailAddress::new(
3391        crate::mailbox::local_machine_name(),
3392        params.locator.harness.as_str(),
3393        &params.locator.session_id,
3394    ) {
3395        Ok(receiver) => receiver,
3396        Err(error) => return refused("delivery_failed", error.to_string()),
3397    };
3398    if params
3399        .voice_for
3400        .as_ref()
3401        .is_some_and(|voice_for| voice_for != &receiver)
3402    {
3403        return refused(
3404            "invalid_sender",
3405            "a voice front must represent the receiving session".into(),
3406        );
3407    }
3408    // A Claude session's registry name must still name only that session.
3409    if params.locator.harness.as_str() == HarnessId::CLAUDE_CODE
3410        && crate::runtime_mail::controlled_runtime(
3411            HarnessId::CLAUDE_CODE,
3412            &params.locator.session_id,
3413        )
3414        .is_none()
3415    {
3416        if let Err(refusal) =
3417            crate::claude_peer::resolve_live_session(&params.homes, &params.locator.session_id)
3418        {
3419            return refused(refusal.reason.as_str(), refusal.message);
3420        }
3421    }
3422    if params.as_user {
3423        return message_as_user(
3424            params,
3425            sender,
3426            receiver,
3427            inbound_controls,
3428            inbound_controls_error,
3429        )
3430        .await;
3431    }
3432    let door = match crate::mail_route::door_for(&params.homes, &receiver) {
3433        Ok(door) => door,
3434        Err(NoDoor::NotRunning | NoDoor::OtherMachine(_)) => {
3435            return refused(
3436                crate::claude_peer::ClaudePeerRefusal::NotLive.as_str(),
3437                format!(
3438                    "no running `{}` session `{}` is reachable; its transcript is persisted only",
3439                    params.locator.harness.as_str(),
3440                    params.locator.session_id
3441                ),
3442            )
3443        }
3444    };
3445    let mut envelope = match crate::mailbox::Envelope::new(
3446        sender.clone(),
3447        params
3448            .sender_name
3449            .clone()
3450            .unwrap_or_else(|| format!("{}@{}", sender.session_id, sender.machine)),
3451        crate::mailbox::MailKind::Peer,
3452        crate::mailbox::ReplyVia::Command,
3453        params.text.clone(),
3454    ) {
3455        Ok(envelope) => envelope,
3456        Err(error) => return refused("delivery_failed", error.to_string()),
3457    };
3458    envelope.in_reply_to = params.in_reply_to.clone();
3459    envelope.voice_for = params.voice_for.clone();
3460    envelope.subject = params.subject.as_deref().map(|value| {
3461        value
3462            .lines()
3463            .next()
3464            .unwrap_or("")
3465            .chars()
3466            .take(200)
3467            .collect()
3468    });
3469    let how = match crate::mail_route::deliver(
3470        &envelope,
3471        &receiver,
3472        &door,
3473        true,
3474        params.notify_when_idle,
3475    )
3476    .await
3477    {
3478        Err(detail) => {
3479            return refused(
3480                crate::claude_peer::ClaudePeerRefusal::DeliveryFailed.as_str(),
3481                detail,
3482            )
3483        }
3484        Ok(Err(Refused::CannotQueueNative)) => unreachable!("sessions.message always wakes"),
3485        Ok(Err(Refused::TooLong(bytes))) => {
3486            return refused(
3487                "too_long",
3488                format!(
3489                    "the message is {bytes} bytes; the limit is {}",
3490                    crate::mail_route::MAX_RELAYED_BYTES
3491                ),
3492            )
3493        }
3494        Ok(Ok(delivered)) => match delivered {
3495            Delivered::Steered => "steered",
3496            Delivered::Started => "started",
3497            Delivered::Native { busy: true } => "next_tool_call",
3498            Delivered::Native { busy: false } => "started",
3499            Delivered::Hooked => "hook",
3500            Delivered::HookWoken => "woken",
3501            Delivered::Queued => "queued",
3502            Delivered::Stored => "stored",
3503            Delivered::Operator => "filed",
3504        },
3505    };
3506    json!({
3507        "delivered_to_bus": !matches!(door, crate::mail_route::Door::Stored),
3508        "message_id": envelope.id,
3509        "reply_to": sender.to_string(),
3510        "target": {
3511            "harness": params.locator.harness.as_str(),
3512            "session_id": params.locator.session_id,
3513            "name": match &door {
3514                crate::mail_route::Door::Native(session) => Some(session.name.clone()),
3515                _ => None,
3516            },
3517        },
3518        "delivery": {"door": door.name(), "how": how},
3519        "inbound_controls": inbound_controls,
3520        "inbound_controls_error": inbound_controls_error,
3521    })
3522}
3523
3524/// `sessions.message` with `as_user`: user mail, through the one door that
3525/// carries the user's authority (a hosted runtime's input, or the session's
3526/// pane when its composer is empty; otherwise it waits in the mailbox).
3527async fn message_as_user(
3528    params: &MessageSessionParams,
3529    sender: crate::mailbox::MailAddress,
3530    receiver: crate::mailbox::MailAddress,
3531    inbound_controls: Value,
3532    inbound_controls_error: Value,
3533) -> Value {
3534    let envelope = match crate::mailbox::Envelope::new(
3535        sender.clone(),
3536        format!("{}@{}", sender.session_id, sender.machine),
3537        crate::mailbox::MailKind::User,
3538        crate::mailbox::ReplyVia::None,
3539        params.text.clone(),
3540    ) {
3541        Ok(envelope) => envelope,
3542        Err(error) => {
3543            return json!({
3544                "delivered_to_bus": false,
3545                "refusal": {"reason": "delivery_failed", "message": error.to_string()},
3546                "inbound_controls": inbound_controls,
3547                "inbound_controls_error": inbound_controls_error,
3548            })
3549        }
3550    };
3551    match crate::mail_route::deliver_user_turn(&params.homes, &envelope, &receiver).await {
3552        Ok(turn) => json!({
3553            "delivered_to_bus": turn != crate::mail_route::UserTurn::Waiting,
3554            "message_id": envelope.id,
3555            "target": {
3556                "harness": params.locator.harness.as_str(),
3557                "session_id": params.locator.session_id,
3558            },
3559            "delivery": {
3560                "door": match turn {
3561                    crate::mail_route::UserTurn::Steered | crate::mail_route::UserTurn::Started => "runtime",
3562                    _ => "pane",
3563                },
3564                "how": turn.as_str(),
3565            },
3566            "inbound_controls": inbound_controls,
3567            "inbound_controls_error": inbound_controls_error,
3568        }),
3569        Err(message) => json!({
3570            "delivered_to_bus": false,
3571            "refusal": {"reason": "no_user_door", "message": message},
3572            "inbound_controls": inbound_controls,
3573            "inbound_controls_error": inbound_controls_error,
3574        }),
3575    }
3576}
3577
3578#[derive(Deserialize)]
3579#[serde(deny_unknown_fields)]
3580struct ActivityUnderParams {
3581    /// Root processes (a terminal pane's shell, say).
3582    pids: Vec<u32>,
3583    #[serde(default)]
3584    homes: crate::HarnessHomes,
3585}
3586
3587/// `harness.v1.sessions.activity_under`: the harness session running beneath
3588/// each root process, with its activity, so a terminal's attention takes a
3589/// harness pane's liveness from the harness's own lifecycle.
3590async fn activity_under_call(params: Value) -> std::result::Result<Value, ServiceError> {
3591    let params = decode::<ActivityUnderParams>(params)?;
3592    if params.pids.len() > 1024 {
3593        return Err(ServiceError::InvalidParams(
3594            "sessions.activity_under accepts at most 1024 pids".into(),
3595        ));
3596    }
3597    let found = crate::session_activity::activity_under(&params.pids, &params.homes)
3598        .await
3599        .map_err(ServiceError::Sdk)?;
3600    Ok(json!({
3601        "activities": found
3602            .into_iter()
3603            .map(|(pid, activity)| json!({"pid": pid, "activity": activity}))
3604            .collect::<Vec<_>>(),
3605    }))
3606}
3607
3608/// Mailbox address of an operator sender (a board, a person's shell) that is
3609/// not itself a harness session: `sc:<machine>:operator:<name>`.
3610fn operator_address(
3611    name: Option<&str>,
3612) -> std::result::Result<crate::mailbox::MailAddress, String> {
3613    let name: String = name
3614        .unwrap_or("supercode")
3615        .trim()
3616        .chars()
3617        .map(|character| {
3618            if character.is_whitespace() || character == '@' {
3619                '-'
3620            } else {
3621                character
3622            }
3623        })
3624        .collect();
3625    if name.is_empty() {
3626        return Err("from_name must not be empty".into());
3627    }
3628    crate::mailbox::MailAddress::new(crate::mailbox::local_machine_name(), "operator", name)
3629        .map_err(|error| error.to_string())
3630}
3631
3632/// `harness.v1.sessions.inbox`: a sender's mailbox, each message with the
3633/// text its reader sees. Unread messages are claimed and acknowledged by this
3634/// read, so a second read does not return them again.
3635fn inbox_call(params: InboxParams) -> std::result::Result<Value, ServiceError> {
3636    let address = match (&params.address, &params.from_name) {
3637        (Some(address), None) => crate::mailbox::MailAddress::parse(address)
3638            .map_err(|error| ServiceError::InvalidParams(error.to_string()))?,
3639        (None, name) => operator_address(name.as_deref()).map_err(ServiceError::InvalidParams)?,
3640        (Some(_), Some(_)) => {
3641            return Err(ServiceError::InvalidParams(
3642                "sessions.inbox takes from_name or address, not both".into(),
3643            ))
3644        }
3645    };
3646    let operation = |error: std::io::Error| ServiceError::Operation(error.to_string());
3647    let mailbox =
3648        crate::mailbox::Mailbox::open(&crate::mailbox::mail_root(), &address).map_err(operation)?;
3649    let claimed = mailbox.claim_unread().map_err(operation)?;
3650    let mut messages: Vec<Value> = Vec::new();
3651    if params.all {
3652        for stored in mailbox.list().map_err(operation)? {
3653            if stored.state == crate::mailbox::MailState::Read {
3654                messages.push(json!({"state": "read", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3655            }
3656        }
3657    }
3658    for stored in &claimed {
3659        messages.push(json!({"state": "unread", "envelope": stored.envelope, "rendered": stored.envelope.render()}));
3660    }
3661    for stored in &claimed {
3662        mailbox.acknowledge(stored).map_err(operation)?;
3663    }
3664    Ok(json!({"address": address.to_string(), "messages": messages}))
3665}
3666
3667/// Source identity of one follow subscription, plus the last lifecycle state
3668/// already reported on it. The follower itself stays purely persistence-facing.
3669// Only the adapter-api poll reads these; the subscription bookkeeping itself is
3670// shared by both builds.
3671#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3672struct FollowedSource {
3673    harness: String,
3674    session_id: String,
3675    reported: Option<String>,
3676}
3677
3678#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
3679struct ActivitySubscription {
3680    locators: Vec<SessionLocator>,
3681    homes: crate::HarnessHomes,
3682    reported: BTreeMap<(String, String), crate::SessionActivity>,
3683}
3684
3685/// Add what makes an indexed row behaviorally equivalent to a discovered row: the attach
3686/// endpoint of a session supercode hosts, and the door a message reaches it by.
3687///
3688/// The durable index owns only persistence metadata. Both are projections: every message/attach
3689/// operation revalidates its authority, so publishing one here never trusts a stale browser-held
3690/// handle. The doors are read once per batch ([`crate::mail_route::LiveSessions`]).
3691fn live_descriptor_value(
3692    session: &SessionDescriptor,
3693    doors: &crate::mail_route::LiveSessions,
3694) -> std::result::Result<Value, ServiceError> {
3695    let mut value = serde_json::to_value(session)
3696        .map_err(|error| ServiceError::Operation(error.to_string()))?;
3697    // A running session with no title of its own is known by its live name.
3698    if value.get("title").is_none_or(Value::is_null) {
3699        if let Some(live) = doors.all().iter().find(|live| {
3700            live.address.harness == session.locator.harness.as_str()
3701                && live.address.session_id == session.locator.session_id
3702        }) {
3703            let name = live.name.split('@').next().unwrap_or(&live.name);
3704            if !name.is_empty() {
3705                value["title"] = json!(name);
3706            }
3707        }
3708    }
3709    if let Some(workspace) = &session.cwd {
3710        let source = LiveRuntimeSource {
3711            harness: session.locator.harness.as_str().to_string(),
3712            session_id: session.locator.session_id.clone(),
3713            workspace: workspace.clone(),
3714        };
3715        if let Some(endpoint) = discover_live_runtime(&source)
3716            .map_err(|error| ServiceError::Operation(error.to_string()))?
3717        {
3718            value["live_endpoint"] = json!(endpoint.as_str());
3719        }
3720    }
3721    // How a message reaches this session right now, chosen by the same
3722    // router every sender uses: `runtime` (supercode controls it), `native`,
3723    // `hook` or `stored`. Absent when no process is running it.
3724    if let Some(door) = doors.door(
3725        session.locator.harness.as_str(),
3726        &session.locator.session_id,
3727    ) {
3728        value["delivery"] = json!(door);
3729    }
3730    // What a session waiting on its user is asking, so every surface can show the question and not only the state.
3731    // Read as `message list` reads it: the transcript's open call, else the prompt on the session's pane.
3732    if let Some(live) = doors.all().iter().find(|live| {
3733        live.address.harness == session.locator.harness.as_str()
3734            && live.address.session_id == session.locator.session_id
3735            && live.status == "waiting"
3736    }) {
3737        if let Some(request) = live.pending_request(&crate::HarnessHomes::default()) {
3738            value["pending_request"] = request;
3739        }
3740    }
3741    Ok(value)
3742}
3743
3744fn live_index_changes(
3745    changes: Vec<crate::session_index::SessionIndexChange>,
3746    homes: &HarnessHomes,
3747) -> std::result::Result<Vec<Value>, ServiceError> {
3748    use crate::session_index::SessionIndexChange;
3749    let doors = crate::mail_route::LiveSessions::read(homes);
3750    changes
3751        .into_iter()
3752        .map(|change| match change {
3753            SessionIndexChange::Added { descriptor } => Ok(json!({
3754                "kind": "added",
3755                "descriptor": live_descriptor_value(&descriptor, &doors)?,
3756            })),
3757            SessionIndexChange::Updated { descriptor } => Ok(json!({
3758                "kind": "updated",
3759                "descriptor": live_descriptor_value(&descriptor, &doors)?,
3760            })),
3761            SessionIndexChange::Removed { key } => Ok(json!({
3762                "kind": "removed",
3763                "key": key,
3764            })),
3765        })
3766        .collect()
3767}
3768
3769fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
3770    use crate::{SessionPresence, SessionTurnState};
3771    match (activity.presence, activity.turn) {
3772        (SessionPresence::Persisted, _) => None,
3773        (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
3774        (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
3775        // Waiting on its user (a question, an approval): the word every other surface uses.
3776        (SessionPresence::Running, SessionTurnState::NeedsInput) => Some("waiting"),
3777        // The normalized activity object can honestly report a live owner even
3778        // when the stock harness never published a turn status. Preserve the
3779        // older field's stricter contract instead of guessing `running`.
3780        (SessionPresence::Running, SessionTurnState::Unknown)
3781            if activity.evidence.native_state.is_none() =>
3782        {
3783            None
3784        }
3785        (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
3786    }
3787}
3788
3789#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
3790#[serde(rename_all = "kebab-case")]
3791enum TransferFormat {
3792    ClaudeCode,
3793    Codex,
3794    #[serde(rename = "opencode", alias = "open-code")]
3795    OpenCode,
3796    Pi,
3797    Grok,
3798    Gemini,
3799    Goose,
3800    /// UNI-18: a Hermes target. Its artifact is the Codex rollout that
3801    /// `hermes sessions import --from codex` reads; `sessions.export` performs
3802    /// that import into the Hermes home.
3803    Hermes,
3804}
3805
3806impl TransferFormat {
3807    fn id(self) -> &'static str {
3808        match self {
3809            Self::ClaudeCode => HarnessId::CLAUDE_CODE,
3810            Self::Codex => HarnessId::CODEX,
3811            Self::OpenCode => HarnessId::OPENCODE,
3812            Self::Pi => HarnessId::PI,
3813            Self::Grok => HarnessId::GROK,
3814            Self::Gemini => HarnessId::GEMINI,
3815            Self::Goose => HarnessId::GOOSE,
3816            Self::Hermes => HarnessId::HERMES,
3817        }
3818    }
3819}
3820
3821impl From<TransferFormat> for SessionFormat {
3822    fn from(value: TransferFormat) -> Self {
3823        match value {
3824            TransferFormat::ClaudeCode => Self::ClaudeCode,
3825            TransferFormat::Codex => Self::Codex,
3826            TransferFormat::OpenCode => Self::OpenCode,
3827            TransferFormat::Pi => Self::Pi,
3828            TransferFormat::Grok => Self::Grok,
3829            TransferFormat::Gemini => Self::Gemini,
3830            TransferFormat::Goose => Self::Goose,
3831            // a Hermes artifact is the Codex rollout Hermes imports
3832            TransferFormat::Hermes => Self::Codex,
3833        }
3834    }
3835}
3836
3837#[derive(Deserialize)]
3838struct ImportSessionParams {
3839    source_harness: TransferFormat,
3840    content: String,
3841}
3842
3843#[derive(Deserialize)]
3844struct ExportSessionParams {
3845    locator: SessionLocator,
3846    target_harness: TransferFormat,
3847}
3848
3849#[derive(Deserialize)]
3850struct ReduceSessionParams {
3851    locator: SessionLocator,
3852    target_harness: TransferFormat,
3853    #[serde(default = "default_keep_last")]
3854    keep_last: usize,
3855}
3856
3857fn default_keep_last() -> usize {
3858    6
3859}
3860
3861#[derive(Deserialize)]
3862struct BranchSessionParams {
3863    locator: SessionLocator,
3864    #[serde(default)]
3865    target_harness: Option<TransferFormat>,
3866}
3867
3868#[derive(Deserialize)]
3869struct HandoffSessionParams {
3870    locator: SessionLocator,
3871    target_harness: TransferFormat,
3872    #[serde(default)]
3873    cwd: Option<PathBuf>,
3874}
3875
3876#[derive(Deserialize)]
3877struct MaterializeSessionParams {
3878    artifact: crate::native_materialize::MaterializeArtifact,
3879    cwd: PathBuf,
3880    /// Where the continuation is written; unset roots are the environment's own, as discovery reads them.
3881    #[serde(default)]
3882    homes: HarnessHomes,
3883}
3884
3885/// A session supercode starts or resumes runs without approval prompts
3886/// (`yolo`) unless the caller asks for the harness's own (`default`).
3887#[derive(Debug, Clone, Copy, Default, Deserialize)]
3888#[serde(rename_all = "snake_case")]
3889enum ResumePolicy {
3890    Default,
3891    #[default]
3892    Yolo,
3893}
3894
3895#[derive(Deserialize)]
3896struct ResumeInstructionsParams {
3897    locator: SessionLocator,
3898    #[serde(default)]
3899    cwd: Option<PathBuf>,
3900    #[serde(default)]
3901    policy: ResumePolicy,
3902}
3903
3904/// `harness.v1.workflow.load` parameters: which harness's board, and its home.
3905#[derive(Deserialize)]
3906struct WorkflowLoadParams {
3907    from: crate::workflow_doors::WorkflowHarness,
3908    home: PathBuf,
3909}
3910
3911/// ONT-4 `harness.v1.orchestration.load` parameters. `flavor` says which layout the
3912/// folder is read as; our own is the default.
3913#[derive(Deserialize)]
3914struct OrchestrationLoadParams {
3915    root: PathBuf,
3916    #[serde(default)]
3917    flavor: crate::orchestration_doors::HomeFlavor,
3918}
3919
3920/// ONT-4 `harness.v1.orchestration.save` parameters. `vault` is merged into the
3921/// home's own secrets; a caller that sends none keeps what is on disk.
3922#[derive(Deserialize)]
3923struct OrchestrationSaveParams {
3924    root: PathBuf,
3925    orchestration: crate::orchestration::Orchestration,
3926    #[serde(default)]
3927    vault: BTreeMap<String, String>,
3928}
3929
3930/// ONT-4 `harness.v1.orchestration.compile` parameters.
3931#[derive(Deserialize)]
3932struct OrchestrationCompileParams {
3933    from: crate::orchestration_doors::OrchestrationHarness,
3934    home: PathBuf,
3935}
3936
3937/// ONT-4 `harness.v1.orchestration.decompile` parameters. `source` is the home the
3938/// orchestration was compiled from: it is re-compiled to recover the io bookkeeping
3939/// that byte reuse and the live-store refusal (UNI-18) are decided from.
3940#[derive(Deserialize)]
3941struct OrchestrationDecompileParams {
3942    to: crate::orchestration_doors::OrchestrationHarness,
3943    orchestration: crate::orchestration::Orchestration,
3944    source: PathBuf,
3945    #[serde(default)]
3946    source_flavor: crate::orchestration_doors::SourceFlavor,
3947    dest: PathBuf,
3948    #[serde(default)]
3949    vault: BTreeMap<String, String>,
3950}
3951
3952/// `harness.v1.orchestration.import` parameters: another harness's home, and the
3953/// folder of ours it becomes.
3954#[derive(Deserialize)]
3955struct OrchestrationImportParams {
3956    from: crate::orchestration_doors::OrchestrationHarness,
3957    home: PathBuf,
3958    into: PathBuf,
3959}
3960
3961/// `harness.v1.orchestration.export` parameters: a folder of ours, and the home of
3962/// another harness it becomes.
3963#[derive(Deserialize)]
3964struct OrchestrationExportParams {
3965    to: crate::orchestration_doors::OrchestrationHarness,
3966    root: PathBuf,
3967    dest: PathBuf,
3968}
3969
3970/// `harness.v1.jobs.get` parameters.
3971#[derive(Deserialize)]
3972struct JobsGetParams {
3973    harness: String,
3974    id: String,
3975    #[serde(default)]
3976    homes: crate::HarnessHomes,
3977}
3978
3979/// ORCH-18: run one mutating job verb through the harness's own CLI.
3980///
3981/// The refusal ladder is deliberate: a harness with no scheduled-job concept
3982/// at all answers with the SAME sentence `jobs.list` gives it, and a harness
3983/// that has jobs but publishes no client-callable verb (Claude Code, whose
3984/// jobs are created by the model inside a session) answers with its own
3985/// reason. Neither is ever a silent no-op.
3986fn mutate_job(
3987    verb: crate::jobs_control::JobVerb,
3988    params: Value,
3989) -> std::result::Result<Value, ServiceError> {
3990    let mutation = decode::<crate::jobs_control::JobMutation>(params)?;
3991    refuse_harness_without_jobs(&mutation.harness, &format!("jobs.{}", verb.as_str()))?;
3992    let outcome = crate::jobs_control::mutate(verb, &mutation).map_err(job_control_error)?;
3993    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
3994}
3995
3996/// ORCH-22: run one mutating skills verb through the harness's own door.
3997///
3998/// The refusal ladder mirrors `jobs.*`: a harness with no skills root at all
3999/// answers with the same sentence `skills.list` gives it, and a harness whose
4000/// door does not publish this verb (OpenClaw has no `skills remove` at the
4001/// pin) answers with its own reason. Neither is ever a silent no-op.
4002fn mutate_skill(
4003    verb: crate::skills_control::SkillVerb,
4004    params: Value,
4005) -> std::result::Result<Value, ServiceError> {
4006    let mutation = decode::<crate::skills_control::SkillMutation>(params)?;
4007    if !crate::skills_control::supports_skill_control(&mutation.harness) {
4008        return Err(ServiceError::UnsupportedAction(format!(
4009            "`{}` has no skills root Volter Harness reads; `skills.{}` is supported for: {}",
4010            mutation.harness,
4011            verb.as_str(),
4012            crate::skills_control::CONTROLLED_SKILL_HARNESSES.join(", ")
4013        )));
4014    }
4015    let outcome =
4016        crate::skills_control::mutate_skill(verb, &mutation).map_err(skill_control_error)?;
4017    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4018}
4019
4020/// The skills twin of [`job_control_error`], with the same mapping rule.
4021fn skill_control_error(error: crate::skills_control::SkillControlError) -> ServiceError {
4022    match error {
4023        crate::skills_control::SkillControlError::Unsupported(message) => {
4024            ServiceError::UnsupportedAction(message)
4025        }
4026        crate::skills_control::SkillControlError::Invalid(message) => {
4027            ServiceError::InvalidParams(message)
4028        }
4029        crate::skills_control::SkillControlError::Failed(message) => {
4030            ServiceError::Operation(message)
4031        }
4032    }
4033}
4034
4035/// ORCH-21: run one mutating profile verb through the harness's own CLI.
4036///
4037/// The refusal ladder mirrors `mutate_job`'s: a harness with no profile
4038/// concept at all answers with the SAME sentence `profiles.list` gives it, and
4039/// a harness that HAS profiles but publishes no client-callable verb (Codex's
4040/// file-authored `[profiles.<name>]` tables, supercode's compiled-in presets)
4041/// answers with its own reason. Neither is ever a silent no-op.
4042fn mutate_profile(
4043    verb: crate::profiles_control::ProfileVerb,
4044    params: Value,
4045) -> std::result::Result<Value, ServiceError> {
4046    let mutation = decode::<crate::profiles_control::ProfileMutation>(params)?;
4047    let outcome =
4048        crate::profiles_control::mutate(verb, &mutation).map_err(profile_control_error)?;
4049    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
4050}
4051
4052/// The same mapping `job_control_error` applies, for the profile noun.
4053fn profile_control_error(error: crate::profiles_control::ProfileControlError) -> ServiceError {
4054    match error {
4055        crate::profiles_control::ProfileControlError::Unsupported(message) => {
4056            ServiceError::UnsupportedAction(message)
4057        }
4058        crate::profiles_control::ProfileControlError::Invalid(message) => {
4059            ServiceError::InvalidParams(message)
4060        }
4061        crate::profiles_control::ProfileControlError::Failed(message) => {
4062            ServiceError::Operation(message)
4063        }
4064    }
4065}
4066
4067/// Map a controlled-tier failure onto the service's error vocabulary. A verb
4068/// the harness lacks is `UnsupportedAction`; a harness verb that RAN and
4069/// failed carries its own stderr through as the operation error.
4070fn job_control_error(error: crate::jobs_control::JobControlError) -> ServiceError {
4071    match error {
4072        crate::jobs_control::JobControlError::Unsupported(message) => {
4073            ServiceError::UnsupportedAction(message)
4074        }
4075        crate::jobs_control::JobControlError::Invalid(message) => {
4076            ServiceError::InvalidParams(message)
4077        }
4078        crate::jobs_control::JobControlError::Failed(message) => ServiceError::Operation(message),
4079    }
4080}
4081
4082/// Map an ORCH-19 controlled-tier failure onto the service's error
4083/// vocabulary. A verb the harness has no door for is `UnsupportedAction`; a
4084/// door that RAN and failed carries the harness's own stderr / HTTP body
4085/// through as the operation error.
4086fn session_control_error(error: crate::SessionControlError) -> ServiceError {
4087    match error {
4088        crate::SessionControlError::Unsupported(message) => {
4089            ServiceError::UnsupportedAction(message)
4090        }
4091        crate::SessionControlError::Invalid(message) => ServiceError::InvalidParams(message),
4092        crate::SessionControlError::Failed(message) => ServiceError::Operation(message),
4093    }
4094}
4095
4096/// A harness without a scheduled-job concept refuses the verb rather than
4097/// answering with an empty list — an absent capability and an empty inventory
4098/// are different answers (the same rule `runtimes.capabilities` applies to
4099/// `steer`).
4100fn refuse_harness_without_jobs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4101    if crate::jobs::supports_jobs(harness) {
4102        return Ok(());
4103    }
4104    Err(ServiceError::UnsupportedAction(format!(
4105        "`{harness}` has no scheduled jobs; `{verb}` is supported for: {}",
4106        crate::jobs::JOB_HARNESSES.join(", ")
4107    )))
4108}
4109
4110/// `harness.v1.runs.get` parameters.
4111#[derive(Deserialize)]
4112struct RunsGetParams {
4113    harness: String,
4114    id: String,
4115    #[serde(default)]
4116    homes: crate::HarnessHomes,
4117}
4118
4119/// A harness with no run store refuses the verb rather than answering with an
4120/// empty history — the same rule `jobs.list` applies. Claude Code lands here
4121/// on purpose: its cron fires are ordinary turns inside the session that
4122/// created the job, so there is no fire record to list.
4123fn refuse_harness_without_runs(harness: &str, verb: &str) -> std::result::Result<(), ServiceError> {
4124    if crate::runs::supports_runs(harness) {
4125        return Ok(());
4126    }
4127    Err(ServiceError::UnsupportedAction(format!(
4128        "`{harness}` keeps no run store; `{verb}` is supported for: {}",
4129        crate::runs::RUN_HARNESSES.join(", ")
4130    )))
4131}
4132
4133#[derive(Serialize)]
4134struct SessionArtifact {
4135    source_harness: HarnessId,
4136    target_harness: &'static str,
4137    session_id: Option<String>,
4138    content: String,
4139    suggested_filename: String,
4140    files: Vec<SessionArtifactFile>,
4141    fidelity: Fidelity,
4142    residue: Vec<String>,
4143}
4144
4145#[derive(Serialize)]
4146struct SessionArtifactFile {
4147    path: String,
4148    content: String,
4149    role: ArtifactFileRole,
4150}
4151
4152#[derive(Serialize)]
4153#[serde(rename_all = "snake_case")]
4154enum ArtifactFileRole {
4155    Primary,
4156    Subagent,
4157    Bundle,
4158    SourceRecovery,
4159}
4160
4161#[derive(Serialize)]
4162struct StructuredLaunch {
4163    cwd: PathBuf,
4164    program: String,
4165    arguments: Vec<String>,
4166    env: BTreeMap<String, String>,
4167}
4168
4169struct HandoffInstructions {
4170    launch: StructuredLaunch,
4171    materialize: Option<StructuredLaunch>,
4172    requires_materialization: bool,
4173    note: String,
4174}
4175
4176#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
4177#[serde(rename_all = "snake_case")]
4178enum HarnessProbeLevel {
4179    #[default]
4180    Passive,
4181    Handshake,
4182}
4183
4184#[derive(Default, Deserialize)]
4185#[serde(default)]
4186struct HarnessInventoryParams {
4187    harness: Option<HarnessId>,
4188    harnesses: Vec<HarnessId>,
4189    workspace: Option<PathBuf>,
4190    probe: HarnessProbeLevel,
4191    include_sessions: bool,
4192    /// Omit subprocess-based `--version` calls when a latency-sensitive UI only needs readiness.
4193    skip_versions: bool,
4194}
4195
4196#[derive(Deserialize)]
4197struct HarnessAuthenticationParams {
4198    harness: HarnessId,
4199}
4200
4201#[derive(Deserialize)]
4202struct BeginHarnessAuthenticationParams {
4203    harness: HarnessId,
4204    #[serde(default = "local_browser_authentication_environment")]
4205    environment: crate::HarnessAuthenticationEnvironment,
4206    #[serde(default)]
4207    method: Option<crate::HarnessAuthenticationMethodId>,
4208    #[serde(default)]
4209    cwd: Option<PathBuf>,
4210}
4211
4212fn local_browser_authentication_environment() -> crate::HarnessAuthenticationEnvironment {
4213    crate::HarnessAuthenticationEnvironment::LocalBrowser
4214}
4215
4216#[derive(Serialize)]
4217struct HarnessInventoryReport {
4218    probe: HarnessProbeLevel,
4219    workspace: Option<PathBuf>,
4220    harnesses: Vec<LocalHarness>,
4221}
4222
4223#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4224#[serde(rename_all = "snake_case")]
4225enum HarnessAuthState {
4226    Ready,
4227    Configured,
4228    Required,
4229    Unknown,
4230}
4231
4232#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4233#[serde(rename_all = "snake_case")]
4234enum HarnessRuntimeState {
4235    Ready,
4236    Degraded,
4237    Unavailable,
4238}
4239
4240#[derive(Serialize)]
4241struct HarnessSessionCounts {
4242    global: Option<usize>,
4243    workspace: Option<usize>,
4244}
4245
4246/// Receipt-backed evidence that a harness has a RUNNING instance right now,
4247/// distinct from being merely installed (UNI-7). Detection is passive and
4248/// default-on: a gateway liveness connect for daemon harnesses, a fresh
4249/// SQLite WAL stamp for store-writer harnesses (precedent: the opencode
4250/// follower's -wal/-shm freshness). Control stays behind per-connection
4251/// grants — this reports observations only.
4252/// ORCH-17: the gateway-health noun on an inventory row. Derived from the
4253/// UNI-7 running-instance probe (Hermes: `state.db-wal` freshness; OpenClaw:
4254/// a TCP connect to the gateway endpoint resolved from its OWN config) plus
4255/// the executable version — never by starting anything.
4256#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
4257#[serde(rename_all = "snake_case")]
4258pub enum GatewayState {
4259    Up,
4260    Down,
4261    Unknown,
4262}
4263
4264/// ORCH-17: `gateway` on a `harness.v1.harnesses.list` row.
4265#[derive(Debug, Clone, Serialize)]
4266pub struct GatewayHealth {
4267    pub state: GatewayState,
4268    /// The endpoint supercode would connect to (OpenClaw: the gateway
4269    /// WebSocket resolved from `openclaw.json`; core harnesses: their
4270    /// declared connect address when one exists). `None` when the harness
4271    /// has no single endpoint (Hermes multiplexes platforms).
4272    #[serde(skip_serializing_if = "Option::is_none")]
4273    pub endpoint: Option<String>,
4274    #[serde(skip_serializing_if = "Option::is_none")]
4275    pub version: Option<String>,
4276    /// What the verdict rests on, or why it is `unknown`.
4277    pub evidence: String,
4278    pub checked_at_ms: u64,
4279}
4280
4281/// OpenClaw's gateway WebSocket endpoint, resolved from its own config the
4282/// way the registry's connect descriptor prescribes (`gateway.url`, else
4283/// `gateway.port`, else the documented default).
4284fn openclaw_gateway_endpoint(home: &Path) -> String {
4285    let config_path = home.join(".openclaw/openclaw.json");
4286    let gateway = std::fs::read_to_string(&config_path)
4287        .ok()
4288        .and_then(|raw| serde_json::from_str::<serde_json::Value>(&raw).ok())
4289        .and_then(|config| config.get("gateway").cloned());
4290    if let Some(url) = gateway
4291        .as_ref()
4292        .and_then(|gateway| gateway.get("url"))
4293        .and_then(serde_json::Value::as_str)
4294    {
4295        return url.to_string();
4296    }
4297    let port = gateway
4298        .as_ref()
4299        .and_then(|gateway| gateway.get("port"))
4300        .and_then(serde_json::Value::as_u64)
4301        .unwrap_or(18789);
4302    format!("ws://127.0.0.1:{port}")
4303}
4304
4305/// Ask Hermes itself (`hermes gateway status`, read-only, ~1 s) whether its
4306/// gateway is up. The command is per-host launchd/systemd text without a JSON
4307/// form at 0.19–0.21; the verdict is read from the lines it prints:
4308/// "supervised by launchd (PID …)" / "is running" → up, "not running" /
4309/// "not installed" → down, anything else → no verdict. `SUPERCODE_HERMES_BIN`
4310/// overrides the executable so a fake can stand in under test.
4311fn hermes_gateway_status() -> Option<(GatewayState, String)> {
4312    let program = crate::harness_command::harness_program(HarnessId::HERMES).ok()?;
4313    let output = std::process::Command::new(&program)
4314        .args(["gateway", "status"])
4315        .stdin(std::process::Stdio::null())
4316        .output()
4317        .ok()?;
4318    let text = format!(
4319        "{}{}",
4320        String::from_utf8_lossy(&output.stdout),
4321        String::from_utf8_lossy(&output.stderr)
4322    );
4323    let verdict = text.lines().find_map(|line| {
4324        let l = line.trim();
4325        if l.contains("supervised by launchd (PID")
4326            || l.contains("supervised by systemd (PID")
4327            || l.contains("Gateway is running")
4328            || l.contains("process is running")
4329        {
4330            Some((GatewayState::Up, format!("`hermes gateway status`: {l}")))
4331        } else if l.contains("not running") || l.contains("not installed") {
4332            Some((GatewayState::Down, format!("`hermes gateway status`: {l}")))
4333        } else {
4334            None
4335        }
4336    });
4337    verdict
4338}
4339
4340fn gateway_health(
4341    id: &str,
4342    installed: bool,
4343    running: Option<&RunningInstance>,
4344    version: Option<&str>,
4345) -> GatewayHealth {
4346    let checked_at_ms = now_epoch_ms();
4347    let home = supercode_interchange::user_home()
4348        .map(std::path::PathBuf::into_os_string)
4349        .map(PathBuf::from);
4350    match id {
4351        HarnessId::HERMES | HarnessId::OPENCLAW => {
4352            let endpoint = (id == HarnessId::OPENCLAW)
4353                .then(|| home.as_deref().map(openclaw_gateway_endpoint))
4354                .flatten();
4355            let (state, evidence) = match running {
4356                Some(instance) => (GatewayState::Up, instance.evidence.clone()),
4357                None if !installed => (
4358                    GatewayState::Unknown,
4359                    format!("`{id}` is not installed; no gateway to probe"),
4360                ),
4361                None if id == HarnessId::HERMES => match hermes_gateway_status() {
4362                    // The harness's own door outranks the WAL heuristic: an idle
4363                    // gateway writes nothing for minutes yet is up.
4364                    Some((state, evidence)) => (state, evidence),
4365                    None => (
4366                        GatewayState::Down,
4367                        "no fresh state.db-wal activity under ~/.hermes and `hermes gateway status` gave no verdict".to_string(),
4368                    ),
4369                },
4370                None => (
4371                    GatewayState::Down,
4372                    format!(
4373                        "no TCP listener at {}",
4374                        endpoint.as_deref().unwrap_or("the gateway endpoint")
4375                    ),
4376                ),
4377            };
4378            GatewayHealth {
4379                state,
4380                endpoint,
4381                version: version.map(str::to_string),
4382                evidence,
4383                checked_at_ms,
4384            }
4385        }
4386        // ORC-7: the orchestrator's gateway IS its daemon, and the daemon's
4387        // own lease file is the record of it. A lease naming a live pid is
4388        // up; a lease whose process is gone is down and says so as a STALE
4389        // lease, never as "no lease"; no lease at all is down. Nothing is
4390        // started, and no port is guessed — the daemon multiplexes adapters
4391        // the way Hermes does, so it has no single endpoint either.
4392        HarnessId::ORCHESTRATOR => {
4393            let root = crate::HarnessHomes::default().orchestrator;
4394            let (state, evidence) = match crate::orchestrator::read_lease(&root) {
4395                Some(lease) if lease.is_live() => (
4396                    GatewayState::Up,
4397                    format!(
4398                        "`{}` names pid {} (started {}), which is live",
4399                        crate::orchestrator::lock_path(&root).display(),
4400                        lease.pid,
4401                        lease.started_at
4402                    ),
4403                ),
4404                Some(lease) => (
4405                    GatewayState::Down,
4406                    format!(
4407                        "stale lease `{}`: pid {} is gone",
4408                        crate::orchestrator::lock_path(&root).display(),
4409                        lease.pid
4410                    ),
4411                ),
4412                None => (
4413                    GatewayState::Down,
4414                    format!(
4415                        "no lease at `{}`; `supercode orchestrator start` writes one",
4416                        crate::orchestrator::lock_path(&root).display()
4417                    ),
4418                ),
4419            };
4420            GatewayHealth {
4421                state,
4422                endpoint: None,
4423                version: version.map(str::to_string),
4424                evidence,
4425                checked_at_ms,
4426            }
4427        }
4428        _ => GatewayHealth {
4429            state: GatewayState::Unknown,
4430            endpoint: None,
4431            version: version.map(str::to_string),
4432            evidence: format!("`{id}` runs per session, not as a gateway"),
4433            checked_at_ms,
4434        },
4435    }
4436}
4437
4438#[derive(Debug, Clone, Serialize)]
4439struct RunningInstance {
4440    /// How the instance was detected.
4441    method: RunningInstanceMethod,
4442    /// The evidence the verdict rests on (endpoint reached / WAL path+age).
4443    evidence: String,
4444    /// Epoch-ms instant the probe executed.
4445    checked_at_ms: u64,
4446}
4447
4448#[derive(Debug, Clone, Copy, Serialize)]
4449#[serde(rename_all = "snake_case")]
4450enum RunningInstanceMethod {
4451    /// A TCP connect to the harness's own configured gateway endpoint
4452    /// succeeded.
4453    GatewayConnect,
4454    /// The harness's session store has an active SQLite WAL (a live writer
4455    /// holds the store open and stamped it recently).
4456    StoreWalActivity,
4457}
4458
4459fn now_epoch_ms() -> u64 {
4460    std::time::SystemTime::now()
4461        .duration_since(std::time::UNIX_EPOCH)
4462        .map(|elapsed| elapsed.as_millis() as u64)
4463        .unwrap_or(0)
4464}
4465
4466/// OpenClaw: the gateway endpoint comes from the harness's OWN config
4467/// (`<home>/.openclaw/openclaw.json` — `gateway.url` or `gateway.port`,
4468/// default port 18789); a successful TCP connect is the running signal.
4469fn probe_openclaw_running(home: &Path) -> Option<RunningInstance> {
4470    let config_path = home.join(".openclaw/openclaw.json");
4471    let text = std::fs::read_to_string(&config_path).ok();
4472    let gateway = text
4473        .as_deref()
4474        .and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
4475        .and_then(|config| config.get("gateway").cloned());
4476    let address = gateway
4477        .as_ref()
4478        .and_then(|gateway| gateway.get("url"))
4479        .and_then(serde_json::Value::as_str)
4480        .and_then(|url| {
4481            url.split("://").nth(1).map(|rest| {
4482                rest.trim_end_matches('/')
4483                    .split('/')
4484                    .next()
4485                    .unwrap_or(rest)
4486                    .to_string()
4487            })
4488        })
4489        .unwrap_or_else(|| {
4490            let port = gateway
4491                .as_ref()
4492                .and_then(|gateway| gateway.get("port"))
4493                .and_then(serde_json::Value::as_u64)
4494                .unwrap_or(18789);
4495            format!("127.0.0.1:{port}")
4496        });
4497    let reachable = std::net::TcpStream::connect_timeout(
4498        &address.parse().ok()?,
4499        std::time::Duration::from_millis(400),
4500    )
4501    .is_ok();
4502    reachable.then(|| RunningInstance {
4503        method: RunningInstanceMethod::GatewayConnect,
4504        evidence: format!(
4505            "gateway endpoint {address} accepted a TCP connect (from {})",
4506            config_path.display()
4507        ),
4508        checked_at_ms: now_epoch_ms(),
4509    })
4510}
4511
4512/// Hermes: `<home>/.hermes/state.db-wal` freshly modified means a live writer
4513/// holds the store open (SQLite WAL exists only while a connection is open;
4514/// a recent stamp distinguishes an active instance from a stale crash
4515/// leftover).
4516fn probe_hermes_running(home: &Path, max_wal_age_ms: u64) -> Option<RunningInstance> {
4517    let wal = home.join(".hermes/state.db-wal");
4518    let modified = std::fs::metadata(&wal).ok()?.modified().ok()?;
4519    let age_ms = std::time::SystemTime::now()
4520        .duration_since(modified)
4521        .map(|age| age.as_millis() as u64)
4522        .unwrap_or(u64::MAX);
4523    (age_ms <= max_wal_age_ms).then(|| RunningInstance {
4524        method: RunningInstanceMethod::StoreWalActivity,
4525        evidence: format!(
4526            "{} stamped {age_ms}ms ago (threshold {max_wal_age_ms}ms)",
4527            wal.display()
4528        ),
4529        checked_at_ms: now_epoch_ms(),
4530    })
4531}
4532
4533/// Default-on running-instance detection for the harnesses that have one.
4534fn probe_running_instance(id: &str) -> Option<RunningInstance> {
4535    let home = supercode_interchange::user_home()
4536        .map(std::path::PathBuf::into_os_string)
4537        .map(PathBuf::from)?;
4538    match id {
4539        HarnessId::OPENCLAW => probe_openclaw_running(&home),
4540        HarnessId::HERMES => probe_hermes_running(&home, 300_000),
4541        _ => None,
4542    }
4543}
4544
4545#[derive(Serialize)]
4546struct LocalHarness {
4547    id: HarnessId,
4548    display_name: String,
4549    supported: bool,
4550    installed: bool,
4551    executable: Option<String>,
4552    version: Option<String>,
4553    auth: HarnessAuthState,
4554    runtime: HarnessRuntimeState,
4555    protocol: String,
4556    capabilities: crate::RuntimeCapabilities,
4557    effective_capabilities: crate::RuntimeCapabilities,
4558    sessions: HarnessSessionCounts,
4559    /// Receipt-backed running-instance detection (None = not detected or the
4560    /// harness has no running-instance concept). Distinct from `installed`.
4561    #[serde(skip_serializing_if = "Option::is_none")]
4562    running: Option<RunningInstance>,
4563    /// ORCH-17: gateway health derived from `running` + the harness's own config.
4564    gateway: GatewayHealth,
4565    reason: Option<String>,
4566    repair: Option<String>,
4567}
4568
4569#[derive(Clone, Deserialize)]
4570struct RuntimeBackendParams {
4571    harness: HarnessId,
4572    #[serde(default)]
4573    protocol: Option<String>,
4574    #[serde(default)]
4575    launch: Option<RuntimeLaunch>,
4576    #[serde(default)]
4577    base_url: Option<String>,
4578    #[serde(default)]
4579    policy: RuntimePolicy,
4580}
4581
4582/// A session supercode starts or resumes runs without approval prompts
4583/// (`yolo`) unless the caller asks for the harness's own (`default`).
4584#[derive(Debug, Clone, Copy, Default, Deserialize)]
4585#[serde(rename_all = "snake_case")]
4586enum RuntimePolicy {
4587    Default,
4588    #[default]
4589    Yolo,
4590}
4591
4592#[derive(Deserialize)]
4593struct RuntimeStartParams {
4594    #[serde(flatten)]
4595    backend: RuntimeBackendParams,
4596    cwd: PathBuf,
4597    /// MCP servers to mount into the new session through the harness's own
4598    /// start door (ORC-6). Backends without such a door ignore them.
4599    #[serde(default)]
4600    mcp_servers: Vec<crate::McpServerLaunch>,
4601    /// The session's approval policy, where the harness's start door takes one (Codex).
4602    #[serde(default)]
4603    approval_policy: Option<String>,
4604}
4605
4606#[derive(Deserialize)]
4607struct RuntimeAttachParams {
4608    #[serde(flatten)]
4609    backend: RuntimeBackendParams,
4610    runtime_id: String,
4611    #[serde(default)]
4612    cwd: Option<PathBuf>,
4613    /// MCP servers to mount into the resumed session (the start door's own
4614    /// field, carried again because a session's tools die with its process).
4615    #[serde(default)]
4616    mcp_servers: Vec<crate::McpServerLaunch>,
4617    /// The session's approval policy, carried again on resume as on start (Codex).
4618    #[serde(default)]
4619    approval_policy: Option<String>,
4620}
4621
4622#[derive(Deserialize)]
4623struct RuntimeConnectionParams {
4624    connection: String,
4625}
4626
4627#[derive(Deserialize)]
4628struct RuntimeInputParams {
4629    connection: String,
4630    text: String,
4631    #[serde(default)]
4632    image_urls: Vec<String>,
4633}
4634
4635const MAX_RUNTIME_IMAGES: usize = 4;
4636const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
4637const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
4638
4639fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
4640    if image_urls.len() > MAX_RUNTIME_IMAGES {
4641        return Err(ServiceError::InvalidParams(format!(
4642            "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
4643        )));
4644    }
4645    let mut total = 0usize;
4646    for url in &image_urls {
4647        if !(url.starts_with("data:image/")
4648            || url.starts_with("https://")
4649            || url.starts_with("http://"))
4650        {
4651            return Err(ServiceError::InvalidParams(
4652                "runtime images must be image data URLs or HTTP(S) URLs".into(),
4653            ));
4654        }
4655        if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
4656            return Err(ServiceError::InvalidParams(format!(
4657                "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
4658            )));
4659        }
4660        total = total.saturating_add(url.len());
4661    }
4662    if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
4663        return Err(ServiceError::InvalidParams(format!(
4664            "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
4665        )));
4666    }
4667    Ok(image_urls)
4668}
4669
4670#[derive(Deserialize)]
4671struct RuntimeRespondParams {
4672    connection: String,
4673    request_id: Value,
4674    response: Value,
4675}
4676
4677fn default_reduction_store_root() -> PathBuf {
4678    if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
4679        return PathBuf::from(root).join("sessions");
4680    }
4681    if let Some(home) = supercode_interchange::user_home().map(std::path::PathBuf::into_os_string) {
4682        return PathBuf::from(home).join(".supercode").join("sessions");
4683    }
4684    PathBuf::from(".supercode").join("sessions")
4685}
4686
4687fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
4688    let mut output = String::new();
4689    for message in messages {
4690        output.push_str(
4691            &serde_json::to_string(message)
4692                .map_err(|error| ServiceError::Operation(error.to_string()))?,
4693        );
4694        output.push('\n');
4695    }
4696    Ok(output)
4697}
4698
4699fn parse_messages_jsonl(
4700    content: &str,
4701) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
4702    content
4703        .lines()
4704        .enumerate()
4705        .filter(|(_, line)| !line.trim().is_empty())
4706        .map(|(index, line)| {
4707            serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
4708                ServiceError::Operation(format!(
4709                    "reduced transcript line {} is invalid: {error}",
4710                    index + 1
4711                ))
4712            })
4713        })
4714        .collect()
4715}
4716
4717fn reduced_bootstrap_prompt(
4718    source: &SessionLocator,
4719    target: TransferFormat,
4720    view_jsonl: &str,
4721    sidecar_path: &Path,
4722    reduction_log_path: &Path,
4723) -> String {
4724    format!(
4725        "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
4726         \n\
4727         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\
4728         \n\
4729         <supercode-reduced-session source-session=\"{source_id}\">\n\
4730         {view_jsonl}\
4731         </supercode-reduced-session>\n\
4732         \n\
4733         Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
4734        source_harness = source.harness.as_str(),
4735        target_harness = target.id(),
4736        sidecar = sidecar_path.display(),
4737        log = reduction_log_path.display(),
4738        source_id = source.session_id,
4739    )
4740}
4741
4742fn session_artifact(
4743    locator: &SessionLocator,
4744    session: &Session,
4745    target: TransferFormat,
4746) -> std::result::Result<SessionArtifact, ServiceError> {
4747    session_artifact_with_id(locator, session, target, None)
4748}
4749
4750fn session_artifact_with_id(
4751    locator: &SessionLocator,
4752    session: &Session,
4753    target: TransferFormat,
4754    target_session_id: Option<&str>,
4755) -> std::result::Result<SessionArtifact, ServiceError> {
4756    let format: SessionFormat = target.into();
4757    let diagonal = format.source() == session.meta.source;
4758    crate::residue_store::store_segments(session);
4759    let has_appended_turns = session
4760        .imported_message_count
4761        .is_some_and(|imported| imported < session.messages.len());
4762    let mut restoration = None;
4763    let content = if let Some(id) = target_session_id {
4764        if diagonal && format != SessionFormat::OpenCode {
4765            session
4766                .to_jsonl_spliced(format, Some(id))
4767                .map_err(operation)?
4768        } else {
4769            let mut rewritten = session.clone();
4770            rewritten.meta.session_id = Some(id.to_string());
4771            rewritten.to_jsonl(format).map_err(operation)?
4772        }
4773    } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
4774        session.raw_verbatim()
4775    } else if diagonal {
4776        session.to_jsonl_spliced(format, None).map_err(operation)?
4777    } else {
4778        // A session that came from `format` before returns its source records verbatim for the
4779        // prefix the residue store holds (docs/plans/portable-residue.md).
4780        match session
4781            .restore_residue(format, crate::residue_store::lookup)
4782            .map_err(operation)?
4783        {
4784            Some((content, report)) => {
4785                restoration = Some(report);
4786                content
4787            }
4788            None => session.to_jsonl(format).map_err(operation)?,
4789        }
4790    };
4791    let stem = sanitize_filename(
4792        target_session_id
4793            .or(session.meta.session_id.as_deref())
4794            .unwrap_or(&locator.session_id),
4795    );
4796    let suggested_filename = if diagonal && target == TransferFormat::Grok {
4797        "chat_history.jsonl".to_string()
4798    } else if target == TransferFormat::Goose {
4799        format!("{stem}.goose.json")
4800    } else {
4801        format!("{stem}.{}.jsonl", target.id())
4802    };
4803    let mut files = vec![SessionArtifactFile {
4804        path: suggested_filename.clone(),
4805        content: content.clone(),
4806        role: ArtifactFileRole::Primary,
4807    }];
4808    if target == TransferFormat::ClaudeCode {
4809        let bundle_stem = Path::new(&suggested_filename)
4810            .file_stem()
4811            .and_then(|stem| stem.to_str())
4812            .unwrap_or(&stem);
4813        let mut child_paths = BTreeSet::new();
4814        for (index, subagent) in session.subagents.iter().enumerate() {
4815            let agent_id = subagent
4816                .meta
4817                .agent_id
4818                .as_deref()
4819                .map(|id| id.strip_prefix("agent-").unwrap_or(id))
4820                .map(sanitize_filename)
4821                .filter(|id| !id.is_empty())
4822                .unwrap_or_else(|| format!("subagent-{}", index + 1));
4823            let child_has_appended_turns = subagent
4824                .imported_message_count
4825                .is_some_and(|imported| imported < subagent.messages.len());
4826            let child_content = if target_session_id.is_none()
4827                && subagent.meta.source == SessionSource::ClaudeCode
4828                && subagent.raw_is_verbatim
4829                && !child_has_appended_turns
4830            {
4831                subagent.raw_verbatim()
4832            } else if subagent.meta.source == SessionSource::ClaudeCode {
4833                subagent
4834                    .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
4835                    .map_err(operation)?
4836            } else {
4837                let mut child = subagent.clone();
4838                if let Some(id) = target_session_id {
4839                    child.meta.session_id = Some(id.to_string());
4840                }
4841                child
4842                    .to_jsonl(SessionFormat::ClaudeCode)
4843                    .map_err(operation)?
4844            };
4845            let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
4846            if !child_paths.insert(path.clone()) {
4847                return Err(ServiceError::Operation(format!(
4848                    "Claude subagent ids collide at artifact path `{path}`"
4849                )));
4850            }
4851            files.push(SessionArtifactFile {
4852                path,
4853                content: child_content,
4854                role: ArtifactFileRole::Subagent,
4855            });
4856        }
4857    }
4858    if diagonal && target == TransferFormat::Grok {
4859        append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
4860    }
4861    if !diagonal || !session.raw_is_verbatim {
4862        files.push(SessionArtifactFile {
4863            path: "recovery/source.supercode.jsonl".into(),
4864            content: session.to_native_jsonl(),
4865            role: ArtifactFileRole::SourceRecovery,
4866        });
4867        for (index, subagent) in session.subagents.iter().enumerate() {
4868            let id = subagent
4869                .meta
4870                .agent_id
4871                .as_deref()
4872                .map(sanitize_filename)
4873                .unwrap_or_else(|| format!("subagent-{}", index + 1));
4874            files.push(SessionArtifactFile {
4875                path: format!("recovery/subagents/{id}.supercode.jsonl"),
4876                content: subagent.to_native_jsonl(),
4877                role: ArtifactFileRole::SourceRecovery,
4878            });
4879        }
4880    }
4881    if !diagonal && session.meta.source == SessionSource::Grok {
4882        append_grok_bundle_files(
4883            locator,
4884            "recovery/grok/",
4885            ArtifactFileRole::SourceRecovery,
4886            &mut files,
4887        )?;
4888    }
4889    let (fidelity, residue) = if diagonal
4890        && target_session_id.is_none()
4891        && session.raw_is_verbatim
4892        && !has_appended_turns
4893    {
4894        (Fidelity::ByteLossless, Vec::new())
4895    } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
4896        (
4897            Fidelity::ValueLossless,
4898            vec![if target_session_id.is_some() {
4899                "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
4900            } else {
4901                "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
4902            }],
4903        )
4904    } else {
4905        match restoration {
4906            Some(report) if report.rendered_messages == 0 => (
4907                Fidelity::ByteLossless,
4908                vec![format!(
4909                    "restored verbatim from this conversation's {} source records in the residue store",
4910                    target.id()
4911                )],
4912            ),
4913            Some(report) => (
4914                Fidelity::Semantic,
4915                vec![format!(
4916                    "{} of {} messages restored verbatim from the residue store; the other {} written by the {} writer",
4917                    report.restored_messages,
4918                    report.restored_messages + report.rendered_messages,
4919                    report.rendered_messages,
4920                    target.id()
4921                )],
4922            ),
4923            None => (
4924                Fidelity::Semantic,
4925                vec!["target schema has no portable slot for every source-native record and metadata field".into()],
4926            ),
4927        }
4928    };
4929    Ok(SessionArtifact {
4930        source_harness: locator.harness.clone(),
4931        target_harness: target.id(),
4932        session_id: target_session_id
4933            .map(str::to_string)
4934            .or_else(|| session.meta.session_id.clone()),
4935        content,
4936        suggested_filename,
4937        files,
4938        fidelity,
4939        residue,
4940    })
4941}
4942
4943fn append_grok_bundle_files(
4944    locator: &SessionLocator,
4945    prefix: &str,
4946    role: ArtifactFileRole,
4947    files: &mut Vec<SessionArtifactFile>,
4948) -> std::result::Result<(), ServiceError> {
4949    let primary = locator.storage.path();
4950    if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
4951        return Err(ServiceError::Operation(format!(
4952            "Grok bundle locator must name chat_history.jsonl, got {}",
4953            primary.display()
4954        )));
4955    }
4956    let parent = primary.parent().ok_or_else(|| {
4957        ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
4958    })?;
4959    for name in ["summary.json", "updates.jsonl"] {
4960        let path = parent.join(name);
4961        let metadata = match std::fs::symlink_metadata(&path) {
4962            Ok(metadata) => metadata,
4963            Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
4964            Err(error) => return Err(ServiceError::Operation(error.to_string())),
4965        };
4966        if metadata.file_type().is_symlink() || !metadata.is_file() {
4967            return Err(ServiceError::Operation(format!(
4968                "refusing non-regular Grok bundle member {}",
4969                path.display()
4970            )));
4971        }
4972        let content = std::fs::read_to_string(&path).map_err(|error| {
4973            ServiceError::Operation(format!(
4974                "Grok bundle member {} is not representable as UTF-8: {error}",
4975                path.display()
4976            ))
4977        })?;
4978        files.push(SessionArtifactFile {
4979            path: format!("{prefix}{name}"),
4980            content,
4981            role: match role {
4982                ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
4983                _ => ArtifactFileRole::SourceRecovery,
4984            },
4985        });
4986    }
4987    Ok(())
4988}
4989
4990fn handoff_artifact(
4991    locator: &SessionLocator,
4992    session: &Session,
4993    target: TransferFormat,
4994) -> std::result::Result<SessionArtifact, ServiceError> {
4995    let target_session_id = target_session_id(target);
4996    session_artifact_with_id(locator, session, target, Some(&target_session_id))
4997}
4998
4999fn target_session_id(target: TransferFormat) -> String {
5000    let uuid = generated_session_id();
5001    match target {
5002        TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
5003        TransferFormat::ClaudeCode
5004        | TransferFormat::Codex
5005        | TransferFormat::Pi
5006        | TransferFormat::Grok
5007        | TransferFormat::Gemini
5008        | TransferFormat::Goose
5009        | TransferFormat::Hermes => uuid,
5010    }
5011}
5012
5013fn sanitize_filename(value: &str) -> String {
5014    let value = value
5015        .chars()
5016        .map(|character| {
5017            if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
5018                character
5019            } else {
5020                '-'
5021            }
5022        })
5023        .collect::<String>();
5024    let value = value.trim_matches('-');
5025    if value.is_empty() {
5026        "session".into()
5027    } else {
5028        value.chars().take(100).collect()
5029    }
5030}
5031
5032fn handoff_instructions(
5033    target: TransferFormat,
5034    session_id: &str,
5035    cwd: &Path,
5036) -> HandoffInstructions {
5037    let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
5038        cwd: cwd.to_path_buf(),
5039        program: program.into(),
5040        arguments,
5041        env: BTreeMap::new(),
5042    };
5043    match target {
5044        TransferFormat::ClaudeCode => HandoffInstructions {
5045            launch: launch("claude", vec!["--resume".into(), session_id.into()]),
5046            materialize: None,
5047            requires_materialization: true,
5048            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(),
5049        },
5050        TransferFormat::Hermes => HandoffInstructions {
5051            launch: launch("hermes", vec!["--resume".into(), session_id.into()]),
5052            materialize: None,
5053            requires_materialization: true,
5054            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(),
5055        },
5056        TransferFormat::Codex => HandoffInstructions {
5057            launch: launch("codex", vec!["resume".into(), session_id.into()]),
5058            materialize: None,
5059            requires_materialization: true,
5060            note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
5061        },
5062        TransferFormat::OpenCode => HandoffInstructions {
5063            launch: launch("opencode", vec!["--session".into(), session_id.into()]),
5064            materialize: Some(launch(
5065                "opencode",
5066                vec!["import".into(), "{artifact_path}".into()],
5067            )),
5068            requires_materialization: true,
5069            note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
5070        },
5071        TransferFormat::Pi => HandoffInstructions {
5072            launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
5073            materialize: None,
5074            requires_materialization: true,
5075            note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
5076        },
5077        TransferFormat::Grok => HandoffInstructions {
5078            launch: launch(
5079                "grok",
5080                vec!["--resume".into(), "{materialized_session_id}".into()],
5081            ),
5082            materialize: None,
5083            requires_materialization: true,
5084            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(),
5085        },
5086        TransferFormat::Gemini => HandoffInstructions {
5087            launch: launch(
5088                "gemini",
5089                vec!["--session-file".into(), "{artifact_path}".into()],
5090            ),
5091            materialize: None,
5092            requires_materialization: true,
5093            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(),
5094        },
5095        TransferFormat::Goose => HandoffInstructions {
5096            launch: launch(
5097                "goose",
5098                vec![
5099                    "session".into(),
5100                    "--resume".into(),
5101                    "--session-id".into(),
5102                    "{imported_session_id}".into(),
5103                ],
5104            ),
5105            materialize: Some(launch(
5106                "goose",
5107                vec!["session".into(), "import".into(), "{artifact_path}".into()],
5108            )),
5109            requires_materialization: true,
5110            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(),
5111        },
5112    }
5113}
5114
5115fn resume_launch(
5116    harness: &str,
5117    session_id: &str,
5118    cwd: &Path,
5119    policy: ResumePolicy,
5120) -> std::result::Result<StructuredLaunch, ServiceError> {
5121    let mut arguments = Vec::new();
5122    let program = match harness {
5123        HarnessId::GROK => {
5124            if matches!(policy, ResumePolicy::Yolo) {
5125                if crate::support::self_sandbox_supported() {
5126                    arguments.extend(["--sandbox".into(), "workspace".into()]);
5127                }
5128                arguments.push("--always-approve".into());
5129            }
5130            arguments.extend(["--resume".into(), session_id.into()]);
5131            "grok"
5132        }
5133        HarnessId::CODEX => {
5134            arguments.extend(crate::startup_prompts::startup_arguments(
5135                harness,
5136                Some(cwd),
5137                &[],
5138                matches!(policy, ResumePolicy::Yolo),
5139            ));
5140            arguments.extend(["resume".into(), session_id.into()]);
5141            "codex"
5142        }
5143        HarnessId::CLAUDE_CODE => {
5144            arguments.extend(crate::startup_prompts::startup_arguments(
5145                harness,
5146                Some(cwd),
5147                &[],
5148                matches!(policy, ResumePolicy::Yolo),
5149            ));
5150            arguments.extend(["--resume".into(), session_id.into()]);
5151            "claude"
5152        }
5153        HarnessId::GEMINI => {
5154            arguments.extend(crate::startup_prompts::startup_arguments(
5155                harness,
5156                Some(cwd),
5157                &[],
5158                matches!(policy, ResumePolicy::Yolo),
5159            ));
5160            arguments.extend(["--resume".into(), session_id.into()]);
5161            "gemini"
5162        }
5163        HarnessId::GOOSE => {
5164            arguments.extend([
5165                "session".into(),
5166                "--resume".into(),
5167                "--session-id".into(),
5168                session_id.into(),
5169            ]);
5170            "goose"
5171        }
5172        HarnessId::PI => {
5173            arguments.extend(crate::startup_prompts::startup_arguments(
5174                harness,
5175                Some(cwd),
5176                &[],
5177                matches!(policy, ResumePolicy::Yolo),
5178            ));
5179            arguments.extend(["--session".into(), session_id.into()]);
5180            "pi"
5181        }
5182        HarnessId::OPENCODE => {
5183            arguments.extend(["--session".into(), session_id.into()]);
5184            "opencode"
5185        }
5186        HarnessId::SUPERCODE => {
5187            if matches!(policy, ResumePolicy::Yolo) {
5188                arguments.push("--dangerous".into());
5189            }
5190            arguments.extend(["resume".into(), session_id.into()]);
5191            "supercode"
5192        }
5193        other => {
5194            return Err(ServiceError::InvalidParams(format!(
5195                "no structured resume launch is registered for harness `{other}`"
5196            )))
5197        }
5198    };
5199    Ok(StructuredLaunch {
5200        cwd: cwd.to_path_buf(),
5201        env: if program == "grok" {
5202            crate::support::grok_home_env()
5203        } else {
5204            BTreeMap::new()
5205        },
5206        program: program.into(),
5207        arguments,
5208    })
5209}
5210
5211/// Stage the resolved gateway credential in a private (0600) file so the
5212/// bridge can read it via `--token-file` — the delivery the real `openclaw
5213/// acp` accepts. One stable file per endpoint (keyed by an address digest,
5214/// no secret material in the name), overwritten on every connect so files
5215/// never accumulate and a rotated token never goes stale on disk.
5216fn openclaw_gateway_token_file(address: &str, secret: &str) -> std::io::Result<PathBuf> {
5217    let digest = blake3::hash(address.as_bytes()).to_hex();
5218    let path = std::env::temp_dir().join(format!(
5219        "supercode-openclaw-gateway-token-{}",
5220        &digest.as_str()[..16]
5221    ));
5222    #[cfg(unix)]
5223    {
5224        use std::io::Write;
5225        use std::os::unix::fs::OpenOptionsExt;
5226        let mut file = std::fs::OpenOptions::new()
5227            .write(true)
5228            .create(true)
5229            .truncate(true)
5230            .mode(0o600)
5231            .open(&path)?;
5232        file.write_all(secret.as_bytes())?;
5233    }
5234    #[cfg(not(unix))]
5235    std::fs::write(&path, secret)?;
5236    Ok(path)
5237}
5238
5239/// Open a connect-mode descriptor: resolve the endpoint address and
5240/// credential from the harness's own config file and build the backend that
5241/// joins the already-running endpoint. Fails closed with a specific
5242/// diagnostic when the config cannot be resolved or the declared protocol has
5243/// no connect-capable client yet.
5244fn open_connect_descriptor(
5245    descriptor: &crate::HarnessSupportDescriptor,
5246    home: &Path,
5247) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5248    let Some(connect) = &descriptor.runtime.connect_launch else {
5249        return Err(ServiceError::InvalidParams(format!(
5250            "harness `{}` has no registered connect-mode launch",
5251            descriptor.id.as_str()
5252        )));
5253    };
5254    let resolved = connect
5255        .resolve(home)
5256        .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
5257    match (descriptor.id.as_str(), connect.protocol.as_str()) {
5258        (HarnessId::OPENCODE, protocol) if protocol.starts_with("opencode-http") => {
5259            let mut backend = OpenCodeRuntimeBackend::connect(&resolved.address);
5260            if let Some(token) = resolved.auth {
5261                backend = backend.with_bearer(token);
5262            }
5263            Ok(Box::new(backend))
5264        }
5265        (HarnessId::OPENCLAW, protocol) if protocol.starts_with("acp") => {
5266            // OpenClaw's own `openclaw acp` binary is the gateway client: a
5267            // stdio ACP bridge that joins the RUNNING gateway at the resolved
5268            // endpoint. Blind-walk finding 2026-08-31: the real bridge does
5269            // NOT honor OPENCLAW_GATEWAY_TOKEN from the environment — the
5270            // credential must arrive via `--token-file` (never bare `--token`
5271            // on argv, where process listings could read it). The env var is
5272            // still set for older bridges that did read it. Requires openclaw
5273            // >= 2026.7: the 2026.2 bridge drops its gateway socket
5274            // mid-prompt and advertises no session resume (executed finding,
5275            // docs/interop/research/openclaw-acp-dialect-2026-08-30.json).
5276            let mut env = BTreeMap::new();
5277            let mut arguments = vec!["acp".into(), "--url".into(), resolved.address.clone()];
5278            if let Some(token) = resolved.auth {
5279                let token_path = openclaw_gateway_token_file(&resolved.address, token.secret())
5280                    .map_err(|error| {
5281                        ServiceError::UnsupportedAction(format!(
5282                            "could not stage the gateway credential for the bridge: {error}"
5283                        ))
5284                    })?;
5285                arguments.push("--token-file".into());
5286                arguments.push(token_path.to_string_lossy().into_owned());
5287                env.insert("OPENCLAW_GATEWAY_TOKEN".to_string(), token.secret().to_string());
5288            }
5289            // The bridge program comes from the descriptor's own default
5290            // launch (the compiled registry pins `openclaw`), so tests can
5291            // substitute an absolute mock-bridge path without touching
5292            // process-global state.
5293            let program = descriptor
5294                .runtime
5295                .default_launch
5296                .as_ref()
5297                .map(|launch| launch.program.clone())
5298                .unwrap_or_else(|| "openclaw".into());
5299            let launch = RuntimeLaunch {
5300                program,
5301                arguments,
5302                env,
5303            };
5304            Ok(Box::new(
5305                crate::AcpRuntimeBackend::new(descriptor.id.clone(), launch)
5306                    .with_resume_support(descriptor.runtime.capabilities.resume_session),
5307            ))
5308        }
5309        _ => Err(ServiceError::UnsupportedAction(format!(
5310            "connect-mode endpoint for `{}` speaks `{}`; joining it needs that protocol's gateway client",
5311            descriptor.id.as_str(),
5312            connect.protocol
5313        ))),
5314    }
5315}
5316
5317/// The registry's connect-mode launch for this harness, honored only when the
5318/// caller supplied neither an explicit launch nor a base URL.
5319fn registry_connect_descriptor(
5320    params: &RuntimeBackendParams,
5321) -> Option<crate::HarnessSupportDescriptor> {
5322    if params.launch.is_some() || params.base_url.is_some() {
5323        return None;
5324    }
5325    harness_support_registry()
5326        .harnesses
5327        .into_iter()
5328        .find(|descriptor| descriptor.id == params.harness)
5329        .filter(|descriptor| descriptor.runtime.connect_launch.is_some())
5330}
5331
5332fn service_home() -> std::result::Result<PathBuf, ServiceError> {
5333    supercode_interchange::user_home()
5334        .map(std::path::PathBuf::into_os_string)
5335        .map(PathBuf::from)
5336        .ok_or_else(|| {
5337            ServiceError::UnsupportedAction(
5338                "connect-mode launches need HOME to locate the harness config".into(),
5339            )
5340        })
5341}
5342
5343/// The doors that open a runtime: each spawns or joins a program and waits on
5344/// that program's protocol handshake before it can answer.
5345pub const RUNTIME_OPEN_METHODS: &[&str] = &[
5346    "harness.v1.runtimes.start",
5347    "harness.v1.runtimes.resume",
5348    "harness.v1.runtimes.attach",
5349    "harness.v1.runtimes.attach_existing",
5350];
5351
5352/// How long a runtime gets to finish opening before its caller is answered an
5353/// error instead. A program that never speaks the protocol at all — the wrong
5354/// binary, a shim that prints usage and waits — never answers the handshake,
5355/// so the wait is unbounded without this.
5356pub const RUNTIME_OPEN_DEADLINE: Duration = Duration::from_secs(60);
5357
5358/// How long a control call on an ALREADY-open runtime — send input, interrupt,
5359/// steer, respond, close — gets before its caller is answered an error
5360/// instead. A live runtime answers these in milliseconds; a wedged one never
5361/// answers at all, and `close` is exactly what a caller reaches for when it
5362/// suspects that.
5363pub const RUNTIME_CONTROL_DEADLINE: Duration = Duration::from_secs(30);
5364
5365/// The doors whose work happens entirely OUTSIDE this service's state once
5366/// its state has been read: probing harnesses, relaying a message into a
5367/// live session, and performing a conversation verb through a harness's own
5368/// CLI / HTTP / store door. Every one of them waits on a child process or a
5369/// network peer. See [`HarnessSessionService::detach`].
5370pub const DETACHED_METHODS: &[&str] = &[
5371    "harness.v1.harnesses.list",
5372    "harness.v1.harnesses.probe",
5373    "harness.v1.sessions.message",
5374    "harness.v1.sessions.new",
5375    "harness.v1.sessions.reset",
5376    "harness.v1.sessions.archive",
5377    "harness.v1.sessions.delete",
5378];
5379
5380/// How long a request moved off a transport's loop gets before its caller is
5381/// answered an error instead. Each of these already bounds its own inner
5382/// waits (a probe's handshake, the relay's send); this is the backstop for
5383/// the ones that do not — a harness CLI that never exits — so no caller waits
5384/// forever on a detached task no one is watching.
5385pub const DETACHED_CALL_DEADLINE: Duration = Duration::from_secs(120);
5386
5387/// How long `sessions.discover` gets before its caller is answered an error
5388/// instead. Discovery reads each harness's own store, and a store on a cold
5389/// or unavailable mount answers at the filesystem's pace rather than its own.
5390///
5391/// Deliberately shorter than the clients' own request deadline (30s): the
5392/// server's answer names the store that did not answer, and it is only read
5393/// if it lands before the client stops listening.
5394pub const SESSION_DISCOVER_DEADLINE: Duration = Duration::from_secs(25);
5395
5396/// Bound one control call on an open runtime by [`RUNTIME_CONTROL_DEADLINE`],
5397/// naming the method and the bound when it blows.
5398async fn within_control_deadline<F: std::future::Future>(
5399    method: &str,
5400    call: F,
5401) -> std::result::Result<F::Output, ServiceError> {
5402    tokio::time::timeout(RUNTIME_CONTROL_DEADLINE, call)
5403        .await
5404        .map_err(|_| {
5405            ServiceError::Operation(format!(
5406                "`{method}` gave up after {}s: the runtime did not answer",
5407                RUNTIME_CONTROL_DEADLINE.as_secs()
5408            ))
5409        })
5410}
5411
5412/// One [`RUNTIME_OPEN_METHODS`] request, parsed but not yet started. See
5413/// [`HarnessSessionService::runtime_open`] for why it exists apart from
5414/// [`HarnessSessionService::handle_async`].
5415pub struct RuntimeOpen {
5416    id: Value,
5417    method: String,
5418    params: Value,
5419}
5420
5421impl RuntimeOpen {
5422    /// Do the waiting: spawn or join the program and complete its handshake,
5423    /// bounded by [`RUNTIME_OPEN_DEADLINE`]. Touches no service state, so this
5424    /// runs on any task.
5425    pub async fn open(self) -> OpenedRuntime {
5426        let Self { id, method, params } = self;
5427        let outcome = open_runtime(&method, params).await;
5428        OpenedRuntime { id, outcome }
5429    }
5430}
5431
5432/// The result of [`RuntimeOpen::open`], ready for
5433/// [`HarnessSessionService::finish_runtime_open`].
5434pub struct OpenedRuntime {
5435    id: Value,
5436    outcome: std::result::Result<OpenRuntime, ServiceError>,
5437}
5438
5439/// One detached request: the half that reads this service's state already
5440/// done, and the half that waits not yet started. See
5441/// [`HarnessSessionService::detach`] and
5442/// [`HarnessSessionService::detach_runtime`].
5443pub struct DetachedCall {
5444    id: Value,
5445    method: String,
5446    work: std::result::Result<Work, ServiceError>,
5447}
5448
5449impl DetachedCall {
5450    /// Do the waiting and answer. Runs on any task: whatever this call needed
5451    /// from the service was taken before it left.
5452    pub async fn run(self) -> DetachedAnswer {
5453        let Self { id, method, work } = self;
5454        match work {
5455            // A call holding a runtime is already bounded by
5456            // RUNTIME_CONTROL_DEADLINE, and its future OWNS that connection:
5457            // a second timeout around it would drop the connection mid-call
5458            // and take down a runtime its caller still has.
5459            Ok(Work::Runtime(work)) => {
5460                let (result, returned) = work.run().await;
5461                DetachedAnswer {
5462                    response: service_response(id, result),
5463                    returned,
5464                }
5465            }
5466            Ok(Work::Free(work)) => {
5467                let result = match tokio::time::timeout(DETACHED_CALL_DEADLINE, work.run()).await {
5468                    Ok(result) => result,
5469                    Err(_) => Err(ServiceError::Operation(format!(
5470                        "`{method}` gave up after {}s: the harness it waits on did not answer",
5471                        DETACHED_CALL_DEADLINE.as_secs()
5472                    ))),
5473                };
5474                DetachedAnswer {
5475                    response: service_response(id, result),
5476                    returned: None,
5477                }
5478            }
5479            Err(error) => DetachedAnswer {
5480                response: service_response(id, Err(error)),
5481                returned: None,
5482            },
5483        }
5484    }
5485}
5486
5487/// One detached call's complete answer, plus whatever it must hand back to
5488/// the service before that answer is written. See
5489/// [`HarnessSessionService::finish_detached`].
5490pub struct DetachedAnswer {
5491    response: Value,
5492    returned: Option<ReturnedRuntime>,
5493}
5494
5495impl DetachedAnswer {
5496    /// The caller's JSON-RPC response, for a transport that owns no service
5497    /// to give a borrowed connection back to.
5498    pub fn into_response(self) -> Value {
5499        self.response
5500    }
5501}
5502
5503/// A connection lent to a detached call, on its way back to the service that
5504/// owns it.
5505pub struct ReturnedRuntime {
5506    connection: String,
5507    runtime: Box<dyn RuntimeConnection>,
5508}
5509
5510/// The waiting half of one detached request: with nothing of the service's
5511/// in hand, or holding a connection the service lent out for the call.
5512enum Work {
5513    Free(DetachedWork),
5514    Runtime(RuntimeWork),
5515}
5516
5517/// The waiting half of one detached request that holds nothing of the
5518/// service's.
5519enum DetachedWork {
5520    /// Probe the selected harnesses: find their executables, ask each its
5521    /// version, and at `probe: handshake` start each one and complete its
5522    /// protocol handshake.
5523    Inventory(InventoryWork),
5524    /// Relay one message into a live session.
5525    Message(MessageSessionParams),
5526    /// Perform one conversation verb through the harness's own CLI, HTTP API,
5527    /// daemon socket, or supercode's own store.
5528    SessionMutation {
5529        verb: crate::SessionVerb,
5530        mutation: crate::SessionMutation,
5531    },
5532}
5533
5534impl DetachedWork {
5535    async fn run(self) -> std::result::Result<Value, ServiceError> {
5536        match self {
5537            Self::Inventory(work) => run_inventory(work).await,
5538            Self::Message(params) => Ok(message_live_session(&params).await),
5539            Self::SessionMutation { verb, mutation } => {
5540                let outcome = run_session_mutation(verb, &mutation).await?;
5541                serde_json::to_value(outcome)
5542                    .map_err(|error| ServiceError::Operation(error.to_string()))
5543            }
5544        }
5545    }
5546}
5547
5548/// One detached call that holds a runtime connection for its whole run.
5549enum RuntimeWork {
5550    /// Tear down a runtime the service has already surrendered.
5551    Close {
5552        runtime: Box<dyn RuntimeConnection>,
5553        process_group: Option<u32>,
5554    },
5555    /// Type one live slash command through a borrowed connection, then give
5556    /// the connection back.
5557    LiveCommand {
5558        connection: String,
5559        runtime: Box<dyn RuntimeConnection>,
5560        verb: crate::SessionVerb,
5561        mutation: crate::SessionMutation,
5562        command: &'static str,
5563        session: String,
5564    },
5565}
5566
5567/// What one [`RuntimeWork`] answers with: the caller's result, and the
5568/// connection to give back when the call only borrowed one.
5569type RuntimeWorkAnswer = (
5570    std::result::Result<Value, ServiceError>,
5571    Option<ReturnedRuntime>,
5572);
5573
5574impl RuntimeWork {
5575    async fn run(self) -> RuntimeWorkAnswer {
5576        match self {
5577            Self::Close {
5578                runtime,
5579                process_group,
5580            } => (close_runtime(runtime, process_group).await, None),
5581            Self::LiveCommand {
5582                connection,
5583                mut runtime,
5584                verb,
5585                mutation,
5586                command,
5587                session,
5588            } => {
5589                let result =
5590                    type_live_command(runtime.as_mut(), verb, &mutation, command, session).await;
5591                (
5592                    result,
5593                    Some(ReturnedRuntime {
5594                        connection,
5595                        runtime,
5596                    }),
5597                )
5598            }
5599        }
5600    }
5601}
5602
5603/// Tear down a runtime already out of the service, within
5604/// [`RUNTIME_CONTROL_DEADLINE`].
5605async fn close_runtime(
5606    mut runtime: Box<dyn RuntimeConnection>,
5607    process_group: Option<u32>,
5608) -> std::result::Result<Value, ServiceError> {
5609    match within_control_deadline("harness.v1.runtimes.close", runtime.close()).await {
5610        Ok(result) => {
5611            result.map_err(operation)?;
5612            Ok(json!({"closed": true}))
5613        }
5614        Err(deadline) => {
5615            // Dropping the handle is not enough: the process that stopped
5616            // answering is held by a task parked on it, so nothing here runs
5617            // its Drop. Signal the group the graceful path would have
5618            // signalled, then say so.
5619            let killed = kill_runtime_process_group(process_group);
5620            drop(runtime);
5621            Ok(json!({
5622                "closed": true,
5623                "killed": killed,
5624                "detail": error_message(deadline),
5625            }))
5626        }
5627    }
5628}
5629
5630/// The conversation a live `sessions.new` / `sessions.reset` acts on: the one
5631/// the request named, or the runtime's own session.
5632fn live_session_name(runtime: &dyn RuntimeConnection, mutation: &crate::SessionMutation) -> String {
5633    mutation
5634        .session
5635        .clone()
5636        .filter(|value| !value.trim().is_empty())
5637        .unwrap_or_else(|| runtime.handle().runtime_id.clone())
5638}
5639
5640/// Type one harness slash command into a live session through the very same
5641/// `send_input` path a human's message takes, within
5642/// [`RUNTIME_CONTROL_DEADLINE`].
5643async fn type_live_command(
5644    runtime: &mut dyn RuntimeConnection,
5645    verb: crate::SessionVerb,
5646    mutation: &crate::SessionMutation,
5647    command: &str,
5648    session: String,
5649) -> std::result::Result<Value, ServiceError> {
5650    within_control_deadline(
5651        &format!("sessions.{}", verb.as_str()),
5652        runtime.send_input(RuntimeInput {
5653            text: command.to_string(),
5654            image_urls: Vec::new(),
5655        }),
5656    )
5657    .await?
5658    .map_err(operation)?;
5659    let outcome = crate::sessions_control::live_outcome(verb, mutation, command, session)
5660        .map_err(session_control_error)?;
5661    serde_json::to_value(outcome).map_err(|error| ServiceError::Operation(error.to_string()))
5662}
5663
5664/// A runtime that is up and whose handshake completed, with what the service
5665/// needs to take ownership of it.
5666enum OpenRuntime {
5667    /// supercode spawned this process, so it also hosts it: a frontend server,
5668    /// a live-runtime registration and a terminal launch of its own.
5669    Hosted {
5670        runtime: Box<dyn RuntimeConnection>,
5671        capabilities: crate::RuntimeCapabilities,
5672        workspace: PathBuf,
5673        /// A new session (`start`): it has no transcript yet, so there is no history to look for.
5674        fresh: bool,
5675    },
5676    /// `attach_existing` joined a process supercode does not own. It is
5677    /// registered as a bare connection and hosts nothing.
5678    Joined { runtime: Box<dyn RuntimeConnection> },
5679}
5680
5681/// Open the runtime one [`RUNTIME_OPEN_METHODS`] request asks for, within
5682/// [`RUNTIME_OPEN_DEADLINE`]. The error a blown deadline answers names the
5683/// method and the bound, so a caller reads why it was cut loose instead of
5684/// waiting on a handshake that is never coming.
5685async fn open_runtime(
5686    method: &str,
5687    params: Value,
5688) -> std::result::Result<OpenRuntime, ServiceError> {
5689    match tokio::time::timeout(
5690        RUNTIME_OPEN_DEADLINE,
5691        open_runtime_unbounded(method, params),
5692    )
5693    .await
5694    {
5695        Ok(result) => result,
5696        Err(_) => Err(ServiceError::Operation(format!(
5697            "`{method}` gave up after {}s: the runtime never finished its protocol handshake",
5698            RUNTIME_OPEN_DEADLINE.as_secs()
5699        ))),
5700    }
5701}
5702
5703async fn open_runtime_unbounded(
5704    method: &str,
5705    params: Value,
5706) -> std::result::Result<OpenRuntime, ServiceError> {
5707    match method {
5708        "harness.v1.runtimes.start" => {
5709            let params = decode::<RuntimeStartParams>(params)?;
5710            let backend = runtime_backend(&params.backend)?;
5711            let capabilities = backend.capabilities();
5712            let workspace = params.cwd.clone();
5713            let runtime = backend
5714                .start(RuntimeStartRequest {
5715                    cwd: params.cwd,
5716                    launch: runtime_launch(&params.backend),
5717                    mcp_servers: params.mcp_servers,
5718                    approval_policy: params.approval_policy,
5719                })
5720                .await
5721                .map_err(operation)?;
5722            Ok(OpenRuntime::Hosted {
5723                runtime,
5724                capabilities,
5725                workspace,
5726                fresh: true,
5727            })
5728        }
5729        "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
5730            let params = decode::<RuntimeAttachParams>(params)?;
5731            let backend = runtime_backend(&params.backend)?;
5732            let capabilities = backend.capabilities();
5733            let workspace = params
5734                .cwd
5735                .clone()
5736                .unwrap_or_else(|| std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")));
5737            let runtime = backend
5738                .attach(RuntimeAttachRequest {
5739                    runtime_id: params.runtime_id,
5740                    cwd: params.cwd,
5741                    launch: runtime_launch(&params.backend),
5742                    mcp_servers: params.mcp_servers,
5743                    approval_policy: params.approval_policy,
5744                })
5745                .await
5746                .map_err(operation)?;
5747            Ok(OpenRuntime::Hosted {
5748                runtime,
5749                capabilities,
5750                workspace,
5751                fresh: false,
5752            })
5753        }
5754        "harness.v1.runtimes.attach_existing" => {
5755            let params = decode::<RuntimeAttachParams>(params)?;
5756            let backend: Box<dyn RuntimeBackend> = match params
5757                .backend
5758                .base_url
5759                .as_deref()
5760                .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
5761            {
5762                Some(endpoint) => {
5763                    #[cfg(not(feature = "adapter-api"))]
5764                    {
5765                        let _ = endpoint;
5766                        return Err(ServiceError::UnsupportedAction(
5767                            "live HTTP attachment adapter is not compiled".into(),
5768                        ));
5769                    }
5770                    #[cfg(feature = "adapter-api")]
5771                    {
5772                        let workspace = params.cwd.clone().ok_or_else(|| {
5773                            ServiceError::InvalidParams(
5774                                "Volter Harness live attach requires the project cwd".into(),
5775                            )
5776                        })?;
5777                        let source = LiveRuntimeSource {
5778                            harness: params.backend.harness.as_str().to_string(),
5779                            session_id: params.runtime_id.clone(),
5780                            workspace,
5781                        };
5782                        let receipt = resolve_live_runtime(&endpoint, &source)
5783                            .map_err(|error| ServiceError::Operation(error.to_string()))?;
5784                        Box::new(SupercodeHttpRuntimeBackend::new(receipt))
5785                    }
5786                }
5787                None => runtime_backend(&params.backend)?,
5788            };
5789            let capabilities = backend.capabilities();
5790            if !capabilities.attach_existing_process {
5791                return Err(ServiceError::Operation(format!(
5792                    "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
5793                    backend.harness().as_str()
5794                )));
5795            }
5796            let runtime = backend
5797                .attach_existing(RuntimeAttachRequest {
5798                    runtime_id: params.runtime_id,
5799                    cwd: params.cwd,
5800                    launch: runtime_launch(&params.backend),
5801                    mcp_servers: params.mcp_servers,
5802                    approval_policy: params.approval_policy,
5803                })
5804                .await
5805                .map_err(operation)?;
5806            Ok(OpenRuntime::Joined { runtime })
5807        }
5808        _ => Err(ServiceError::MethodNotFound),
5809    }
5810}
5811
5812/// Wrap one service outcome in its JSON-RPC 2.0 envelope.
5813fn service_response(id: Value, result: std::result::Result<Value, ServiceError>) -> Value {
5814    match result {
5815        Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
5816        Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
5817        Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
5818        Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
5819        Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
5820        Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
5821    }
5822}
5823
5824fn runtime_backend(
5825    params: &RuntimeBackendParams,
5826) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
5827    if let Some(descriptor) = registry_connect_descriptor(params) {
5828        return open_connect_descriptor(&descriptor, &service_home()?);
5829    }
5830    if params.protocol.as_deref() == Some("acp") {
5831        let launch = params
5832            .launch
5833            .clone()
5834            .or_else(|| {
5835                harness_support_registry()
5836                    .harnesses
5837                    .into_iter()
5838                    .find(|harness| harness.id == params.harness)
5839                    .filter(|harness| {
5840                        harness.runtime.implementation == ImplementationKind::GenericProtocol
5841                            && harness.runtime.protocol.starts_with("acp")
5842                    })
5843                    .and_then(|harness| harness.runtime.default_launch)
5844            })
5845            .ok_or_else(|| {
5846                ServiceError::InvalidParams(
5847                    "an ACP runtime requires `launch` unless the harness has a registered default"
5848                        .into(),
5849                )
5850            })?;
5851        let resume_session = harness_support_registry()
5852            .harnesses
5853            .into_iter()
5854            .find(|harness| harness.id == params.harness)
5855            .is_some_and(|harness| harness.runtime.capabilities.resume_session);
5856        return Ok(Box::new(
5857            AcpRuntimeBackend::new(params.harness.clone(), launch)
5858                .with_resume_support(resume_session),
5859        ));
5860    }
5861    let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
5862        HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
5863        HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
5864        HarnessId::PI => Box::new(PiRuntimeBackend::new()),
5865        HarnessId::OPENCODE => match &params.base_url {
5866            Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
5867            None => Box::new(OpenCodeRuntimeBackend::new()),
5868        },
5869        harness => {
5870            let descriptor = harness_support_registry()
5871                .harnesses
5872                .into_iter()
5873                .find(|descriptor| descriptor.id.as_str() == harness)
5874                .filter(|descriptor| {
5875                    descriptor.runtime.implementation == ImplementationKind::GenericProtocol
5876                        && descriptor.runtime.protocol.starts_with("acp")
5877                });
5878            let Some(descriptor) = descriptor else {
5879                return Err(ServiceError::InvalidParams(format!(
5880                    "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
5881                )));
5882            };
5883            let resume = descriptor.runtime.capabilities.resume_session;
5884            Box::new(
5885                AcpRuntimeBackend::new(
5886                    descriptor.id,
5887                    descriptor
5888                        .runtime
5889                        .default_launch
5890                        .expect("generic ACP registry entry includes its launch"),
5891                )
5892                .with_resume_support(resume),
5893            )
5894        }
5895    };
5896    Ok(backend)
5897}
5898
5899fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
5900    if let Some(launch) = &params.launch {
5901        return Some(launch.clone());
5902    }
5903    if !matches!(params.policy, RuntimePolicy::Yolo) {
5904        return None;
5905    }
5906    let launch = match params.harness.as_str() {
5907        HarnessId::GROK => RuntimeLaunch {
5908            program: "grok".into(),
5909            arguments: {
5910                let mut arguments: Vec<String> = Vec::new();
5911                if crate::support::self_sandbox_supported() {
5912                    arguments.extend(["--sandbox".into(), "workspace".into()]);
5913                }
5914                arguments.extend([
5915                    "--always-approve".into(),
5916                    "agent".into(),
5917                    "--no-leader".into(),
5918                    "stdio".into(),
5919                ]);
5920                arguments
5921            },
5922            env: crate::support::grok_env(),
5923        },
5924        HarnessId::CODEX => RuntimeLaunch {
5925            program: "codex".into(),
5926            arguments: vec![
5927                "--dangerously-bypass-approvals-and-sandbox".into(),
5928                "--dangerously-bypass-hook-trust".into(),
5929                "app-server".into(),
5930            ],
5931            env: BTreeMap::new(),
5932        },
5933        HarnessId::CLAUDE_CODE => RuntimeLaunch {
5934            program: "claude".into(),
5935            arguments: vec![
5936                "--dangerously-skip-permissions".into(),
5937                "--print".into(),
5938                "--input-format".into(),
5939                "stream-json".into(),
5940                "--output-format".into(),
5941                "stream-json".into(),
5942                "--verbose".into(),
5943            ],
5944            env: BTreeMap::new(),
5945        },
5946        HarnessId::PI => RuntimeLaunch {
5947            program: "pi".into(),
5948            arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
5949            env: BTreeMap::new(),
5950        },
5951        HarnessId::OPENCODE => RuntimeLaunch {
5952            program: "opencode".into(),
5953            arguments: vec!["serve".into()],
5954            env: BTreeMap::new(),
5955        },
5956        HarnessId::GEMINI => RuntimeLaunch {
5957            program: "gemini".into(),
5958            arguments: vec!["--acp".into(), "--yolo".into()],
5959            env: BTreeMap::new(),
5960        },
5961        HarnessId::GOOSE => RuntimeLaunch {
5962            program: "goose".into(),
5963            arguments: vec!["acp".into()],
5964            env: BTreeMap::new(),
5965        },
5966        HarnessId::SUPERCODE => RuntimeLaunch {
5967            program: "supercode".into(),
5968            arguments: vec!["acp".into(), "--dangerous".into()],
5969            env: BTreeMap::new(),
5970        },
5971        _ => return None,
5972    };
5973    Some(launch)
5974}
5975
5976/// Disposable harness state for a no-prompt readiness probe. Merely opening
5977/// several stock CLIs writes a session header or migrates configuration, so a
5978/// handshake must never point at the user's real home. Authentication files
5979/// are copied into the private temporary home; all writes disappear with the
5980/// guard after the connection closes.
5981struct IsolatedProbeHome {
5982    launch: RuntimeLaunch,
5983    root: PathBuf,
5984}
5985
5986impl IsolatedProbeHome {
5987    fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
5988        let root = std::env::temp_dir().join(format!(
5989            "supercode-harness-probe-{harness}-{}",
5990            generated_session_id()
5991        ));
5992        std::fs::create_dir_all(&root)?;
5993        set_private_dir_permissions(&root)?;
5994
5995        if let Some(source_home) = supercode_interchange::user_home()
5996            .map(std::path::PathBuf::into_os_string)
5997            .map(PathBuf::from)
5998        {
5999            for relative in probe_auth_files(harness) {
6000                copy_probe_file(&source_home, &root, relative)?;
6001            }
6002        }
6003        // supercode reads its own config home ($SUPERCODE_HOME, else
6004        // $XDG_CONFIG_HOME/supercode, else ~/.config/supercode), not a fixed
6005        // place under HOME: a login kept under XDG_CONFIG_HOME probed as
6006        // "no API key found" while `supercode run` answered.
6007        if harness == HarnessId::SUPERCODE {
6008            let config_home = crate::agent::global_instructions_dir();
6009            for file in ["config.toml", "credentials.toml"] {
6010                copy_probe_path(
6011                    &config_home.join(file),
6012                    &root.join(".config/supercode").join(file),
6013                )?;
6014            }
6015        }
6016        configure_isolated_probe_auth(harness, &root)?;
6017
6018        let root_text = root.to_string_lossy().into_owned();
6019        for (key, value) in [
6020            ("HOME", root_text.clone()),
6021            (
6022                "XDG_CACHE_HOME",
6023                root.join(".cache").to_string_lossy().into_owned(),
6024            ),
6025            (
6026                "XDG_CONFIG_HOME",
6027                root.join(".config").to_string_lossy().into_owned(),
6028            ),
6029            (
6030                "XDG_DATA_HOME",
6031                root.join(".local/share").to_string_lossy().into_owned(),
6032            ),
6033        ] {
6034            launch.env.insert(key.into(), value);
6035        }
6036        let scoped = match harness {
6037            HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
6038            HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
6039            HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
6040            HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
6041            HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
6042            HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
6043            _ => None,
6044        };
6045        if let Some((key, value)) = scoped {
6046            launch
6047                .env
6048                .insert(key.into(), value.to_string_lossy().into_owned());
6049        }
6050        Ok(Self { launch, root })
6051    }
6052
6053    fn cleanup(&self) -> std::io::Result<()> {
6054        match std::fs::remove_dir_all(&self.root) {
6055            Ok(()) => Ok(()),
6056            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
6057            Err(error) => Err(error),
6058        }
6059    }
6060}
6061
6062impl Drop for IsolatedProbeHome {
6063    fn drop(&mut self) {
6064        let _ = self.cleanup();
6065    }
6066}
6067
6068fn probe_auth_files(harness: &str) -> &'static [&'static str] {
6069    match harness {
6070        HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
6071        // The gateway endpoint + token live in openclaw's own config; without
6072        // it the isolated probe dials the default endpoint unauthenticated
6073        // (PARITY-24 finding 2026-08-31).
6074        HarnessId::OPENCLAW => &[".openclaw/openclaw.json"],
6075        HarnessId::CODEX => &[".codex/auth.json"],
6076        HarnessId::GEMINI => &[
6077            ".gemini/google_accounts.json",
6078            ".gemini/oauth_creds.json",
6079            ".gemini/settings.json",
6080        ],
6081        HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
6082        HarnessId::OPENCODE => &[
6083            ".config/opencode/auth.json",
6084            ".local/share/opencode/auth.json",
6085        ],
6086        HarnessId::PI => &[".pi/agent/auth.json"],
6087        // Hermes keeps its provider selection in config.yaml, its OAuth
6088        // credential pool in auth.json, and API keys in .env; without them
6089        // the isolated probe sees "No LLM provider configured" for a
6090        // hermes that answers fine from the user's real home.
6091        HarnessId::HERMES => &[".hermes/config.yaml", ".hermes/auth.json", ".hermes/.env"],
6092        _ => &[],
6093    }
6094}
6095
6096fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
6097    copy_probe_path(&source_home.join(relative), &probe_home.join(relative))
6098}
6099
6100fn copy_probe_path(source: &Path, destination: &Path) -> std::io::Result<()> {
6101    if !source.is_file() {
6102        return Ok(());
6103    }
6104    if let Some(parent) = destination.parent() {
6105        std::fs::create_dir_all(parent)?;
6106        set_private_dir_permissions(parent)?;
6107    }
6108    std::fs::copy(source, destination)?;
6109    set_private_file_permissions(destination)
6110}
6111
6112fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
6113    if harness != HarnessId::GEMINI {
6114        return Ok(());
6115    }
6116    let oauth = probe_home.join(".gemini/oauth_creds.json");
6117    if !oauth.is_file() {
6118        return Ok(());
6119    }
6120    let settings_path = probe_home.join(".gemini/settings.json");
6121    let mut settings = std::fs::read_to_string(&settings_path)
6122        .ok()
6123        .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
6124        .unwrap_or_else(|| json!({}));
6125    settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
6126    std::fs::write(
6127        &settings_path,
6128        serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
6129    )?;
6130    set_private_file_permissions(&settings_path)
6131}
6132
6133#[cfg(unix)]
6134fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
6135    use std::os::unix::fs::PermissionsExt;
6136    std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
6137}
6138
6139#[cfg(not(unix))]
6140fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
6141    Ok(())
6142}
6143
6144#[cfg(unix)]
6145fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
6146    use std::os::unix::fs::PermissionsExt;
6147    std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
6148}
6149
6150#[cfg(not(unix))]
6151fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
6152    Ok(())
6153}
6154
6155fn find_executable(program: &str) -> Option<PathBuf> {
6156    let candidate = PathBuf::from(program);
6157    if candidate.components().count() > 1 {
6158        return candidate.is_file().then_some(candidate);
6159    }
6160    let path = std::env::var_os("PATH")?;
6161    for directory in std::env::split_paths(&path) {
6162        let candidate = directory.join(program);
6163        if candidate.is_file() {
6164            return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6165        }
6166        #[cfg(windows)]
6167        {
6168            for extension in ["exe", "cmd", "bat"] {
6169                let candidate = directory.join(format!("{program}.{extension}"));
6170                if candidate.is_file() {
6171                    return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
6172                }
6173            }
6174        }
6175    }
6176    None
6177}
6178
6179async fn executable_version(executable: &Path) -> Option<String> {
6180    let mut command = tokio::process::Command::new(executable);
6181    command
6182        .arg("--version")
6183        .stdin(std::process::Stdio::null())
6184        .stdout(std::process::Stdio::piped())
6185        .stderr(std::process::Stdio::piped())
6186        .kill_on_drop(true);
6187    let output = tokio::time::timeout(Duration::from_secs(3), command.output())
6188        .await
6189        .ok()?
6190        .ok()?;
6191    let stdout = String::from_utf8_lossy(&output.stdout);
6192    let stderr = String::from_utf8_lossy(&output.stderr);
6193    stdout
6194        .lines()
6195        .chain(stderr.lines())
6196        .map(str::trim)
6197        .find(|line| !line.is_empty())
6198        .map(|line| truncate_text(line, 200))
6199}
6200
6201pub(crate) fn auth_evidence(harness: &str) -> bool {
6202    let env_names: &[&str] = match harness {
6203        HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
6204        HarnessId::CODEX => &["OPENAI_API_KEY"],
6205        HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6206        HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
6207        HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
6208        HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
6209        HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
6210        _ => &[],
6211    };
6212    if env_names
6213        .iter()
6214        .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
6215    {
6216        return true;
6217    }
6218    let Some(home) = supercode_interchange::user_home()
6219        .map(std::path::PathBuf::into_os_string)
6220        .map(PathBuf::from)
6221    else {
6222        return false;
6223    };
6224    let files: Vec<PathBuf> = match harness {
6225        HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
6226        HarnessId::CODEX => vec![home.join(".codex/auth.json")],
6227        HarnessId::OPENCODE => vec![
6228            home.join(".local/share/opencode/auth.json"),
6229            home.join(".config/opencode/auth.json"),
6230        ],
6231        HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
6232        HarnessId::GROK => vec![home.join(".grok/auth.json")],
6233        HarnessId::GEMINI => vec![
6234            home.join(".gemini/oauth_creds.json"),
6235            home.join(".gemini/google_accounts.json"),
6236        ],
6237        HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
6238        HarnessId::HERMES => vec![home.join(".hermes/auth.json"), home.join(".hermes/.env")],
6239        _ => Vec::new(),
6240    };
6241    if files.into_iter().any(|path| {
6242        std::fs::metadata(path)
6243            .map(|metadata| metadata.is_file() && metadata.len() > 2)
6244            .unwrap_or(false)
6245    }) {
6246        return true;
6247    }
6248    // macOS keeps Claude Code's OAuth login in the Keychain, so
6249    // `.claude/.credentials.json` never exists there and the file probe above
6250    // reports a signed-in install as unauthenticated forever. A completed
6251    // login also writes an `oauthAccount` record into `~/.claude.json` on
6252    // every platform — file-based, prompt-free evidence (querying the
6253    // Keychain itself from an unsigned daemon can raise a UI prompt).
6254    if harness == HarnessId::CLAUDE_CODE {
6255        return std::fs::read_to_string(home.join(".claude.json"))
6256            .map(|text| text.contains("\"oauthAccount\""))
6257            .unwrap_or(false);
6258    }
6259    false
6260}
6261
6262fn looks_like_auth_error(message: &str) -> bool {
6263    let message = message.to_ascii_lowercase();
6264    [
6265        "auth",
6266        "login",
6267        "sign in",
6268        "sign-in",
6269        "credential",
6270        "unauthorized",
6271        "forbidden",
6272        "token",
6273    ]
6274    .iter()
6275    .any(|needle| message.contains(needle))
6276}
6277
6278fn unavailable_capabilities() -> crate::RuntimeCapabilities {
6279    crate::RuntimeCapabilities {
6280        start_session: false,
6281        resume_session: false,
6282        attach_existing_process: false,
6283        send_input: false,
6284        stream_events: false,
6285        interrupt: false,
6286        steer: false,
6287        respond_to_requests: false,
6288    }
6289}
6290
6291fn truncate_text(text: &str, max_chars: usize) -> String {
6292    let mut chars = text.chars();
6293    let truncated = chars.by_ref().take(max_chars).collect::<String>();
6294    if chars.next().is_some() {
6295        format!("{truncated}…")
6296    } else {
6297        truncated
6298    }
6299}
6300
6301/// The process group a runtime's own handle names, when it names one.
6302///
6303/// Every adapter that spawns a local process spawns it as its own group
6304/// leader (`Command::process_group(0)`), so the endpoint's pid IS the group
6305/// id. A runtime reached over HTTP, or one supercode joined rather than
6306/// spawned, names no group here and is left alone.
6307fn runtime_process_group(handle: &crate::RuntimeHandle) -> Option<u32> {
6308    match &handle.endpoint {
6309        crate::RuntimeEndpoint::LocalProcess { pid, .. } => *pid,
6310        crate::RuntimeEndpoint::Http { .. } => None,
6311    }
6312}
6313
6314/// SIGKILL a wedged runtime's whole process group, reporting whether there
6315/// was one to signal. This is the same group teardown a graceful `close`
6316/// performs; it runs here only when the graceful path blew its deadline,
6317/// because the task parked on the unanswered call still owns the process
6318/// handle and so no `Drop` of ours can reach it.
6319fn kill_runtime_process_group(process_group: Option<u32>) -> bool {
6320    match process_group {
6321        #[cfg(unix)]
6322        Some(pid) => {
6323            crate::lsp::kill_process_group(pid);
6324            true
6325        }
6326        #[cfg(not(unix))]
6327        Some(_) => false,
6328        None => false,
6329    }
6330}
6331
6332fn error_message(error: ServiceError) -> String {
6333    match error {
6334        ServiceError::InvalidParams(message)
6335        | ServiceError::Operation(message)
6336        | ServiceError::UnsupportedAction(message) => message,
6337        ServiceError::MethodNotFound => "runtime adapter is not available".into(),
6338        ServiceError::Sdk(error) => error.to_string(),
6339    }
6340}
6341
6342#[derive(Debug)]
6343enum ServiceError {
6344    InvalidParams(String),
6345    MethodNotFound,
6346    UnsupportedAction(String),
6347    Operation(String),
6348    Sdk(SdkError),
6349}
6350
6351fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
6352    match error {
6353        ServiceError::InvalidParams(message) => {
6354            SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
6355        }
6356        ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
6357            SdkError::unsupported(operation)
6358        }
6359        ServiceError::Operation(message) => {
6360            let code = if message.contains("already in progress") {
6361                SdkErrorCode::Busy
6362            } else if message.contains("not supported by this runtime") {
6363                SdkErrorCode::UnsupportedAction
6364            } else if message.contains("unknown runtime connection") {
6365                SdkErrorCode::NotFound
6366            } else {
6367                SdkErrorCode::Execution
6368            };
6369            SdkError::new(code, operation, message)
6370        }
6371        ServiceError::Sdk(error) => error,
6372    }
6373}
6374
6375fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
6376    let error_code = error.code();
6377    let code = match error_code {
6378        SdkErrorCode::Unauthenticated => -32030,
6379        SdkErrorCode::Unauthorized => -32031,
6380        SdkErrorCode::ControllerRequired => -32032,
6381        SdkErrorCode::LeaseExpired => -32033,
6382        SdkErrorCode::InvalidArgument => -32602,
6383        SdkErrorCode::NotFound => -32004,
6384        SdkErrorCode::Busy => -32000,
6385        SdkErrorCode::UnsupportedAction => -32020,
6386        SdkErrorCode::Execution => -32002,
6387        SdkErrorCode::Transport => -32003,
6388    };
6389    json!({
6390        "jsonrpc": "2.0",
6391        "id": id,
6392        "error": {
6393            "code": code,
6394            "name": error_code,
6395            "operation": error.operation(),
6396            "message": error.to_string(),
6397        },
6398    })
6399}
6400
6401fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
6402    serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
6403}
6404
6405fn operation(error: impl Into<crate::Error>) -> ServiceError {
6406    let error = error.into();
6407    match error {
6408        crate::Error::Sdk(error) => ServiceError::Sdk(error),
6409        error => ServiceError::Operation(error.to_string()),
6410    }
6411}
6412
6413/// ORCH-12 `harness.v1.memory.show|search` params. `homes` is the same
6414/// storage-root override every read-only method accepts, so a caller can
6415/// point the read at a fixture home without touching the real ones.
6416#[derive(Debug, Clone, Deserialize, Default)]
6417#[serde(default)]
6418struct MemoryRequest {
6419    /// Harness whose store is read. Required.
6420    harness: Option<String>,
6421    /// The needle, required by `search`.
6422    query: Option<String>,
6423    /// Hermes profile, OpenClaw agent, or Claude Code project.
6424    profile: Option<String>,
6425    /// Claude Code session id selecting a project store (`show` only).
6426    session: Option<String>,
6427    /// Include each document's whole text (`show` only).
6428    full: bool,
6429    /// Treat `query` as a regular expression (`search` only).
6430    regex: bool,
6431    /// Working tree whose project store is read.
6432    cwd: Option<std::path::PathBuf>,
6433    /// Storage roots to read.
6434    homes: crate::HarnessHomes,
6435}
6436
6437/// Read the memory noun. A harness with no memory store fails with
6438/// `UnsupportedAction` (RPC `-32020`), never an empty list.
6439fn memory_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6440    let request = decode::<MemoryRequest>(params)?;
6441    let harness = request
6442        .harness
6443        .clone()
6444        .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6445    let to_service = |error: crate::memory::MemoryError| match error {
6446        crate::memory::MemoryError::UnsupportedHarness { .. }
6447        | crate::memory::MemoryError::SessionNotScoped { .. } => {
6448            ServiceError::UnsupportedAction(error.to_string())
6449        }
6450        other => ServiceError::InvalidParams(other.to_string()),
6451    };
6452    match method {
6453        "harness.v1.memory.show" => {
6454            let documents = crate::memory::show_memory(&crate::memory::MemoryQuery {
6455                harness,
6456                profile: request.profile,
6457                session: request.session,
6458                full: request.full,
6459                cwd: request.cwd,
6460                homes: request.homes,
6461            })
6462            .map_err(to_service)?;
6463            Ok(json!({
6464                "schema": crate::memory::MEMORY_SCHEMA,
6465                "documents": documents,
6466            }))
6467        }
6468        "harness.v1.memory.search" => {
6469            let query = request
6470                .query
6471                .ok_or_else(|| ServiceError::InvalidParams("`query` is required".into()))?;
6472            let matches = crate::memory::search_memory(&crate::memory::MemorySearchQuery {
6473                harness,
6474                query,
6475                profile: request.profile,
6476                regex: request.regex,
6477                cwd: request.cwd,
6478                homes: request.homes,
6479            })
6480            .map_err(to_service)?;
6481            Ok(json!({
6482                "schema": crate::memory::MEMORY_SCHEMA,
6483                "matches": matches,
6484            }))
6485        }
6486        _ => Err(ServiceError::MethodNotFound),
6487    }
6488}
6489
6490/// ORCH-10 `harness.v1.profiles.list|get` params. `homes` is the same
6491/// storage-root override every read-only method accepts, so a caller can
6492/// point the read at a fixture home without touching the real ones.
6493#[derive(Debug, Clone, Deserialize)]
6494#[serde(default)]
6495struct ProfilesQuery {
6496    /// Restrict the listing to one harness. `get` requires it.
6497    harness: Option<String>,
6498    /// Profile name, required by `get`.
6499    name: Option<String>,
6500    /// Storage roots to read.
6501    homes: crate::HarnessHomes,
6502}
6503
6504impl Default for ProfilesQuery {
6505    fn default() -> Self {
6506        Self {
6507            harness: None,
6508            name: None,
6509            homes: crate::HarnessHomes::default(),
6510        }
6511    }
6512}
6513
6514/// Read the profile noun. A harness with no profile concept fails with
6515/// `UnsupportedAction` (RPC `-32020`), never an empty list.
6516fn profiles_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6517    let query = decode::<ProfilesQuery>(params)?;
6518    let to_service = |error: crate::profiles::ProfileError| match error {
6519        crate::profiles::ProfileError::UnsupportedHarness { .. } => {
6520            ServiceError::UnsupportedAction(error.to_string())
6521        }
6522        crate::profiles::ProfileError::NotFound { .. } => {
6523            ServiceError::InvalidParams(error.to_string())
6524        }
6525    };
6526    match method {
6527        "harness.v1.profiles.list" => {
6528            let profiles = crate::profiles::list_profiles(&query.homes, query.harness.as_deref())
6529                .map_err(to_service)?;
6530            Ok(json!({
6531                "schema": crate::profiles::PROFILES_SCHEMA,
6532                "profiles": profiles,
6533            }))
6534        }
6535        "harness.v1.profiles.get" => {
6536            let harness = query
6537                .harness
6538                .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6539            let name = query
6540                .name
6541                .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6542            let profile =
6543                crate::profiles::get_profile(&query.homes, &harness, &name).map_err(to_service)?;
6544            Ok(json!({
6545                "schema": crate::profiles::PROFILES_SCHEMA,
6546                "profile": profile,
6547            }))
6548        }
6549        _ => Err(ServiceError::MethodNotFound),
6550    }
6551}
6552
6553/// ORCH-14 `harness.v1.channels.list|status` params, the same storage-root
6554/// override every read-only method accepts so a caller can point the read at
6555/// a fixture home without touching the real ones.
6556#[derive(Debug, Clone, Deserialize)]
6557#[serde(default)]
6558struct ChannelsQuery {
6559    /// Restrict the listing to one harness. `status` requires it.
6560    harness: Option<String>,
6561    /// Channel name, required by `status`.
6562    name: Option<String>,
6563    /// Storage roots to read.
6564    homes: crate::HarnessHomes,
6565}
6566
6567impl Default for ChannelsQuery {
6568    fn default() -> Self {
6569        Self {
6570            harness: None,
6571            name: None,
6572            homes: crate::HarnessHomes::default(),
6573        }
6574    }
6575}
6576
6577/// Read the channel noun. A harness with no channel concept fails with
6578/// `UnsupportedAction` (RPC `-32020`), never an empty list. No row carries a
6579/// token, key or secret — see `crate::channels` "Secrecy".
6580#[derive(Debug, Clone, Deserialize)]
6581#[serde(default)]
6582struct RoutesQuery {
6583    harness: Option<String>,
6584    /// Restrict to routes targeting one profile / agent.
6585    profile: Option<String>,
6586    homes: crate::HarnessHomes,
6587}
6588
6589impl Default for RoutesQuery {
6590    fn default() -> Self {
6591        Self {
6592            harness: None,
6593            profile: None,
6594            homes: crate::HarnessHomes::default(),
6595        }
6596    }
6597}
6598
6599#[derive(Debug, Clone, Deserialize)]
6600#[serde(default)]
6601struct TriggersQuery {
6602    harness: Option<String>,
6603    homes: crate::HarnessHomes,
6604}
6605
6606impl Default for TriggersQuery {
6607    fn default() -> Self {
6608        Self {
6609            harness: None,
6610            homes: crate::HarnessHomes::default(),
6611        }
6612    }
6613}
6614
6615fn triggers_call(params: Value) -> std::result::Result<Value, ServiceError> {
6616    let query = decode::<TriggersQuery>(params)?;
6617    let triggers = crate::triggers::list_triggers(&query.homes, query.harness.as_deref())
6618        .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6619    Ok(json!({
6620        "schema": crate::triggers::TRIGGERS_SCHEMA,
6621        "triggers": triggers,
6622    }))
6623}
6624
6625fn routes_call(params: Value) -> std::result::Result<Value, ServiceError> {
6626    let query = decode::<RoutesQuery>(params)?;
6627    let routes = crate::routes::list_routes(
6628        &query.homes,
6629        query.harness.as_deref(),
6630        query.profile.as_deref(),
6631    )
6632    .map_err(|error| ServiceError::UnsupportedAction(error.to_string()))?;
6633    Ok(json!({
6634        "schema": crate::routes::ROUTES_SCHEMA,
6635        "routes": routes,
6636    }))
6637}
6638
6639fn channels_call(method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
6640    let query = decode::<ChannelsQuery>(params)?;
6641    let to_service = |error: crate::channels::ChannelError| match error {
6642        crate::channels::ChannelError::UnsupportedHarness { .. } => {
6643            ServiceError::UnsupportedAction(error.to_string())
6644        }
6645        crate::channels::ChannelError::NotFound { .. } => {
6646            ServiceError::InvalidParams(error.to_string())
6647        }
6648    };
6649    match method {
6650        "harness.v1.channels.list" => {
6651            let channels = crate::channels::list_channels(&query.homes, query.harness.as_deref())
6652                .map_err(to_service)?;
6653            Ok(json!({
6654                "schema": crate::channels::CHANNELS_SCHEMA,
6655                "channels": channels,
6656            }))
6657        }
6658        "harness.v1.channels.status" => {
6659            let harness = query
6660                .harness
6661                .ok_or_else(|| ServiceError::InvalidParams("`harness` is required".into()))?;
6662            let name = query
6663                .name
6664                .ok_or_else(|| ServiceError::InvalidParams("`name` is required".into()))?;
6665            let channel = crate::channels::channel_status(&query.homes, &harness, &name)
6666                .map_err(to_service)?;
6667            Ok(json!({
6668                "schema": crate::channels::CHANNELS_SCHEMA,
6669                "channel": channel,
6670            }))
6671        }
6672        _ => Err(ServiceError::MethodNotFound),
6673    }
6674}
6675
6676fn rpc_error(id: Value, code: i64, message: &str) -> Value {
6677    json!({
6678        "jsonrpc": "2.0",
6679        "id": id,
6680        "error": {"code": code, "message": message},
6681    })
6682}