Skip to main content

supercode_harness/
harness_service.rs

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