Skip to main content

supercode_harness/
harness_service.rs

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