Skip to main content

supercode_harness/
harness_service.rs

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