Skip to main content

supercode_harness/
harness_service.rs

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