Skip to main content

supercode_harness/
harness_service.rs

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