Skip to main content

supercode_harness/
harness_service.rs

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