Skip to main content

supercode/
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::time::Duration;
10
11use serde::{Deserialize, Serialize};
12use serde_json::{json, Value};
13
14use crate::runtime::generated_session_id;
15#[cfg(feature = "adapter-api")]
16use crate::runtime::{HostedHarnessConnection, HostedHarnessRuntime};
17use crate::sdk::{
18    discover_sessions, load_session, load_session_with_fidelity, SdkCapabilities, SdkError,
19    SdkErrorCode, SdkEvent, SdkOperation, SdkRequest, SdkRuntimeEvent, SdkService,
20};
21use crate::watch::normalized_session_json;
22use crate::Fidelity;
23#[cfg(feature = "adapter-api")]
24use crate::SupercodeHttpRuntimeBackend;
25use crate::{
26    discover_live_runtime, harness_support_registry, AcpRuntimeBackend, ClaudeCodeRuntimeBackend,
27    CodexRuntimeBackend, DiscoveryQuery, HarnessCatalog, HarnessId, ImplementationKind,
28    LiveRuntimeEndpoint, LiveRuntimeSource, OpenCodeRuntimeBackend, PiRuntimeBackend,
29    RuntimeAttachRequest, RuntimeBackend, RuntimeConnection, RuntimeInput, RuntimeLaunch,
30    RuntimeStartRequest, Session, SessionFollower, SessionFormat, SessionLocator, SessionSource,
31};
32#[cfg(feature = "adapter-api")]
33use crate::{register_live_runtime, resolve_live_runtime, LiveRuntimeRegistration};
34
35/// Protocol namespace implemented by this service.
36pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
37/// Notification method emitted for followed-session changes.
38pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
39/// Notification method emitted for live runtime events.
40pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
41
42/// Stateful persisted-session service. Each instance owns its follow
43/// subscriptions; discovery and loading remain read-only.
44pub struct HarnessSessionService {
45    catalog: HarnessCatalog,
46    followers: BTreeMap<String, SessionFollower>,
47    followed_sources: BTreeMap<String, FollowedSource>,
48    next_subscription: u64,
49    runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
50    terminal_launches: BTreeMap<String, StructuredLaunch>,
51    runtime_sequences: BTreeMap<String, u64>,
52    next_runtime: u64,
53}
54
55impl Default for HarnessSessionService {
56    fn default() -> Self {
57        Self::new()
58    }
59}
60
61impl HarnessSessionService {
62    /// Create an empty service instance.
63    pub fn new() -> Self {
64        Self {
65            catalog: HarnessCatalog::new(),
66            followers: BTreeMap::new(),
67            followed_sources: BTreeMap::new(),
68            next_subscription: 1,
69            runtimes: BTreeMap::new(),
70            terminal_launches: BTreeMap::new(),
71            runtime_sequences: BTreeMap::new(),
72            next_runtime: 1,
73        }
74    }
75
76    /// Handle one JSON-RPC 2.0 request and return one JSON-RPC response.
77    #[cfg(feature = "adapter-api")]
78    pub fn handle(&mut self, request: Value) -> Value {
79        let id = request.get("id").cloned().unwrap_or(Value::Null);
80        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
81            return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
82        }
83        let Some(method) = request.get("method").and_then(Value::as_str) else {
84            return rpc_error(id, -32600, "request is missing `method`");
85        };
86        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
87        match self.call(method, params) {
88            Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
89            Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
90            Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
91            Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
92            Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
93            Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
94        }
95    }
96
97    /// Handle either a persisted-session request or an asynchronous live
98    /// runtime request.
99    #[cfg(feature = "adapter-api")]
100    pub async fn handle_async(&mut self, request: Value) -> Value {
101        let method = request
102            .get("method")
103            .and_then(Value::as_str)
104            .unwrap_or_default();
105        if matches!(
106            method,
107            "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
108        ) {
109            let id = request.get("id").cloned().unwrap_or(Value::Null);
110            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
111                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
112            }
113            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
114            return match self.inventory_call(method, params).await {
115                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
116                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
117                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
118                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
119                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
120                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
121            };
122        }
123        if method == "harness.v1.sessions.message" {
124            let id = request.get("id").cloned().unwrap_or(Value::Null);
125            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
126                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
127            }
128            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
129            return match self.message_call(params).await {
130                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
131                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
132                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
133                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
134                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
135                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
136            };
137        }
138        if let Some(operation) = SdkOperation::from_method(method) {
139            let id = request.get("id").cloned().unwrap_or(Value::Null);
140            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
141                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
142            }
143            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
144            return match self.execute(SdkRequest { operation, params }).await {
145                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
146                Err(error) => sdk_rpc_error(id, &error),
147            };
148        }
149        if !method.starts_with("harness.v1.runtimes.") {
150            return self.handle(request);
151        }
152        let id = request.get("id").cloned().unwrap_or(Value::Null);
153        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
154            return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
155        }
156        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
157        match self.runtime_call(method, params).await {
158            Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
159            Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
160            Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
161            Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
162            Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
163            Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
164        }
165    }
166
167    /// Poll all active subscriptions once and return zero or more JSON-RPC
168    /// notifications. Recoverable follower errors are delivered as events.
169    #[cfg(feature = "adapter-api")]
170    pub fn poll(&mut self) -> Vec<Value> {
171        let mut notifications = Vec::new();
172        for (subscription, follower) in &mut self.followers {
173            match follower.poll() {
174                Ok(Some(event)) => notifications.push(json!({
175                    "jsonrpc": "2.0",
176                    "method": SESSION_EVENT_METHOD,
177                    "params": {
178                        "subscription": subscription,
179                        "event": event.to_json(),
180                    }
181                })),
182                Ok(None) => {}
183                Err(error) => notifications.push(json!({
184                    "jsonrpc": "2.0",
185                    "method": SESSION_EVENT_METHOD,
186                    "params": {
187                        "subscription": subscription,
188                        "event": {
189                            "type": "watch_error",
190                            "recoverable": true,
191                            "message": error.to_string(),
192                        },
193                    }
194                })),
195            }
196        }
197        notifications
198    }
199
200    /// Report each followed session's live-runtime lifecycle state on that
201    /// session's own subscription, emitting only when the state changes.
202    ///
203    /// A growing transcript is not evidence that an agent is working, so the
204    /// state comes from the live-runtime registry and nowhere else. A followed
205    /// session with no registered Supercode runtime — a harness running outside
206    /// Supercode — reports `persisted`, which says plainly that its activity is
207    /// unknown rather than guessing at it. These events carry no sequence
208    /// number and no transcript content; they never interleave with the
209    /// content follower's sequenced stream.
210    #[cfg(feature = "adapter-api")]
211    pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
212        let registry = crate::LocalRuntimeRegistry::new();
213        let authorization = crate::RuntimeAuthorization::observer();
214        let mut notifications = Vec::new();
215        for (subscription, source) in &mut self.followed_sources {
216            let state = match registry
217                .source_state(&source.harness, &source.session_id, &authorization)
218                .await
219            {
220                Ok(Some(state)) => state,
221                Ok(None) => crate::RuntimeRegistryState::Persisted,
222                // A failed registry read is not evidence of a state change.
223                Err(_) => continue,
224            };
225            if source.reported.as_deref() == Some(state.as_str()) {
226                continue;
227            }
228            source.reported = Some(state.as_str().to_string());
229            notifications.push(json!({
230                "jsonrpc": "2.0",
231                "method": SESSION_EVENT_METHOD,
232                "params": {
233                    "subscription": subscription,
234                    "event": {"type": "runtime_state", "state": state.as_str()},
235                },
236            }));
237        }
238        notifications
239    }
240
241    /// Non-blockingly sample one event from every connected live runtime.
242    #[cfg(feature = "adapter-api")]
243    pub async fn poll_runtimes(&mut self) -> Vec<Value> {
244        self.poll_sdk_events()
245            .await
246            .into_iter()
247            .map(|(connection, runtime_event)| {
248                json!({
249                    "jsonrpc": "2.0",
250                    "method": RUNTIME_EVENT_METHOD,
251                    "params": {
252                        "connection": connection,
253                        "session_id": runtime_event.session_id,
254                        "sequence": runtime_event.event.sequence,
255                        "event": {
256                            "kind": runtime_event.event.kind,
257                            "payload": runtime_event.event.payload,
258                        },
259                    },
260                })
261            })
262            .collect()
263    }
264
265    async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
266        let mut events = Vec::new();
267        let mut closed = Vec::new();
268        for (connection, runtime) in &mut self.runtimes {
269            let session_id = runtime.handle().runtime_id.clone();
270            match tokio::time::timeout(Duration::from_millis(1), runtime.next_event()).await {
271                Ok(Ok(Some(event))) => {
272                    let terminal = event.kind == "transport_closed";
273                    let next_sequence = self
274                        .runtime_sequences
275                        .entry(session_id.clone())
276                        .or_insert(0);
277                    let sequence = event.sequence.unwrap_or_else(|| {
278                        *next_sequence = next_sequence.saturating_add(1);
279                        *next_sequence
280                    });
281                    *next_sequence = (*next_sequence).max(sequence);
282                    events.push((
283                        connection.clone(),
284                        SdkRuntimeEvent {
285                            session_id: session_id.clone(),
286                            event: SdkEvent {
287                                sequence,
288                                kind: event.kind,
289                                payload: event.payload,
290                            },
291                        },
292                    ));
293                    if terminal {
294                        closed.push(connection.clone());
295                    }
296                }
297                Ok(Ok(None)) => {
298                    let sequence = self
299                        .runtime_sequences
300                        .entry(session_id.clone())
301                        .or_insert(0);
302                    *sequence = sequence.saturating_add(1);
303                    events.push((
304                        connection.clone(),
305                        SdkRuntimeEvent {
306                            session_id,
307                            event: SdkEvent {
308                                sequence: *sequence,
309                                kind: "transport_closed".into(),
310                                payload: json!({"message": "Harness runtime transport closed."}),
311                            },
312                        },
313                    ));
314                    closed.push(connection.clone());
315                }
316                Err(_) => {}
317                Ok(Err(error)) => {
318                    let sequence = self
319                        .runtime_sequences
320                        .entry(session_id.clone())
321                        .or_insert(0);
322                    *sequence = sequence.saturating_add(1);
323                    events.push((
324                        connection.clone(),
325                        SdkRuntimeEvent {
326                            session_id,
327                            event: SdkEvent {
328                                sequence: *sequence,
329                                kind: "transport_error".into(),
330                                payload: json!({"message": error.to_string(), "terminal": true}),
331                            },
332                        },
333                    ));
334                    closed.push(connection.clone());
335                }
336            }
337        }
338        for connection in closed {
339            if let Some(runtime) = self.runtimes.remove(&connection) {
340                self.runtime_sequences.remove(&runtime.handle().runtime_id);
341            }
342            self.terminal_launches.remove(&connection);
343        }
344        events
345    }
346
347    fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
348        match method {
349            "harness.v1.capabilities" => Ok(json!({
350                "version": HARNESS_SERVICE_VERSION,
351                "sdk": self.capabilities(),
352                "methods": [
353                    "harness.v1.support.report",
354                    "harness.v1.harnesses.list",
355                    "harness.v1.harnesses.probe",
356                    "harness.v1.sessions.discover",
357                    "harness.v1.sessions.load",
358                    "harness.v1.sessions.follow",
359                    "harness.v1.sessions.unfollow",
360                    "harness.v1.sessions.message",
361                    "harness.v1.sessions.import",
362                    "harness.v1.sessions.export",
363                    "harness.v1.sessions.translate",
364                    "harness.v1.sessions.branch",
365                    "harness.v1.sessions.handoff",
366                    "harness.v1.sessions.resume_instructions",
367                    "harness.v1.runtimes.capabilities",
368                    "harness.v1.runtimes.start",
369                    "harness.v1.runtimes.resume",
370                    "harness.v1.runtimes.attach_existing",
371                    "harness.v1.runtimes.attach",
372                    "harness.v1.runtimes.send_input",
373                    "harness.v1.runtimes.interrupt",
374                    "harness.v1.runtimes.steer",
375                    "harness.v1.runtimes.respond",
376                    "harness.v1.runtimes.terminal_instructions",
377                    "harness.v1.runtimes.close",
378                ],
379                "notifications": [SESSION_EVENT_METHOD, RUNTIME_EVENT_METHOD],
380                "harnesses": harness_support_registry()
381                    .harnesses
382                    .into_iter()
383                    .map(|harness| harness.id)
384                    .collect::<Vec<_>>(),
385            })),
386            "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
387                .map_err(|error| ServiceError::Operation(error.to_string())),
388            "harness.v1.sessions.discover" => {
389                let query = decode::<DiscoveryQuery>(params)?;
390                let sessions = discover_sessions(&query).map_err(operation)?;
391                // Claude Code is the one harness that publishes its RUNNING
392                // sessions. The registry is read once per discovery and joined
393                // by session id; every record in it has already survived a
394                // `kill(pid, 0)` liveness check inside `read_registry`.
395                let peers = if sessions
396                    .iter()
397                    .any(|session| session.locator.harness.as_str() == HarnessId::CLAUDE_CODE)
398                {
399                    crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(
400                        &query.homes,
401                    ))
402                } else {
403                    Vec::new()
404                };
405                let sessions = sessions
406                    .into_iter()
407                    .map(|session| {
408                        let mut value = serde_json::to_value(&session)
409                            .map_err(|error| ServiceError::Operation(error.to_string()))?;
410                        if let Some(workspace) = &session.cwd {
411                            let source = LiveRuntimeSource {
412                                harness: session.locator.harness.as_str().to_string(),
413                                session_id: session.locator.session_id.clone(),
414                                workspace: workspace.clone(),
415                            };
416                            if let Some(endpoint) = discover_live_runtime(&source)
417                                .map_err(|error| ServiceError::Operation(error.to_string()))?
418                            {
419                                value["live_endpoint"] = json!(endpoint.as_str());
420                            }
421                        }
422                        if let Some(peer) = peers.iter().find(|peer| {
423                            session.locator.harness.as_str() == HarnessId::CLAUDE_CODE
424                                && peer.session_id == session.locator.session_id
425                        }) {
426                            // A Supercode-hosted runtime is the richer
427                            // attachment, so it keeps the endpoint slot; the
428                            // peer endpoint fills it only when nothing else did.
429                            if value.get("live_endpoint").is_none() {
430                                value["live_endpoint"] = json!(peer.endpoint().as_str());
431                            }
432                            if let Some(status) = peer.status {
433                                value["live_status"] = json!(status.as_str());
434                            }
435                        }
436                        Ok(value)
437                    })
438                    .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
439                Ok(json!({"sessions": sessions}))
440            }
441            "harness.v1.sessions.load" => {
442                let params = decode::<LocatorParams>(params)?;
443                load_session_with_fidelity(&params.locator, params.read_fidelity())
444                    .map(|session| json!({"session": normalized_session_json(&session)}))
445                    .map_err(operation)
446            }
447            "harness.v1.sessions.follow" => {
448                let params = decode::<LocatorParams>(params)?;
449                let mut follower = self
450                    .catalog
451                    .follow_with_fidelity(&params.locator, params.read_fidelity())
452                    .map_err(operation)?;
453                let initial = follower
454                    .poll()
455                    .map_err(operation)?
456                    .map(|event| event.to_json());
457                let subscription = format!("sub-{}", self.next_subscription);
458                self.next_subscription += 1;
459                self.followers.insert(subscription.clone(), follower);
460                self.followed_sources.insert(
461                    subscription.clone(),
462                    FollowedSource {
463                        harness: params.locator.harness.as_str().to_string(),
464                        session_id: params.locator.session_id.clone(),
465                        reported: None,
466                    },
467                );
468                Ok(json!({"subscription": subscription, "initial": initial}))
469            }
470            "harness.v1.sessions.unfollow" => {
471                let params = decode::<UnfollowParams>(params)?;
472                self.followed_sources.remove(&params.subscription);
473                Ok(json!({
474                    "removed": self.followers.remove(&params.subscription).is_some()
475                }))
476            }
477            "harness.v1.sessions.import" => {
478                let params = decode::<ImportSessionParams>(params)?;
479                let session = Session::load_str(&params.content, params.source_harness.into())
480                    .map_err(operation)?;
481                Ok(json!({"session": normalized_session_json(&session)}))
482            }
483            "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
484                let params = decode::<ExportSessionParams>(params)?;
485                let session = load_session(&params.locator).map_err(operation)?;
486                let artifact = session_artifact(&params.locator, &session, params.target_harness)?;
487                Ok(json!({"artifact": artifact}))
488            }
489            "harness.v1.sessions.branch" => {
490                let params = decode::<BranchSessionParams>(params)?;
491                let session = load_session(&params.locator).map_err(operation)?;
492                let storage = params.locator.storage.path().display().to_string();
493                let bootstrap_prompt = format!(
494                    "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.",
495                    params.locator.harness.as_str(), params.locator.session_id, storage
496                );
497                let artifact = params
498                    .target_harness
499                    .map(|target| session_artifact(&params.locator, &session, target))
500                    .transpose()?;
501                Ok(json!({
502                    "parent": params.locator,
503                    "session": normalized_session_json(&session),
504                    "bootstrap_prompt": bootstrap_prompt,
505                    "artifact": artifact,
506                }))
507            }
508            "harness.v1.sessions.handoff" => {
509                let params = decode::<HandoffSessionParams>(params)?;
510                let session = load_session(&params.locator).map_err(operation)?;
511                let cwd = params
512                    .cwd
513                    .or_else(|| session.meta.cwd.clone())
514                    .unwrap_or_else(|| PathBuf::from("."));
515                let artifact =
516                    handoff_artifact(&params.locator, &session, params.target_harness, &cwd)?;
517                let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
518                    ServiceError::Operation(
519                        "handoff artifact omitted target session identity".into(),
520                    )
521                })?;
522                let instructions =
523                    handoff_instructions(params.target_harness, target_session_id, &cwd);
524                Ok(json!({
525                    "artifact": artifact,
526                    "launch": instructions.launch,
527                    "materialize": instructions.materialize,
528                    "requires_materialization": instructions.requires_materialization,
529                    "note": instructions.note,
530                }))
531            }
532            "harness.v1.sessions.resume_instructions" => {
533                let params = decode::<ResumeInstructionsParams>(params)?;
534                let session = load_session(&params.locator).map_err(operation)?;
535                let cwd = params
536                    .cwd
537                    .or(session.meta.cwd)
538                    .unwrap_or_else(|| PathBuf::from("."));
539                let launch = resume_launch(
540                    params.locator.harness.as_str(),
541                    &params.locator.session_id,
542                    &cwd,
543                    params.policy,
544                )?;
545                Ok(json!({"launch": launch}))
546            }
547            _ => Err(ServiceError::MethodNotFound),
548        }
549    }
550
551    async fn runtime_call(
552        &mut self,
553        method: &str,
554        params: Value,
555    ) -> std::result::Result<Value, ServiceError> {
556        match method {
557            "harness.v1.runtimes.capabilities" => {
558                let params = decode::<RuntimeBackendParams>(params)?;
559                let backend = runtime_backend(&params)?;
560                Ok(json!({
561                    "harness": backend.harness(),
562                    "capabilities": backend.capabilities(),
563                }))
564            }
565            "harness.v1.runtimes.start" => {
566                let params = decode::<RuntimeStartParams>(params)?;
567                let backend = runtime_backend(&params.backend)?;
568                let capabilities = backend.capabilities();
569                let workspace = params.cwd.clone();
570                let runtime = backend
571                    .start(RuntimeStartRequest {
572                        cwd: params.cwd,
573                        launch: runtime_launch(&params.backend),
574                    })
575                    .await
576                    .map_err(operation)?;
577                self.insert_hosted_runtime(runtime, capabilities, workspace)
578                    .await
579            }
580            "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
581                let params = decode::<RuntimeAttachParams>(params)?;
582                let backend = runtime_backend(&params.backend)?;
583                let capabilities = backend.capabilities();
584                let workspace = params.cwd.clone().unwrap_or_else(|| {
585                    std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))
586                });
587                let runtime = backend
588                    .attach(RuntimeAttachRequest {
589                        runtime_id: params.runtime_id,
590                        cwd: params.cwd,
591                        launch: runtime_launch(&params.backend),
592                    })
593                    .await
594                    .map_err(operation)?;
595                self.insert_hosted_runtime(runtime, capabilities, workspace)
596                    .await
597            }
598            "harness.v1.runtimes.attach_existing" => {
599                let params = decode::<RuntimeAttachParams>(params)?;
600                let backend: Box<dyn RuntimeBackend> = match params
601                    .backend
602                    .base_url
603                    .as_deref()
604                    .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
605                {
606                    Some(endpoint) => {
607                        #[cfg(not(feature = "adapter-api"))]
608                        {
609                            let _ = endpoint;
610                            return Err(ServiceError::UnsupportedAction(
611                                "live HTTP attachment adapter is not compiled".into(),
612                            ));
613                        }
614                        #[cfg(feature = "adapter-api")]
615                        {
616                            let workspace = params.cwd.clone().ok_or_else(|| {
617                                ServiceError::InvalidParams(
618                                    "Supercode live attach requires the project cwd".into(),
619                                )
620                            })?;
621                            let source = LiveRuntimeSource {
622                                harness: params.backend.harness.as_str().to_string(),
623                                session_id: params.runtime_id.clone(),
624                                workspace,
625                            };
626                            let receipt = resolve_live_runtime(&endpoint, &source)
627                                .map_err(|error| ServiceError::Operation(error.to_string()))?;
628                            Box::new(SupercodeHttpRuntimeBackend::new(receipt))
629                        }
630                    }
631                    None => runtime_backend(&params.backend)?,
632                };
633                if !backend.capabilities().attach_existing_process {
634                    return Err(ServiceError::Operation(format!(
635                        "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
636                        backend.harness().as_str()
637                    )));
638                }
639                let runtime = backend
640                    .attach_existing(RuntimeAttachRequest {
641                        runtime_id: params.runtime_id,
642                        cwd: params.cwd,
643                        launch: runtime_launch(&params.backend),
644                    })
645                    .await
646                    .map_err(operation)?;
647                self.insert_runtime(runtime)
648            }
649            "harness.v1.runtimes.send_input" => {
650                let params = decode::<RuntimeInputParams>(params)?;
651                let runtime = self.runtime_mut(&params.connection)?;
652                let turn_id = runtime
653                    .send_input(RuntimeInput { text: params.text })
654                    .await
655                    .map_err(operation)?;
656                Ok(json!({"turn_id": turn_id}))
657            }
658            "harness.v1.runtimes.interrupt" => {
659                let params = decode::<RuntimeConnectionParams>(params)?;
660                self.runtime_mut(&params.connection)?
661                    .interrupt()
662                    .await
663                    .map_err(operation)?;
664                Ok(json!({}))
665            }
666            "harness.v1.runtimes.steer" => Err(ServiceError::UnsupportedAction(
667                "steer is not supported by this harness-native runtime adapter".into(),
668            )),
669            "harness.v1.runtimes.respond" => {
670                let params = decode::<RuntimeRespondParams>(params)?;
671                self.runtime_mut(&params.connection)?
672                    .respond(params.request_id, params.response)
673                    .await
674                    .map_err(operation)?;
675                Ok(json!({}))
676            }
677            "harness.v1.runtimes.terminal_instructions" => {
678                let params = decode::<RuntimeConnectionParams>(params)?;
679                let launch = self
680                    .terminal_launches
681                    .get(&params.connection)
682                    .ok_or_else(|| {
683                        ServiceError::Operation(
684                            "this runtime is not hosted for terminal attachment".into(),
685                        )
686                    })?;
687                Ok(json!({"launch":launch}))
688            }
689            "harness.v1.runtimes.close" => {
690                let params = decode::<RuntimeConnectionParams>(params)?;
691                let Some(mut runtime) = self.runtimes.remove(&params.connection) else {
692                    return Err(ServiceError::InvalidParams(format!(
693                        "unknown runtime connection `{}`",
694                        params.connection
695                    )));
696                };
697                self.terminal_launches.remove(&params.connection);
698                self.runtime_sequences.remove(&runtime.handle().runtime_id);
699                runtime.close().await.map_err(operation)?;
700                Ok(json!({"closed": true}))
701            }
702            _ => Err(ServiceError::MethodNotFound),
703        }
704    }
705
706    /// Deliver one message into a session that is running right now.
707    #[cfg(feature = "adapter-api")]
708    async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
709        let params = decode::<MessageSessionParams>(params)?;
710        Ok(message_live_session(&params, &crate::claude_peer::ProcessCourierRunner).await)
711    }
712
713    fn insert_runtime(
714        &mut self,
715        runtime: Box<dyn RuntimeConnection>,
716    ) -> std::result::Result<Value, ServiceError> {
717        let connection = format!("runtime-{}", self.next_runtime);
718        self.next_runtime += 1;
719        let handle = runtime.handle().clone();
720        self.runtime_sequences
721            .entry(handle.runtime_id.clone())
722            .or_insert(0);
723        self.runtimes.insert(connection.clone(), runtime);
724        Ok(json!({"connection": connection, "handle": handle}))
725    }
726
727    #[cfg(feature = "adapter-api")]
728    async fn insert_hosted_runtime(
729        &mut self,
730        runtime: Box<dyn RuntimeConnection>,
731        capabilities: crate::RuntimeCapabilities,
732        workspace: PathBuf,
733    ) -> std::result::Result<Value, ServiceError> {
734        let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities);
735        let token: std::sync::Arc<str> = crate::server::generate_token().into();
736        let server = crate::server::run_frontend_http(
737            host.clone(),
738            host.frontend_sender(),
739            "127.0.0.1:0",
740            token.clone(),
741        )
742        .await
743        .map_err(|error| ServiceError::Operation(error.to_string()))?;
744        let source = LiveRuntimeSource {
745            harness: connection.handle().harness.as_str().to_string(),
746            session_id: connection.handle().runtime_id.clone(),
747            workspace: workspace.clone(),
748        };
749        let registration = register_live_runtime(
750            connection.handle().runtime_id.clone(),
751            source.clone(),
752            format!("http://{}", server.address()),
753            token.to_string(),
754        )
755        .map_err(|error| ServiceError::Operation(error.to_string()))?;
756        let endpoint = registration.endpoint().to_string();
757        let launch = StructuredLaunch {
758            cwd: workspace,
759            // Pin attachment to the executable hosting this runtime. A bare
760            // `supercode` could resolve to an older global install whose CLI
761            // does not understand the receipt it is being asked to open.
762            program: std::env::current_exe()
763                .ok()
764                .map(|path| path.to_string_lossy().into_owned())
765                .unwrap_or_else(|| "supercode".into()),
766            arguments: vec![
767                "harness".into(),
768                "attach".into(),
769                "--endpoint".into(),
770                endpoint,
771                "--harness".into(),
772                source.harness,
773                "--session".into(),
774                source.session_id,
775            ],
776            env: BTreeMap::new(),
777        };
778        let lease = HostedRuntimeLease {
779            connection,
780            _host: host,
781            _registration: registration,
782            _server: server,
783        };
784        let opened = self.insert_runtime(Box::new(lease))?;
785        let connection_id = opened["connection"]
786            .as_str()
787            .expect("insert_runtime returns a connection id")
788            .to_string();
789        self.terminal_launches.insert(connection_id, launch);
790        Ok(opened)
791    }
792
793    #[cfg(not(feature = "adapter-api"))]
794    async fn insert_hosted_runtime(
795        &mut self,
796        runtime: Box<dyn RuntimeConnection>,
797        _capabilities: crate::RuntimeCapabilities,
798        _workspace: PathBuf,
799    ) -> std::result::Result<Value, ServiceError> {
800        self.insert_runtime(runtime)
801    }
802
803    fn runtime_mut(
804        &mut self,
805        connection: &str,
806    ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
807        self.runtimes.get_mut(connection).ok_or_else(|| {
808            ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
809        })
810    }
811
812    async fn inventory_call(
813        &self,
814        method: &str,
815        params: Value,
816    ) -> std::result::Result<Value, ServiceError> {
817        let mut params = decode::<HarnessInventoryParams>(params)?;
818        if method == "harness.v1.harnesses.probe" {
819            let harness = params.harness.take().ok_or_else(|| {
820                ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
821            })?;
822            params.harnesses = vec![harness];
823        }
824        let selected = params
825            .harnesses
826            .iter()
827            .map(HarnessId::as_str)
828            .collect::<std::collections::BTreeSet<_>>();
829        let supported = harness_support_registry()
830            .harnesses
831            .into_iter()
832            .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
833            .collect::<Vec<_>>();
834        if !params.harnesses.is_empty() && supported.len() != selected.len() {
835            let known = supported
836                .iter()
837                .map(|harness| harness.id.as_str())
838                .collect::<std::collections::BTreeSet<_>>();
839            let missing = params
840                .harnesses
841                .iter()
842                .filter(|id| !known.contains(id.as_str()))
843                .map(HarnessId::as_str)
844                .collect::<Vec<_>>();
845            return Err(ServiceError::InvalidParams(format!(
846                "unknown harness(es): {}",
847                missing.join(", ")
848            )));
849        }
850        let global_counts = params
851            .include_sessions
852            .then(|| self.session_counts(None, &params.harnesses));
853        let workspace_counts = params.include_sessions.then(|| {
854            params
855                .workspace
856                .as_deref()
857                .map(|workspace| self.session_counts(Some(workspace), &params.harnesses))
858        });
859        let probes = supported.into_iter().map(|descriptor| {
860            let global = global_counts
861                .as_ref()
862                .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
863            let workspace = workspace_counts
864                .as_ref()
865                .and_then(Option::as_ref)
866                .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
867            self.probe_harness(descriptor, &params, global, workspace)
868        });
869        let harnesses = futures::future::join_all(probes).await;
870        serde_json::to_value(HarnessInventoryReport {
871            probe: params.probe,
872            workspace: params.workspace,
873            harnesses,
874        })
875        .map_err(|error| ServiceError::Operation(error.to_string()))
876    }
877
878    async fn probe_harness(
879        &self,
880        descriptor: crate::HarnessSupportDescriptor,
881        params: &HarnessInventoryParams,
882        global: Option<usize>,
883        workspace: Option<usize>,
884    ) -> LocalHarness {
885        let launch = descriptor.runtime.default_launch.as_ref();
886        let executable = launch.and_then(|launch| find_executable(&launch.program));
887        let installed = executable.is_some();
888        let version = match executable.as_deref() {
889            Some(path) => executable_version(path).await,
890            None => None,
891        };
892        let configured = auth_evidence(descriptor.id.as_str());
893        let mut auth = if configured {
894            HarnessAuthState::Configured
895        } else {
896            HarnessAuthState::Unknown
897        };
898        let mut runtime = if installed {
899            HarnessRuntimeState::Degraded
900        } else {
901            HarnessRuntimeState::Unavailable
902        };
903        let mut reason = (!installed).then(|| {
904            format!(
905                "{} is supported but `{}` was not found on PATH",
906                descriptor.display_name,
907                launch
908                    .map(|launch| launch.program.as_str())
909                    .unwrap_or("executable")
910            )
911        });
912        let mut repair = (!installed).then(|| {
913            format!(
914                "Install {} and ensure `{}` is on PATH.",
915                descriptor.display_name,
916                launch
917                    .map(|launch| launch.program.as_str())
918                    .unwrap_or("its executable")
919            )
920        });
921
922        if installed && params.probe == HarnessProbeLevel::Handshake {
923            let backend_params = RuntimeBackendParams {
924                harness: descriptor.id.clone(),
925                protocol: None,
926                launch: None,
927                base_url: None,
928                policy: RuntimePolicy::Default,
929            };
930            match runtime_backend(&backend_params) {
931                Ok(backend) => {
932                    let cwd = params
933                        .workspace
934                        .clone()
935                        .or_else(|| std::env::current_dir().ok())
936                        .unwrap_or_else(|| PathBuf::from("."));
937                    match tokio::time::timeout(
938                        Duration::from_secs(30),
939                        backend.start(RuntimeStartRequest { cwd, launch: None }),
940                    )
941                    .await
942                    {
943                        Ok(Ok(mut connection)) => {
944                            match stabilize_handshake(connection.as_mut()).await {
945                                Ok(()) => {
946                                    auth = HarnessAuthState::Ready;
947                                    runtime = HarnessRuntimeState::Ready;
948                                    reason = Some(
949                                        "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
950                                            .into(),
951                                    );
952                                    repair = None;
953                                }
954                                Err(message) => {
955                                    auth = if looks_like_auth_error(&message) {
956                                        HarnessAuthState::Required
957                                    } else if configured {
958                                        HarnessAuthState::Configured
959                                    } else {
960                                        HarnessAuthState::Unknown
961                                    };
962                                    reason = Some(format!(
963                                        "No-prompt runtime handshake became unhealthy during startup: {message}"
964                                    ));
965                                    repair = Some(if auth == HarnessAuthState::Required {
966                                        format!(
967                                            "Run `{}` interactively once and complete sign-in, then probe again.",
968                                            launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
969                                        )
970                                    } else {
971                                        "Run the harness directly to inspect its startup failure, then probe again."
972                                            .into()
973                                    });
974                                }
975                            }
976                            let _ =
977                                tokio::time::timeout(Duration::from_secs(3), connection.close())
978                                    .await;
979                        }
980                        Ok(Err(error)) => {
981                            let message = truncate_text(&error.to_string(), 500);
982                            auth = if looks_like_auth_error(&message) {
983                                HarnessAuthState::Required
984                            } else if configured {
985                                HarnessAuthState::Configured
986                            } else {
987                                HarnessAuthState::Unknown
988                            };
989                            reason = Some(format!("No-prompt runtime handshake failed: {message}"));
990                            repair = Some(if auth == HarnessAuthState::Required {
991                                format!(
992                                    "Run `{}` interactively once and complete sign-in, then probe again.",
993                                    launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
994                                )
995                            } else {
996                                "Check the harness installation and run the handshake probe again."
997                                    .into()
998                            });
999                        }
1000                        Err(_) => {
1001                            reason = Some(
1002                                "No-prompt runtime handshake timed out after 30 seconds.".into(),
1003                            );
1004                            repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
1005                        }
1006                    }
1007                }
1008                Err(error) => {
1009                    reason = Some(error_message(error));
1010                }
1011            }
1012        } else if installed && configured {
1013            reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
1014        } else if installed {
1015            reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
1016            repair =
1017                Some(format!(
1018                "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
1019                launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1020            ));
1021        }
1022
1023        let effective_capabilities = if installed {
1024            descriptor.runtime.capabilities.clone()
1025        } else {
1026            unavailable_capabilities()
1027        };
1028        LocalHarness {
1029            id: descriptor.id,
1030            display_name: descriptor.display_name,
1031            supported: true,
1032            installed,
1033            executable: executable.map(|path| path.to_string_lossy().into_owned()),
1034            version,
1035            auth,
1036            runtime,
1037            protocol: descriptor.runtime.protocol,
1038            capabilities: descriptor.runtime.capabilities,
1039            effective_capabilities,
1040            sessions: HarnessSessionCounts { global, workspace },
1041            reason,
1042            repair,
1043        }
1044    }
1045
1046    fn session_counts(
1047        &self,
1048        workspace: Option<&Path>,
1049        harnesses: &[HarnessId],
1050    ) -> BTreeMap<String, usize> {
1051        let mut counts = BTreeMap::new();
1052        for session in self
1053            .catalog
1054            .discover(&DiscoveryQuery {
1055                workspace: workspace.map(Path::to_path_buf),
1056                harnesses: harnesses.to_vec(),
1057                ..DiscoveryQuery::default()
1058            })
1059            .unwrap_or_default()
1060        {
1061            *counts
1062                .entry(session.locator.harness.as_str().to_string())
1063                .or_insert(0) += 1;
1064        }
1065        counts
1066    }
1067}
1068
1069#[async_trait::async_trait]
1070impl SdkService for HarnessSessionService {
1071    fn capabilities(&self) -> SdkCapabilities {
1072        SdkCapabilities::default()
1073    }
1074
1075    async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
1076        if request.operation == SdkOperation::Events {
1077            let events = self
1078                .poll_sdk_events()
1079                .await
1080                .into_iter()
1081                .map(|(_, event)| event)
1082                .collect::<Vec<_>>();
1083            return serde_json::to_value(events).map_err(|error| {
1084                SdkError::new(
1085                    SdkErrorCode::Execution,
1086                    request.operation,
1087                    error.to_string(),
1088                )
1089            });
1090        }
1091        let method = request
1092            .operation
1093            .method()
1094            .ok_or_else(|| SdkError::unsupported(request.operation))?;
1095        let result = match request.operation {
1096            SdkOperation::Discover | SdkOperation::Load | SdkOperation::Export => {
1097                self.call(method, request.params)
1098            }
1099            SdkOperation::Start
1100            | SdkOperation::Resume
1101            | SdkOperation::Input
1102            | SdkOperation::Interrupt
1103            | SdkOperation::Steer
1104            | SdkOperation::Respond
1105            | SdkOperation::Close => self.runtime_call(method, request.params).await,
1106            SdkOperation::Events => unreachable!("handled before method dispatch"),
1107        };
1108        result.map_err(|error| sdk_error(request.operation, error))
1109    }
1110
1111    async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
1112        Ok(self
1113            .poll_sdk_events()
1114            .await
1115            .into_iter()
1116            .map(|(_, event)| event)
1117            .collect())
1118    }
1119}
1120
1121#[cfg(feature = "adapter-api")]
1122struct HostedRuntimeLease {
1123    connection: HostedHarnessConnection,
1124    _host: std::sync::Arc<HostedHarnessRuntime>,
1125    _registration: LiveRuntimeRegistration,
1126    _server: crate::server::FrontendHttpServer,
1127}
1128
1129#[async_trait::async_trait]
1130#[cfg(feature = "adapter-api")]
1131impl RuntimeConnection for HostedRuntimeLease {
1132    fn handle(&self) -> &crate::RuntimeHandle {
1133        self.connection.handle()
1134    }
1135
1136    async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
1137        self.connection.send_input(input).await
1138    }
1139
1140    async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
1141        self.connection.next_event().await
1142    }
1143
1144    async fn interrupt(&mut self) -> crate::Result<()> {
1145        self.connection.interrupt().await
1146    }
1147
1148    async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
1149        self.connection.respond(request_id, response).await
1150    }
1151
1152    async fn close(&mut self) -> crate::Result<()> {
1153        self.connection.close().await
1154    }
1155}
1156
1157async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
1158    let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
1159    loop {
1160        let now = tokio::time::Instant::now();
1161        if now >= deadline {
1162            return Ok(());
1163        }
1164        match tokio::time::timeout(deadline - now, connection.next_event()).await {
1165            Err(_) => return Ok(()),
1166            Ok(Ok(Some(event))) => {
1167                if let Some(message) = handshake_event_failure(&event) {
1168                    return Err(truncate_text(&message, 500));
1169                }
1170            }
1171            Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
1172            Ok(Err(error)) => return Err(error.to_string()),
1173        }
1174    }
1175}
1176
1177fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
1178    let detail = event
1179        .payload
1180        .get("message")
1181        .or_else(|| event.payload.get("line"))
1182        .and_then(Value::as_str)
1183        .unwrap_or(event.kind.as_str());
1184    match event.kind.as_str() {
1185        "transport_closed" => Some("runtime transport closed during startup".into()),
1186        "transport_error" => Some(format!("runtime transport error: {detail}")),
1187        "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
1188        // Stderr is retained as a runtime event, but is not transport health.
1189        // Grok, for example, can log an AuthorizationRequired error from an
1190        // optional background worker while its ACP session continues to send
1191        // updates and complete prompts normally.
1192        _ => None,
1193    }
1194}
1195
1196#[derive(Deserialize)]
1197struct LocatorParams {
1198    locator: SessionLocator,
1199    /// Optional fidelity for the READ surfaces (`sessions.load`,
1200    /// `sessions.follow`).
1201    ///
1202    /// Omitted means [`Fidelity::Semantic`]: these two methods only ever
1203    /// produce a read-only view, and a compacted or resumed-across-files
1204    /// transcript — the everyday shape of a long Claude Code session — has no
1205    /// losslessly reconstructable record graph, so refusing to render it made
1206    /// the mirror unusable rather than accurate. A caller that intends to
1207    /// CONTINUE from what it reads asks for a lossless level explicitly and
1208    /// gets the strict refusal back. Every other method (export, translate,
1209    /// branch, handoff, resume_instructions) is lossless-only and has no
1210    /// such knob.
1211    #[serde(default)]
1212    fidelity: Option<Fidelity>,
1213}
1214
1215impl LocatorParams {
1216    fn read_fidelity(&self) -> Fidelity {
1217        self.fidelity.unwrap_or(Fidelity::Semantic)
1218    }
1219}
1220
1221#[derive(Deserialize)]
1222struct UnfollowParams {
1223    subscription: String,
1224}
1225
1226#[derive(Deserialize)]
1227struct MessageSessionParams {
1228    locator: SessionLocator,
1229    text: String,
1230    /// Same storage roots discovery accepts, so a caller (and a test) can
1231    /// point the live-session registry somewhere other than `$HOME`.
1232    #[serde(default)]
1233    homes: crate::HarnessHomes,
1234}
1235
1236/// Deliver `text` into a session that is running right now, or say why not.
1237///
1238/// A refusal is a RESULT, not a JSON-RPC error: "that session is persisted
1239/// only" is an answer about the session, which a mirror renders next to the
1240/// transcript, and this service's error envelope carries no structured data
1241/// field a machine-readable reason could survive in.
1242///
1243/// `delivered_to_bus` is the honest ceiling of what the courier proves. The
1244/// message reached the receiving session's inbox; whether that session ever
1245/// reads it is governed by ITS OWN inbound controls (`crossSessionInbound`,
1246/// approval dialogs), which Supercode neither sees nor overrides.
1247#[cfg(feature = "adapter-api")]
1248async fn message_live_session(
1249    params: &MessageSessionParams,
1250    runner: &dyn crate::claude_peer::CourierRunner,
1251) -> Value {
1252    if params.locator.harness.as_str() != HarnessId::CLAUDE_CODE {
1253        return json!({
1254            "delivered_to_bus": false,
1255            "refusal": {
1256                "reason": crate::claude_peer::ClaudePeerRefusal::HarnessUnsupported.as_str(),
1257                "message": format!(
1258                    "`{}` does not publish a live-session registry; only claude-code sessions can be messaged in place",
1259                    params.locator.harness.as_str()
1260                ),
1261            },
1262        });
1263    }
1264    match crate::claude_peer::message_claude_peer(
1265        &params.homes,
1266        &params.locator.session_id,
1267        &params.text,
1268        runner,
1269    )
1270    .await
1271    {
1272        Ok(delivery) => json!({
1273            "delivered_to_bus": true,
1274            "target": {
1275                "session_id": delivery.target.session_id,
1276                "name": delivery.target.name,
1277                "pid": delivery.target.pid,
1278                "cwd": delivery.target.cwd,
1279                "status": delivery.target.status.map(|status| status.as_str()),
1280            },
1281            "courier": {
1282                "model": crate::claude_peer::COURIER_MODEL,
1283                "report": delivery.courier_report,
1284            },
1285        }),
1286        Err(refusal) => json!({
1287            "delivered_to_bus": false,
1288            "refusal": {"reason": refusal.reason.as_str(), "message": refusal.message},
1289        }),
1290    }
1291}
1292
1293/// Source identity of one follow subscription, plus the last lifecycle state
1294/// already reported on it. The follower itself stays purely persistence-facing.
1295// Only the adapter-api poll reads these; the subscription bookkeeping itself is
1296// shared by both builds.
1297#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
1298struct FollowedSource {
1299    harness: String,
1300    session_id: String,
1301    reported: Option<String>,
1302}
1303
1304#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
1305#[serde(rename_all = "kebab-case")]
1306enum TransferFormat {
1307    ClaudeCode,
1308    Codex,
1309    #[serde(rename = "opencode", alias = "open-code")]
1310    OpenCode,
1311    Pi,
1312    Grok,
1313}
1314
1315impl TransferFormat {
1316    fn id(self) -> &'static str {
1317        match self {
1318            Self::ClaudeCode => HarnessId::CLAUDE_CODE,
1319            Self::Codex => HarnessId::CODEX,
1320            Self::OpenCode => HarnessId::OPENCODE,
1321            Self::Pi => HarnessId::PI,
1322            Self::Grok => HarnessId::GROK,
1323        }
1324    }
1325}
1326
1327impl From<TransferFormat> for SessionFormat {
1328    fn from(value: TransferFormat) -> Self {
1329        match value {
1330            TransferFormat::ClaudeCode => Self::ClaudeCode,
1331            TransferFormat::Codex => Self::Codex,
1332            TransferFormat::OpenCode => Self::OpenCode,
1333            TransferFormat::Pi => Self::Pi,
1334            TransferFormat::Grok => Self::Grok,
1335        }
1336    }
1337}
1338
1339#[derive(Deserialize)]
1340struct ImportSessionParams {
1341    source_harness: TransferFormat,
1342    content: String,
1343}
1344
1345#[derive(Deserialize)]
1346struct ExportSessionParams {
1347    locator: SessionLocator,
1348    target_harness: TransferFormat,
1349}
1350
1351#[derive(Deserialize)]
1352struct BranchSessionParams {
1353    locator: SessionLocator,
1354    #[serde(default)]
1355    target_harness: Option<TransferFormat>,
1356}
1357
1358#[derive(Deserialize)]
1359struct HandoffSessionParams {
1360    locator: SessionLocator,
1361    target_harness: TransferFormat,
1362    #[serde(default)]
1363    cwd: Option<PathBuf>,
1364}
1365
1366#[derive(Debug, Clone, Copy, Default, Deserialize)]
1367#[serde(rename_all = "snake_case")]
1368enum ResumePolicy {
1369    #[default]
1370    Default,
1371    Yolo,
1372}
1373
1374#[derive(Deserialize)]
1375struct ResumeInstructionsParams {
1376    locator: SessionLocator,
1377    #[serde(default)]
1378    cwd: Option<PathBuf>,
1379    #[serde(default)]
1380    policy: ResumePolicy,
1381}
1382
1383#[derive(Serialize)]
1384struct SessionArtifact {
1385    source_harness: HarnessId,
1386    target_harness: &'static str,
1387    session_id: Option<String>,
1388    content: String,
1389    suggested_filename: String,
1390    files: Vec<SessionArtifactFile>,
1391    fidelity: Fidelity,
1392    residue: Vec<String>,
1393}
1394
1395#[derive(Serialize)]
1396struct SessionArtifactFile {
1397    path: String,
1398    content: String,
1399    role: ArtifactFileRole,
1400}
1401
1402#[derive(Serialize)]
1403#[serde(rename_all = "snake_case")]
1404enum ArtifactFileRole {
1405    Primary,
1406    Subagent,
1407    Bundle,
1408    SourceRecovery,
1409}
1410
1411#[derive(Serialize)]
1412struct StructuredLaunch {
1413    cwd: PathBuf,
1414    program: String,
1415    arguments: Vec<String>,
1416    env: BTreeMap<String, String>,
1417}
1418
1419struct HandoffInstructions {
1420    launch: StructuredLaunch,
1421    materialize: Option<StructuredLaunch>,
1422    requires_materialization: bool,
1423    note: String,
1424}
1425
1426#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
1427#[serde(rename_all = "snake_case")]
1428enum HarnessProbeLevel {
1429    #[default]
1430    Passive,
1431    Handshake,
1432}
1433
1434#[derive(Default, Deserialize)]
1435#[serde(default)]
1436struct HarnessInventoryParams {
1437    harness: Option<HarnessId>,
1438    harnesses: Vec<HarnessId>,
1439    workspace: Option<PathBuf>,
1440    probe: HarnessProbeLevel,
1441    include_sessions: bool,
1442}
1443
1444#[derive(Serialize)]
1445struct HarnessInventoryReport {
1446    probe: HarnessProbeLevel,
1447    workspace: Option<PathBuf>,
1448    harnesses: Vec<LocalHarness>,
1449}
1450
1451#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1452#[serde(rename_all = "snake_case")]
1453enum HarnessAuthState {
1454    Ready,
1455    Configured,
1456    Required,
1457    Unknown,
1458}
1459
1460#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1461#[serde(rename_all = "snake_case")]
1462enum HarnessRuntimeState {
1463    Ready,
1464    Degraded,
1465    Unavailable,
1466}
1467
1468#[derive(Serialize)]
1469struct HarnessSessionCounts {
1470    global: Option<usize>,
1471    workspace: Option<usize>,
1472}
1473
1474#[derive(Serialize)]
1475struct LocalHarness {
1476    id: HarnessId,
1477    display_name: String,
1478    supported: bool,
1479    installed: bool,
1480    executable: Option<String>,
1481    version: Option<String>,
1482    auth: HarnessAuthState,
1483    runtime: HarnessRuntimeState,
1484    protocol: String,
1485    capabilities: crate::RuntimeCapabilities,
1486    effective_capabilities: crate::RuntimeCapabilities,
1487    sessions: HarnessSessionCounts,
1488    reason: Option<String>,
1489    repair: Option<String>,
1490}
1491
1492#[derive(Clone, Deserialize)]
1493struct RuntimeBackendParams {
1494    harness: HarnessId,
1495    #[serde(default)]
1496    protocol: Option<String>,
1497    #[serde(default)]
1498    launch: Option<RuntimeLaunch>,
1499    #[serde(default)]
1500    base_url: Option<String>,
1501    #[serde(default)]
1502    policy: RuntimePolicy,
1503}
1504
1505#[derive(Debug, Clone, Copy, Default, Deserialize)]
1506#[serde(rename_all = "snake_case")]
1507enum RuntimePolicy {
1508    #[default]
1509    Default,
1510    Yolo,
1511}
1512
1513#[derive(Deserialize)]
1514struct RuntimeStartParams {
1515    #[serde(flatten)]
1516    backend: RuntimeBackendParams,
1517    cwd: PathBuf,
1518}
1519
1520#[derive(Deserialize)]
1521struct RuntimeAttachParams {
1522    #[serde(flatten)]
1523    backend: RuntimeBackendParams,
1524    runtime_id: String,
1525    #[serde(default)]
1526    cwd: Option<PathBuf>,
1527}
1528
1529#[derive(Deserialize)]
1530struct RuntimeConnectionParams {
1531    connection: String,
1532}
1533
1534#[derive(Deserialize)]
1535struct RuntimeInputParams {
1536    connection: String,
1537    text: String,
1538}
1539
1540#[derive(Deserialize)]
1541struct RuntimeRespondParams {
1542    connection: String,
1543    request_id: Value,
1544    response: Value,
1545}
1546
1547fn session_artifact(
1548    locator: &SessionLocator,
1549    session: &Session,
1550    target: TransferFormat,
1551) -> std::result::Result<SessionArtifact, ServiceError> {
1552    session_artifact_with_id(locator, session, target, None)
1553}
1554
1555fn session_artifact_with_id(
1556    locator: &SessionLocator,
1557    session: &Session,
1558    target: TransferFormat,
1559    target_session_id: Option<&str>,
1560) -> std::result::Result<SessionArtifact, ServiceError> {
1561    let format: SessionFormat = target.into();
1562    let diagonal = format.source() == session.meta.source;
1563    let has_appended_turns = session
1564        .imported_message_count
1565        .is_some_and(|imported| imported < session.messages.len());
1566    let content = if let Some(id) = target_session_id {
1567        if diagonal && format != SessionFormat::OpenCode {
1568            session
1569                .to_jsonl_spliced(format, Some(id))
1570                .map_err(operation)?
1571        } else {
1572            let mut rewritten = session.clone();
1573            rewritten.meta.session_id = Some(id.to_string());
1574            rewritten.to_jsonl(format).map_err(operation)?
1575        }
1576    } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
1577        session.raw_verbatim()
1578    } else if diagonal {
1579        session.to_jsonl_spliced(format, None).map_err(operation)?
1580    } else {
1581        session.to_jsonl(format).map_err(operation)?
1582    };
1583    let stem = sanitize_filename(
1584        target_session_id
1585            .or(session.meta.session_id.as_deref())
1586            .unwrap_or(&locator.session_id),
1587    );
1588    let suggested_filename = if diagonal && target == TransferFormat::Grok {
1589        "chat_history.jsonl".to_string()
1590    } else {
1591        format!("{stem}.{}.jsonl", target.id())
1592    };
1593    let mut files = vec![SessionArtifactFile {
1594        path: suggested_filename.clone(),
1595        content: content.clone(),
1596        role: ArtifactFileRole::Primary,
1597    }];
1598    if target == TransferFormat::ClaudeCode {
1599        let bundle_stem = Path::new(&suggested_filename)
1600            .file_stem()
1601            .and_then(|stem| stem.to_str())
1602            .unwrap_or(&stem);
1603        let mut child_paths = BTreeSet::new();
1604        for (index, subagent) in session.subagents.iter().enumerate() {
1605            let agent_id = subagent
1606                .meta
1607                .agent_id
1608                .as_deref()
1609                .map(|id| id.strip_prefix("agent-").unwrap_or(id))
1610                .map(sanitize_filename)
1611                .filter(|id| !id.is_empty())
1612                .unwrap_or_else(|| format!("subagent-{}", index + 1));
1613            let child_has_appended_turns = subagent
1614                .imported_message_count
1615                .is_some_and(|imported| imported < subagent.messages.len());
1616            let child_content = if target_session_id.is_none()
1617                && subagent.meta.source == SessionSource::ClaudeCode
1618                && subagent.raw_is_verbatim
1619                && !child_has_appended_turns
1620            {
1621                subagent.raw_verbatim()
1622            } else if subagent.meta.source == SessionSource::ClaudeCode {
1623                subagent
1624                    .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
1625                    .map_err(operation)?
1626            } else {
1627                let mut child = subagent.clone();
1628                if let Some(id) = target_session_id {
1629                    child.meta.session_id = Some(id.to_string());
1630                }
1631                child
1632                    .to_jsonl(SessionFormat::ClaudeCode)
1633                    .map_err(operation)?
1634            };
1635            let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
1636            if !child_paths.insert(path.clone()) {
1637                return Err(ServiceError::Operation(format!(
1638                    "Claude subagent ids collide at artifact path `{path}`"
1639                )));
1640            }
1641            files.push(SessionArtifactFile {
1642                path,
1643                content: child_content,
1644                role: ArtifactFileRole::Subagent,
1645            });
1646        }
1647    }
1648    if diagonal && target == TransferFormat::Grok {
1649        append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
1650    }
1651    if !diagonal || !session.raw_is_verbatim {
1652        files.push(SessionArtifactFile {
1653            path: "recovery/source.supercode.jsonl".into(),
1654            content: session.to_native_jsonl(),
1655            role: ArtifactFileRole::SourceRecovery,
1656        });
1657        for (index, subagent) in session.subagents.iter().enumerate() {
1658            let id = subagent
1659                .meta
1660                .agent_id
1661                .as_deref()
1662                .map(sanitize_filename)
1663                .unwrap_or_else(|| format!("subagent-{}", index + 1));
1664            files.push(SessionArtifactFile {
1665                path: format!("recovery/subagents/{id}.supercode.jsonl"),
1666                content: subagent.to_native_jsonl(),
1667                role: ArtifactFileRole::SourceRecovery,
1668            });
1669        }
1670    }
1671    if !diagonal && session.meta.source == SessionSource::Grok {
1672        append_grok_bundle_files(
1673            locator,
1674            "recovery/grok/",
1675            ArtifactFileRole::SourceRecovery,
1676            &mut files,
1677        )?;
1678    }
1679    let (fidelity, residue) = if diagonal
1680        && target_session_id.is_none()
1681        && session.raw_is_verbatim
1682        && !has_appended_turns
1683    {
1684        (Fidelity::ByteLossless, Vec::new())
1685    } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
1686        (
1687            Fidelity::ValueLossless,
1688            vec![if target_session_id.is_some() {
1689                "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
1690            } else {
1691                "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
1692            }],
1693        )
1694    } else {
1695        (
1696            Fidelity::Semantic,
1697            vec!["target schema has no portable slot for every source-native record and metadata field".into()],
1698        )
1699    };
1700    Ok(SessionArtifact {
1701        source_harness: locator.harness.clone(),
1702        target_harness: target.id(),
1703        session_id: target_session_id
1704            .map(str::to_string)
1705            .or_else(|| session.meta.session_id.clone()),
1706        content,
1707        suggested_filename,
1708        files,
1709        fidelity,
1710        residue,
1711    })
1712}
1713
1714fn append_grok_bundle_files(
1715    locator: &SessionLocator,
1716    prefix: &str,
1717    role: ArtifactFileRole,
1718    files: &mut Vec<SessionArtifactFile>,
1719) -> std::result::Result<(), ServiceError> {
1720    let primary = locator.storage.path();
1721    if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
1722        return Err(ServiceError::Operation(format!(
1723            "Grok bundle locator must name chat_history.jsonl, got {}",
1724            primary.display()
1725        )));
1726    }
1727    let parent = primary.parent().ok_or_else(|| {
1728        ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
1729    })?;
1730    for name in ["summary.json", "updates.jsonl"] {
1731        let path = parent.join(name);
1732        let metadata = match std::fs::symlink_metadata(&path) {
1733            Ok(metadata) => metadata,
1734            Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
1735            Err(error) => return Err(ServiceError::Operation(error.to_string())),
1736        };
1737        if metadata.file_type().is_symlink() || !metadata.is_file() {
1738            return Err(ServiceError::Operation(format!(
1739                "refusing non-regular Grok bundle member {}",
1740                path.display()
1741            )));
1742        }
1743        let content = std::fs::read_to_string(&path).map_err(|error| {
1744            ServiceError::Operation(format!(
1745                "Grok bundle member {} is not representable as UTF-8: {error}",
1746                path.display()
1747            ))
1748        })?;
1749        files.push(SessionArtifactFile {
1750            path: format!("{prefix}{name}"),
1751            content,
1752            role: match role {
1753                ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
1754                _ => ArtifactFileRole::SourceRecovery,
1755            },
1756        });
1757    }
1758    Ok(())
1759}
1760
1761fn handoff_artifact(
1762    locator: &SessionLocator,
1763    session: &Session,
1764    target: TransferFormat,
1765    cwd: &Path,
1766) -> std::result::Result<SessionArtifact, ServiceError> {
1767    if target != TransferFormat::Grok {
1768        let target_session_id = target_session_id(target);
1769        return session_artifact_with_id(locator, session, target, Some(&target_session_id));
1770    }
1771
1772    // Stock Grok's importer accepts Claude/Codex transcripts and materializes its own
1773    // multi-file session bundle. A synthesized Grok chat_history.jsonl alone is not a
1774    // resumable handoff because updates.jsonl is the authoritative restore log.
1775    let mut importable = session.clone();
1776    // The Claude importer validates sessionId as a UUID. Source harness identities
1777    // are not portable (OpenCode, for example, uses `ses_...`), and a handoff must
1778    // not overwrite an existing target session when the source already uses UUIDs.
1779    // Mint a distinct target identity and still bind the importer-returned ID at
1780    // launch time because the importer remains the authority on materialization.
1781    importable.meta.session_id = Some(target_session_id(TransferFormat::ClaudeCode));
1782    importable.meta.cwd = Some(if cwd.is_absolute() {
1783        cwd.to_path_buf()
1784    } else {
1785        std::env::current_dir()
1786            .map_err(|error| ServiceError::Operation(error.to_string()))?
1787            .join(cwd)
1788    });
1789    let content = importable
1790        .to_jsonl(SessionFormat::ClaudeCode)
1791        .map_err(operation)?;
1792    let stem = sanitize_filename(
1793        importable
1794            .meta
1795            .session_id
1796            .as_deref()
1797            .unwrap_or(&locator.session_id),
1798    );
1799    let suggested_filename = format!("{stem}.grok-import.claude-code.jsonl");
1800    Ok(SessionArtifact {
1801        source_harness: locator.harness.clone(),
1802        // This names the artifact's actual wire format. The requested handoff target
1803        // remains Grok; its official importer is the materialization boundary.
1804        target_harness: TransferFormat::ClaudeCode.id(),
1805        session_id: importable.meta.session_id.clone(),
1806        content: content.clone(),
1807        suggested_filename: suggested_filename.clone(),
1808        files: vec![SessionArtifactFile {
1809            path: suggested_filename,
1810            content,
1811            role: ArtifactFileRole::Primary,
1812        }],
1813        fidelity: Fidelity::Semantic,
1814        residue: vec!["Grok's stock importer accepts a Claude Code transcript, not a complete Grok updates/session bundle".into()],
1815    })
1816}
1817
1818fn target_session_id(target: TransferFormat) -> String {
1819    let uuid = generated_session_id();
1820    match target {
1821        TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
1822        TransferFormat::ClaudeCode
1823        | TransferFormat::Codex
1824        | TransferFormat::Pi
1825        | TransferFormat::Grok => uuid,
1826    }
1827}
1828
1829fn sanitize_filename(value: &str) -> String {
1830    let value = value
1831        .chars()
1832        .map(|character| {
1833            if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
1834                character
1835            } else {
1836                '-'
1837            }
1838        })
1839        .collect::<String>();
1840    let value = value.trim_matches('-');
1841    if value.is_empty() {
1842        "session".into()
1843    } else {
1844        value.chars().take(100).collect()
1845    }
1846}
1847
1848fn handoff_instructions(
1849    target: TransferFormat,
1850    session_id: &str,
1851    cwd: &Path,
1852) -> HandoffInstructions {
1853    let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
1854        cwd: cwd.to_path_buf(),
1855        program: program.into(),
1856        arguments,
1857        env: BTreeMap::new(),
1858    };
1859    match target {
1860        TransferFormat::ClaudeCode => HandoffInstructions {
1861            launch: launch("claude", vec!["--resume".into(), session_id.into()]),
1862            materialize: None,
1863            requires_materialization: true,
1864            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(),
1865        },
1866        TransferFormat::Codex => HandoffInstructions {
1867            launch: launch("codex", vec!["resume".into(), session_id.into()]),
1868            materialize: None,
1869            requires_materialization: true,
1870            note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
1871        },
1872        TransferFormat::OpenCode => HandoffInstructions {
1873            launch: launch("opencode", vec!["--session".into(), session_id.into()]),
1874            materialize: Some(launch(
1875                "opencode",
1876                vec!["import".into(), "{artifact_path}".into()],
1877            )),
1878            requires_materialization: true,
1879            note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
1880        },
1881        TransferFormat::Pi => HandoffInstructions {
1882            launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
1883            materialize: None,
1884            requires_materialization: true,
1885            note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
1886        },
1887        TransferFormat::Grok => HandoffInstructions {
1888            launch: launch(
1889                "grok",
1890                vec![
1891                    "--resume".into(),
1892                    "{imported_session_id}".into(),
1893                    "--fork-session".into(),
1894                ],
1895            ),
1896            materialize: Some(launch(
1897                "grok",
1898                vec!["import".into(), "--json".into(), "{artifact_path}".into()],
1899            )),
1900            requires_materialization: true,
1901            note: "The artifact is Claude Code JSONL for Grok's official importer. Write it to a file, run the materialize command, read sessionId from its NDJSON outcome=imported record, replace {imported_session_id} in the launch arguments, then launch a writable fork of the imported session.".into(),
1902        },
1903    }
1904}
1905
1906fn resume_launch(
1907    harness: &str,
1908    session_id: &str,
1909    cwd: &Path,
1910    policy: ResumePolicy,
1911) -> std::result::Result<StructuredLaunch, ServiceError> {
1912    let mut arguments = Vec::new();
1913    let program = match harness {
1914        HarnessId::GROK => {
1915            if matches!(policy, ResumePolicy::Yolo) {
1916                arguments.extend([
1917                    "--sandbox".into(),
1918                    "workspace".into(),
1919                    "--always-approve".into(),
1920                ]);
1921            }
1922            arguments.extend(["--resume".into(), session_id.into()]);
1923            "grok"
1924        }
1925        HarnessId::CODEX => {
1926            if matches!(policy, ResumePolicy::Yolo) {
1927                arguments.extend([
1928                    "--dangerously-bypass-approvals-and-sandbox".into(),
1929                    "--dangerously-bypass-hook-trust".into(),
1930                ]);
1931            }
1932            arguments.extend(["resume".into(), session_id.into()]);
1933            "codex"
1934        }
1935        HarnessId::CLAUDE_CODE => {
1936            if matches!(policy, ResumePolicy::Yolo) {
1937                arguments.push("--dangerously-skip-permissions".into());
1938            }
1939            arguments.extend(["--resume".into(), session_id.into()]);
1940            "claude"
1941        }
1942        HarnessId::PI => {
1943            if matches!(policy, ResumePolicy::Yolo) {
1944                arguments.push("--approve".into());
1945            }
1946            arguments.extend(["--session".into(), session_id.into()]);
1947            "pi"
1948        }
1949        HarnessId::OPENCODE => {
1950            arguments.extend(["--session".into(), session_id.into()]);
1951            "opencode"
1952        }
1953        other => {
1954            return Err(ServiceError::InvalidParams(format!(
1955                "no structured resume launch is registered for harness `{other}`"
1956            )))
1957        }
1958    };
1959    Ok(StructuredLaunch {
1960        cwd: cwd.to_path_buf(),
1961        program: program.into(),
1962        arguments,
1963        env: BTreeMap::new(),
1964    })
1965}
1966
1967fn runtime_backend(
1968    params: &RuntimeBackendParams,
1969) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
1970    if params.protocol.as_deref() == Some("acp") {
1971        let launch = params
1972            .launch
1973            .clone()
1974            .or_else(|| {
1975                harness_support_registry()
1976                    .harnesses
1977                    .into_iter()
1978                    .find(|harness| harness.id == params.harness)
1979                    .filter(|harness| {
1980                        harness.runtime.implementation == ImplementationKind::GenericProtocol
1981                            && harness.runtime.protocol.starts_with("acp")
1982                    })
1983                    .and_then(|harness| harness.runtime.default_launch)
1984            })
1985            .ok_or_else(|| {
1986                ServiceError::InvalidParams(
1987                    "an ACP runtime requires `launch` unless the harness has a registered default"
1988                        .into(),
1989                )
1990            })?;
1991        let resume_session = params.harness.as_str() == HarnessId::GROK;
1992        return Ok(Box::new(
1993            AcpRuntimeBackend::new(params.harness.clone(), launch)
1994                .with_resume_support(resume_session),
1995        ));
1996    }
1997    let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
1998        HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
1999        HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
2000        HarnessId::PI => Box::new(PiRuntimeBackend::new()),
2001        HarnessId::OPENCODE => match &params.base_url {
2002            Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
2003            None => Box::new(OpenCodeRuntimeBackend::new()),
2004        },
2005        HarnessId::GROK => {
2006            let descriptor = harness_support_registry()
2007                .harnesses
2008                .into_iter()
2009                .find(|harness| harness.id.as_str() == HarnessId::GROK)
2010                .expect("Grok descriptor is part of the canonical registry");
2011            Box::new(
2012                AcpRuntimeBackend::new(
2013                    descriptor.id,
2014                    descriptor
2015                        .runtime
2016                        .default_launch
2017                        .expect("Grok registry includes its ACP launch"),
2018                )
2019                .with_resume_support(true),
2020            )
2021        }
2022        other => {
2023            return Err(ServiceError::InvalidParams(format!(
2024                "no runtime adapter for harness `{other}`; use protocol `acp` with a launch command"
2025            )))
2026        }
2027    };
2028    Ok(backend)
2029}
2030
2031fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
2032    if let Some(launch) = &params.launch {
2033        return Some(launch.clone());
2034    }
2035    if !matches!(params.policy, RuntimePolicy::Yolo) {
2036        return None;
2037    }
2038    let launch = match params.harness.as_str() {
2039        HarnessId::GROK => RuntimeLaunch {
2040            program: "grok".into(),
2041            arguments: vec![
2042                "--sandbox".into(),
2043                "workspace".into(),
2044                "--always-approve".into(),
2045                "agent".into(),
2046                "--no-leader".into(),
2047                "stdio".into(),
2048            ],
2049            env: BTreeMap::from([("GROK_AGENT_DASHBOARD".into(), "0".into())]),
2050        },
2051        HarnessId::CODEX => RuntimeLaunch {
2052            program: "codex".into(),
2053            arguments: vec![
2054                "--dangerously-bypass-approvals-and-sandbox".into(),
2055                "--dangerously-bypass-hook-trust".into(),
2056                "app-server".into(),
2057            ],
2058            env: BTreeMap::new(),
2059        },
2060        HarnessId::CLAUDE_CODE => RuntimeLaunch {
2061            program: "claude".into(),
2062            arguments: vec![
2063                "--dangerously-skip-permissions".into(),
2064                "--print".into(),
2065                "--input-format".into(),
2066                "stream-json".into(),
2067                "--output-format".into(),
2068                "stream-json".into(),
2069                "--verbose".into(),
2070            ],
2071            env: BTreeMap::new(),
2072        },
2073        HarnessId::PI => RuntimeLaunch {
2074            program: "pi".into(),
2075            arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
2076            env: BTreeMap::new(),
2077        },
2078        HarnessId::OPENCODE => RuntimeLaunch {
2079            program: "opencode".into(),
2080            arguments: vec!["serve".into()],
2081            env: BTreeMap::new(),
2082        },
2083        _ => return None,
2084    };
2085    Some(launch)
2086}
2087
2088fn find_executable(program: &str) -> Option<PathBuf> {
2089    let candidate = PathBuf::from(program);
2090    if candidate.components().count() > 1 {
2091        return candidate.is_file().then_some(candidate);
2092    }
2093    let path = std::env::var_os("PATH")?;
2094    for directory in std::env::split_paths(&path) {
2095        let candidate = directory.join(program);
2096        if candidate.is_file() {
2097            return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
2098        }
2099        #[cfg(windows)]
2100        {
2101            for extension in ["exe", "cmd", "bat"] {
2102                let candidate = directory.join(format!("{program}.{extension}"));
2103                if candidate.is_file() {
2104                    return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
2105                }
2106            }
2107        }
2108    }
2109    None
2110}
2111
2112async fn executable_version(executable: &Path) -> Option<String> {
2113    let mut command = tokio::process::Command::new(executable);
2114    command
2115        .arg("--version")
2116        .stdin(std::process::Stdio::null())
2117        .stdout(std::process::Stdio::piped())
2118        .stderr(std::process::Stdio::piped())
2119        .kill_on_drop(true);
2120    let output = tokio::time::timeout(Duration::from_secs(3), command.output())
2121        .await
2122        .ok()?
2123        .ok()?;
2124    let stdout = String::from_utf8_lossy(&output.stdout);
2125    let stderr = String::from_utf8_lossy(&output.stderr);
2126    stdout
2127        .lines()
2128        .chain(stderr.lines())
2129        .map(str::trim)
2130        .find(|line| !line.is_empty())
2131        .map(|line| truncate_text(line, 200))
2132}
2133
2134fn auth_evidence(harness: &str) -> bool {
2135    let env_names: &[&str] = match harness {
2136        HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
2137        HarnessId::CODEX => &["OPENAI_API_KEY"],
2138        HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
2139        HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
2140        HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
2141        _ => &[],
2142    };
2143    if env_names
2144        .iter()
2145        .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
2146    {
2147        return true;
2148    }
2149    let Some(home) = std::env::var_os("HOME").map(PathBuf::from) else {
2150        return false;
2151    };
2152    let files: Vec<PathBuf> = match harness {
2153        HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
2154        HarnessId::CODEX => vec![home.join(".codex/auth.json")],
2155        HarnessId::OPENCODE => vec![
2156            home.join(".local/share/opencode/auth.json"),
2157            home.join(".config/opencode/auth.json"),
2158        ],
2159        HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
2160        HarnessId::GROK => vec![home.join(".grok/auth.json")],
2161        _ => Vec::new(),
2162    };
2163    if files.into_iter().any(|path| {
2164        std::fs::metadata(path)
2165            .map(|metadata| metadata.is_file() && metadata.len() > 2)
2166            .unwrap_or(false)
2167    }) {
2168        return true;
2169    }
2170    // macOS keeps Claude Code's OAuth login in the Keychain, so
2171    // `.claude/.credentials.json` never exists there and the file probe above
2172    // reports a signed-in install as unauthenticated forever. A completed
2173    // login also writes an `oauthAccount` record into `~/.claude.json` on
2174    // every platform — file-based, prompt-free evidence (querying the
2175    // Keychain itself from an unsigned daemon can raise a UI prompt).
2176    if harness == HarnessId::CLAUDE_CODE {
2177        return std::fs::read_to_string(home.join(".claude.json"))
2178            .map(|text| text.contains("\"oauthAccount\""))
2179            .unwrap_or(false);
2180    }
2181    false
2182}
2183
2184fn looks_like_auth_error(message: &str) -> bool {
2185    let message = message.to_ascii_lowercase();
2186    [
2187        "auth",
2188        "login",
2189        "sign in",
2190        "sign-in",
2191        "credential",
2192        "unauthorized",
2193        "forbidden",
2194        "token",
2195    ]
2196    .iter()
2197    .any(|needle| message.contains(needle))
2198}
2199
2200fn unavailable_capabilities() -> crate::RuntimeCapabilities {
2201    crate::RuntimeCapabilities {
2202        start_session: false,
2203        resume_session: false,
2204        attach_existing_process: false,
2205        send_input: false,
2206        stream_events: false,
2207        interrupt: false,
2208        respond_to_requests: false,
2209    }
2210}
2211
2212fn truncate_text(text: &str, max_chars: usize) -> String {
2213    let mut chars = text.chars();
2214    let truncated = chars.by_ref().take(max_chars).collect::<String>();
2215    if chars.next().is_some() {
2216        format!("{truncated}…")
2217    } else {
2218        truncated
2219    }
2220}
2221
2222fn error_message(error: ServiceError) -> String {
2223    match error {
2224        ServiceError::InvalidParams(message)
2225        | ServiceError::Operation(message)
2226        | ServiceError::UnsupportedAction(message) => message,
2227        ServiceError::MethodNotFound => "runtime adapter is not available".into(),
2228        ServiceError::Sdk(error) => error.to_string(),
2229    }
2230}
2231
2232enum ServiceError {
2233    InvalidParams(String),
2234    MethodNotFound,
2235    UnsupportedAction(String),
2236    Operation(String),
2237    Sdk(SdkError),
2238}
2239
2240fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
2241    match error {
2242        ServiceError::InvalidParams(message) => {
2243            SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
2244        }
2245        ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
2246            SdkError::unsupported(operation)
2247        }
2248        ServiceError::Operation(message) => {
2249            let code = if message.contains("already in progress") {
2250                SdkErrorCode::Busy
2251            } else if message.contains("not supported by this runtime") {
2252                SdkErrorCode::UnsupportedAction
2253            } else if message.contains("unknown runtime connection") {
2254                SdkErrorCode::NotFound
2255            } else {
2256                SdkErrorCode::Execution
2257            };
2258            SdkError::new(code, operation, message)
2259        }
2260        ServiceError::Sdk(error) => error,
2261    }
2262}
2263
2264fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
2265    let error_code = error.code();
2266    let code = match error_code {
2267        SdkErrorCode::Unauthenticated => -32030,
2268        SdkErrorCode::Unauthorized => -32031,
2269        SdkErrorCode::ControllerRequired => -32032,
2270        SdkErrorCode::LeaseExpired => -32033,
2271        SdkErrorCode::InvalidArgument => -32602,
2272        SdkErrorCode::NotFound => -32004,
2273        SdkErrorCode::Busy => -32000,
2274        SdkErrorCode::UnsupportedAction => -32020,
2275        SdkErrorCode::Execution => -32002,
2276        SdkErrorCode::Transport => -32003,
2277    };
2278    json!({
2279        "jsonrpc": "2.0",
2280        "id": id,
2281        "error": {
2282            "code": code,
2283            "name": error_code,
2284            "operation": error.operation(),
2285            "message": error.to_string(),
2286        },
2287    })
2288}
2289
2290fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
2291    serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
2292}
2293
2294fn operation(error: crate::Error) -> ServiceError {
2295    match error {
2296        crate::Error::Sdk(error) => ServiceError::Sdk(error),
2297        error => ServiceError::Operation(error.to_string()),
2298    }
2299}
2300
2301fn rpc_error(id: Value, code: i64, message: &str) -> Value {
2302    json!({
2303        "jsonrpc": "2.0",
2304        "id": id,
2305        "error": {"code": code, "message": message},
2306    })
2307}
2308
2309#[cfg(test)]
2310mod tests {
2311    use super::*;
2312    use crate::{HarnessEvent, HarnessId, RuntimeEndpoint, RuntimeHandle, StorageLocator};
2313    use async_trait::async_trait;
2314    use std::path::PathBuf;
2315
2316    struct EndingRuntime {
2317        handle: RuntimeHandle,
2318        event: Option<HarnessEvent>,
2319    }
2320
2321    #[async_trait]
2322    impl RuntimeConnection for EndingRuntime {
2323        fn handle(&self) -> &RuntimeHandle {
2324            &self.handle
2325        }
2326
2327        async fn send_input(&mut self, _input: RuntimeInput) -> crate::Result<Option<String>> {
2328            unreachable!("ending runtime does not accept input")
2329        }
2330
2331        async fn next_event(&mut self) -> crate::Result<Option<HarnessEvent>> {
2332            Ok(self.event.take())
2333        }
2334
2335        async fn interrupt(&mut self) -> crate::Result<()> {
2336            Ok(())
2337        }
2338
2339        async fn respond(&mut self, _request_id: Value, _response: Value) -> crate::Result<()> {
2340            Ok(())
2341        }
2342
2343        async fn close(&mut self) -> crate::Result<()> {
2344            Ok(())
2345        }
2346    }
2347
2348    fn ending_runtime(event: Option<HarnessEvent>) -> Box<dyn RuntimeConnection> {
2349        Box::new(EndingRuntime {
2350            handle: RuntimeHandle {
2351                harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2352                runtime_id: "ending-session".into(),
2353                endpoint: RuntimeEndpoint::LocalProcess {
2354                    pid: None,
2355                    command: vec!["ending-runtime".into()],
2356                    protocol: "test".into(),
2357                },
2358            },
2359            event,
2360        })
2361    }
2362
2363    fn request(id: u64, method: &str, params: Value) -> Value {
2364        json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})
2365    }
2366
2367    fn pi_locator() -> SessionLocator {
2368        SessionLocator {
2369            harness: HarnessId::from(HarnessId::PI),
2370            session_id: "1e6f2a3b-0000-4000-8000-000000000001".into(),
2371            storage: StorageLocator::File {
2372                path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2373                    .join("tests/fixtures/pi_session.jsonl"),
2374            },
2375        }
2376    }
2377
2378    fn opencode_locator() -> SessionLocator {
2379        let session_id = "ses_fixtureAAAAAAAAAAAAAAA1";
2380        SessionLocator {
2381            harness: HarnessId::from(HarnessId::OPENCODE),
2382            session_id: session_id.into(),
2383            storage: StorageLocator::Sqlite {
2384                path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2385                    .join("tests/fixtures/opencode_fixture/opencode.db"),
2386                selector: session_id.into(),
2387            },
2388        }
2389    }
2390
2391    fn grok_locator() -> SessionLocator {
2392        SessionLocator {
2393            harness: HarnessId::from(HarnessId::GROK),
2394            session_id: "73c09283-4b33-41fa-90f1-0bcb0f7be523".into(),
2395            storage: StorageLocator::File {
2396                path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2397                    .join("tests/fixtures/grok_session/chat_history.jsonl"),
2398            },
2399        }
2400    }
2401
2402    #[test]
2403    fn capabilities_are_explicit_and_versioned() {
2404        let mut service = HarnessSessionService::new();
2405        let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
2406        assert_eq!(response["result"]["version"], HARNESS_SERVICE_VERSION);
2407        assert_eq!(
2408            response["result"]["sdk"]["schema_version"],
2409            crate::SDK_SCHEMA_VERSION
2410        );
2411        assert_eq!(
2412            response["result"]["sdk"]["operations"]
2413                .as_array()
2414                .unwrap()
2415                .len(),
2416            SdkOperation::ALL.len()
2417        );
2418        assert_eq!(response["result"]["harnesses"].as_array().unwrap().len(), 5);
2419        assert!(response["result"]["harnesses"]
2420            .as_array()
2421            .unwrap()
2422            .iter()
2423            .any(|harness| harness == HarnessId::GROK));
2424    }
2425
2426    #[test]
2427    fn handshake_health_uses_protocol_liveness_not_stderr_severity() {
2428        let noisy_stderr = crate::HarnessEvent {
2429            sequence: None,
2430            kind: "transport_stderr".into(),
2431            payload: json!({"line": "ERROR optional worker AuthorizationRequired"}),
2432        };
2433        assert_eq!(handshake_event_failure(&noisy_stderr), None);
2434
2435        let closed = crate::HarnessEvent {
2436            sequence: None,
2437            kind: "transport_closed".into(),
2438            payload: json!({}),
2439        };
2440        assert!(handshake_event_failure(&closed).is_some());
2441    }
2442
2443    #[tokio::test]
2444    async fn runtime_eof_is_notified_and_removed_for_raw_and_explicit_close() {
2445        let mut service = HarnessSessionService::new();
2446        service
2447            .runtimes
2448            .insert("raw-eof".into(), ending_runtime(None));
2449        service.runtimes.insert(
2450            "explicit-close".into(),
2451            ending_runtime(Some(HarnessEvent {
2452                sequence: None,
2453                kind: "transport_closed".into(),
2454                payload: json!({"message": "native transport exited"}),
2455            })),
2456        );
2457
2458        let notifications = service.poll_runtimes().await;
2459
2460        assert_eq!(notifications.len(), 2);
2461        assert!(notifications
2462            .iter()
2463            .all(|notification| { notification["params"]["event"]["kind"] == "transport_closed" }));
2464        assert!(notifications.iter().all(|notification| {
2465            notification["params"]["session_id"] == "ending-session"
2466                && notification["params"]["connection"].is_string()
2467        }));
2468        let mut sequences = notifications
2469            .iter()
2470            .filter_map(|notification| notification["params"]["sequence"].as_u64())
2471            .collect::<Vec<_>>();
2472        sequences.sort_unstable();
2473        assert_eq!(sequences, vec![1, 2]);
2474        assert!(service.runtimes.is_empty());
2475    }
2476
2477    #[tokio::test]
2478    async fn sdk_facade_returns_named_unsupported_actions() {
2479        let mut service = HarnessSessionService::new();
2480        let error = service
2481            .execute(SdkRequest {
2482                operation: SdkOperation::Steer,
2483                params: json!({"connection": "runtime-1", "text": "go left"}),
2484            })
2485            .await
2486            .unwrap_err();
2487        assert_eq!(error.code(), SdkErrorCode::UnsupportedAction);
2488        assert_eq!(error.operation(), Some(SdkOperation::Steer));
2489
2490        let response = service
2491            .handle_async(request(
2492                7,
2493                "harness.v1.runtimes.steer",
2494                json!({"connection": "runtime-1", "text": "go left"}),
2495            ))
2496            .await;
2497        assert_eq!(response["error"]["name"], "unsupported_action");
2498        assert_eq!(response["error"]["operation"], "steer");
2499    }
2500
2501    #[test]
2502    fn support_report_and_grok_default_binding_share_the_registry() {
2503        let mut service = HarnessSessionService::new();
2504        let response = service.handle(request(1, "harness.v1.support.report", json!({})));
2505        assert_eq!(response["result"]["schema"], crate::SUPPORT_REGISTRY_SCHEMA);
2506        let params = RuntimeBackendParams {
2507            harness: HarnessId::from(HarnessId::GROK),
2508            protocol: None,
2509            launch: None,
2510            base_url: None,
2511            policy: RuntimePolicy::Default,
2512        };
2513        let backend = match runtime_backend(&params) {
2514            Ok(backend) => backend,
2515            Err(_) => panic!("Grok should bind through its registered ACP launch"),
2516        };
2517        assert_eq!(backend.harness().as_str(), HarnessId::GROK);
2518        assert!(backend.capabilities().start_session);
2519        let registered = harness_support_registry()
2520            .harnesses
2521            .into_iter()
2522            .find(|harness| harness.id.as_str() == HarnessId::GROK)
2523            .and_then(|harness| harness.runtime.default_launch)
2524            .unwrap();
2525        assert!(!registered
2526            .arguments
2527            .iter()
2528            .any(|argument| argument == "--always-approve"));
2529        assert!(runtime_launch(&params).is_none());
2530
2531        let yolo = RuntimeBackendParams {
2532            policy: RuntimePolicy::Yolo,
2533            ..params
2534        };
2535        assert!(runtime_launch(&yolo)
2536            .unwrap()
2537            .arguments
2538            .iter()
2539            .any(|argument| argument == "--always-approve"));
2540
2541        let mismatched_protocol = RuntimeBackendParams {
2542            harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2543            protocol: Some("acp".into()),
2544            launch: None,
2545            base_url: None,
2546            policy: RuntimePolicy::Default,
2547        };
2548        assert!(runtime_backend(&mismatched_protocol).is_err());
2549    }
2550
2551    #[test]
2552    fn load_follow_and_unfollow_share_the_same_locator() {
2553        let mut service = HarnessSessionService::new();
2554        let locator = pi_locator();
2555        let loaded = service.handle(request(
2556            1,
2557            "harness.v1.sessions.load",
2558            json!({"locator": locator}),
2559        ));
2560        assert_eq!(
2561            loaded["result"]["session"]["session_id"],
2562            locator.session_id
2563        );
2564
2565        let followed = service.handle(request(
2566            2,
2567            "harness.v1.sessions.follow",
2568            json!({"locator": locator}),
2569        ));
2570        assert_eq!(followed["result"]["subscription"], "sub-1");
2571        assert_eq!(followed["result"]["initial"]["type"], "session_snapshot");
2572        assert!(service.poll().is_empty());
2573
2574        let unfollowed = service.handle(request(
2575            3,
2576            "harness.v1.sessions.unfollow",
2577            json!({"subscription": "sub-1"}),
2578        ));
2579        assert_eq!(unfollowed["result"]["removed"], true);
2580    }
2581
2582    #[test]
2583    fn import_translate_branch_and_handoff_use_typed_artifacts() {
2584        let mut service = HarnessSessionService::new();
2585        let locator = pi_locator();
2586        let translated = service.handle(request(
2587            1,
2588            "harness.v1.sessions.translate",
2589            json!({"locator": locator, "target_harness": "grok"}),
2590        ));
2591        assert_eq!(translated["result"]["artifact"]["source_harness"], "pi");
2592        assert_eq!(translated["result"]["artifact"]["target_harness"], "grok");
2593        assert!(translated["result"]["artifact"]["content"]
2594            .as_str()
2595            .is_some_and(|content| !content.is_empty()));
2596
2597        for target in ["opencode", "open-code"] {
2598            let opencode = service.handle(request(
2599                6,
2600                "harness.v1.sessions.translate",
2601                json!({"locator": locator, "target_harness": target}),
2602            ));
2603            assert_eq!(opencode["result"]["artifact"]["target_harness"], "opencode");
2604        }
2605
2606        let imported = service.handle(request(
2607            2,
2608            "harness.v1.sessions.import",
2609            json!({
2610                "source_harness": "grok",
2611                "content": translated["result"]["artifact"]["content"],
2612            }),
2613        ));
2614        assert_eq!(imported["result"]["session"]["source"], "grok");
2615
2616        let branched = service.handle(request(
2617            3,
2618            "harness.v1.sessions.branch",
2619            json!({"locator": locator, "target_harness": "codex"}),
2620        ));
2621        assert_eq!(branched["result"]["parent"]["harness"], "pi");
2622        assert!(branched["result"]["bootstrap_prompt"]
2623            .as_str()
2624            .unwrap()
2625            .contains("frozen parent transcript"));
2626        assert_eq!(branched["result"]["artifact"]["target_harness"], "codex");
2627
2628        let handoff = service.handle(request(
2629            4,
2630            "harness.v1.sessions.handoff",
2631            json!({"locator": locator, "target_harness": "pi", "cwd": "/tmp/project"}),
2632        ));
2633        assert_eq!(handoff["result"]["launch"]["program"], "pi");
2634        assert_eq!(handoff["result"]["launch"]["cwd"], "/tmp/project");
2635        assert_eq!(handoff["result"]["requires_materialization"], true);
2636
2637        let resumed = service.handle(request(
2638            5,
2639            "harness.v1.sessions.resume_instructions",
2640            json!({"locator": locator, "cwd": "/tmp/project", "policy": "yolo"}),
2641        ));
2642        assert_eq!(resumed["result"]["launch"]["program"], "pi");
2643        assert_eq!(resumed["result"]["launch"]["arguments"][0], "--approve");
2644    }
2645
2646    #[test]
2647    fn read_surfaces_view_a_severed_claude_graph_while_transfer_still_refuses_it() {
2648        let temp = std::env::temp_dir().join(format!(
2649            "supercode-severed-view-{}-{}",
2650            std::process::id(),
2651            generated_session_id()
2652        ));
2653        std::fs::create_dir_all(&temp).unwrap();
2654        let path = temp.join("severed.jsonl");
2655        // A live record whose parent was pruned — what a compacted or
2656        // resumed-across-files Claude Code session looks like on disk.
2657        std::fs::write(
2658            &path,
2659            concat!(
2660                r#"{"type":"user","uuid":"orphan-u","parentUuid":null,"message":{"role":"user","content":"stranded prompt"}}"#,
2661                "\n",
2662                r#"{"type":"assistant","uuid":"live-a","parentUuid":"pruned","message":{"id":"m","role":"assistant","content":[{"type":"text","text":"live answer"}]}}"#,
2663                "\n",
2664            ),
2665        )
2666        .unwrap();
2667        let locator = SessionLocator {
2668            harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2669            session_id: "severed".into(),
2670            storage: StorageLocator::File { path },
2671        };
2672        let mut service = HarnessSessionService::new();
2673
2674        let viewed = service.handle(request(
2675            1,
2676            "harness.v1.sessions.load",
2677            json!({"locator": locator}),
2678        ));
2679        let session = &viewed["result"]["session"];
2680        assert_eq!(session["fidelity"], "semantic");
2681        assert_eq!(session["messages"].as_array().unwrap().len(), 2);
2682        assert!(session["residue"].as_array().unwrap().iter().any(|entry| {
2683            entry
2684                .as_str()
2685                .is_some_and(|entry| entry.contains("live-a") && entry.contains("pruned"))
2686        }));
2687
2688        // Asking a READ surface for a lossless reconstruction gets the strict
2689        // refusal back, unchanged.
2690        let strict = service.handle(request(
2691            2,
2692            "harness.v1.sessions.load",
2693            json!({"locator": locator, "fidelity": "byte_lossless"}),
2694        ));
2695        assert!(strict["error"]["message"]
2696            .as_str()
2697            .unwrap()
2698            .contains("cannot reconstruct lossless Claude continuation"));
2699
2700        // Transfer/continuation surfaces have no view mode at all.
2701        let translated = service.handle(request(
2702            3,
2703            "harness.v1.sessions.translate",
2704            json!({"locator": locator, "target_harness": "codex"}),
2705        ));
2706        assert!(translated["error"]["message"]
2707            .as_str()
2708            .unwrap()
2709            .contains("cannot reconstruct lossless Claude continuation"));
2710        let resumed = service.handle(request(
2711            4,
2712            "harness.v1.sessions.resume_instructions",
2713            json!({"locator": locator}),
2714        ));
2715        assert!(resumed["error"]["message"]
2716            .as_str()
2717            .unwrap()
2718            .contains("cannot reconstruct lossless Claude continuation"));
2719
2720        let _ = std::fs::remove_dir_all(&temp);
2721    }
2722
2723    #[test]
2724    fn diagonal_artifacts_preserve_claude_subagents_and_grok_bundle_members() {
2725        let temp = std::env::temp_dir().join(format!(
2726            "supercode-harness-artifact-{}-{}",
2727            std::process::id(),
2728            generated_session_id()
2729        ));
2730        let main_path = temp.join("parent.jsonl");
2731        let subagent_path = temp.join("parent/subagents/agent-child.jsonl");
2732        std::fs::create_dir_all(subagent_path.parent().unwrap()).unwrap();
2733        let fixture = std::fs::read_to_string(
2734            PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2735                .join("tests/fixtures/claude_code_session.jsonl"),
2736        )
2737        .unwrap();
2738        let parent = fixture.trim_end_matches('\n');
2739        let child = fixture.trim_end_matches('\n');
2740        std::fs::write(&main_path, parent).unwrap();
2741        std::fs::write(&subagent_path, child).unwrap();
2742        let locator = SessionLocator {
2743            harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2744            session_id: "213bb148-51ea-453f-9206-f8b4b1168547".into(),
2745            storage: StorageLocator::File {
2746                path: main_path.clone(),
2747            },
2748        };
2749        let mut service = HarnessSessionService::new();
2750        let claude = service.handle(request(
2751            1,
2752            "harness.v1.sessions.translate",
2753            json!({"locator": locator, "target_harness": "claude-code"}),
2754        ));
2755        let artifact = &claude["result"]["artifact"];
2756        assert_eq!(artifact["fidelity"], "byte_lossless");
2757        assert_eq!(artifact["content"], parent);
2758        let files = artifact["files"].as_array().unwrap();
2759        assert!(files.iter().any(|file| {
2760            file["role"] == "subagent"
2761                && file["path"]
2762                    .as_str()
2763                    .is_some_and(|path| path.ends_with("/subagents/agent-child.jsonl"))
2764                && file["content"] == child
2765        }));
2766        assert!(!artifact["content"].as_str().unwrap().ends_with('\n'));
2767
2768        let grok = service.handle(request(
2769            2,
2770            "harness.v1.sessions.translate",
2771            json!({"locator": grok_locator(), "target_harness": "grok"}),
2772        ));
2773        let files = grok["result"]["artifact"]["files"].as_array().unwrap();
2774        for name in ["summary.json", "updates.jsonl"] {
2775            let expected = std::fs::read_to_string(
2776                PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2777                    .join("tests/fixtures/grok_session")
2778                    .join(name),
2779            )
2780            .unwrap();
2781            assert!(files.iter().any(|file| {
2782                file["path"] == name && file["role"] == "bundle" && file["content"] == expected
2783            }));
2784        }
2785        std::fs::remove_dir_all(temp).ok();
2786    }
2787
2788    #[test]
2789    fn every_non_grok_handoff_mints_and_uses_a_fresh_target_identity() {
2790        let mut service = HarnessSessionService::new();
2791        let source = pi_locator();
2792        for (target, format) in [
2793            ("claude-code", SessionFormat::ClaudeCode),
2794            ("codex", SessionFormat::Codex),
2795            ("opencode", SessionFormat::OpenCode),
2796            ("pi", SessionFormat::Pi),
2797        ] {
2798            let result = service.handle(request(
2799                1,
2800                "harness.v1.sessions.handoff",
2801                json!({"locator": source, "target_harness": target, "cwd": "/tmp/project"}),
2802            ));
2803            let artifact = &result["result"]["artifact"];
2804            let target_id = artifact["session_id"].as_str().unwrap();
2805            assert_ne!(target_id, source.session_id, "{target}");
2806            let parsed = Session::load_str(artifact["content"].as_str().unwrap(), format).unwrap();
2807            assert_eq!(
2808                parsed.meta.session_id.as_deref(),
2809                Some(target_id),
2810                "{target}"
2811            );
2812            if target != "pi" {
2813                assert!(result["result"]["launch"]["arguments"]
2814                    .as_array()
2815                    .unwrap()
2816                    .iter()
2817                    .any(|argument| argument == target_id));
2818            }
2819            if target == "opencode" {
2820                assert!(target_id.starts_with("ses_"));
2821                fn assert_session_ids(value: &Value, target_id: &str) {
2822                    match value {
2823                        Value::Object(fields) => {
2824                            if let Some(session_id) = fields.get("sessionID") {
2825                                assert_eq!(session_id, target_id);
2826                            }
2827                            for child in fields.values() {
2828                                assert_session_ids(child, target_id);
2829                            }
2830                        }
2831                        Value::Array(values) => {
2832                            for child in values {
2833                                assert_session_ids(child, target_id);
2834                            }
2835                        }
2836                        _ => {}
2837                    }
2838                }
2839                let document: Value =
2840                    serde_json::from_str(artifact["content"].as_str().unwrap()).unwrap();
2841                assert_session_ids(&document, target_id);
2842            }
2843        }
2844
2845        let first = service.handle(request(
2846            2,
2847            "harness.v1.sessions.handoff",
2848            json!({"locator": source, "target_harness": "codex"}),
2849        ));
2850        let second = service.handle(request(
2851            3,
2852            "harness.v1.sessions.handoff",
2853            json!({"locator": source, "target_harness": "codex"}),
2854        ));
2855        assert_ne!(
2856            first["result"]["artifact"]["session_id"],
2857            second["result"]["artifact"]["session_id"]
2858        );
2859    }
2860
2861    #[test]
2862    fn grok_handoff_uses_the_official_importer_contract() {
2863        let mut service = HarnessSessionService::new();
2864        let source = opencode_locator();
2865        let response = service.handle(request(
2866            1,
2867            "harness.v1.sessions.handoff",
2868            json!({
2869                "locator": source,
2870                "target_harness": "grok",
2871                "cwd": "/tmp/grok-handoff-project",
2872            }),
2873        ));
2874        let result = &response["result"];
2875
2876        // The target is Grok, but the artifact truthfully names the Claude Code wire
2877        // format accepted by Grok's official importer. Raw Grok chat_history JSONL is
2878        // not a complete stock-resumable bundle.
2879        assert_eq!(result["artifact"]["target_harness"], "claude-code");
2880        assert!(result["artifact"]["suggested_filename"]
2881            .as_str()
2882            .unwrap()
2883            .ends_with(".grok-import.claude-code.jsonl"));
2884        let artifact = Session::load_str(
2885            result["artifact"]["content"].as_str().unwrap(),
2886            SessionFormat::ClaudeCode,
2887        )
2888        .unwrap();
2889        assert_eq!(
2890            artifact.meta.cwd.as_deref(),
2891            Some(Path::new("/tmp/grok-handoff-project"))
2892        );
2893        let target_session_id = artifact.meta.session_id.as_deref().unwrap();
2894        assert_eq!(target_session_id.len(), 36);
2895        assert_eq!(target_session_id.as_bytes()[14], b'4');
2896        assert_ne!(target_session_id, opencode_locator().session_id);
2897        assert_eq!(
2898            result["artifact"]["session_id"],
2899            artifact.meta.session_id.as_deref().unwrap()
2900        );
2901
2902        assert_eq!(
2903            result["materialize"]["arguments"],
2904            json!(["import", "--json", "{artifact_path}"])
2905        );
2906        assert_eq!(
2907            result["launch"]["arguments"],
2908            json!(["--resume", "{imported_session_id}", "--fork-session"])
2909        );
2910        assert!(result["note"]
2911            .as_str()
2912            .unwrap()
2913            .contains("outcome=imported"));
2914        assert!(!result["launch"]["arguments"]
2915            .as_array()
2916            .unwrap()
2917            .iter()
2918            .any(|argument| argument == &opencode_locator().session_id));
2919    }
2920
2921    #[tokio::test]
2922    async fn inventory_rejects_unknown_harnesses_and_runtime_attach_is_honest() {
2923        let mut service = HarnessSessionService::new();
2924        let inventory = service
2925            .handle_async(request(
2926                1,
2927                "harness.v1.harnesses.list",
2928                json!({"harnesses": ["missing"]}),
2929            ))
2930            .await;
2931        assert_eq!(inventory["error"]["code"], -32602);
2932
2933        let attached = service
2934            .handle_async(request(
2935                2,
2936                "harness.v1.runtimes.attach_existing",
2937                json!({"harness": "codex", "runtime_id": "thread-1"}),
2938            ))
2939            .await;
2940        assert_eq!(attached["error"]["code"], -32000);
2941        assert!(attached["error"]["message"]
2942            .as_str()
2943            .unwrap()
2944            .contains("runtimes.resume"));
2945    }
2946
2947    #[test]
2948    fn invalid_params_and_unknown_methods_use_json_rpc_errors() {
2949        let mut service = HarnessSessionService::new();
2950        let invalid = service.handle(request(1, "harness.v1.sessions.load", json!({})));
2951        assert_eq!(invalid["error"]["code"], -32602);
2952        let unknown = service.handle(request(2, "harness.v1.unknown", json!({})));
2953        assert_eq!(unknown["error"]["code"], -32601);
2954    }
2955
2956    #[cfg(unix)]
2957    #[tokio::test]
2958    // The test mutates process-wide harness environment and deliberately
2959    // holds the global test lock until every async runtime operation ends.
2960    #[allow(clippy::await_holding_lock)]
2961    async fn async_service_drives_a_generic_acp_runtime() {
2962        let _environment_guard = crate::live_runtime::test_environment_lock();
2963        let script = r#"
2964            i=0
2965            while IFS= read -r line; do
2966              i=$((i + 1))
2967              case "$i" in
2968                1) printf '%s\n' '{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentCapabilities":{},"authMethods":[]}}' ;;
2969                2) printf '%s\n' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"svc_acp"}}' ;;
2970                3)
2971                  printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"ok"}}}}'
2972                  printf '%s\n' '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
2973                  ;;
2974                4)
2975                  printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"from terminal"}}}}'
2976                  printf '%s\n' '{"jsonrpc":"2.0","id":4,"result":{"stopReason":"end_turn"}}'
2977                  ;;
2978              esac
2979            done
2980        "#;
2981        let mut service = HarnessSessionService::new();
2982        let started = service
2983            .handle_async(request(
2984                1,
2985                "harness.v1.runtimes.start",
2986                json!({
2987                    "harness": "codex",
2988                    "protocol": "acp",
2989                    "cwd": std::env::current_dir().unwrap(),
2990                    "launch": {"program": "/bin/sh", "arguments": ["-c", script], "env": {}},
2991                }),
2992            ))
2993            .await;
2994        assert_eq!(started["result"]["connection"], "runtime-1");
2995        assert_eq!(started["result"]["handle"]["runtime_id"], "svc_acp");
2996
2997        let terminal = service
2998            .handle_async(request(
2999                9,
3000                "harness.v1.runtimes.terminal_instructions",
3001                json!({"connection":"runtime-1"}),
3002            ))
3003            .await;
3004        let arguments = terminal["result"]["launch"]["arguments"]
3005            .as_array()
3006            .expect("hosted runtime should return terminal arguments");
3007        let endpoint_index = arguments
3008            .iter()
3009            .position(|value| value == "--endpoint")
3010            .expect("terminal command should use an opaque endpoint");
3011        let endpoint = LiveRuntimeEndpoint::parse(
3012            arguments[endpoint_index + 1]
3013                .as_str()
3014                .expect("endpoint argument should be text"),
3015        )
3016        .unwrap();
3017        assert!(!terminal.to_string().contains("Bearer"));
3018        let workspace = std::env::current_dir().unwrap();
3019        let receipt = resolve_live_runtime(
3020            &endpoint,
3021            &LiveRuntimeSource {
3022                harness: "codex".into(),
3023                session_id: "svc_acp".into(),
3024                workspace,
3025            },
3026        )
3027        .unwrap();
3028        let remote = crate::HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
3029            .await
3030            .unwrap();
3031        let mut attachment = crate::FrontendRuntime::attach(remote.as_ref(), 100)
3032            .await
3033            .unwrap();
3034
3035        let sent = service
3036            .handle_async(request(
3037                2,
3038                "harness.v1.runtimes.send_input",
3039                json!({"connection": "runtime-1", "text": "hi"}),
3040            ))
3041            .await;
3042        assert_eq!(sent["result"]["turn_id"], "3");
3043
3044        let mut events = Vec::new();
3045        for _ in 0..20 {
3046            events.extend(service.poll_runtimes().await);
3047            if events.len() >= 2 {
3048                break;
3049            }
3050            tokio::time::sleep(Duration::from_millis(2)).await;
3051        }
3052        assert!(events
3053            .iter()
3054            .any(|event| { event["params"]["event"]["kind"] == "session/update" }));
3055        assert!(events.iter().any(|event| {
3056            event["params"]["event"]["kind"] == "supercode/acp_request_completed"
3057        }));
3058
3059        let saw_editor_reply = tokio::time::timeout(Duration::from_secs(2), async {
3060            loop {
3061                let event = attachment.next_event().await.unwrap();
3062                if event.kind == "text_delta" && event.payload["text"] == "ok" {
3063                    break;
3064                }
3065            }
3066        })
3067        .await;
3068        assert!(
3069            saw_editor_reply.is_ok(),
3070            "terminal should observe the editor-driven turn"
3071        );
3072
3073        crate::FrontendRuntime::submit(remote.as_ref(), "DRIVE FROM TERMINAL".into())
3074            .await
3075            .unwrap();
3076        let saw_terminal_reply = tokio::time::timeout(Duration::from_secs(2), async {
3077            loop {
3078                let event = attachment.next_event().await.unwrap();
3079                if event.kind == "text_delta" && event.payload["text"] == "from terminal" {
3080                    break;
3081                }
3082            }
3083        })
3084        .await;
3085        assert!(
3086            saw_terminal_reply.is_ok(),
3087            "terminal should drive the same runtime"
3088        );
3089
3090        let closed = service
3091            .handle_async(request(
3092                3,
3093                "harness.v1.runtimes.close",
3094                json!({"connection": "runtime-1"}),
3095            ))
3096            .await;
3097        assert_eq!(closed["result"]["closed"], true);
3098    }
3099}