Skip to main content

supercode_harness/
harness_service.rs

1//! Versioned, language-neutral service over persisted harness sessions.
2//!
3//! The service is transport-agnostic: [`HarnessSessionService::handle`] accepts
4//! one JSON-RPC value and [`HarnessSessionService::poll`] produces subscription
5//! notifications. The CLI exposes those primitives as NDJSON over stdio.
6
7use std::collections::{BTreeMap, BTreeSet};
8use std::path::{Path, PathBuf};
9use std::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_session_page, load_session, load_session_with_fidelity, SdkCapabilities, SdkError,
19    SdkErrorCode, SdkEvent, SdkOperation, SdkRequest, SdkRuntimeEvent, SdkService,
20};
21use crate::watch::{bound_session_view, message_json, 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, Role,
29    RuntimeAttachRequest, RuntimeBackend, RuntimeConnection, RuntimeInput, RuntimeLaunch,
30    RuntimeStartRequest, Session, SessionFollower, SessionFormat, SessionLocator, SessionSource,
31};
32use crate::{reduce, tokens};
33#[cfg(feature = "adapter-api")]
34use crate::{register_live_runtime, resolve_live_runtime, LiveRuntimeRegistration};
35
36/// Protocol namespace implemented by this service.
37pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
38/// Notification method emitted for followed-session changes.
39pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
40/// Notification method emitted for live runtime events.
41pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
42
43/// Stateful persisted-session service. Each instance owns its follow
44/// subscriptions; discovery and loading remain read-only.
45pub struct HarnessSessionService {
46    catalog: HarnessCatalog,
47    followers: BTreeMap<String, SessionFollower>,
48    followed_sources: BTreeMap<String, FollowedSource>,
49    next_subscription: u64,
50    runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
51    terminal_launches: BTreeMap<String, StructuredLaunch>,
52    runtime_sequences: BTreeMap<String, u64>,
53    next_runtime: u64,
54    reduction_store_root: Option<PathBuf>,
55}
56
57impl Default for HarnessSessionService {
58    fn default() -> Self {
59        Self::new()
60    }
61}
62
63impl HarnessSessionService {
64    /// Create an empty service instance.
65    pub fn new() -> Self {
66        Self {
67            catalog: HarnessCatalog::new(),
68            followers: BTreeMap::new(),
69            followed_sources: BTreeMap::new(),
70            next_subscription: 1,
71            runtimes: BTreeMap::new(),
72            terminal_launches: BTreeMap::new(),
73            runtime_sequences: BTreeMap::new(),
74            next_runtime: 1,
75            reduction_store_root: None,
76        }
77    }
78
79    /// Override the trusted, service-owned store used for durable reduction
80    /// bundles. Embedders and tests use this to keep all writes inside an
81    /// explicitly selected root; the CLI otherwise uses the normal
82    /// `$SUPERCODE_HOME/sessions` location.
83    pub fn with_reduction_store_root(mut self, root: impl Into<PathBuf>) -> Self {
84        self.reduction_store_root = Some(root.into());
85        self
86    }
87
88    /// Handle one JSON-RPC 2.0 request and return one JSON-RPC response.
89    #[cfg(feature = "adapter-api")]
90    pub fn handle(&mut self, request: Value) -> Value {
91        let id = request.get("id").cloned().unwrap_or(Value::Null);
92        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
93            return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
94        }
95        let Some(method) = request.get("method").and_then(Value::as_str) else {
96            return rpc_error(id, -32600, "request is missing `method`");
97        };
98        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
99        match self.call(method, params) {
100            Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
101            Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
102            Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
103            Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
104            Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
105            Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
106        }
107    }
108
109    /// Handle either a persisted-session request or an asynchronous live
110    /// runtime request.
111    #[cfg(feature = "adapter-api")]
112    pub async fn handle_async(&mut self, request: Value) -> Value {
113        let method = request
114            .get("method")
115            .and_then(Value::as_str)
116            .unwrap_or_default();
117        if matches!(
118            method,
119            "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
120        ) {
121            let id = request.get("id").cloned().unwrap_or(Value::Null);
122            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
123                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
124            }
125            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
126            return match self.inventory_call(method, params).await {
127                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
128                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
129                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
130                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
131                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
132                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
133            };
134        }
135        if method == "harness.v1.sessions.message" {
136            let id = request.get("id").cloned().unwrap_or(Value::Null);
137            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
138                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
139            }
140            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
141            return match self.message_call(params).await {
142                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
143                Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
144                Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
145                Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
146                Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
147                Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
148            };
149        }
150        if let Some(operation) = SdkOperation::from_method(method) {
151            let id = request.get("id").cloned().unwrap_or(Value::Null);
152            if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
153                return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
154            }
155            let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
156            return match self.execute(SdkRequest { operation, params }).await {
157                Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
158                Err(error) => sdk_rpc_error(id, &error),
159            };
160        }
161        if !method.starts_with("harness.v1.runtimes.") {
162            return self.handle(request);
163        }
164        let id = request.get("id").cloned().unwrap_or(Value::Null);
165        if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
166            return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
167        }
168        let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
169        match self.runtime_call(method, params).await {
170            Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
171            Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
172            Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
173            Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
174            Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
175            Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
176        }
177    }
178
179    /// Poll all active subscriptions once and return zero or more JSON-RPC
180    /// notifications. Recoverable follower errors are delivered as events.
181    #[cfg(feature = "adapter-api")]
182    pub fn poll(&mut self) -> Vec<Value> {
183        let mut notifications = Vec::new();
184        for (subscription, follower) in &mut self.followers {
185            match follower.poll() {
186                Ok(Some(event)) => notifications.push(json!({
187                    "jsonrpc": "2.0",
188                    "method": SESSION_EVENT_METHOD,
189                    "params": {
190                        "subscription": subscription,
191                        "event": event.to_json(),
192                    }
193                })),
194                Ok(None) => {}
195                Err(error) => notifications.push(json!({
196                    "jsonrpc": "2.0",
197                    "method": SESSION_EVENT_METHOD,
198                    "params": {
199                        "subscription": subscription,
200                        "event": {
201                            "type": "watch_error",
202                            "recoverable": true,
203                            "message": error.to_string(),
204                        },
205                    }
206                })),
207            }
208        }
209        notifications
210    }
211
212    /// Report each followed session's live-runtime lifecycle state on that
213    /// session's own subscription, emitting only when the state changes.
214    ///
215    /// A growing transcript is not evidence that an agent is working, so the
216    /// state comes from the live-runtime registry and nowhere else. A followed
217    /// session with no registered Supercode runtime — a harness running outside
218    /// Supercode — reports `persisted`, which says plainly that its activity is
219    /// unknown rather than guessing at it. These events carry no sequence
220    /// number and no transcript content; they never interleave with the
221    /// content follower's sequenced stream.
222    #[cfg(feature = "adapter-api")]
223    pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
224        let registry = crate::LocalRuntimeRegistry::new();
225        let authorization = crate::RuntimeAuthorization::observer();
226        let mut notifications = Vec::new();
227        for (subscription, source) in &mut self.followed_sources {
228            let state = match registry
229                .source_state(&source.harness, &source.session_id, &authorization)
230                .await
231            {
232                Ok(Some(state)) => state,
233                Ok(None) => crate::RuntimeRegistryState::Persisted,
234                // A failed registry read is not evidence of a state change.
235                Err(_) => continue,
236            };
237            if source.reported.as_deref() == Some(state.as_str()) {
238                continue;
239            }
240            source.reported = Some(state.as_str().to_string());
241            notifications.push(json!({
242                "jsonrpc": "2.0",
243                "method": SESSION_EVENT_METHOD,
244                "params": {
245                    "subscription": subscription,
246                    "event": {"type": "runtime_state", "state": state.as_str()},
247                },
248            }));
249        }
250        notifications
251    }
252
253    /// Non-blockingly sample one event from every connected live runtime.
254    #[cfg(feature = "adapter-api")]
255    pub async fn poll_runtimes(&mut self) -> Vec<Value> {
256        self.poll_sdk_events()
257            .await
258            .into_iter()
259            .map(|(connection, runtime_event)| {
260                json!({
261                    "jsonrpc": "2.0",
262                    "method": RUNTIME_EVENT_METHOD,
263                    "params": {
264                        "connection": connection,
265                        "session_id": runtime_event.session_id,
266                        "sequence": runtime_event.event.sequence,
267                        "event": {
268                            "kind": runtime_event.event.kind,
269                            "payload": runtime_event.event.payload,
270                        },
271                    },
272                })
273            })
274            .collect()
275    }
276
277    async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
278        let mut events = Vec::new();
279        let mut closed = Vec::new();
280        for (connection, runtime) in &mut self.runtimes {
281            let session_id = runtime.handle().runtime_id.clone();
282            match tokio::time::timeout(Duration::from_millis(1), runtime.next_event()).await {
283                Ok(Ok(Some(event))) => {
284                    let terminal = event.kind == "transport_closed";
285                    let next_sequence = self
286                        .runtime_sequences
287                        .entry(session_id.clone())
288                        .or_insert(0);
289                    let sequence = event.sequence.unwrap_or_else(|| {
290                        *next_sequence = next_sequence.saturating_add(1);
291                        *next_sequence
292                    });
293                    *next_sequence = (*next_sequence).max(sequence);
294                    events.push((
295                        connection.clone(),
296                        SdkRuntimeEvent {
297                            session_id: session_id.clone(),
298                            event: SdkEvent {
299                                sequence,
300                                kind: event.kind,
301                                payload: event.payload,
302                            },
303                        },
304                    ));
305                    if terminal {
306                        closed.push(connection.clone());
307                    }
308                }
309                Ok(Ok(None)) => {
310                    let sequence = self
311                        .runtime_sequences
312                        .entry(session_id.clone())
313                        .or_insert(0);
314                    *sequence = sequence.saturating_add(1);
315                    events.push((
316                        connection.clone(),
317                        SdkRuntimeEvent {
318                            session_id,
319                            event: SdkEvent {
320                                sequence: *sequence,
321                                kind: "transport_closed".into(),
322                                payload: json!({"message": "Harness runtime transport closed."}),
323                            },
324                        },
325                    ));
326                    closed.push(connection.clone());
327                }
328                Err(_) => {}
329                Ok(Err(error)) => {
330                    let sequence = self
331                        .runtime_sequences
332                        .entry(session_id.clone())
333                        .or_insert(0);
334                    *sequence = sequence.saturating_add(1);
335                    events.push((
336                        connection.clone(),
337                        SdkRuntimeEvent {
338                            session_id,
339                            event: SdkEvent {
340                                sequence: *sequence,
341                                kind: "transport_error".into(),
342                                payload: json!({"message": error.to_string(), "terminal": true}),
343                            },
344                        },
345                    ));
346                    closed.push(connection.clone());
347                }
348            }
349        }
350        for connection in closed {
351            if let Some(runtime) = self.runtimes.remove(&connection) {
352                self.runtime_sequences.remove(&runtime.handle().runtime_id);
353            }
354            self.terminal_launches.remove(&connection);
355        }
356        events
357    }
358
359    fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
360        match method {
361            "harness.v1.capabilities" => Ok(json!({
362                "version": HARNESS_SERVICE_VERSION,
363                "sdk": self.capabilities(),
364                "methods": [
365                    "harness.v1.support.report",
366                    "harness.v1.harnesses.list",
367                    "harness.v1.harnesses.probe",
368                    "harness.v1.sessions.discover",
369                    "harness.v1.sessions.load",
370                    "harness.v1.sessions.follow",
371                    "harness.v1.sessions.unfollow",
372                    "harness.v1.sessions.message",
373                    "harness.v1.sessions.import",
374                    "harness.v1.sessions.export",
375                    "harness.v1.sessions.translate",
376                    "harness.v1.sessions.reduce",
377                    "harness.v1.sessions.branch",
378                    "harness.v1.sessions.handoff",
379                    "harness.v1.sessions.resume_instructions",
380                    "harness.v1.runtimes.capabilities",
381                    "harness.v1.runtimes.start",
382                    "harness.v1.runtimes.resume",
383                    "harness.v1.runtimes.attach_existing",
384                    "harness.v1.runtimes.attach",
385                    "harness.v1.runtimes.send_input",
386                    "harness.v1.runtimes.interrupt",
387                    "harness.v1.runtimes.steer",
388                    "harness.v1.runtimes.respond",
389                    "harness.v1.runtimes.terminal_instructions",
390                    "harness.v1.runtimes.close",
391                ],
392                "notifications": [SESSION_EVENT_METHOD, RUNTIME_EVENT_METHOD],
393                "harnesses": harness_support_registry()
394                    .harnesses
395                    .into_iter()
396                    .map(|harness| harness.id)
397                    .collect::<Vec<_>>(),
398            })),
399            "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
400                .map_err(|error| ServiceError::Operation(error.to_string())),
401            "harness.v1.sessions.discover" => {
402                let query = decode::<DiscoveryQuery>(params)?;
403                let page = discover_session_page(&query).map_err(operation)?;
404                // Claude Code is the one harness that publishes its RUNNING
405                // sessions. The registry is read once per discovery and joined
406                // by session id; every record in it has already survived a
407                // `kill(pid, 0)` liveness check inside `read_registry`.
408                let peers = if page
409                    .sessions
410                    .iter()
411                    .any(|session| session.locator.harness.as_str() == HarnessId::CLAUDE_CODE)
412                {
413                    crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(
414                        &query.homes,
415                    ))
416                } else {
417                    Vec::new()
418                };
419                let sessions = page
420                    .sessions
421                    .into_iter()
422                    .map(|session| {
423                        let mut value = serde_json::to_value(&session)
424                            .map_err(|error| ServiceError::Operation(error.to_string()))?;
425                        if let Some(workspace) = &session.cwd {
426                            let source = LiveRuntimeSource {
427                                harness: session.locator.harness.as_str().to_string(),
428                                session_id: session.locator.session_id.clone(),
429                                workspace: workspace.clone(),
430                            };
431                            if let Some(endpoint) = discover_live_runtime(&source)
432                                .map_err(|error| ServiceError::Operation(error.to_string()))?
433                            {
434                                value["live_endpoint"] = json!(endpoint.as_str());
435                            }
436                        }
437                        if let Some(peer) = peers.iter().find(|peer| {
438                            session.locator.harness.as_str() == HarnessId::CLAUDE_CODE
439                                && peer.session_id == session.locator.session_id
440                        }) {
441                            // A Supercode-hosted runtime is the richer
442                            // attachment, so it keeps the endpoint slot; the
443                            // peer endpoint fills it only when nothing else did.
444                            if value.get("live_endpoint").is_none() {
445                                value["live_endpoint"] = json!(peer.endpoint().as_str());
446                            }
447                            if let Some(status) = peer.status {
448                                value["live_status"] = json!(status.as_str());
449                            }
450                        }
451                        Ok(value)
452                    })
453                    .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
454                Ok(json!({"sessions": sessions, "next_cursor": page.next_cursor}))
455            }
456            "harness.v1.sessions.load" => {
457                let params = decode::<LoadSessionParams>(params)?;
458                if let Some(options) = &params.options {
459                    options.validate()?;
460                    return load_session(&params.read.locator)
461                        .map(|session| projected_session_result(&session, options))
462                        .map_err(operation);
463                }
464                let mut session = if params.read.display_history() {
465                    self.catalog
466                        .load_display_view(
467                            &params.read.locator,
468                            params.read.read_fidelity(),
469                            params.read.tail_messages().unwrap_or(500),
470                        )
471                        .map_err(crate::Error::from)
472                } else if params.read.include_subagents() {
473                    load_session_with_fidelity(&params.read.locator, params.read.read_fidelity())
474                } else {
475                    self.catalog
476                        .load_parent_with_fidelity(
477                            &params.read.locator,
478                            params.read.read_fidelity(),
479                        )
480                        .map_err(crate::Error::from)
481                }
482                .map_err(operation)?;
483                params.read.bound_session(&mut session);
484                Ok(json!({"session": normalized_session_json(&session)}))
485            }
486            "harness.v1.sessions.follow" => {
487                let params = decode::<LocatorParams>(params)?;
488                let mut follower = self
489                    .catalog
490                    .follow_read_view(
491                        &params.locator,
492                        params.read_fidelity(),
493                        params.include_subagents(),
494                        params.tail_messages(),
495                        params.max_message_chars(),
496                        params.display_history(),
497                    )
498                    .map_err(operation)?;
499                let initial = follower
500                    .poll()
501                    .map_err(operation)?
502                    .map(|event| event.to_json());
503                let subscription = format!("sub-{}", self.next_subscription);
504                self.next_subscription += 1;
505                self.followers.insert(subscription.clone(), follower);
506                self.followed_sources.insert(
507                    subscription.clone(),
508                    FollowedSource {
509                        harness: params.locator.harness.as_str().to_string(),
510                        session_id: params.locator.session_id.clone(),
511                        reported: None,
512                    },
513                );
514                Ok(json!({"subscription": subscription, "initial": initial}))
515            }
516            "harness.v1.sessions.unfollow" => {
517                let params = decode::<UnfollowParams>(params)?;
518                self.followed_sources.remove(&params.subscription);
519                Ok(json!({
520                    "removed": self.followers.remove(&params.subscription).is_some()
521                }))
522            }
523            "harness.v1.sessions.import" => {
524                let params = decode::<ImportSessionParams>(params)?;
525                let session = Session::load_str(&params.content, params.source_harness.into())
526                    .map_err(operation)?;
527                Ok(json!({"session": normalized_session_json(&session)}))
528            }
529            "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
530                let params = decode::<ExportSessionParams>(params)?;
531                let session = load_session(&params.locator).map_err(operation)?;
532                let artifact = session_artifact(&params.locator, &session, params.target_harness)?;
533                Ok(json!({"artifact": artifact}))
534            }
535            "harness.v1.sessions.reduce" => {
536                let params = decode::<ReduceSessionParams>(params)?;
537                self.reduce_session(params)
538            }
539            "harness.v1.sessions.branch" => {
540                let params = decode::<BranchSessionParams>(params)?;
541                let session = load_session(&params.locator).map_err(operation)?;
542                let storage = params.locator.storage.path().display().to_string();
543                let bootstrap_prompt = format!(
544                    "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.",
545                    params.locator.harness.as_str(), params.locator.session_id, storage
546                );
547                let artifact = params
548                    .target_harness
549                    .map(|target| session_artifact(&params.locator, &session, target))
550                    .transpose()?;
551                Ok(json!({
552                    "parent": params.locator,
553                    "session": normalized_session_json(&session),
554                    "bootstrap_prompt": bootstrap_prompt,
555                    "artifact": artifact,
556                }))
557            }
558            "harness.v1.sessions.handoff" => {
559                let params = decode::<HandoffSessionParams>(params)?;
560                let session = load_session(&params.locator).map_err(operation)?;
561                let cwd = params
562                    .cwd
563                    .or_else(|| session.meta.cwd.clone())
564                    .unwrap_or_else(|| PathBuf::from("."));
565                let artifact =
566                    handoff_artifact(&params.locator, &session, params.target_harness, &cwd)?;
567                let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
568                    ServiceError::Operation(
569                        "handoff artifact omitted target session identity".into(),
570                    )
571                })?;
572                let instructions =
573                    handoff_instructions(params.target_harness, target_session_id, &cwd);
574                Ok(json!({
575                    "artifact": artifact,
576                    "launch": instructions.launch,
577                    "materialize": instructions.materialize,
578                    "requires_materialization": instructions.requires_materialization,
579                    "note": instructions.note,
580                }))
581            }
582            "harness.v1.sessions.resume_instructions" => {
583                let params = decode::<ResumeInstructionsParams>(params)?;
584                let session = load_session(&params.locator).map_err(operation)?;
585                let cwd = params
586                    .cwd
587                    .or(session.meta.cwd)
588                    .unwrap_or_else(|| PathBuf::from("."));
589                let launch = resume_launch(
590                    params.locator.harness.as_str(),
591                    &params.locator.session_id,
592                    &cwd,
593                    params.policy,
594                )?;
595                Ok(json!({"launch": launch}))
596            }
597            _ => Err(ServiceError::MethodNotFound),
598        }
599    }
600
601    fn reduce_session(
602        &self,
603        params: ReduceSessionParams,
604    ) -> std::result::Result<Value, ServiceError> {
605        let session = load_session(&params.locator).map_err(operation)?;
606        if session.messages.is_empty() {
607            return Err(ServiceError::InvalidParams(
608                "cannot reduce an empty session".into(),
609            ));
610        }
611        let keep_last = params.keep_last.clamp(1, 128);
612        let policy = reduce::ReductionPolicy {
613            clear_turns_older_than: Some(keep_last),
614            ..Default::default()
615        };
616        let (view, log) =
617            reduce::project_messages(&session.messages, &policy, &reduce::ReductionLog::default());
618        if log.reductions.is_empty() {
619            return Err(ServiceError::UnsupportedAction(format!(
620                "session `{}` is already too small for a meaningful reversible reduction",
621                params.locator.session_id
622            )));
623        }
624        let source_tokens = tokens::estimate_view_tokens(&session.messages);
625        let reduced_tokens = tokens::estimate_view_tokens(&view);
626        if reduced_tokens >= source_tokens {
627            return Err(ServiceError::UnsupportedAction(format!(
628                "session `{}` has no token-reducing reversible projection",
629                params.locator.session_id
630            )));
631        }
632
633        let store_root = self
634            .reduction_store_root
635            .clone()
636            .unwrap_or_else(default_reduction_store_root);
637        let store = crate::SessionStore::open(&store_root).map_err(operation)?;
638        let rescue_id = format!("rescue-{}", generated_session_id());
639        let imported = session
640            .imported_message_count
641            .unwrap_or(session.messages.len())
642            .min(session.messages.len());
643        let sidecar_jsonl = session.to_native_jsonl_v2(&session.messages[imported..]);
644        let view_jsonl = messages_jsonl(&view)?;
645        let title = format!(
646            "Reduced {} continuation from {}",
647            params.target_harness.id(),
648            params.locator.session_id
649        );
650
651        // Durability order is intentional: the full source of truth lands
652        // before either object that can refer to it. A crash may leave an
653        // unused sidecar, but can never leave a reduced view whose originals
654        // were not durably written first.
655        store
656            .save_sidecar(&rescue_id, &sidecar_jsonl)
657            .map_err(operation)?;
658        store
659            .save_reduction_log(&rescue_id, &log)
660            .map_err(operation)?;
661        store
662            .save(&rescue_id, &title, &view_jsonl)
663            .map_err(operation)?;
664
665        let source_bytes = serde_json::to_vec(&session.messages)
666            .map_err(|error| ServiceError::Operation(error.to_string()))?
667            .len() as u64;
668        let reduced_bytes = serde_json::to_vec(&view)
669            .map_err(|error| ServiceError::Operation(error.to_string()))?
670            .len() as u64;
671        store
672            .set_reduction_stats(
673                &rescue_id,
674                &title,
675                source_bytes,
676                reduced_bytes,
677                log.reductions.len() as u32,
678            )
679            .map_err(operation)?;
680
681        // The receipt is issued only after a real disk reload. This proves
682        // the exact files another process will consume, not the convenient
683        // in-memory values that produced them.
684        let reloaded_sidecar = store
685            .load_sidecar(&rescue_id)
686            .map_err(operation)?
687            .ok_or_else(|| ServiceError::Operation("reduction sidecar disappeared".into()))?;
688        let reloaded_sidecar = Session::from_sidecar_str(&reloaded_sidecar).map_err(operation)?;
689        let reloaded_log = store
690            .load_reduction_log(&rescue_id)
691            .map_err(operation)?
692            .ok_or_else(|| ServiceError::Operation("reduction log disappeared".into()))?;
693        let reloaded_view = parse_messages_jsonl(&store.load(&rescue_id).map_err(operation)?)?;
694        reduce::verify_log(&reloaded_log, &reloaded_sidecar).map_err(operation)?;
695        // `sc.reduction` is deliberately in-memory-only metadata: it must
696        // never leak onto a provider-facing transcript. Reapplying the
697        // durable log to the durable sidecar restores those ids. Comparing
698        // its wire form with the transcript reloaded above proves that the
699        // persisted view is exactly the deterministic projection before we
700        // use the restamped form for inversion.
701        let (restamped_view, restamped_log) =
702            reduce::project_messages(&reloaded_sidecar.messages, &policy, &reloaded_log);
703        if messages_jsonl(&restamped_view)? != messages_jsonl(&reloaded_view)? {
704            return Err(ServiceError::Operation(
705                "persisted reduction view does not match its durable log and sidecar".into(),
706            ));
707        }
708        if restamped_log != reloaded_log {
709            return Err(ServiceError::Operation(
710                "reapplying the durable reduction log changed its identity".into(),
711            ));
712        }
713        let inverted =
714            reduce::invert(&restamped_view, &reloaded_log, &reloaded_sidecar).map_err(operation)?;
715        if inverted != session.messages {
716            return Err(ServiceError::Operation(
717                "reduction inversion did not restore the source messages byte-exactly".into(),
718            ));
719        }
720
721        let ratio = source_tokens as f64 / reduced_tokens.max(1) as f64;
722        let sidecar_path = store.sidecar_path(&rescue_id);
723        let reduction_log_path = store.reduction_log_path(&rescue_id).map_err(operation)?;
724        let bootstrap_prompt = reduced_bootstrap_prompt(
725            &params.locator,
726            params.target_harness,
727            &view_jsonl,
728            &sidecar_path,
729            &reduction_log_path,
730        );
731        let mut reduced_session = session.clone();
732        reduced_session.meta.session_id = Some(rescue_id.clone());
733        reduced_session.messages = view;
734
735        Ok(json!({
736            "session": normalized_session_json(&reduced_session),
737            "bootstrap_prompt": bootstrap_prompt,
738            "receipt": {
739                "id": rescue_id,
740                "sidecar_id": rescue_id,
741                "source_harness": params.locator.harness,
742                "target_harness": params.target_harness.id(),
743                "source_tokens": source_tokens,
744                "reduced_tokens": reduced_tokens,
745                "ratio": ratio,
746                "source_bytes": source_bytes,
747                "reduced_bytes": reduced_bytes,
748                "reductions": reloaded_log.reductions.len(),
749                "sidecar_path": sidecar_path,
750                "reduction_log_path": reduction_log_path,
751                "verified": true,
752                "reversible": true,
753            }
754        }))
755    }
756
757    async fn runtime_call(
758        &mut self,
759        method: &str,
760        params: Value,
761    ) -> std::result::Result<Value, ServiceError> {
762        match method {
763            "harness.v1.runtimes.capabilities" => {
764                let params = decode::<RuntimeBackendParams>(params)?;
765                let backend = runtime_backend(&params)?;
766                Ok(json!({
767                    "harness": backend.harness(),
768                    "capabilities": backend.capabilities(),
769                }))
770            }
771            "harness.v1.runtimes.start" => {
772                let params = decode::<RuntimeStartParams>(params)?;
773                let backend = runtime_backend(&params.backend)?;
774                let capabilities = backend.capabilities();
775                let workspace = params.cwd.clone();
776                let runtime = backend
777                    .start(RuntimeStartRequest {
778                        cwd: params.cwd,
779                        launch: runtime_launch(&params.backend),
780                    })
781                    .await
782                    .map_err(operation)?;
783                self.insert_hosted_runtime(runtime, capabilities, workspace)
784                    .await
785            }
786            "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
787                let params = decode::<RuntimeAttachParams>(params)?;
788                let backend = runtime_backend(&params.backend)?;
789                let capabilities = backend.capabilities();
790                let workspace = params.cwd.clone().unwrap_or_else(|| {
791                    std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))
792                });
793                let runtime = backend
794                    .attach(RuntimeAttachRequest {
795                        runtime_id: params.runtime_id,
796                        cwd: params.cwd,
797                        launch: runtime_launch(&params.backend),
798                    })
799                    .await
800                    .map_err(operation)?;
801                self.insert_hosted_runtime(runtime, capabilities, workspace)
802                    .await
803            }
804            "harness.v1.runtimes.attach_existing" => {
805                let params = decode::<RuntimeAttachParams>(params)?;
806                let backend: Box<dyn RuntimeBackend> = match params
807                    .backend
808                    .base_url
809                    .as_deref()
810                    .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
811                {
812                    Some(endpoint) => {
813                        #[cfg(not(feature = "adapter-api"))]
814                        {
815                            let _ = endpoint;
816                            return Err(ServiceError::UnsupportedAction(
817                                "live HTTP attachment adapter is not compiled".into(),
818                            ));
819                        }
820                        #[cfg(feature = "adapter-api")]
821                        {
822                            let workspace = params.cwd.clone().ok_or_else(|| {
823                                ServiceError::InvalidParams(
824                                    "Supercode live attach requires the project cwd".into(),
825                                )
826                            })?;
827                            let source = LiveRuntimeSource {
828                                harness: params.backend.harness.as_str().to_string(),
829                                session_id: params.runtime_id.clone(),
830                                workspace,
831                            };
832                            let receipt = resolve_live_runtime(&endpoint, &source)
833                                .map_err(|error| ServiceError::Operation(error.to_string()))?;
834                            Box::new(SupercodeHttpRuntimeBackend::new(receipt))
835                        }
836                    }
837                    None => runtime_backend(&params.backend)?,
838                };
839                if !backend.capabilities().attach_existing_process {
840                    return Err(ServiceError::Operation(format!(
841                        "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
842                        backend.harness().as_str()
843                    )));
844                }
845                let runtime = backend
846                    .attach_existing(RuntimeAttachRequest {
847                        runtime_id: params.runtime_id,
848                        cwd: params.cwd,
849                        launch: runtime_launch(&params.backend),
850                    })
851                    .await
852                    .map_err(operation)?;
853                self.insert_runtime(runtime)
854            }
855            "harness.v1.runtimes.send_input" => {
856                let params = decode::<RuntimeInputParams>(params)?;
857                let runtime = self.runtime_mut(&params.connection)?;
858                let turn_id = runtime
859                    .send_input(RuntimeInput { text: params.text })
860                    .await
861                    .map_err(operation)?;
862                Ok(json!({"turn_id": turn_id}))
863            }
864            "harness.v1.runtimes.interrupt" => {
865                let params = decode::<RuntimeConnectionParams>(params)?;
866                self.runtime_mut(&params.connection)?
867                    .interrupt()
868                    .await
869                    .map_err(operation)?;
870                Ok(json!({}))
871            }
872            "harness.v1.runtimes.steer" => Err(ServiceError::UnsupportedAction(
873                "steer is not supported by this harness-native runtime adapter".into(),
874            )),
875            "harness.v1.runtimes.respond" => {
876                let params = decode::<RuntimeRespondParams>(params)?;
877                self.runtime_mut(&params.connection)?
878                    .respond(params.request_id, params.response)
879                    .await
880                    .map_err(operation)?;
881                Ok(json!({}))
882            }
883            "harness.v1.runtimes.terminal_instructions" => {
884                let params = decode::<RuntimeConnectionParams>(params)?;
885                let launch = self
886                    .terminal_launches
887                    .get(&params.connection)
888                    .ok_or_else(|| {
889                        ServiceError::Operation(
890                            "this runtime is not hosted for terminal attachment".into(),
891                        )
892                    })?;
893                Ok(json!({"launch":launch}))
894            }
895            "harness.v1.runtimes.close" => {
896                let params = decode::<RuntimeConnectionParams>(params)?;
897                let Some(mut runtime) = self.runtimes.remove(&params.connection) else {
898                    return Err(ServiceError::InvalidParams(format!(
899                        "unknown runtime connection `{}`",
900                        params.connection
901                    )));
902                };
903                self.terminal_launches.remove(&params.connection);
904                self.runtime_sequences.remove(&runtime.handle().runtime_id);
905                runtime.close().await.map_err(operation)?;
906                Ok(json!({"closed": true}))
907            }
908            _ => Err(ServiceError::MethodNotFound),
909        }
910    }
911
912    /// Deliver one message into a session that is running right now.
913    #[cfg(feature = "adapter-api")]
914    async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
915        let params = decode::<MessageSessionParams>(params)?;
916        Ok(message_live_session(&params, &crate::claude_peer::ProcessCourierRunner).await)
917    }
918
919    fn insert_runtime(
920        &mut self,
921        runtime: Box<dyn RuntimeConnection>,
922    ) -> std::result::Result<Value, ServiceError> {
923        let connection = format!("runtime-{}", self.next_runtime);
924        self.next_runtime += 1;
925        let handle = runtime.handle().clone();
926        self.runtime_sequences
927            .entry(handle.runtime_id.clone())
928            .or_insert(0);
929        self.runtimes.insert(connection.clone(), runtime);
930        Ok(json!({"connection": connection, "handle": handle}))
931    }
932
933    #[cfg(feature = "adapter-api")]
934    async fn insert_hosted_runtime(
935        &mut self,
936        runtime: Box<dyn RuntimeConnection>,
937        capabilities: crate::RuntimeCapabilities,
938        workspace: PathBuf,
939    ) -> std::result::Result<Value, ServiceError> {
940        let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities);
941        let token: std::sync::Arc<str> = crate::server::generate_token().into();
942        let server = crate::server::run_frontend_http(
943            host.clone(),
944            host.frontend_sender(),
945            "127.0.0.1:0",
946            token.clone(),
947        )
948        .await
949        .map_err(|error| ServiceError::Operation(error.to_string()))?;
950        let source = LiveRuntimeSource {
951            harness: connection.handle().harness.as_str().to_string(),
952            session_id: connection.handle().runtime_id.clone(),
953            workspace: workspace.clone(),
954        };
955        let registration = register_live_runtime(
956            connection.handle().runtime_id.clone(),
957            source.clone(),
958            format!("http://{}", server.address()),
959            token.to_string(),
960        )
961        .map_err(|error| ServiceError::Operation(error.to_string()))?;
962        let endpoint = registration.endpoint().to_string();
963        let launch = StructuredLaunch {
964            cwd: workspace,
965            // Pin attachment to the executable hosting this runtime. A bare
966            // `supercode` could resolve to an older global install whose CLI
967            // does not understand the receipt it is being asked to open.
968            program: std::env::current_exe()
969                .ok()
970                .map(|path| path.to_string_lossy().into_owned())
971                .unwrap_or_else(|| "supercode".into()),
972            arguments: vec![
973                "harness".into(),
974                "attach".into(),
975                "--endpoint".into(),
976                endpoint,
977                "--harness".into(),
978                source.harness,
979                "--session".into(),
980                source.session_id,
981            ],
982            env: BTreeMap::new(),
983        };
984        let lease = HostedRuntimeLease {
985            connection,
986            _host: host,
987            _registration: registration,
988            _server: server,
989        };
990        let opened = self.insert_runtime(Box::new(lease))?;
991        let connection_id = opened["connection"]
992            .as_str()
993            .expect("insert_runtime returns a connection id")
994            .to_string();
995        self.terminal_launches.insert(connection_id, launch);
996        Ok(opened)
997    }
998
999    #[cfg(not(feature = "adapter-api"))]
1000    async fn insert_hosted_runtime(
1001        &mut self,
1002        runtime: Box<dyn RuntimeConnection>,
1003        _capabilities: crate::RuntimeCapabilities,
1004        _workspace: PathBuf,
1005    ) -> std::result::Result<Value, ServiceError> {
1006        self.insert_runtime(runtime)
1007    }
1008
1009    fn runtime_mut(
1010        &mut self,
1011        connection: &str,
1012    ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
1013        self.runtimes.get_mut(connection).ok_or_else(|| {
1014            ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
1015        })
1016    }
1017
1018    async fn inventory_call(
1019        &self,
1020        method: &str,
1021        params: Value,
1022    ) -> std::result::Result<Value, ServiceError> {
1023        let mut params = decode::<HarnessInventoryParams>(params)?;
1024        if method == "harness.v1.harnesses.probe" {
1025            let harness = params.harness.take().ok_or_else(|| {
1026                ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
1027            })?;
1028            params.harnesses = vec![harness];
1029        }
1030        let selected = params
1031            .harnesses
1032            .iter()
1033            .map(HarnessId::as_str)
1034            .collect::<std::collections::BTreeSet<_>>();
1035        let supported = harness_support_registry()
1036            .harnesses
1037            .into_iter()
1038            .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
1039            .collect::<Vec<_>>();
1040        if !params.harnesses.is_empty() && supported.len() != selected.len() {
1041            let known = supported
1042                .iter()
1043                .map(|harness| harness.id.as_str())
1044                .collect::<std::collections::BTreeSet<_>>();
1045            let missing = params
1046                .harnesses
1047                .iter()
1048                .filter(|id| !known.contains(id.as_str()))
1049                .map(HarnessId::as_str)
1050                .collect::<Vec<_>>();
1051            return Err(ServiceError::InvalidParams(format!(
1052                "unknown harness(es): {}",
1053                missing.join(", ")
1054            )));
1055        }
1056        let global_counts = params
1057            .include_sessions
1058            .then(|| self.session_counts(None, &params.harnesses));
1059        let workspace_counts = params.include_sessions.then(|| {
1060            params
1061                .workspace
1062                .as_deref()
1063                .map(|workspace| self.session_counts(Some(workspace), &params.harnesses))
1064        });
1065        let probes = supported.into_iter().map(|descriptor| {
1066            let global = global_counts
1067                .as_ref()
1068                .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
1069            let workspace = workspace_counts
1070                .as_ref()
1071                .and_then(Option::as_ref)
1072                .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
1073            self.probe_harness(descriptor, &params, global, workspace)
1074        });
1075        let harnesses = futures::future::join_all(probes).await;
1076        serde_json::to_value(HarnessInventoryReport {
1077            probe: params.probe,
1078            workspace: params.workspace,
1079            harnesses,
1080        })
1081        .map_err(|error| ServiceError::Operation(error.to_string()))
1082    }
1083
1084    async fn probe_harness(
1085        &self,
1086        descriptor: crate::HarnessSupportDescriptor,
1087        params: &HarnessInventoryParams,
1088        global: Option<usize>,
1089        workspace: Option<usize>,
1090    ) -> LocalHarness {
1091        let launch = descriptor.runtime.default_launch.as_ref();
1092        let executable = launch.and_then(|launch| find_executable(&launch.program));
1093        let installed = executable.is_some();
1094        let version = if params.skip_versions {
1095            None
1096        } else {
1097            match executable.as_deref() {
1098                Some(path) => executable_version(path).await,
1099                None => None,
1100            }
1101        };
1102        let configured = auth_evidence(descriptor.id.as_str());
1103        let mut auth = if configured {
1104            HarnessAuthState::Configured
1105        } else {
1106            HarnessAuthState::Unknown
1107        };
1108        let mut runtime = if installed {
1109            HarnessRuntimeState::Degraded
1110        } else {
1111            HarnessRuntimeState::Unavailable
1112        };
1113        let mut reason = (!installed).then(|| {
1114            format!(
1115                "{} is supported but `{}` was not found on PATH",
1116                descriptor.display_name,
1117                launch
1118                    .map(|launch| launch.program.as_str())
1119                    .unwrap_or("executable")
1120            )
1121        });
1122        let mut repair = (!installed).then(|| {
1123            format!(
1124                "Install {} and ensure `{}` is on PATH.",
1125                descriptor.display_name,
1126                launch
1127                    .map(|launch| launch.program.as_str())
1128                    .unwrap_or("its executable")
1129            )
1130        });
1131
1132        if installed && params.probe == HarnessProbeLevel::Handshake {
1133            let backend_params = RuntimeBackendParams {
1134                harness: descriptor.id.clone(),
1135                protocol: None,
1136                launch: None,
1137                base_url: None,
1138                policy: RuntimePolicy::Default,
1139            };
1140            match runtime_backend(&backend_params) {
1141                Ok(backend) => {
1142                    let cwd = params
1143                        .workspace
1144                        .clone()
1145                        .or_else(|| std::env::current_dir().ok())
1146                        .unwrap_or_else(|| PathBuf::from("."));
1147                    let isolated = descriptor
1148                        .runtime
1149                        .default_launch
1150                        .clone()
1151                        .and_then(|launch| {
1152                            IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok()
1153                        });
1154                    let Some(isolated) = isolated else {
1155                        reason = Some(
1156                            "No-prompt runtime handshake could not create its isolated harness home."
1157                                .into(),
1158                        );
1159                        repair = Some(
1160                            "Check temporary-directory permissions, then run the handshake probe again."
1161                                .into(),
1162                        );
1163                        return LocalHarness {
1164                            id: descriptor.id,
1165                            display_name: descriptor.display_name,
1166                            supported: true,
1167                            installed,
1168                            executable: executable.map(|path| path.to_string_lossy().into_owned()),
1169                            version,
1170                            auth,
1171                            runtime,
1172                            protocol: descriptor.runtime.protocol,
1173                            capabilities: descriptor.runtime.capabilities.clone(),
1174                            effective_capabilities: descriptor.runtime.capabilities,
1175                            sessions: HarnessSessionCounts { global, workspace },
1176                            reason,
1177                            repair,
1178                        };
1179                    };
1180                    match tokio::time::timeout(
1181                        Duration::from_secs(30),
1182                        backend.start(RuntimeStartRequest {
1183                            cwd,
1184                            launch: Some(isolated.launch.clone()),
1185                        }),
1186                    )
1187                    .await
1188                    {
1189                        Ok(Ok(mut connection)) => {
1190                            match stabilize_handshake(connection.as_mut()).await {
1191                                Ok(()) => {
1192                                    auth = HarnessAuthState::Ready;
1193                                    runtime = HarnessRuntimeState::Ready;
1194                                    reason = Some(
1195                                        "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
1196                                            .into(),
1197                                    );
1198                                    repair = None;
1199                                }
1200                                Err(message) => {
1201                                    auth = if looks_like_auth_error(&message) {
1202                                        HarnessAuthState::Required
1203                                    } else if configured {
1204                                        HarnessAuthState::Configured
1205                                    } else {
1206                                        HarnessAuthState::Unknown
1207                                    };
1208                                    reason = Some(format!(
1209                                        "No-prompt runtime handshake became unhealthy during startup: {message}"
1210                                    ));
1211                                    repair = Some(if auth == HarnessAuthState::Required {
1212                                        format!(
1213                                            "Run `{}` interactively once and complete sign-in, then probe again.",
1214                                            launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1215                                        )
1216                                    } else {
1217                                        "Run the harness directly to inspect its startup failure, then probe again."
1218                                            .into()
1219                                    });
1220                                }
1221                            }
1222                            let _ =
1223                                tokio::time::timeout(Duration::from_secs(3), connection.close())
1224                                    .await;
1225                        }
1226                        Ok(Err(error)) => {
1227                            let message = truncate_text(&error.to_string(), 500);
1228                            auth = if looks_like_auth_error(&message) {
1229                                HarnessAuthState::Required
1230                            } else if configured {
1231                                HarnessAuthState::Configured
1232                            } else {
1233                                HarnessAuthState::Unknown
1234                            };
1235                            reason = Some(format!("No-prompt runtime handshake failed: {message}"));
1236                            repair = Some(if auth == HarnessAuthState::Required {
1237                                format!(
1238                                    "Run `{}` interactively once and complete sign-in, then probe again.",
1239                                    launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1240                                )
1241                            } else {
1242                                "Check the harness installation and run the handshake probe again."
1243                                    .into()
1244                            });
1245                        }
1246                        Err(_) => {
1247                            reason = Some(
1248                                "No-prompt runtime handshake timed out after 30 seconds.".into(),
1249                            );
1250                            repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
1251                        }
1252                    }
1253                    // Keep the isolated home alive through process teardown.
1254                    // Otherwise the compiler may release the last meaningful
1255                    // use after cloning `launch`, and a still-starting CLI can
1256                    // recreate its state directory after Drop removed it.
1257                    // Some Node-based launchers finish a short asynchronous
1258                    // installation-id write just after their parent process
1259                    // is reaped. Remove once immediately, allow that bounded
1260                    // writer to settle, then perform the authoritative pass.
1261                    let _ = isolated.cleanup();
1262                    tokio::time::sleep(Duration::from_millis(250)).await;
1263                    if let Err(error) = isolated.cleanup() {
1264                        auth = if configured {
1265                            HarnessAuthState::Configured
1266                        } else {
1267                            HarnessAuthState::Unknown
1268                        };
1269                        runtime = HarnessRuntimeState::Degraded;
1270                        reason = Some(format!(
1271                            "No-prompt runtime handshake could not remove its isolated harness home: {error}"
1272                        ));
1273                        repair = Some(
1274                            "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
1275                                .into(),
1276                        );
1277                    }
1278                }
1279                Err(error) => {
1280                    reason = Some(error_message(error));
1281                }
1282            }
1283        } else if installed && configured {
1284            reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
1285        } else if installed {
1286            reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
1287            repair =
1288                Some(format!(
1289                "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
1290                launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1291            ));
1292        }
1293
1294        let effective_capabilities = if installed {
1295            descriptor.runtime.capabilities.clone()
1296        } else {
1297            unavailable_capabilities()
1298        };
1299        LocalHarness {
1300            id: descriptor.id,
1301            display_name: descriptor.display_name,
1302            supported: true,
1303            installed,
1304            executable: executable.map(|path| path.to_string_lossy().into_owned()),
1305            version,
1306            auth,
1307            runtime,
1308            protocol: descriptor.runtime.protocol,
1309            capabilities: descriptor.runtime.capabilities,
1310            effective_capabilities,
1311            sessions: HarnessSessionCounts { global, workspace },
1312            reason,
1313            repair,
1314        }
1315    }
1316
1317    fn session_counts(
1318        &self,
1319        workspace: Option<&Path>,
1320        harnesses: &[HarnessId],
1321    ) -> BTreeMap<String, usize> {
1322        let mut counts = BTreeMap::new();
1323        for session in self
1324            .catalog
1325            .discover(&DiscoveryQuery {
1326                workspace: workspace.map(Path::to_path_buf),
1327                harnesses: harnesses.to_vec(),
1328                ..DiscoveryQuery::default()
1329            })
1330            .unwrap_or_default()
1331        {
1332            *counts
1333                .entry(session.locator.harness.as_str().to_string())
1334                .or_insert(0) += 1;
1335        }
1336        counts
1337    }
1338}
1339
1340#[async_trait::async_trait]
1341impl SdkService for HarnessSessionService {
1342    fn capabilities(&self) -> SdkCapabilities {
1343        SdkCapabilities::default()
1344    }
1345
1346    async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
1347        if request.operation == SdkOperation::Events {
1348            let events = self
1349                .poll_sdk_events()
1350                .await
1351                .into_iter()
1352                .map(|(_, event)| event)
1353                .collect::<Vec<_>>();
1354            return serde_json::to_value(events).map_err(|error| {
1355                SdkError::new(
1356                    SdkErrorCode::Execution,
1357                    request.operation,
1358                    error.to_string(),
1359                )
1360            });
1361        }
1362        let method = request
1363            .operation
1364            .method()
1365            .ok_or_else(|| SdkError::unsupported(request.operation))?;
1366        let result = match request.operation {
1367            SdkOperation::Discover | SdkOperation::Load | SdkOperation::Export => {
1368                self.call(method, request.params)
1369            }
1370            SdkOperation::Start
1371            | SdkOperation::Resume
1372            | SdkOperation::Input
1373            | SdkOperation::Interrupt
1374            | SdkOperation::Steer
1375            | SdkOperation::Respond
1376            | SdkOperation::Close => self.runtime_call(method, request.params).await,
1377            SdkOperation::Events => unreachable!("handled before method dispatch"),
1378        };
1379        result.map_err(|error| sdk_error(request.operation, error))
1380    }
1381
1382    async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
1383        Ok(self
1384            .poll_sdk_events()
1385            .await
1386            .into_iter()
1387            .map(|(_, event)| event)
1388            .collect())
1389    }
1390}
1391
1392#[cfg(feature = "adapter-api")]
1393struct HostedRuntimeLease {
1394    connection: HostedHarnessConnection,
1395    _host: std::sync::Arc<HostedHarnessRuntime>,
1396    _registration: LiveRuntimeRegistration,
1397    _server: crate::server::FrontendHttpServer,
1398}
1399
1400#[async_trait::async_trait]
1401#[cfg(feature = "adapter-api")]
1402impl RuntimeConnection for HostedRuntimeLease {
1403    fn handle(&self) -> &crate::RuntimeHandle {
1404        self.connection.handle()
1405    }
1406
1407    async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
1408        self.connection.send_input(input).await
1409    }
1410
1411    async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
1412        self.connection.next_event().await
1413    }
1414
1415    async fn interrupt(&mut self) -> crate::Result<()> {
1416        self.connection.interrupt().await
1417    }
1418
1419    async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
1420        self.connection.respond(request_id, response).await
1421    }
1422
1423    async fn close(&mut self) -> crate::Result<()> {
1424        self.connection.close().await
1425    }
1426}
1427
1428async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
1429    let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
1430    loop {
1431        let now = tokio::time::Instant::now();
1432        if now >= deadline {
1433            return Ok(());
1434        }
1435        match tokio::time::timeout(deadline - now, connection.next_event()).await {
1436            Err(_) => return Ok(()),
1437            Ok(Ok(Some(event))) => {
1438                if let Some(message) = handshake_event_failure(&event) {
1439                    return Err(truncate_text(&message, 500));
1440                }
1441            }
1442            Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
1443            Ok(Err(error)) => return Err(error.to_string()),
1444        }
1445    }
1446}
1447
1448fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
1449    let detail = event
1450        .payload
1451        .get("message")
1452        .or_else(|| event.payload.get("line"))
1453        .and_then(Value::as_str)
1454        .unwrap_or(event.kind.as_str());
1455    match event.kind.as_str() {
1456        "transport_closed" => Some("runtime transport closed during startup".into()),
1457        "transport_error" => Some(format!("runtime transport error: {detail}")),
1458        "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
1459        // Stderr is retained as a runtime event, but is not transport health.
1460        // Grok, for example, can log an AuthorizationRequired error from an
1461        // optional background worker while its ACP session continues to send
1462        // updates and complete prompts normally.
1463        _ => None,
1464    }
1465}
1466
1467fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
1468    let total_messages = session.messages.len();
1469    let (offset, end) = projected_message_window(total_messages, options);
1470    json!({
1471        "session": projected_session_json(session, options),
1472        "summary": projected_session_summary(session, options),
1473        "window": {
1474            "has_more": offset > 0 || end < total_messages,
1475            "has_newer": end < total_messages,
1476            "has_older": offset > 0,
1477            "newer_items": normalized_item_count(&session.messages[end..]),
1478            "offset": offset,
1479            "older_items": normalized_item_count(&session.messages[..offset]),
1480            "returned": end.saturating_sub(offset),
1481            "total_messages": total_messages,
1482        }
1483    })
1484}
1485
1486fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
1487    messages
1488        .iter()
1489        .map(|message| {
1490            let conversation = usize::from(
1491                matches!(message.role, Role::Assistant | Role::User)
1492                    && message_has_content(message),
1493            );
1494            let tool_result =
1495                usize::from(message.role == Role::Tool && message_has_content(message));
1496            conversation + tool_result + message.tool_calls().len()
1497        })
1498        .sum()
1499}
1500
1501fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
1502    let mut conversational = session.messages.iter().filter(|message| {
1503        matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
1504    });
1505    let first_message = conversational.clone().next();
1506    let last_message = conversational.next_back();
1507    let mut assistant = session
1508        .messages
1509        .iter()
1510        .filter(|message| message.role == Role::Assistant && message_has_content(message));
1511    let first_assistant_message = assistant.clone().next();
1512    let last_assistant_message = assistant.next_back();
1513    let end_of_turn = session
1514        .messages
1515        .iter()
1516        .rev()
1517        .find(|message| message.role != Role::System)
1518        .is_some_and(|message| {
1519            message.role == Role::Assistant
1520                && message_has_content(message)
1521                && message.tool_calls().is_empty()
1522        });
1523    let project = |message: Option<&crate::ChatMessage>| {
1524        message.map(|message| project_inline_media(message_json(message), options))
1525    };
1526    json!({
1527        "end_of_turn": end_of_turn,
1528        "first_assistant_message": project(first_assistant_message),
1529        "first_message": project(first_message),
1530        "last_assistant_message": project(last_assistant_message),
1531        "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
1532        "last_message": project(last_message),
1533    })
1534}
1535
1536fn message_has_content(message: &crate::ChatMessage) -> bool {
1537    message
1538        .content
1539        .as_deref()
1540        .is_some_and(|content| !content.trim().is_empty())
1541        || message
1542            .content_parts
1543            .as_ref()
1544            .is_some_and(|parts| !parts.is_empty())
1545}
1546
1547fn message_text(message: &crate::ChatMessage) -> String {
1548    if let Some(content) = &message.content {
1549        return content.clone();
1550    }
1551    message
1552        .content_parts
1553        .as_ref()
1554        .into_iter()
1555        .flatten()
1556        .filter_map(|part| part.get("text").and_then(Value::as_str))
1557        .collect::<Vec<_>>()
1558        .join("\n")
1559}
1560
1561fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
1562    let (offset, end) = projected_message_window(session.messages.len(), options);
1563    let messages = session.messages[offset..end]
1564        .iter()
1565        .map(|message| project_inline_media(message_json(message), options))
1566        .collect::<Vec<_>>();
1567    let subagents = if options.include_subagents.unwrap_or(true) {
1568        // The reported window describes the top-level transcript. Applying it
1569        // recursively would silently truncate subagents without returning a
1570        // window for each child. Keep their histories complete while carrying
1571        // the caller's media policy through the tree.
1572        let subagent_options = SessionLoadOptions {
1573            message_limit: None,
1574            message_offset: None,
1575            message_tail: None,
1576            ..options.clone()
1577        };
1578        session
1579            .subagents
1580            .iter()
1581            .map(|subagent| projected_session_json(subagent, &subagent_options))
1582            .collect::<Vec<_>>()
1583    } else {
1584        Vec::new()
1585    };
1586    json!({
1587        "source": match session.meta.source {
1588            SessionSource::ClaudeCode => "claude_code",
1589            SessionSource::Codex => "codex",
1590            SessionSource::Gemini => "gemini",
1591            SessionSource::Goose => "goose",
1592            SessionSource::Grok => "grok",
1593            SessionSource::Native => "native",
1594            SessionSource::OpenCode => "opencode",
1595            SessionSource::Pi => "pi",
1596        },
1597        "session_id": session.meta.session_id,
1598        "model": session.meta.model,
1599        "cwd": session.meta.cwd,
1600        "system_prompt": session.meta.system_prompt,
1601        "agent_id": session.meta.agent_id,
1602        "parent_tool_use_id": session.meta.parent_tool_use_id,
1603        "lineage": session.meta.lineage,
1604        "messages": messages,
1605        "subagents": subagents,
1606        "raw_record_count": session.raw.len(),
1607        "parse_error_lines": session.parse_error_lines,
1608    })
1609}
1610
1611fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
1612    if let Some(tail) = options.message_tail {
1613        return (total.saturating_sub(tail), total);
1614    }
1615    let offset = options.message_offset.unwrap_or(0).min(total);
1616    let end = options
1617        .message_limit
1618        .map(|limit| offset.saturating_add(limit).min(total))
1619        .unwrap_or(total);
1620    (offset, end)
1621}
1622
1623fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
1624    let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
1625        return message;
1626    };
1627    for part in parts {
1628        let Some(url) = part
1629            .get("image_url")
1630            .and_then(|image| image.get("url"))
1631            .and_then(Value::as_str)
1632        else {
1633            continue;
1634        };
1635        let Some(rest) = url.strip_prefix("data:") else {
1636            continue;
1637        };
1638        let Some((media_type, encoded)) = rest.split_once(";base64,") else {
1639            continue;
1640        };
1641        let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
1642        let decoded_bytes = encoded.len().saturating_mul(3) / 4;
1643        let decoded_bytes = decoded_bytes.saturating_sub(padding);
1644        let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
1645            || options
1646                .max_inline_media_bytes
1647                .is_some_and(|limit| decoded_bytes > limit);
1648        if should_elide {
1649            *part = json!({
1650                "type": "media_reference",
1651                "media_type": media_type,
1652                "encoding": "base64",
1653                "encoded_bytes": encoded.len(),
1654                "decoded_bytes": decoded_bytes,
1655                "omitted": true,
1656            });
1657        }
1658    }
1659    message
1660}
1661
1662#[derive(Deserialize)]
1663struct LocatorParams {
1664    locator: SessionLocator,
1665    /// Optional fidelity for the READ surfaces (`sessions.load`,
1666    /// `sessions.follow`).
1667    ///
1668    /// Omitted means [`Fidelity::Semantic`]: these two methods only ever
1669    /// produce a read-only view, and a compacted or resumed-across-files
1670    /// transcript — the everyday shape of a long Claude Code session — has no
1671    /// losslessly reconstructable record graph, so refusing to render it made
1672    /// the mirror unusable rather than accurate. A caller that intends to
1673    /// CONTINUE from what it reads asks for a lossless level explicitly and
1674    /// gets the strict refusal back. Every other method (export, translate,
1675    /// branch, handoff, resume_instructions) is lossless-only and has no
1676    /// such knob.
1677    #[serde(default)]
1678    fidelity: Option<Fidelity>,
1679    /// Optional bounded frontend projection. Absent preserves the historical
1680    /// complete-session read contract.
1681    #[serde(default)]
1682    view: Option<SessionReadView>,
1683}
1684
1685#[derive(Deserialize)]
1686struct SessionReadView {
1687    /// Number of trailing normalized messages to return. Zero is treated as
1688    /// one so a caller cannot accidentally request an unbounded empty mode.
1689    #[serde(default)]
1690    tail_messages: Option<usize>,
1691    /// Whether Claude Code child transcripts belong in this view. The
1692    /// frontend default is false; the legacy no-view path remains true.
1693    #[serde(default)]
1694    include_subagents: bool,
1695    /// Preserve human-visible native history across model-context compaction.
1696    #[serde(default)]
1697    display_history: bool,
1698    /// Bound each individual text field so a single tool result cannot turn a
1699    /// small message window into a hundred-megabyte RPC response.
1700    #[serde(default)]
1701    max_message_chars: Option<usize>,
1702}
1703
1704impl LocatorParams {
1705    fn read_fidelity(&self) -> Fidelity {
1706        self.fidelity.unwrap_or(Fidelity::Semantic)
1707    }
1708
1709    fn include_subagents(&self) -> bool {
1710        self.view
1711            .as_ref()
1712            .map(|view| view.include_subagents)
1713            .unwrap_or(true)
1714    }
1715
1716    fn tail_messages(&self) -> Option<usize> {
1717        self.view
1718            .as_ref()
1719            .and_then(|view| view.tail_messages)
1720            .map(|limit| limit.clamp(1, 5_000))
1721    }
1722
1723    fn display_history(&self) -> bool {
1724        self.view.as_ref().is_some_and(|view| view.display_history)
1725    }
1726
1727    fn max_message_chars(&self) -> Option<usize> {
1728        self.view
1729            .as_ref()
1730            .and_then(|view| view.max_message_chars)
1731            .map(|limit| limit.clamp(256, 64_000))
1732    }
1733
1734    fn bound_session(&self, session: &mut Session) {
1735        bound_session_view(session, self.tail_messages(), self.max_message_chars());
1736    }
1737}
1738
1739#[derive(Debug, Clone, Copy, Default, Deserialize)]
1740#[serde(rename_all = "snake_case")]
1741enum InlineMediaMode {
1742    #[default]
1743    Full,
1744    Metadata,
1745}
1746
1747#[derive(Debug, Clone, Default, Deserialize)]
1748#[serde(default)]
1749struct SessionLoadOptions {
1750    include_subagents: Option<bool>,
1751    inline_media: InlineMediaMode,
1752    max_inline_media_bytes: Option<usize>,
1753    message_limit: Option<usize>,
1754    message_offset: Option<usize>,
1755    message_tail: Option<usize>,
1756}
1757
1758impl SessionLoadOptions {
1759    fn validate(&self) -> std::result::Result<(), ServiceError> {
1760        if self.message_tail.is_some()
1761            && (self.message_limit.is_some() || self.message_offset.is_some())
1762        {
1763            return Err(ServiceError::InvalidParams(
1764                "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
1765                    .into(),
1766            ));
1767        }
1768        Ok(())
1769    }
1770}
1771
1772#[derive(Deserialize)]
1773struct LoadSessionParams {
1774    #[serde(flatten)]
1775    read: LocatorParams,
1776    #[serde(default)]
1777    options: Option<SessionLoadOptions>,
1778}
1779
1780#[derive(Deserialize)]
1781struct UnfollowParams {
1782    subscription: String,
1783}
1784
1785#[derive(Deserialize)]
1786struct MessageSessionParams {
1787    locator: SessionLocator,
1788    text: String,
1789    /// Same storage roots discovery accepts, so a caller (and a test) can
1790    /// point the live-session registry somewhere other than `$HOME`.
1791    #[serde(default)]
1792    homes: crate::HarnessHomes,
1793}
1794
1795/// Deliver `text` into a session that is running right now, or say why not.
1796///
1797/// A refusal is a RESULT, not a JSON-RPC error: "that session is persisted
1798/// only" is an answer about the session, which a mirror renders next to the
1799/// transcript, and this service's error envelope carries no structured data
1800/// field a machine-readable reason could survive in.
1801///
1802/// `delivered_to_bus` is the honest ceiling of what the courier proves. The
1803/// message reached the receiving session's inbox; whether that session ever
1804/// reads it is governed by ITS OWN inbound controls (`crossSessionInbound`,
1805/// approval dialogs), which Supercode neither sees nor overrides.
1806#[cfg(feature = "adapter-api")]
1807async fn message_live_session(
1808    params: &MessageSessionParams,
1809    runner: &dyn crate::claude_peer::CourierRunner,
1810) -> Value {
1811    if params.locator.harness.as_str() != HarnessId::CLAUDE_CODE {
1812        return json!({
1813            "delivered_to_bus": false,
1814            "refusal": {
1815                "reason": crate::claude_peer::ClaudePeerRefusal::HarnessUnsupported.as_str(),
1816                "message": format!(
1817                    "`{}` does not publish a live-session registry; only claude-code sessions can be messaged in place",
1818                    params.locator.harness.as_str()
1819                ),
1820            },
1821        });
1822    }
1823    match crate::claude_peer::message_claude_peer(
1824        &params.homes,
1825        &params.locator.session_id,
1826        &params.text,
1827        runner,
1828    )
1829    .await
1830    {
1831        Ok(delivery) => json!({
1832            "delivered_to_bus": true,
1833            "target": {
1834                "session_id": delivery.target.session_id,
1835                "name": delivery.target.name,
1836                "pid": delivery.target.pid,
1837                "cwd": delivery.target.cwd,
1838                "status": delivery.target.status.map(|status| status.as_str()),
1839            },
1840            "courier": {
1841                "model": crate::claude_peer::COURIER_MODEL,
1842                "report": delivery.courier_report,
1843            },
1844        }),
1845        Err(refusal) => json!({
1846            "delivered_to_bus": false,
1847            "refusal": {"reason": refusal.reason.as_str(), "message": refusal.message},
1848        }),
1849    }
1850}
1851
1852/// Source identity of one follow subscription, plus the last lifecycle state
1853/// already reported on it. The follower itself stays purely persistence-facing.
1854// Only the adapter-api poll reads these; the subscription bookkeeping itself is
1855// shared by both builds.
1856#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
1857struct FollowedSource {
1858    harness: String,
1859    session_id: String,
1860    reported: Option<String>,
1861}
1862
1863#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
1864#[serde(rename_all = "kebab-case")]
1865enum TransferFormat {
1866    ClaudeCode,
1867    Codex,
1868    #[serde(rename = "opencode", alias = "open-code")]
1869    OpenCode,
1870    Pi,
1871    Grok,
1872    Gemini,
1873    Goose,
1874}
1875
1876impl TransferFormat {
1877    fn id(self) -> &'static str {
1878        match self {
1879            Self::ClaudeCode => HarnessId::CLAUDE_CODE,
1880            Self::Codex => HarnessId::CODEX,
1881            Self::OpenCode => HarnessId::OPENCODE,
1882            Self::Pi => HarnessId::PI,
1883            Self::Grok => HarnessId::GROK,
1884            Self::Gemini => HarnessId::GEMINI,
1885            Self::Goose => HarnessId::GOOSE,
1886        }
1887    }
1888}
1889
1890impl From<TransferFormat> for SessionFormat {
1891    fn from(value: TransferFormat) -> Self {
1892        match value {
1893            TransferFormat::ClaudeCode => Self::ClaudeCode,
1894            TransferFormat::Codex => Self::Codex,
1895            TransferFormat::OpenCode => Self::OpenCode,
1896            TransferFormat::Pi => Self::Pi,
1897            TransferFormat::Grok => Self::Grok,
1898            TransferFormat::Gemini => Self::Gemini,
1899            TransferFormat::Goose => Self::Goose,
1900        }
1901    }
1902}
1903
1904#[derive(Deserialize)]
1905struct ImportSessionParams {
1906    source_harness: TransferFormat,
1907    content: String,
1908}
1909
1910#[derive(Deserialize)]
1911struct ExportSessionParams {
1912    locator: SessionLocator,
1913    target_harness: TransferFormat,
1914}
1915
1916#[derive(Deserialize)]
1917struct ReduceSessionParams {
1918    locator: SessionLocator,
1919    target_harness: TransferFormat,
1920    #[serde(default = "default_keep_last")]
1921    keep_last: usize,
1922}
1923
1924fn default_keep_last() -> usize {
1925    6
1926}
1927
1928#[derive(Deserialize)]
1929struct BranchSessionParams {
1930    locator: SessionLocator,
1931    #[serde(default)]
1932    target_harness: Option<TransferFormat>,
1933}
1934
1935#[derive(Deserialize)]
1936struct HandoffSessionParams {
1937    locator: SessionLocator,
1938    target_harness: TransferFormat,
1939    #[serde(default)]
1940    cwd: Option<PathBuf>,
1941}
1942
1943#[derive(Debug, Clone, Copy, Default, Deserialize)]
1944#[serde(rename_all = "snake_case")]
1945enum ResumePolicy {
1946    #[default]
1947    Default,
1948    Yolo,
1949}
1950
1951#[derive(Deserialize)]
1952struct ResumeInstructionsParams {
1953    locator: SessionLocator,
1954    #[serde(default)]
1955    cwd: Option<PathBuf>,
1956    #[serde(default)]
1957    policy: ResumePolicy,
1958}
1959
1960#[derive(Serialize)]
1961struct SessionArtifact {
1962    source_harness: HarnessId,
1963    target_harness: &'static str,
1964    session_id: Option<String>,
1965    content: String,
1966    suggested_filename: String,
1967    files: Vec<SessionArtifactFile>,
1968    fidelity: Fidelity,
1969    residue: Vec<String>,
1970}
1971
1972#[derive(Serialize)]
1973struct SessionArtifactFile {
1974    path: String,
1975    content: String,
1976    role: ArtifactFileRole,
1977}
1978
1979#[derive(Serialize)]
1980#[serde(rename_all = "snake_case")]
1981enum ArtifactFileRole {
1982    Primary,
1983    Subagent,
1984    Bundle,
1985    SourceRecovery,
1986}
1987
1988#[derive(Serialize)]
1989struct StructuredLaunch {
1990    cwd: PathBuf,
1991    program: String,
1992    arguments: Vec<String>,
1993    env: BTreeMap<String, String>,
1994}
1995
1996struct HandoffInstructions {
1997    launch: StructuredLaunch,
1998    materialize: Option<StructuredLaunch>,
1999    requires_materialization: bool,
2000    note: String,
2001}
2002
2003#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
2004#[serde(rename_all = "snake_case")]
2005enum HarnessProbeLevel {
2006    #[default]
2007    Passive,
2008    Handshake,
2009}
2010
2011#[derive(Default, Deserialize)]
2012#[serde(default)]
2013struct HarnessInventoryParams {
2014    harness: Option<HarnessId>,
2015    harnesses: Vec<HarnessId>,
2016    workspace: Option<PathBuf>,
2017    probe: HarnessProbeLevel,
2018    include_sessions: bool,
2019    /// Omit subprocess-based `--version` calls when a latency-sensitive UI only needs readiness.
2020    skip_versions: bool,
2021}
2022
2023#[derive(Serialize)]
2024struct HarnessInventoryReport {
2025    probe: HarnessProbeLevel,
2026    workspace: Option<PathBuf>,
2027    harnesses: Vec<LocalHarness>,
2028}
2029
2030#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
2031#[serde(rename_all = "snake_case")]
2032enum HarnessAuthState {
2033    Ready,
2034    Configured,
2035    Required,
2036    Unknown,
2037}
2038
2039#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
2040#[serde(rename_all = "snake_case")]
2041enum HarnessRuntimeState {
2042    Ready,
2043    Degraded,
2044    Unavailable,
2045}
2046
2047#[derive(Serialize)]
2048struct HarnessSessionCounts {
2049    global: Option<usize>,
2050    workspace: Option<usize>,
2051}
2052
2053#[derive(Serialize)]
2054struct LocalHarness {
2055    id: HarnessId,
2056    display_name: String,
2057    supported: bool,
2058    installed: bool,
2059    executable: Option<String>,
2060    version: Option<String>,
2061    auth: HarnessAuthState,
2062    runtime: HarnessRuntimeState,
2063    protocol: String,
2064    capabilities: crate::RuntimeCapabilities,
2065    effective_capabilities: crate::RuntimeCapabilities,
2066    sessions: HarnessSessionCounts,
2067    reason: Option<String>,
2068    repair: Option<String>,
2069}
2070
2071#[derive(Clone, Deserialize)]
2072struct RuntimeBackendParams {
2073    harness: HarnessId,
2074    #[serde(default)]
2075    protocol: Option<String>,
2076    #[serde(default)]
2077    launch: Option<RuntimeLaunch>,
2078    #[serde(default)]
2079    base_url: Option<String>,
2080    #[serde(default)]
2081    policy: RuntimePolicy,
2082}
2083
2084#[derive(Debug, Clone, Copy, Default, Deserialize)]
2085#[serde(rename_all = "snake_case")]
2086enum RuntimePolicy {
2087    #[default]
2088    Default,
2089    Yolo,
2090}
2091
2092#[derive(Deserialize)]
2093struct RuntimeStartParams {
2094    #[serde(flatten)]
2095    backend: RuntimeBackendParams,
2096    cwd: PathBuf,
2097}
2098
2099#[derive(Deserialize)]
2100struct RuntimeAttachParams {
2101    #[serde(flatten)]
2102    backend: RuntimeBackendParams,
2103    runtime_id: String,
2104    #[serde(default)]
2105    cwd: Option<PathBuf>,
2106}
2107
2108#[derive(Deserialize)]
2109struct RuntimeConnectionParams {
2110    connection: String,
2111}
2112
2113#[derive(Deserialize)]
2114struct RuntimeInputParams {
2115    connection: String,
2116    text: String,
2117}
2118
2119#[derive(Deserialize)]
2120struct RuntimeRespondParams {
2121    connection: String,
2122    request_id: Value,
2123    response: Value,
2124}
2125
2126fn default_reduction_store_root() -> PathBuf {
2127    if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
2128        return PathBuf::from(root).join("sessions");
2129    }
2130    if let Some(home) = std::env::var_os("HOME") {
2131        return PathBuf::from(home).join(".supercode").join("sessions");
2132    }
2133    PathBuf::from(".supercode").join("sessions")
2134}
2135
2136fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
2137    let mut output = String::new();
2138    for message in messages {
2139        output.push_str(
2140            &serde_json::to_string(message)
2141                .map_err(|error| ServiceError::Operation(error.to_string()))?,
2142        );
2143        output.push('\n');
2144    }
2145    Ok(output)
2146}
2147
2148fn parse_messages_jsonl(
2149    content: &str,
2150) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
2151    content
2152        .lines()
2153        .enumerate()
2154        .filter(|(_, line)| !line.trim().is_empty())
2155        .map(|(index, line)| {
2156            serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
2157                ServiceError::Operation(format!(
2158                    "reduced transcript line {} is invalid: {error}",
2159                    index + 1
2160                ))
2161            })
2162        })
2163        .collect()
2164}
2165
2166fn reduced_bootstrap_prompt(
2167    source: &SessionLocator,
2168    target: TransferFormat,
2169    view_jsonl: &str,
2170    sidecar_path: &Path,
2171    reduction_log_path: &Path,
2172) -> String {
2173    format!(
2174        "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
2175         \n\
2176         The bounded working transcript is below. Treat reduction markers as transparent placeholders, not missing work. If a detail behind a marker is needed, use ordinary file-reading/search tools against the full Supercode sidecar at `{sidecar}` and its reduction index at `{log}`. Do not guess hidden content. Both files were reloaded and verified before this continuation was issued.\n\
2177         \n\
2178         <supercode-reduced-session source-session=\"{source_id}\">\n\
2179         {view_jsonl}\
2180         </supercode-reduced-session>\n\
2181         \n\
2182         Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
2183        source_harness = source.harness.as_str(),
2184        target_harness = target.id(),
2185        sidecar = sidecar_path.display(),
2186        log = reduction_log_path.display(),
2187        source_id = source.session_id,
2188    )
2189}
2190
2191fn session_artifact(
2192    locator: &SessionLocator,
2193    session: &Session,
2194    target: TransferFormat,
2195) -> std::result::Result<SessionArtifact, ServiceError> {
2196    session_artifact_with_id(locator, session, target, None)
2197}
2198
2199fn session_artifact_with_id(
2200    locator: &SessionLocator,
2201    session: &Session,
2202    target: TransferFormat,
2203    target_session_id: Option<&str>,
2204) -> std::result::Result<SessionArtifact, ServiceError> {
2205    let format: SessionFormat = target.into();
2206    let diagonal = format.source() == session.meta.source;
2207    let has_appended_turns = session
2208        .imported_message_count
2209        .is_some_and(|imported| imported < session.messages.len());
2210    let content = if let Some(id) = target_session_id {
2211        if diagonal && format != SessionFormat::OpenCode {
2212            session
2213                .to_jsonl_spliced(format, Some(id))
2214                .map_err(operation)?
2215        } else {
2216            let mut rewritten = session.clone();
2217            rewritten.meta.session_id = Some(id.to_string());
2218            rewritten.to_jsonl(format).map_err(operation)?
2219        }
2220    } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
2221        session.raw_verbatim()
2222    } else if diagonal {
2223        session.to_jsonl_spliced(format, None).map_err(operation)?
2224    } else {
2225        session.to_jsonl(format).map_err(operation)?
2226    };
2227    let stem = sanitize_filename(
2228        target_session_id
2229            .or(session.meta.session_id.as_deref())
2230            .unwrap_or(&locator.session_id),
2231    );
2232    let suggested_filename = if diagonal && target == TransferFormat::Grok {
2233        "chat_history.jsonl".to_string()
2234    } else if target == TransferFormat::Goose {
2235        format!("{stem}.goose.json")
2236    } else {
2237        format!("{stem}.{}.jsonl", target.id())
2238    };
2239    let mut files = vec![SessionArtifactFile {
2240        path: suggested_filename.clone(),
2241        content: content.clone(),
2242        role: ArtifactFileRole::Primary,
2243    }];
2244    if target == TransferFormat::ClaudeCode {
2245        let bundle_stem = Path::new(&suggested_filename)
2246            .file_stem()
2247            .and_then(|stem| stem.to_str())
2248            .unwrap_or(&stem);
2249        let mut child_paths = BTreeSet::new();
2250        for (index, subagent) in session.subagents.iter().enumerate() {
2251            let agent_id = subagent
2252                .meta
2253                .agent_id
2254                .as_deref()
2255                .map(|id| id.strip_prefix("agent-").unwrap_or(id))
2256                .map(sanitize_filename)
2257                .filter(|id| !id.is_empty())
2258                .unwrap_or_else(|| format!("subagent-{}", index + 1));
2259            let child_has_appended_turns = subagent
2260                .imported_message_count
2261                .is_some_and(|imported| imported < subagent.messages.len());
2262            let child_content = if target_session_id.is_none()
2263                && subagent.meta.source == SessionSource::ClaudeCode
2264                && subagent.raw_is_verbatim
2265                && !child_has_appended_turns
2266            {
2267                subagent.raw_verbatim()
2268            } else if subagent.meta.source == SessionSource::ClaudeCode {
2269                subagent
2270                    .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
2271                    .map_err(operation)?
2272            } else {
2273                let mut child = subagent.clone();
2274                if let Some(id) = target_session_id {
2275                    child.meta.session_id = Some(id.to_string());
2276                }
2277                child
2278                    .to_jsonl(SessionFormat::ClaudeCode)
2279                    .map_err(operation)?
2280            };
2281            let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
2282            if !child_paths.insert(path.clone()) {
2283                return Err(ServiceError::Operation(format!(
2284                    "Claude subagent ids collide at artifact path `{path}`"
2285                )));
2286            }
2287            files.push(SessionArtifactFile {
2288                path,
2289                content: child_content,
2290                role: ArtifactFileRole::Subagent,
2291            });
2292        }
2293    }
2294    if diagonal && target == TransferFormat::Grok {
2295        append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
2296    }
2297    if !diagonal || !session.raw_is_verbatim {
2298        files.push(SessionArtifactFile {
2299            path: "recovery/source.supercode.jsonl".into(),
2300            content: session.to_native_jsonl(),
2301            role: ArtifactFileRole::SourceRecovery,
2302        });
2303        for (index, subagent) in session.subagents.iter().enumerate() {
2304            let id = subagent
2305                .meta
2306                .agent_id
2307                .as_deref()
2308                .map(sanitize_filename)
2309                .unwrap_or_else(|| format!("subagent-{}", index + 1));
2310            files.push(SessionArtifactFile {
2311                path: format!("recovery/subagents/{id}.supercode.jsonl"),
2312                content: subagent.to_native_jsonl(),
2313                role: ArtifactFileRole::SourceRecovery,
2314            });
2315        }
2316    }
2317    if !diagonal && session.meta.source == SessionSource::Grok {
2318        append_grok_bundle_files(
2319            locator,
2320            "recovery/grok/",
2321            ArtifactFileRole::SourceRecovery,
2322            &mut files,
2323        )?;
2324    }
2325    let (fidelity, residue) = if diagonal
2326        && target_session_id.is_none()
2327        && session.raw_is_verbatim
2328        && !has_appended_turns
2329    {
2330        (Fidelity::ByteLossless, Vec::new())
2331    } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
2332        (
2333            Fidelity::ValueLossless,
2334            vec![if target_session_id.is_some() {
2335                "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
2336            } else {
2337                "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
2338            }],
2339        )
2340    } else {
2341        (
2342            Fidelity::Semantic,
2343            vec!["target schema has no portable slot for every source-native record and metadata field".into()],
2344        )
2345    };
2346    Ok(SessionArtifact {
2347        source_harness: locator.harness.clone(),
2348        target_harness: target.id(),
2349        session_id: target_session_id
2350            .map(str::to_string)
2351            .or_else(|| session.meta.session_id.clone()),
2352        content,
2353        suggested_filename,
2354        files,
2355        fidelity,
2356        residue,
2357    })
2358}
2359
2360fn append_grok_bundle_files(
2361    locator: &SessionLocator,
2362    prefix: &str,
2363    role: ArtifactFileRole,
2364    files: &mut Vec<SessionArtifactFile>,
2365) -> std::result::Result<(), ServiceError> {
2366    let primary = locator.storage.path();
2367    if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
2368        return Err(ServiceError::Operation(format!(
2369            "Grok bundle locator must name chat_history.jsonl, got {}",
2370            primary.display()
2371        )));
2372    }
2373    let parent = primary.parent().ok_or_else(|| {
2374        ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
2375    })?;
2376    for name in ["summary.json", "updates.jsonl"] {
2377        let path = parent.join(name);
2378        let metadata = match std::fs::symlink_metadata(&path) {
2379            Ok(metadata) => metadata,
2380            Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
2381            Err(error) => return Err(ServiceError::Operation(error.to_string())),
2382        };
2383        if metadata.file_type().is_symlink() || !metadata.is_file() {
2384            return Err(ServiceError::Operation(format!(
2385                "refusing non-regular Grok bundle member {}",
2386                path.display()
2387            )));
2388        }
2389        let content = std::fs::read_to_string(&path).map_err(|error| {
2390            ServiceError::Operation(format!(
2391                "Grok bundle member {} is not representable as UTF-8: {error}",
2392                path.display()
2393            ))
2394        })?;
2395        files.push(SessionArtifactFile {
2396            path: format!("{prefix}{name}"),
2397            content,
2398            role: match role {
2399                ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
2400                _ => ArtifactFileRole::SourceRecovery,
2401            },
2402        });
2403    }
2404    Ok(())
2405}
2406
2407fn handoff_artifact(
2408    locator: &SessionLocator,
2409    session: &Session,
2410    target: TransferFormat,
2411    cwd: &Path,
2412) -> std::result::Result<SessionArtifact, ServiceError> {
2413    if target != TransferFormat::Grok {
2414        let target_session_id = target_session_id(target);
2415        return session_artifact_with_id(locator, session, target, Some(&target_session_id));
2416    }
2417
2418    // Stock Grok's importer accepts Claude/Codex transcripts and materializes its own
2419    // multi-file session bundle. A synthesized Grok chat_history.jsonl alone is not a
2420    // resumable handoff because updates.jsonl is the authoritative restore log.
2421    let mut importable = session.clone();
2422    // The Claude importer validates sessionId as a UUID. Source harness identities
2423    // are not portable (OpenCode, for example, uses `ses_...`), and a handoff must
2424    // not overwrite an existing target session when the source already uses UUIDs.
2425    // Mint a distinct target identity and still bind the importer-returned ID at
2426    // launch time because the importer remains the authority on materialization.
2427    importable.meta.session_id = Some(target_session_id(TransferFormat::ClaudeCode));
2428    importable.meta.cwd = Some(if cwd.is_absolute() {
2429        cwd.to_path_buf()
2430    } else {
2431        std::env::current_dir()
2432            .map_err(|error| ServiceError::Operation(error.to_string()))?
2433            .join(cwd)
2434    });
2435    let content = importable
2436        .to_jsonl(SessionFormat::ClaudeCode)
2437        .map_err(operation)?;
2438    let stem = sanitize_filename(
2439        importable
2440            .meta
2441            .session_id
2442            .as_deref()
2443            .unwrap_or(&locator.session_id),
2444    );
2445    let suggested_filename = format!("{stem}.grok-import.claude-code.jsonl");
2446    Ok(SessionArtifact {
2447        source_harness: locator.harness.clone(),
2448        // This names the artifact's actual wire format. The requested handoff target
2449        // remains Grok; its official importer is the materialization boundary.
2450        target_harness: TransferFormat::ClaudeCode.id(),
2451        session_id: importable.meta.session_id.clone(),
2452        content: content.clone(),
2453        suggested_filename: suggested_filename.clone(),
2454        files: vec![SessionArtifactFile {
2455            path: suggested_filename,
2456            content,
2457            role: ArtifactFileRole::Primary,
2458        }],
2459        fidelity: Fidelity::Semantic,
2460        residue: vec!["Grok's stock importer accepts a Claude Code transcript, not a complete Grok updates/session bundle".into()],
2461    })
2462}
2463
2464fn target_session_id(target: TransferFormat) -> String {
2465    let uuid = generated_session_id();
2466    match target {
2467        TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
2468        TransferFormat::ClaudeCode
2469        | TransferFormat::Codex
2470        | TransferFormat::Pi
2471        | TransferFormat::Grok
2472        | TransferFormat::Gemini
2473        | TransferFormat::Goose => uuid,
2474    }
2475}
2476
2477fn sanitize_filename(value: &str) -> String {
2478    let value = value
2479        .chars()
2480        .map(|character| {
2481            if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
2482                character
2483            } else {
2484                '-'
2485            }
2486        })
2487        .collect::<String>();
2488    let value = value.trim_matches('-');
2489    if value.is_empty() {
2490        "session".into()
2491    } else {
2492        value.chars().take(100).collect()
2493    }
2494}
2495
2496fn handoff_instructions(
2497    target: TransferFormat,
2498    session_id: &str,
2499    cwd: &Path,
2500) -> HandoffInstructions {
2501    let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
2502        cwd: cwd.to_path_buf(),
2503        program: program.into(),
2504        arguments,
2505        env: BTreeMap::new(),
2506    };
2507    match target {
2508        TransferFormat::ClaudeCode => HandoffInstructions {
2509            launch: launch("claude", vec!["--resume".into(), session_id.into()]),
2510            materialize: None,
2511            requires_materialization: true,
2512            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(),
2513        },
2514        TransferFormat::Codex => HandoffInstructions {
2515            launch: launch("codex", vec!["resume".into(), session_id.into()]),
2516            materialize: None,
2517            requires_materialization: true,
2518            note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
2519        },
2520        TransferFormat::OpenCode => HandoffInstructions {
2521            launch: launch("opencode", vec!["--session".into(), session_id.into()]),
2522            materialize: Some(launch(
2523                "opencode",
2524                vec!["import".into(), "{artifact_path}".into()],
2525            )),
2526            requires_materialization: true,
2527            note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
2528        },
2529        TransferFormat::Pi => HandoffInstructions {
2530            launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
2531            materialize: None,
2532            requires_materialization: true,
2533            note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
2534        },
2535        TransferFormat::Grok => HandoffInstructions {
2536            launch: launch(
2537                "grok",
2538                vec![
2539                    "--resume".into(),
2540                    "{imported_session_id}".into(),
2541                    "--fork-session".into(),
2542                ],
2543            ),
2544            materialize: Some(launch(
2545                "grok",
2546                vec!["import".into(), "--json".into(), "{artifact_path}".into()],
2547            )),
2548            requires_materialization: true,
2549            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(),
2550        },
2551        TransferFormat::Gemini => HandoffInstructions {
2552            launch: launch(
2553                "gemini",
2554                vec!["--session-file".into(), "{artifact_path}".into()],
2555            ),
2556            materialize: None,
2557            requires_materialization: true,
2558            note: "Write the Gemini JSONL artifact to a file and replace {artifact_path}; Gemini imports it into the current project's chat store before opening the continuation.".into(),
2559        },
2560        TransferFormat::Goose => HandoffInstructions {
2561            launch: launch(
2562                "goose",
2563                vec![
2564                    "session".into(),
2565                    "--resume".into(),
2566                    "--session-id".into(),
2567                    "{imported_session_id}".into(),
2568                ],
2569            ),
2570            materialize: Some(launch(
2571                "goose",
2572                vec!["session".into(), "import".into(), "{artifact_path}".into()],
2573            )),
2574            requires_materialization: true,
2575            note: "Write the Goose JSON artifact to a file, run the materialize command, read the imported session id from its output, replace {imported_session_id}, then resume that native Goose session.".into(),
2576        },
2577    }
2578}
2579
2580fn resume_launch(
2581    harness: &str,
2582    session_id: &str,
2583    cwd: &Path,
2584    policy: ResumePolicy,
2585) -> std::result::Result<StructuredLaunch, ServiceError> {
2586    let mut arguments = Vec::new();
2587    let program = match harness {
2588        HarnessId::GROK => {
2589            if matches!(policy, ResumePolicy::Yolo) {
2590                arguments.extend([
2591                    "--sandbox".into(),
2592                    "workspace".into(),
2593                    "--always-approve".into(),
2594                ]);
2595            }
2596            arguments.extend(["--resume".into(), session_id.into()]);
2597            "grok"
2598        }
2599        HarnessId::CODEX => {
2600            if matches!(policy, ResumePolicy::Yolo) {
2601                arguments.extend([
2602                    "--dangerously-bypass-approvals-and-sandbox".into(),
2603                    "--dangerously-bypass-hook-trust".into(),
2604                ]);
2605            }
2606            arguments.extend(["resume".into(), session_id.into()]);
2607            "codex"
2608        }
2609        HarnessId::CLAUDE_CODE => {
2610            if matches!(policy, ResumePolicy::Yolo) {
2611                arguments.push("--dangerously-skip-permissions".into());
2612            }
2613            arguments.extend(["--resume".into(), session_id.into()]);
2614            "claude"
2615        }
2616        HarnessId::GEMINI => {
2617            if matches!(policy, ResumePolicy::Yolo) {
2618                arguments.push("--yolo".into());
2619            }
2620            arguments.extend(["--resume".into(), session_id.into()]);
2621            "gemini"
2622        }
2623        HarnessId::GOOSE => {
2624            arguments.extend([
2625                "session".into(),
2626                "--resume".into(),
2627                "--session-id".into(),
2628                session_id.into(),
2629            ]);
2630            "goose"
2631        }
2632        HarnessId::PI => {
2633            if matches!(policy, ResumePolicy::Yolo) {
2634                arguments.push("--approve".into());
2635            }
2636            arguments.extend(["--session".into(), session_id.into()]);
2637            "pi"
2638        }
2639        HarnessId::OPENCODE => {
2640            arguments.extend(["--session".into(), session_id.into()]);
2641            "opencode"
2642        }
2643        HarnessId::SUPERCODE => {
2644            if matches!(policy, ResumePolicy::Yolo) {
2645                arguments.push("--dangerous".into());
2646            }
2647            arguments.extend(["resume".into(), session_id.into()]);
2648            "supercode"
2649        }
2650        other => {
2651            return Err(ServiceError::InvalidParams(format!(
2652                "no structured resume launch is registered for harness `{other}`"
2653            )))
2654        }
2655    };
2656    Ok(StructuredLaunch {
2657        cwd: cwd.to_path_buf(),
2658        program: program.into(),
2659        arguments,
2660        env: BTreeMap::new(),
2661    })
2662}
2663
2664fn runtime_backend(
2665    params: &RuntimeBackendParams,
2666) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
2667    if params.protocol.as_deref() == Some("acp") {
2668        let launch = params
2669            .launch
2670            .clone()
2671            .or_else(|| {
2672                harness_support_registry()
2673                    .harnesses
2674                    .into_iter()
2675                    .find(|harness| harness.id == params.harness)
2676                    .filter(|harness| {
2677                        harness.runtime.implementation == ImplementationKind::GenericProtocol
2678                            && harness.runtime.protocol.starts_with("acp")
2679                    })
2680                    .and_then(|harness| harness.runtime.default_launch)
2681            })
2682            .ok_or_else(|| {
2683                ServiceError::InvalidParams(
2684                    "an ACP runtime requires `launch` unless the harness has a registered default"
2685                        .into(),
2686                )
2687            })?;
2688        let resume_session = harness_support_registry()
2689            .harnesses
2690            .into_iter()
2691            .find(|harness| harness.id == params.harness)
2692            .is_some_and(|harness| harness.runtime.capabilities.resume_session);
2693        return Ok(Box::new(
2694            AcpRuntimeBackend::new(params.harness.clone(), launch)
2695                .with_resume_support(resume_session),
2696        ));
2697    }
2698    let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
2699        HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
2700        HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
2701        HarnessId::PI => Box::new(PiRuntimeBackend::new()),
2702        HarnessId::OPENCODE => match &params.base_url {
2703            Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
2704            None => Box::new(OpenCodeRuntimeBackend::new()),
2705        },
2706        harness => {
2707            let descriptor = harness_support_registry()
2708                .harnesses
2709                .into_iter()
2710                .find(|descriptor| descriptor.id.as_str() == harness)
2711                .filter(|descriptor| {
2712                    descriptor.runtime.implementation == ImplementationKind::GenericProtocol
2713                        && descriptor.runtime.protocol.starts_with("acp")
2714                });
2715            let Some(descriptor) = descriptor else {
2716                return Err(ServiceError::InvalidParams(format!(
2717                    "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
2718                )));
2719            };
2720            let resume = descriptor.runtime.capabilities.resume_session;
2721            Box::new(
2722                AcpRuntimeBackend::new(
2723                    descriptor.id,
2724                    descriptor
2725                        .runtime
2726                        .default_launch
2727                        .expect("generic ACP registry entry includes its launch"),
2728                )
2729                .with_resume_support(resume),
2730            )
2731        }
2732    };
2733    Ok(backend)
2734}
2735
2736fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
2737    if let Some(launch) = &params.launch {
2738        return Some(launch.clone());
2739    }
2740    if !matches!(params.policy, RuntimePolicy::Yolo) {
2741        return None;
2742    }
2743    let launch = match params.harness.as_str() {
2744        HarnessId::GROK => RuntimeLaunch {
2745            program: "grok".into(),
2746            arguments: vec![
2747                "--sandbox".into(),
2748                "workspace".into(),
2749                "--always-approve".into(),
2750                "agent".into(),
2751                "--no-leader".into(),
2752                "stdio".into(),
2753            ],
2754            env: BTreeMap::from([("GROK_AGENT_DASHBOARD".into(), "0".into())]),
2755        },
2756        HarnessId::CODEX => RuntimeLaunch {
2757            program: "codex".into(),
2758            arguments: vec![
2759                "--dangerously-bypass-approvals-and-sandbox".into(),
2760                "--dangerously-bypass-hook-trust".into(),
2761                "app-server".into(),
2762            ],
2763            env: BTreeMap::new(),
2764        },
2765        HarnessId::CLAUDE_CODE => RuntimeLaunch {
2766            program: "claude".into(),
2767            arguments: vec![
2768                "--dangerously-skip-permissions".into(),
2769                "--print".into(),
2770                "--input-format".into(),
2771                "stream-json".into(),
2772                "--output-format".into(),
2773                "stream-json".into(),
2774                "--verbose".into(),
2775            ],
2776            env: BTreeMap::new(),
2777        },
2778        HarnessId::PI => RuntimeLaunch {
2779            program: "pi".into(),
2780            arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
2781            env: BTreeMap::new(),
2782        },
2783        HarnessId::OPENCODE => RuntimeLaunch {
2784            program: "opencode".into(),
2785            arguments: vec!["serve".into()],
2786            env: BTreeMap::new(),
2787        },
2788        HarnessId::GEMINI => RuntimeLaunch {
2789            program: "gemini".into(),
2790            arguments: vec!["--acp".into(), "--yolo".into()],
2791            env: BTreeMap::new(),
2792        },
2793        HarnessId::GOOSE => RuntimeLaunch {
2794            program: "goose".into(),
2795            arguments: vec!["acp".into()],
2796            env: BTreeMap::new(),
2797        },
2798        HarnessId::SUPERCODE => RuntimeLaunch {
2799            program: "supercode".into(),
2800            arguments: vec!["acp".into(), "--dangerous".into()],
2801            env: BTreeMap::new(),
2802        },
2803        _ => return None,
2804    };
2805    Some(launch)
2806}
2807
2808/// Disposable harness state for a no-prompt readiness probe. Merely opening
2809/// several stock CLIs writes a session header or migrates configuration, so a
2810/// handshake must never point at the user's real home. Authentication files
2811/// are copied into the private temporary home; all writes disappear with the
2812/// guard after the connection closes.
2813struct IsolatedProbeHome {
2814    launch: RuntimeLaunch,
2815    root: PathBuf,
2816}
2817
2818impl IsolatedProbeHome {
2819    fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
2820        let root = std::env::temp_dir().join(format!(
2821            "supercode-harness-probe-{harness}-{}",
2822            generated_session_id()
2823        ));
2824        std::fs::create_dir_all(&root)?;
2825        set_private_dir_permissions(&root)?;
2826
2827        if let Some(source_home) = std::env::var_os("HOME").map(PathBuf::from) {
2828            for relative in probe_auth_files(harness) {
2829                copy_probe_file(&source_home, &root, relative)?;
2830            }
2831        }
2832        configure_isolated_probe_auth(harness, &root)?;
2833
2834        let root_text = root.to_string_lossy().into_owned();
2835        for (key, value) in [
2836            ("HOME", root_text.clone()),
2837            (
2838                "XDG_CACHE_HOME",
2839                root.join(".cache").to_string_lossy().into_owned(),
2840            ),
2841            (
2842                "XDG_CONFIG_HOME",
2843                root.join(".config").to_string_lossy().into_owned(),
2844            ),
2845            (
2846                "XDG_DATA_HOME",
2847                root.join(".local/share").to_string_lossy().into_owned(),
2848            ),
2849        ] {
2850            launch.env.insert(key.into(), value);
2851        }
2852        let scoped = match harness {
2853            HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
2854            HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
2855            HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
2856            HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
2857            HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
2858            HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
2859            _ => None,
2860        };
2861        if let Some((key, value)) = scoped {
2862            launch
2863                .env
2864                .insert(key.into(), value.to_string_lossy().into_owned());
2865        }
2866        Ok(Self { launch, root })
2867    }
2868
2869    fn cleanup(&self) -> std::io::Result<()> {
2870        match std::fs::remove_dir_all(&self.root) {
2871            Ok(()) => Ok(()),
2872            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
2873            Err(error) => Err(error),
2874        }
2875    }
2876}
2877
2878impl Drop for IsolatedProbeHome {
2879    fn drop(&mut self) {
2880        let _ = self.cleanup();
2881    }
2882}
2883
2884fn probe_auth_files(harness: &str) -> &'static [&'static str] {
2885    match harness {
2886        HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
2887        HarnessId::CODEX => &[".codex/auth.json"],
2888        HarnessId::GEMINI => &[
2889            ".gemini/google_accounts.json",
2890            ".gemini/oauth_creds.json",
2891            ".gemini/settings.json",
2892        ],
2893        HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
2894        HarnessId::OPENCODE => &[
2895            ".config/opencode/auth.json",
2896            ".local/share/opencode/auth.json",
2897        ],
2898        HarnessId::PI => &[".pi/agent/auth.json"],
2899        HarnessId::SUPERCODE => &[
2900            ".config/supercode/config.toml",
2901            ".config/supercode/credentials.toml",
2902        ],
2903        _ => &[],
2904    }
2905}
2906
2907fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
2908    let source = source_home.join(relative);
2909    if !source.is_file() {
2910        return Ok(());
2911    }
2912    let destination = probe_home.join(relative);
2913    if let Some(parent) = destination.parent() {
2914        std::fs::create_dir_all(parent)?;
2915        set_private_dir_permissions(parent)?;
2916    }
2917    std::fs::copy(source, &destination)?;
2918    set_private_file_permissions(&destination)
2919}
2920
2921fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
2922    if harness != HarnessId::GEMINI {
2923        return Ok(());
2924    }
2925    let oauth = probe_home.join(".gemini/oauth_creds.json");
2926    if !oauth.is_file() {
2927        return Ok(());
2928    }
2929    let settings_path = probe_home.join(".gemini/settings.json");
2930    let mut settings = std::fs::read_to_string(&settings_path)
2931        .ok()
2932        .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
2933        .unwrap_or_else(|| json!({}));
2934    settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
2935    std::fs::write(
2936        &settings_path,
2937        serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
2938    )?;
2939    set_private_file_permissions(&settings_path)
2940}
2941
2942#[cfg(unix)]
2943fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
2944    use std::os::unix::fs::PermissionsExt;
2945    std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
2946}
2947
2948#[cfg(not(unix))]
2949fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
2950    Ok(())
2951}
2952
2953#[cfg(unix)]
2954fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
2955    use std::os::unix::fs::PermissionsExt;
2956    std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
2957}
2958
2959#[cfg(not(unix))]
2960fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
2961    Ok(())
2962}
2963
2964fn find_executable(program: &str) -> Option<PathBuf> {
2965    let candidate = PathBuf::from(program);
2966    if candidate.components().count() > 1 {
2967        return candidate.is_file().then_some(candidate);
2968    }
2969    let path = std::env::var_os("PATH")?;
2970    for directory in std::env::split_paths(&path) {
2971        let candidate = directory.join(program);
2972        if candidate.is_file() {
2973            return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
2974        }
2975        #[cfg(windows)]
2976        {
2977            for extension in ["exe", "cmd", "bat"] {
2978                let candidate = directory.join(format!("{program}.{extension}"));
2979                if candidate.is_file() {
2980                    return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
2981                }
2982            }
2983        }
2984    }
2985    None
2986}
2987
2988async fn executable_version(executable: &Path) -> Option<String> {
2989    let mut command = tokio::process::Command::new(executable);
2990    command
2991        .arg("--version")
2992        .stdin(std::process::Stdio::null())
2993        .stdout(std::process::Stdio::piped())
2994        .stderr(std::process::Stdio::piped())
2995        .kill_on_drop(true);
2996    let output = tokio::time::timeout(Duration::from_secs(3), command.output())
2997        .await
2998        .ok()?
2999        .ok()?;
3000    let stdout = String::from_utf8_lossy(&output.stdout);
3001    let stderr = String::from_utf8_lossy(&output.stderr);
3002    stdout
3003        .lines()
3004        .chain(stderr.lines())
3005        .map(str::trim)
3006        .find(|line| !line.is_empty())
3007        .map(|line| truncate_text(line, 200))
3008}
3009
3010fn auth_evidence(harness: &str) -> bool {
3011    let env_names: &[&str] = match harness {
3012        HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
3013        HarnessId::CODEX => &["OPENAI_API_KEY"],
3014        HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
3015        HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
3016        HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
3017        HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
3018        HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
3019        _ => &[],
3020    };
3021    if env_names
3022        .iter()
3023        .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
3024    {
3025        return true;
3026    }
3027    let Some(home) = std::env::var_os("HOME").map(PathBuf::from) else {
3028        return false;
3029    };
3030    let files: Vec<PathBuf> = match harness {
3031        HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
3032        HarnessId::CODEX => vec![home.join(".codex/auth.json")],
3033        HarnessId::OPENCODE => vec![
3034            home.join(".local/share/opencode/auth.json"),
3035            home.join(".config/opencode/auth.json"),
3036        ],
3037        HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
3038        HarnessId::GROK => vec![home.join(".grok/auth.json")],
3039        HarnessId::GEMINI => vec![
3040            home.join(".gemini/oauth_creds.json"),
3041            home.join(".gemini/google_accounts.json"),
3042        ],
3043        HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
3044        _ => Vec::new(),
3045    };
3046    if files.into_iter().any(|path| {
3047        std::fs::metadata(path)
3048            .map(|metadata| metadata.is_file() && metadata.len() > 2)
3049            .unwrap_or(false)
3050    }) {
3051        return true;
3052    }
3053    // macOS keeps Claude Code's OAuth login in the Keychain, so
3054    // `.claude/.credentials.json` never exists there and the file probe above
3055    // reports a signed-in install as unauthenticated forever. A completed
3056    // login also writes an `oauthAccount` record into `~/.claude.json` on
3057    // every platform — file-based, prompt-free evidence (querying the
3058    // Keychain itself from an unsigned daemon can raise a UI prompt).
3059    if harness == HarnessId::CLAUDE_CODE {
3060        return std::fs::read_to_string(home.join(".claude.json"))
3061            .map(|text| text.contains("\"oauthAccount\""))
3062            .unwrap_or(false);
3063    }
3064    false
3065}
3066
3067fn looks_like_auth_error(message: &str) -> bool {
3068    let message = message.to_ascii_lowercase();
3069    [
3070        "auth",
3071        "login",
3072        "sign in",
3073        "sign-in",
3074        "credential",
3075        "unauthorized",
3076        "forbidden",
3077        "token",
3078    ]
3079    .iter()
3080    .any(|needle| message.contains(needle))
3081}
3082
3083fn unavailable_capabilities() -> crate::RuntimeCapabilities {
3084    crate::RuntimeCapabilities {
3085        start_session: false,
3086        resume_session: false,
3087        attach_existing_process: false,
3088        send_input: false,
3089        stream_events: false,
3090        interrupt: false,
3091        respond_to_requests: false,
3092    }
3093}
3094
3095fn truncate_text(text: &str, max_chars: usize) -> String {
3096    let mut chars = text.chars();
3097    let truncated = chars.by_ref().take(max_chars).collect::<String>();
3098    if chars.next().is_some() {
3099        format!("{truncated}…")
3100    } else {
3101        truncated
3102    }
3103}
3104
3105fn error_message(error: ServiceError) -> String {
3106    match error {
3107        ServiceError::InvalidParams(message)
3108        | ServiceError::Operation(message)
3109        | ServiceError::UnsupportedAction(message) => message,
3110        ServiceError::MethodNotFound => "runtime adapter is not available".into(),
3111        ServiceError::Sdk(error) => error.to_string(),
3112    }
3113}
3114
3115#[derive(Debug)]
3116enum ServiceError {
3117    InvalidParams(String),
3118    MethodNotFound,
3119    UnsupportedAction(String),
3120    Operation(String),
3121    Sdk(SdkError),
3122}
3123
3124fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
3125    match error {
3126        ServiceError::InvalidParams(message) => {
3127            SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
3128        }
3129        ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
3130            SdkError::unsupported(operation)
3131        }
3132        ServiceError::Operation(message) => {
3133            let code = if message.contains("already in progress") {
3134                SdkErrorCode::Busy
3135            } else if message.contains("not supported by this runtime") {
3136                SdkErrorCode::UnsupportedAction
3137            } else if message.contains("unknown runtime connection") {
3138                SdkErrorCode::NotFound
3139            } else {
3140                SdkErrorCode::Execution
3141            };
3142            SdkError::new(code, operation, message)
3143        }
3144        ServiceError::Sdk(error) => error,
3145    }
3146}
3147
3148fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
3149    let error_code = error.code();
3150    let code = match error_code {
3151        SdkErrorCode::Unauthenticated => -32030,
3152        SdkErrorCode::Unauthorized => -32031,
3153        SdkErrorCode::ControllerRequired => -32032,
3154        SdkErrorCode::LeaseExpired => -32033,
3155        SdkErrorCode::InvalidArgument => -32602,
3156        SdkErrorCode::NotFound => -32004,
3157        SdkErrorCode::Busy => -32000,
3158        SdkErrorCode::UnsupportedAction => -32020,
3159        SdkErrorCode::Execution => -32002,
3160        SdkErrorCode::Transport => -32003,
3161    };
3162    json!({
3163        "jsonrpc": "2.0",
3164        "id": id,
3165        "error": {
3166            "code": code,
3167            "name": error_code,
3168            "operation": error.operation(),
3169            "message": error.to_string(),
3170        },
3171    })
3172}
3173
3174fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
3175    serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
3176}
3177
3178fn operation(error: impl Into<crate::Error>) -> ServiceError {
3179    let error = error.into();
3180    match error {
3181        crate::Error::Sdk(error) => ServiceError::Sdk(error),
3182        error => ServiceError::Operation(error.to_string()),
3183    }
3184}
3185
3186fn rpc_error(id: Value, code: i64, message: &str) -> Value {
3187    json!({
3188        "jsonrpc": "2.0",
3189        "id": id,
3190        "error": {"code": code, "message": message},
3191    })
3192}
3193
3194#[cfg(test)]
3195mod tests {
3196    use super::*;
3197    use crate::{HarnessEvent, HarnessId, RuntimeEndpoint, RuntimeHandle, StorageLocator};
3198    use async_trait::async_trait;
3199    use std::io::Write;
3200    use std::path::PathBuf;
3201    use std::time::Instant;
3202
3203    struct EndingRuntime {
3204        handle: RuntimeHandle,
3205        event: Option<HarnessEvent>,
3206    }
3207
3208    #[async_trait]
3209    impl RuntimeConnection for EndingRuntime {
3210        fn handle(&self) -> &RuntimeHandle {
3211            &self.handle
3212        }
3213
3214        async fn send_input(&mut self, _input: RuntimeInput) -> crate::Result<Option<String>> {
3215            unreachable!("ending runtime does not accept input")
3216        }
3217
3218        async fn next_event(&mut self) -> crate::Result<Option<HarnessEvent>> {
3219            Ok(self.event.take())
3220        }
3221
3222        async fn interrupt(&mut self) -> crate::Result<()> {
3223            Ok(())
3224        }
3225
3226        async fn respond(&mut self, _request_id: Value, _response: Value) -> crate::Result<()> {
3227            Ok(())
3228        }
3229
3230        async fn close(&mut self) -> crate::Result<()> {
3231            Ok(())
3232        }
3233    }
3234
3235    fn ending_runtime(event: Option<HarnessEvent>) -> Box<dyn RuntimeConnection> {
3236        Box::new(EndingRuntime {
3237            handle: RuntimeHandle {
3238                harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3239                runtime_id: "ending-session".into(),
3240                endpoint: RuntimeEndpoint::LocalProcess {
3241                    pid: None,
3242                    command: vec!["ending-runtime".into()],
3243                    protocol: "test".into(),
3244                },
3245            },
3246            event,
3247        })
3248    }
3249
3250    fn request(id: u64, method: &str, params: Value) -> Value {
3251        json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})
3252    }
3253
3254    fn pi_locator() -> SessionLocator {
3255        SessionLocator {
3256            harness: HarnessId::from(HarnessId::PI),
3257            session_id: "1e6f2a3b-0000-4000-8000-000000000001".into(),
3258            storage: StorageLocator::File {
3259                path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3260                    .join("tests/fixtures/pi_session.jsonl"),
3261            },
3262        }
3263    }
3264
3265    fn opencode_locator() -> SessionLocator {
3266        let session_id = "ses_fixtureAAAAAAAAAAAAAAA1";
3267        SessionLocator {
3268            harness: HarnessId::from(HarnessId::OPENCODE),
3269            session_id: session_id.into(),
3270            storage: StorageLocator::Sqlite {
3271                path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3272                    .join("tests/fixtures/opencode_fixture/opencode.db"),
3273                selector: session_id.into(),
3274            },
3275        }
3276    }
3277
3278    fn grok_locator() -> SessionLocator {
3279        SessionLocator {
3280            harness: HarnessId::from(HarnessId::GROK),
3281            session_id: "73c09283-4b33-41fa-90f1-0bcb0f7be523".into(),
3282            storage: StorageLocator::File {
3283                path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3284                    .join("tests/fixtures/grok_session/chat_history.jsonl"),
3285            },
3286        }
3287    }
3288
3289    #[test]
3290    fn capabilities_are_explicit_and_versioned() {
3291        let mut service = HarnessSessionService::new();
3292        let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
3293        assert_eq!(response["result"]["version"], HARNESS_SERVICE_VERSION);
3294        assert_eq!(
3295            response["result"]["sdk"]["schema_version"],
3296            crate::SDK_SCHEMA_VERSION
3297        );
3298        assert_eq!(
3299            response["result"]["sdk"]["operations"]
3300                .as_array()
3301                .unwrap()
3302                .len(),
3303            SdkOperation::ALL.len()
3304        );
3305        assert_eq!(response["result"]["harnesses"].as_array().unwrap().len(), 8);
3306        assert!(response["result"]["harnesses"]
3307            .as_array()
3308            .unwrap()
3309            .iter()
3310            .any(|harness| harness == HarnessId::GROK));
3311        assert!(response["result"]["harnesses"]
3312            .as_array()
3313            .unwrap()
3314            .iter()
3315            .any(|harness| harness == HarnessId::GOOSE));
3316    }
3317
3318    #[test]
3319    fn handshake_health_uses_protocol_liveness_not_stderr_severity() {
3320        let noisy_stderr = crate::HarnessEvent {
3321            sequence: None,
3322            kind: "transport_stderr".into(),
3323            payload: json!({"line": "ERROR optional worker AuthorizationRequired"}),
3324        };
3325        assert_eq!(handshake_event_failure(&noisy_stderr), None);
3326
3327        let closed = crate::HarnessEvent {
3328            sequence: None,
3329            kind: "transport_closed".into(),
3330            payload: json!({}),
3331        };
3332        assert!(handshake_event_failure(&closed).is_some());
3333    }
3334
3335    #[tokio::test]
3336    async fn runtime_eof_is_notified_and_removed_for_raw_and_explicit_close() {
3337        let mut service = HarnessSessionService::new();
3338        service
3339            .runtimes
3340            .insert("raw-eof".into(), ending_runtime(None));
3341        service.runtimes.insert(
3342            "explicit-close".into(),
3343            ending_runtime(Some(HarnessEvent {
3344                sequence: None,
3345                kind: "transport_closed".into(),
3346                payload: json!({"message": "native transport exited"}),
3347            })),
3348        );
3349
3350        let notifications = service.poll_runtimes().await;
3351
3352        assert_eq!(notifications.len(), 2);
3353        assert!(notifications
3354            .iter()
3355            .all(|notification| { notification["params"]["event"]["kind"] == "transport_closed" }));
3356        assert!(notifications.iter().all(|notification| {
3357            notification["params"]["session_id"] == "ending-session"
3358                && notification["params"]["connection"].is_string()
3359        }));
3360        let mut sequences = notifications
3361            .iter()
3362            .filter_map(|notification| notification["params"]["sequence"].as_u64())
3363            .collect::<Vec<_>>();
3364        sequences.sort_unstable();
3365        assert_eq!(sequences, vec![1, 2]);
3366        assert!(service.runtimes.is_empty());
3367    }
3368
3369    #[tokio::test]
3370    async fn sdk_facade_returns_named_unsupported_actions() {
3371        let mut service = HarnessSessionService::new();
3372        let error = service
3373            .execute(SdkRequest {
3374                operation: SdkOperation::Steer,
3375                params: json!({"connection": "runtime-1", "text": "go left"}),
3376            })
3377            .await
3378            .unwrap_err();
3379        assert_eq!(error.code(), SdkErrorCode::UnsupportedAction);
3380        assert_eq!(error.operation(), Some(SdkOperation::Steer));
3381
3382        let response = service
3383            .handle_async(request(
3384                7,
3385                "harness.v1.runtimes.steer",
3386                json!({"connection": "runtime-1", "text": "go left"}),
3387            ))
3388            .await;
3389        assert_eq!(response["error"]["name"], "unsupported_action");
3390        assert_eq!(response["error"]["operation"], "steer");
3391    }
3392
3393    #[test]
3394    fn support_report_and_grok_default_binding_share_the_registry() {
3395        let mut service = HarnessSessionService::new();
3396        let response = service.handle(request(1, "harness.v1.support.report", json!({})));
3397        assert_eq!(response["result"]["schema"], crate::SUPPORT_REGISTRY_SCHEMA);
3398        let params = RuntimeBackendParams {
3399            harness: HarnessId::from(HarnessId::GROK),
3400            protocol: None,
3401            launch: None,
3402            base_url: None,
3403            policy: RuntimePolicy::Default,
3404        };
3405        let backend = match runtime_backend(&params) {
3406            Ok(backend) => backend,
3407            Err(_) => panic!("Grok should bind through its registered ACP launch"),
3408        };
3409        assert_eq!(backend.harness().as_str(), HarnessId::GROK);
3410        assert!(backend.capabilities().start_session);
3411        let registered = harness_support_registry()
3412            .harnesses
3413            .into_iter()
3414            .find(|harness| harness.id.as_str() == HarnessId::GROK)
3415            .and_then(|harness| harness.runtime.default_launch)
3416            .unwrap();
3417        assert!(!registered
3418            .arguments
3419            .iter()
3420            .any(|argument| argument == "--always-approve"));
3421        assert!(runtime_launch(&params).is_none());
3422
3423        let yolo = RuntimeBackendParams {
3424            policy: RuntimePolicy::Yolo,
3425            ..params
3426        };
3427        assert!(runtime_launch(&yolo)
3428            .unwrap()
3429            .arguments
3430            .iter()
3431            .any(|argument| argument == "--always-approve"));
3432
3433        let mismatched_protocol = RuntimeBackendParams {
3434            harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3435            protocol: Some("acp".into()),
3436            launch: None,
3437            base_url: None,
3438            policy: RuntimePolicy::Default,
3439        };
3440        assert!(runtime_backend(&mismatched_protocol).is_err());
3441    }
3442
3443    #[test]
3444    fn load_follow_and_unfollow_share_the_same_locator() {
3445        let mut service = HarnessSessionService::new();
3446        let locator = pi_locator();
3447        let loaded = service.handle(request(
3448            1,
3449            "harness.v1.sessions.load",
3450            json!({"locator": locator}),
3451        ));
3452        assert_eq!(
3453            loaded["result"]["session"]["session_id"],
3454            locator.session_id
3455        );
3456
3457        let followed = service.handle(request(
3458            2,
3459            "harness.v1.sessions.follow",
3460            json!({"locator": locator}),
3461        ));
3462        assert_eq!(followed["result"]["subscription"], "sub-1");
3463        assert_eq!(followed["result"]["initial"]["type"], "session_snapshot");
3464        assert!(service.poll().is_empty());
3465
3466        let unfollowed = service.handle(request(
3467            3,
3468            "harness.v1.sessions.unfollow",
3469            json!({"subscription": "sub-1"}),
3470        ));
3471        assert_eq!(unfollowed["result"]["removed"], true);
3472    }
3473
3474    #[test]
3475    fn bounded_read_view_excludes_subagents_and_keeps_only_the_tail() {
3476        let temp = std::env::temp_dir().join(format!(
3477            "supercode-bounded-view-{}-{}",
3478            std::process::id(),
3479            generated_session_id()
3480        ));
3481        let path = temp.join("parent.jsonl");
3482        let subagents = temp.join("parent/subagents");
3483        std::fs::create_dir_all(&subagents).unwrap();
3484        let long_last = "x".repeat(300);
3485        let parent_records = [
3486            json!({"type":"user","uuid":"u1","parentUuid":null,"message":{"role":"user","content":"first"}}),
3487            json!({"type":"assistant","uuid":"a1","parentUuid":"u1","message":{"role":"assistant","content":[{"type":"text","text":"middle"}]}}),
3488            json!({"type":"user","uuid":"u2","parentUuid":"a1","message":{"role":"user","content":long_last}}),
3489        ];
3490        std::fs::write(
3491            &path,
3492            format!(
3493                "{}\n",
3494                parent_records
3495                    .iter()
3496                    .map(Value::to_string)
3497                    .collect::<Vec<_>>()
3498                    .join("\n")
3499            ),
3500        )
3501        .unwrap();
3502        std::fs::write(
3503            subagents.join("agent-child.jsonl"),
3504            concat!(
3505                r#"{"type":"user","uuid":"cu","parentUuid":null,"agentId":"child","message":{"role":"user","content":"child work"}}"#,
3506                "\n",
3507            ),
3508        )
3509        .unwrap();
3510        let locator = SessionLocator {
3511            harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3512            session_id: "parent".into(),
3513            storage: StorageLocator::File { path },
3514        };
3515        let mut service = HarnessSessionService::new();
3516
3517        let complete = service.handle(request(
3518            1,
3519            "harness.v1.sessions.load",
3520            json!({"locator": locator}),
3521        ));
3522        assert_eq!(
3523            complete["result"]["session"]["subagents"]
3524                .as_array()
3525                .unwrap()
3526                .len(),
3527            1
3528        );
3529
3530        let bounded = service.handle(request(
3531            2,
3532            "harness.v1.sessions.load",
3533            json!({
3534                "locator": locator,
3535                "view": {
3536                    "tail_messages": 1,
3537                    "max_message_chars": 256,
3538                    "include_subagents": false
3539                },
3540            }),
3541        ));
3542        let session = &bounded["result"]["session"];
3543        assert!(session["subagents"].as_array().unwrap().is_empty());
3544        assert_eq!(session["messages"].as_array().unwrap().len(), 1);
3545        assert_eq!(
3546            session["messages"][0]["content"],
3547            format!("{}\n…", "x".repeat(256))
3548        );
3549
3550        let followed = service.handle(request(
3551            3,
3552            "harness.v1.sessions.follow",
3553            json!({
3554                "locator": locator,
3555                "view": {
3556                    "tail_messages": 1,
3557                    "max_message_chars": 256,
3558                    "include_subagents": false
3559                },
3560            }),
3561        ));
3562        let initial = &followed["result"]["initial"]["session"];
3563        assert!(initial["subagents"].as_array().unwrap().is_empty());
3564        assert_eq!(initial["messages"].as_array().unwrap().len(), 1);
3565
3566        let _ = std::fs::remove_dir_all(&temp);
3567    }
3568
3569    #[test]
3570    fn forty_megabyte_display_load_is_bounded_and_prompt() {
3571        let temp = std::env::temp_dir().join(format!(
3572            "supercode-large-display-view-{}-{}",
3573            std::process::id(),
3574            generated_session_id()
3575        ));
3576        std::fs::create_dir_all(&temp).unwrap();
3577        let path = temp.join("rollout.jsonl");
3578        let mut file = std::io::BufWriter::new(std::fs::File::create(&path).unwrap());
3579        writeln!(
3580            file,
3581            r#"{{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{{"id":"large-display","cwd":"/tmp"}}}}"#
3582        )
3583        .unwrap();
3584        let padding = "x".repeat(80 * 1024);
3585        for index in 0..512 {
3586            let marker = if index == 0 {
3587                "OLDEST-SHOULD-NOT-LOAD"
3588            } else if index == 511 {
3589                "LATEST-MUST-LOAD"
3590            } else {
3591                "bulk"
3592            };
3593            writeln!(
3594                file,
3595                "{}",
3596                json!({
3597                    "timestamp": "2026-01-01T00:00:01Z",
3598                    "type": "response_item",
3599                    "payload": {
3600                        "type": "message",
3601                        "role": "assistant",
3602                        "content": [{"type": "output_text", "text": format!("{marker}:{padding}")}],
3603                    },
3604                })
3605            )
3606            .unwrap();
3607        }
3608        file.flush().unwrap();
3609        drop(file);
3610        assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
3611
3612        let locator = SessionLocator {
3613            harness: HarnessId::from(HarnessId::CODEX),
3614            session_id: "large-display".into(),
3615            storage: StorageLocator::File { path },
3616        };
3617        let started = Instant::now();
3618        let response = HarnessSessionService::new().handle(request(
3619            1,
3620            "harness.v1.sessions.load",
3621            json!({
3622                "locator": locator,
3623                "view": {
3624                    "tail_messages": 500,
3625                    "max_message_chars": 1024,
3626                    "include_subagents": false,
3627                    "display_history": true,
3628                },
3629            }),
3630        ));
3631        let elapsed = started.elapsed();
3632        let wire = response.to_string();
3633        eprintln!(
3634            "bounded 40 MiB display load: {elapsed:?}, {} response bytes",
3635            wire.len()
3636        );
3637        assert!(response.get("error").is_none(), "{response:#}");
3638        assert!(wire.contains("LATEST-MUST-LOAD"));
3639        assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
3640        assert!(
3641            wire.len() < 2 * 1024 * 1024,
3642            "bounded wire was {} bytes",
3643            wire.len()
3644        );
3645        assert!(
3646            elapsed.as_secs_f64() < 3.0,
3647            "bounded 40 MiB load took {elapsed:?}"
3648        );
3649
3650        let _ = std::fs::remove_dir_all(&temp);
3651    }
3652
3653    #[test]
3654    fn forty_megabyte_goose_store_display_load_reads_only_the_tail() {
3655        let temp = std::env::temp_dir().join(format!(
3656            "supercode-large-goose-view-{}-{}",
3657            std::process::id(),
3658            generated_session_id()
3659        ));
3660        std::fs::create_dir_all(&temp).unwrap();
3661        let path = temp.join("sessions.db");
3662        let connection = rusqlite::Connection::open(&path).unwrap();
3663        connection
3664            .execute_batch(
3665                "CREATE TABLE sessions (
3666                    id TEXT PRIMARY KEY, name TEXT NOT NULL, working_dir TEXT NOT NULL,
3667                    created_at TEXT NOT NULL, updated_at TEXT NOT NULL,
3668                    session_type TEXT NOT NULL, extension_data TEXT,
3669                    goose_mode TEXT NOT NULL, provider_name TEXT, model_config_json TEXT,
3670                    archived_at TEXT
3671                 );
3672                 CREATE TABLE messages (
3673                    id INTEGER PRIMARY KEY, session_id TEXT NOT NULL, message_id TEXT,
3674                    role TEXT NOT NULL, content_json TEXT NOT NULL,
3675                    created_timestamp INTEGER NOT NULL, metadata_json TEXT
3676                 );",
3677            )
3678            .unwrap();
3679        connection
3680            .execute(
3681                "INSERT INTO sessions VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL)",
3682                rusqlite::params![
3683                    "goose-large",
3684                    "Large Goose session",
3685                    "/tmp",
3686                    "2026-01-01 00:00:00",
3687                    "2026-01-01 00:00:02",
3688                    "user",
3689                    "{}",
3690                    "auto",
3691                    "anthropic",
3692                    r#"{"model_name":"claude-sonnet"}"#,
3693                ],
3694            )
3695            .unwrap();
3696        let old_content = serde_json::to_string(&vec![json!({
3697            "type": "text",
3698            "text": format!("OLDEST-SHOULD-NOT-LOAD:{}", "x".repeat(40 * 1024 * 1024)),
3699        })])
3700        .unwrap();
3701        connection
3702            .execute(
3703                "INSERT INTO messages VALUES (1, ?1, 'old', 'user', ?2, 1, '{}')",
3704                rusqlite::params!["goose-large", old_content],
3705            )
3706            .unwrap();
3707        connection
3708            .execute(
3709                "INSERT INTO messages VALUES (2, ?1, 'new', 'assistant', ?2, 2, '{}')",
3710                rusqlite::params![
3711                    "goose-large",
3712                    r#"[{"type":"text","text":"LATEST-MUST-LOAD"}]"#
3713                ],
3714            )
3715            .unwrap();
3716        drop(connection);
3717        assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
3718
3719        let locator = SessionLocator {
3720            harness: HarnessId::from(HarnessId::GOOSE),
3721            session_id: "goose-large".into(),
3722            storage: StorageLocator::Sqlite {
3723                path,
3724                selector: "goose-large".into(),
3725            },
3726        };
3727        let started = Instant::now();
3728        let response = HarnessSessionService::new().handle(request(
3729            1,
3730            "harness.v1.sessions.load",
3731            json!({
3732                "locator": locator,
3733                "view": {
3734                    "tail_messages": 1,
3735                    "max_message_chars": 1024,
3736                    "include_subagents": false,
3737                    "display_history": true,
3738                },
3739            }),
3740        ));
3741        let elapsed = started.elapsed();
3742        let wire = response.to_string();
3743        eprintln!(
3744            "bounded 40 MiB Goose display load: {elapsed:?}, {} response bytes",
3745            wire.len()
3746        );
3747        assert!(response.get("error").is_none(), "{response:#}");
3748        assert!(wire.contains("LATEST-MUST-LOAD"));
3749        assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
3750        assert!(
3751            wire.len() < 64 * 1024,
3752            "bounded wire was {} bytes",
3753            wire.len()
3754        );
3755        assert!(
3756            elapsed.as_secs_f64() < 1.0,
3757            "bounded Goose load took {elapsed:?}"
3758        );
3759
3760        let _ = std::fs::remove_dir_all(&temp);
3761    }
3762
3763    #[test]
3764    fn display_view_keeps_codex_assistant_history_across_compaction() {
3765        let temp = std::env::temp_dir().join(format!(
3766            "supercode-codex-display-view-{}-{}",
3767            std::process::id(),
3768            generated_session_id()
3769        ));
3770        std::fs::create_dir_all(&temp).unwrap();
3771        let path = temp.join("rollout.jsonl");
3772        std::fs::write(
3773            &path,
3774            concat!(
3775                r#"{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{"id":"codex-display","cwd":"/tmp"}}"#,
3776                "\n",
3777                r#"{"timestamp":"2026-01-01T00:00:01Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"old prompt"}]}}"#,
3778                "\n",
3779                r#"{"timestamp":"2026-01-01T00:00:02Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"old answer"}]}}"#,
3780                "\n",
3781                r#"{"timestamp":"2026-01-01T00:00:03Z","type":"compacted","payload":{"replacement_history":[{"type":"message","role":"user","content":[{"type":"input_text","text":"old prompt"}]},{"type":"compaction","encrypted_content":"opaque"}]}}"#,
3782                "\n",
3783                r#"{"timestamp":"2026-01-01T00:00:04Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"new prompt"}]}}"#,
3784                "\n",
3785                r#"{"timestamp":"2026-01-01T00:00:05Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"new answer"}]}}"#,
3786                "\n",
3787            ),
3788        )
3789        .unwrap();
3790        let locator = SessionLocator {
3791            harness: HarnessId::from(HarnessId::CODEX),
3792            session_id: "codex-display".into(),
3793            storage: StorageLocator::File { path },
3794        };
3795        let mut service = HarnessSessionService::new();
3796
3797        let continuation = service.handle(request(
3798            1,
3799            "harness.v1.sessions.load",
3800            json!({"locator": locator}),
3801        ));
3802        let continuation_text = continuation["result"]["session"]["messages"].to_string();
3803        assert!(!continuation_text.contains("old answer"));
3804
3805        let display = service.handle(request(
3806            2,
3807            "harness.v1.sessions.load",
3808            json!({
3809                "locator": locator,
3810                "view": {
3811                    "tail_messages": 10,
3812                    "include_subagents": false,
3813                    "display_history": true,
3814                },
3815            }),
3816        ));
3817        let display_text = display["result"]["session"]["messages"].to_string();
3818        assert!(display_text.contains("old prompt"));
3819        assert!(display_text.contains("old answer"));
3820        assert!(display_text.contains("new prompt"));
3821        assert!(display_text.contains("new answer"));
3822
3823        let _ = std::fs::remove_dir_all(&temp);
3824    }
3825
3826    #[test]
3827    fn load_supports_bounded_windows_and_media_metadata() {
3828        let mut service = HarnessSessionService::new();
3829        let locator = pi_locator();
3830        let bounded = service.handle(request(
3831            1,
3832            "harness.v1.sessions.load",
3833            json!({
3834                "locator": locator,
3835                "options": {
3836                    "include_subagents": false,
3837                    "message_limit": 2,
3838                    "message_offset": 1
3839                }
3840            }),
3841        ));
3842        assert_eq!(bounded["result"]["window"]["offset"], 1);
3843        assert_eq!(bounded["result"]["window"]["returned"], 2);
3844        assert!(bounded["result"]["summary"]["first_message"].is_object());
3845        assert!(bounded["result"]["summary"]["last_message"].is_object());
3846        assert_eq!(
3847            bounded["result"]["session"]["messages"]
3848                .as_array()
3849                .unwrap()
3850                .len(),
3851            2
3852        );
3853        assert!(bounded["result"]["session"]["subagents"]
3854            .as_array()
3855            .unwrap()
3856            .is_empty());
3857
3858        let tail = service.handle(request(
3859            2,
3860            "harness.v1.sessions.load",
3861            json!({"locator": locator, "options": {"message_tail": 1}}),
3862        ));
3863        assert_eq!(tail["result"]["window"]["returned"], 1);
3864        assert_eq!(tail["result"]["window"]["has_more"], true);
3865        assert_eq!(tail["result"]["window"]["has_older"], true);
3866        assert!(tail["result"]["window"]["older_items"].as_u64().unwrap() > 0);
3867        assert!(tail["result"]["summary"]["first_message"].is_object());
3868
3869        let metadata_only = service.handle(request(
3870            3,
3871            "harness.v1.sessions.load",
3872            json!({"locator": locator, "options": {"inline_media": "metadata"}}),
3873        ));
3874        assert!(metadata_only["result"]["session"]
3875            .to_string()
3876            .contains("media_reference"));
3877        assert!(!metadata_only["result"]["session"]
3878            .to_string()
3879            .contains("data:image/"));
3880    }
3881
3882    #[test]
3883    fn import_translate_branch_and_handoff_use_typed_artifacts() {
3884        let mut service = HarnessSessionService::new();
3885        let locator = pi_locator();
3886        let translated = service.handle(request(
3887            1,
3888            "harness.v1.sessions.translate",
3889            json!({"locator": locator, "target_harness": "grok"}),
3890        ));
3891        assert_eq!(translated["result"]["artifact"]["source_harness"], "pi");
3892        assert_eq!(translated["result"]["artifact"]["target_harness"], "grok");
3893        assert!(translated["result"]["artifact"]["content"]
3894            .as_str()
3895            .is_some_and(|content| !content.is_empty()));
3896
3897        for target in ["opencode", "open-code"] {
3898            let opencode = service.handle(request(
3899                6,
3900                "harness.v1.sessions.translate",
3901                json!({"locator": locator, "target_harness": target}),
3902            ));
3903            assert_eq!(opencode["result"]["artifact"]["target_harness"], "opencode");
3904        }
3905        let goose = service.handle(request(
3906            7,
3907            "harness.v1.sessions.translate",
3908            json!({"locator": locator, "target_harness": "goose"}),
3909        ));
3910        assert_eq!(goose["result"]["artifact"]["target_harness"], "goose");
3911        assert!(serde_json::from_str::<Value>(
3912            goose["result"]["artifact"]["content"].as_str().unwrap()
3913        )
3914        .unwrap()["conversation"]
3915            .is_array());
3916
3917        let imported = service.handle(request(
3918            2,
3919            "harness.v1.sessions.import",
3920            json!({
3921                "source_harness": "grok",
3922                "content": translated["result"]["artifact"]["content"],
3923            }),
3924        ));
3925        assert_eq!(imported["result"]["session"]["source"], "grok");
3926
3927        let branched = service.handle(request(
3928            3,
3929            "harness.v1.sessions.branch",
3930            json!({"locator": locator, "target_harness": "codex"}),
3931        ));
3932        assert_eq!(branched["result"]["parent"]["harness"], "pi");
3933        assert!(branched["result"]["bootstrap_prompt"]
3934            .as_str()
3935            .unwrap()
3936            .contains("frozen parent transcript"));
3937        assert_eq!(branched["result"]["artifact"]["target_harness"], "codex");
3938
3939        let handoff = service.handle(request(
3940            4,
3941            "harness.v1.sessions.handoff",
3942            json!({"locator": locator, "target_harness": "pi", "cwd": "/tmp/project"}),
3943        ));
3944        assert_eq!(handoff["result"]["launch"]["program"], "pi");
3945        assert_eq!(handoff["result"]["launch"]["cwd"], "/tmp/project");
3946        assert_eq!(handoff["result"]["requires_materialization"], true);
3947
3948        let goose_handoff = service.handle(request(
3949            8,
3950            "harness.v1.sessions.handoff",
3951            json!({"locator": locator, "target_harness": "goose", "cwd": "/tmp/project"}),
3952        ));
3953        assert_eq!(goose_handoff["result"]["launch"]["program"], "goose");
3954        assert_eq!(
3955            goose_handoff["result"]["materialize"]["arguments"],
3956            json!(["session", "import", "{artifact_path}"])
3957        );
3958
3959        let resumed = service.handle(request(
3960            5,
3961            "harness.v1.sessions.resume_instructions",
3962            json!({"locator": locator, "cwd": "/tmp/project", "policy": "yolo"}),
3963        ));
3964        assert_eq!(resumed["result"]["launch"]["program"], "pi");
3965        assert_eq!(resumed["result"]["launch"]["arguments"][0], "--approve");
3966    }
3967
3968    #[test]
3969    fn reduce_persists_and_reloads_a_byte_exact_reversible_bundle() {
3970        let temp = std::env::temp_dir().join(format!(
3971            "supercode-service-reduce-{}-{}",
3972            std::process::id(),
3973            generated_session_id()
3974        ));
3975        let source_path = temp.join("source.jsonl");
3976        let store_root = temp.join("store");
3977        std::fs::create_dir_all(&temp).unwrap();
3978
3979        let mut records = vec![json!({
3980            "timestamp": "2026-01-01T00:00:00Z",
3981            "type": "session_meta",
3982            "payload": {"id": "codex-reduce", "cwd": "/tmp/project"},
3983        })];
3984        for turn in 0..16 {
3985            records.push(json!({
3986                "timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 1),
3987                "type": "response_item",
3988                "payload": {
3989                    "type": "message",
3990                    "role": "user",
3991                    "content": [{
3992                        "type": "input_text",
3993                        "text": format!("request {turn}: {}", "context ".repeat(80)),
3994                    }],
3995                },
3996            }));
3997            records.push(json!({
3998                "timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 2),
3999                "type": "response_item",
4000                "payload": {
4001                    "type": "message",
4002                    "role": "assistant",
4003                    "content": [{
4004                        "type": "output_text",
4005                        "text": format!("answer {turn}: {}", "implementation detail ".repeat(80)),
4006                    }],
4007                },
4008            }));
4009        }
4010        let source = format!(
4011            "{}\n",
4012            records
4013                .iter()
4014                .map(Value::to_string)
4015                .collect::<Vec<_>>()
4016                .join("\n")
4017        );
4018        std::fs::write(&source_path, &source).unwrap();
4019        let locator = SessionLocator {
4020            harness: HarnessId::from(HarnessId::CODEX),
4021            session_id: "codex-reduce".into(),
4022            storage: StorageLocator::File {
4023                path: source_path.clone(),
4024            },
4025        };
4026        let original = load_session(&locator).unwrap();
4027        let mut service =
4028            HarnessSessionService::new().with_reduction_store_root(store_root.clone());
4029
4030        let response = service.handle(request(
4031            1,
4032            "harness.v1.sessions.reduce",
4033            json!({
4034                "locator": locator,
4035                "target_harness": "claude-code",
4036                "keep_last": 4,
4037            }),
4038        ));
4039        assert!(response.get("error").is_none(), "{response:#}");
4040        let receipt = &response["result"]["receipt"];
4041        assert_eq!(receipt["source_harness"], "codex");
4042        assert_eq!(receipt["target_harness"], "claude-code");
4043        assert_eq!(receipt["verified"], true);
4044        assert_eq!(receipt["reversible"], true);
4045        assert!(receipt["reductions"].as_u64().unwrap() > 0);
4046        assert!(
4047            receipt["source_tokens"].as_u64().unwrap()
4048                > receipt["reduced_tokens"].as_u64().unwrap()
4049        );
4050        assert!(receipt["ratio"].as_f64().unwrap() > 1.0);
4051        assert!(response["result"]["bootstrap_prompt"]
4052            .as_str()
4053            .unwrap()
4054            .contains("Do not guess hidden content"));
4055
4056        let rescue_id = receipt["id"].as_str().unwrap();
4057        let store = crate::SessionStore::open(&store_root).unwrap();
4058        let sidecar =
4059            Session::from_sidecar_str(&store.load_sidecar(rescue_id).unwrap().unwrap()).unwrap();
4060        let log = store.load_reduction_log(rescue_id).unwrap().unwrap();
4061        let persisted_view = parse_messages_jsonl(&store.load(rescue_id).unwrap()).unwrap();
4062        let policy = reduce::ReductionPolicy {
4063            clear_turns_older_than: Some(4),
4064            ..Default::default()
4065        };
4066        let (restamped_view, reapplied_log) =
4067            reduce::project_messages(&sidecar.messages, &policy, &log);
4068        assert_eq!(
4069            messages_jsonl(&persisted_view).unwrap(),
4070            messages_jsonl(&restamped_view).unwrap()
4071        );
4072        assert_eq!(reapplied_log, log);
4073        reduce::verify_log(&log, &sidecar).unwrap();
4074        assert_eq!(
4075            reduce::invert(&restamped_view, &log, &sidecar).unwrap(),
4076            original.messages
4077        );
4078        assert_eq!(std::fs::read_to_string(&source_path).unwrap(), source);
4079
4080        std::fs::remove_dir_all(temp).ok();
4081    }
4082
4083    #[test]
4084    fn read_surfaces_view_a_severed_claude_graph_while_transfer_still_refuses_it() {
4085        let temp = std::env::temp_dir().join(format!(
4086            "supercode-severed-view-{}-{}",
4087            std::process::id(),
4088            generated_session_id()
4089        ));
4090        std::fs::create_dir_all(&temp).unwrap();
4091        let path = temp.join("severed.jsonl");
4092        // A live record whose parent was pruned — what a compacted or
4093        // resumed-across-files Claude Code session looks like on disk.
4094        std::fs::write(
4095            &path,
4096            concat!(
4097                r#"{"type":"user","uuid":"orphan-u","parentUuid":null,"message":{"role":"user","content":"stranded prompt"}}"#,
4098                "\n",
4099                r#"{"type":"assistant","uuid":"live-a","parentUuid":"pruned","message":{"id":"m","role":"assistant","content":[{"type":"text","text":"live answer"}]}}"#,
4100                "\n",
4101            ),
4102        )
4103        .unwrap();
4104        let locator = SessionLocator {
4105            harness: HarnessId::from(HarnessId::CLAUDE_CODE),
4106            session_id: "severed".into(),
4107            storage: StorageLocator::File { path },
4108        };
4109        let mut service = HarnessSessionService::new();
4110
4111        let viewed = service.handle(request(
4112            1,
4113            "harness.v1.sessions.load",
4114            json!({"locator": locator}),
4115        ));
4116        let session = &viewed["result"]["session"];
4117        assert_eq!(session["fidelity"], "semantic");
4118        assert_eq!(session["messages"].as_array().unwrap().len(), 2);
4119        assert!(session["residue"].as_array().unwrap().iter().any(|entry| {
4120            entry
4121                .as_str()
4122                .is_some_and(|entry| entry.contains("live-a") && entry.contains("pruned"))
4123        }));
4124
4125        // Asking a READ surface for a lossless reconstruction gets the strict
4126        // refusal back, unchanged.
4127        let strict = service.handle(request(
4128            2,
4129            "harness.v1.sessions.load",
4130            json!({"locator": locator, "fidelity": "byte_lossless"}),
4131        ));
4132        assert!(strict["error"]["message"]
4133            .as_str()
4134            .unwrap()
4135            .contains("cannot reconstruct lossless Claude continuation"));
4136
4137        // Transfer/continuation surfaces have no view mode at all.
4138        let translated = service.handle(request(
4139            3,
4140            "harness.v1.sessions.translate",
4141            json!({"locator": locator, "target_harness": "codex"}),
4142        ));
4143        assert!(translated["error"]["message"]
4144            .as_str()
4145            .unwrap()
4146            .contains("cannot reconstruct lossless Claude continuation"));
4147        let resumed = service.handle(request(
4148            4,
4149            "harness.v1.sessions.resume_instructions",
4150            json!({"locator": locator}),
4151        ));
4152        assert!(resumed["error"]["message"]
4153            .as_str()
4154            .unwrap()
4155            .contains("cannot reconstruct lossless Claude continuation"));
4156
4157        let _ = std::fs::remove_dir_all(&temp);
4158    }
4159
4160    #[test]
4161    fn structured_resume_launches_cover_gemini_goose_and_supercode() {
4162        let gemini = resume_launch(
4163            HarnessId::GEMINI,
4164            "gemini-session",
4165            Path::new("/tmp/project"),
4166            ResumePolicy::Yolo,
4167        )
4168        .unwrap_or_else(|_| panic!("Gemini resume launch must be registered"));
4169        assert_eq!(gemini.program, "gemini");
4170        assert_eq!(gemini.arguments, ["--yolo", "--resume", "gemini-session"]);
4171
4172        let goose = resume_launch(
4173            HarnessId::GOOSE,
4174            "goose-session",
4175            Path::new("/tmp/project"),
4176            ResumePolicy::Yolo,
4177        )
4178        .unwrap_or_else(|_| panic!("Goose resume launch must be registered"));
4179        assert_eq!(goose.program, "goose");
4180        assert_eq!(
4181            goose.arguments,
4182            ["session", "--resume", "--session-id", "goose-session"]
4183        );
4184
4185        let supercode = resume_launch(
4186            HarnessId::SUPERCODE,
4187            "supercode-session",
4188            Path::new("/tmp/project"),
4189            ResumePolicy::Yolo,
4190        )
4191        .unwrap_or_else(|_| panic!("Supercode resume launch must be registered"));
4192        assert_eq!(supercode.program, "supercode");
4193        assert_eq!(
4194            supercode.arguments,
4195            ["--dangerous", "resume", "supercode-session"]
4196        );
4197    }
4198
4199    #[test]
4200    fn diagonal_artifacts_preserve_claude_subagents_and_grok_bundle_members() {
4201        let temp = std::env::temp_dir().join(format!(
4202            "supercode-harness-artifact-{}-{}",
4203            std::process::id(),
4204            generated_session_id()
4205        ));
4206        let main_path = temp.join("parent.jsonl");
4207        let subagent_path = temp.join("parent/subagents/agent-child.jsonl");
4208        std::fs::create_dir_all(subagent_path.parent().unwrap()).unwrap();
4209        let fixture = std::fs::read_to_string(
4210            PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4211                .join("tests/fixtures/claude_code_session.jsonl"),
4212        )
4213        .unwrap();
4214        let parent = fixture.trim_end_matches('\n');
4215        let child = fixture.trim_end_matches('\n');
4216        std::fs::write(&main_path, parent).unwrap();
4217        std::fs::write(&subagent_path, child).unwrap();
4218        let locator = SessionLocator {
4219            harness: HarnessId::from(HarnessId::CLAUDE_CODE),
4220            session_id: "213bb148-51ea-453f-9206-f8b4b1168547".into(),
4221            storage: StorageLocator::File {
4222                path: main_path.clone(),
4223            },
4224        };
4225        let mut service = HarnessSessionService::new();
4226        let claude = service.handle(request(
4227            1,
4228            "harness.v1.sessions.translate",
4229            json!({"locator": locator, "target_harness": "claude-code"}),
4230        ));
4231        let artifact = &claude["result"]["artifact"];
4232        assert_eq!(artifact["fidelity"], "byte_lossless");
4233        assert_eq!(artifact["content"], parent);
4234        let files = artifact["files"].as_array().unwrap();
4235        assert!(files.iter().any(|file| {
4236            file["role"] == "subagent"
4237                && file["path"]
4238                    .as_str()
4239                    .is_some_and(|path| path.ends_with("/subagents/agent-child.jsonl"))
4240                && file["content"] == child
4241        }));
4242        assert!(!artifact["content"].as_str().unwrap().ends_with('\n'));
4243
4244        let grok = service.handle(request(
4245            2,
4246            "harness.v1.sessions.translate",
4247            json!({"locator": grok_locator(), "target_harness": "grok"}),
4248        ));
4249        let files = grok["result"]["artifact"]["files"].as_array().unwrap();
4250        for name in ["summary.json", "updates.jsonl"] {
4251            let expected = std::fs::read_to_string(
4252                PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4253                    .join("tests/fixtures/grok_session")
4254                    .join(name),
4255            )
4256            .unwrap();
4257            assert!(files.iter().any(|file| {
4258                file["path"] == name && file["role"] == "bundle" && file["content"] == expected
4259            }));
4260        }
4261        std::fs::remove_dir_all(temp).ok();
4262    }
4263
4264    #[test]
4265    fn every_non_grok_handoff_mints_and_uses_a_fresh_target_identity() {
4266        let mut service = HarnessSessionService::new();
4267        let source = pi_locator();
4268        for (target, format) in [
4269            ("claude-code", SessionFormat::ClaudeCode),
4270            ("codex", SessionFormat::Codex),
4271            ("opencode", SessionFormat::OpenCode),
4272            ("pi", SessionFormat::Pi),
4273        ] {
4274            let result = service.handle(request(
4275                1,
4276                "harness.v1.sessions.handoff",
4277                json!({"locator": source, "target_harness": target, "cwd": "/tmp/project"}),
4278            ));
4279            let artifact = &result["result"]["artifact"];
4280            let target_id = artifact["session_id"].as_str().unwrap();
4281            assert_ne!(target_id, source.session_id, "{target}");
4282            let parsed = Session::load_str(artifact["content"].as_str().unwrap(), format).unwrap();
4283            assert_eq!(
4284                parsed.meta.session_id.as_deref(),
4285                Some(target_id),
4286                "{target}"
4287            );
4288            if target != "pi" {
4289                assert!(result["result"]["launch"]["arguments"]
4290                    .as_array()
4291                    .unwrap()
4292                    .iter()
4293                    .any(|argument| argument == target_id));
4294            }
4295            if target == "opencode" {
4296                assert!(target_id.starts_with("ses_"));
4297                fn assert_session_ids(value: &Value, target_id: &str) {
4298                    match value {
4299                        Value::Object(fields) => {
4300                            if let Some(session_id) = fields.get("sessionID") {
4301                                assert_eq!(session_id, target_id);
4302                            }
4303                            for child in fields.values() {
4304                                assert_session_ids(child, target_id);
4305                            }
4306                        }
4307                        Value::Array(values) => {
4308                            for child in values {
4309                                assert_session_ids(child, target_id);
4310                            }
4311                        }
4312                        _ => {}
4313                    }
4314                }
4315                let document: Value =
4316                    serde_json::from_str(artifact["content"].as_str().unwrap()).unwrap();
4317                assert_session_ids(&document, target_id);
4318            }
4319        }
4320
4321        let first = service.handle(request(
4322            2,
4323            "harness.v1.sessions.handoff",
4324            json!({"locator": source, "target_harness": "codex"}),
4325        ));
4326        let second = service.handle(request(
4327            3,
4328            "harness.v1.sessions.handoff",
4329            json!({"locator": source, "target_harness": "codex"}),
4330        ));
4331        assert_ne!(
4332            first["result"]["artifact"]["session_id"],
4333            second["result"]["artifact"]["session_id"]
4334        );
4335    }
4336
4337    #[test]
4338    fn grok_handoff_uses_the_official_importer_contract() {
4339        let mut service = HarnessSessionService::new();
4340        let source = opencode_locator();
4341        let response = service.handle(request(
4342            1,
4343            "harness.v1.sessions.handoff",
4344            json!({
4345                "locator": source,
4346                "target_harness": "grok",
4347                "cwd": "/tmp/grok-handoff-project",
4348            }),
4349        ));
4350        let result = &response["result"];
4351
4352        // The target is Grok, but the artifact truthfully names the Claude Code wire
4353        // format accepted by Grok's official importer. Raw Grok chat_history JSONL is
4354        // not a complete stock-resumable bundle.
4355        assert_eq!(result["artifact"]["target_harness"], "claude-code");
4356        assert!(result["artifact"]["suggested_filename"]
4357            .as_str()
4358            .unwrap()
4359            .ends_with(".grok-import.claude-code.jsonl"));
4360        let artifact = Session::load_str(
4361            result["artifact"]["content"].as_str().unwrap(),
4362            SessionFormat::ClaudeCode,
4363        )
4364        .unwrap();
4365        assert_eq!(
4366            artifact.meta.cwd.as_deref(),
4367            Some(Path::new("/tmp/grok-handoff-project"))
4368        );
4369        let target_session_id = artifact.meta.session_id.as_deref().unwrap();
4370        assert_eq!(target_session_id.len(), 36);
4371        assert_eq!(target_session_id.as_bytes()[14], b'4');
4372        assert_ne!(target_session_id, opencode_locator().session_id);
4373        assert_eq!(
4374            result["artifact"]["session_id"],
4375            artifact.meta.session_id.as_deref().unwrap()
4376        );
4377
4378        assert_eq!(
4379            result["materialize"]["arguments"],
4380            json!(["import", "--json", "{artifact_path}"])
4381        );
4382        assert_eq!(
4383            result["launch"]["arguments"],
4384            json!(["--resume", "{imported_session_id}", "--fork-session"])
4385        );
4386        assert!(result["note"]
4387            .as_str()
4388            .unwrap()
4389            .contains("outcome=imported"));
4390        assert!(!result["launch"]["arguments"]
4391            .as_array()
4392            .unwrap()
4393            .iter()
4394            .any(|argument| argument == &opencode_locator().session_id));
4395    }
4396
4397    #[tokio::test]
4398    async fn inventory_rejects_unknown_harnesses_and_runtime_attach_is_honest() {
4399        let mut service = HarnessSessionService::new();
4400        let inventory = service
4401            .handle_async(request(
4402                1,
4403                "harness.v1.harnesses.list",
4404                json!({"harnesses": ["missing"]}),
4405            ))
4406            .await;
4407        assert_eq!(inventory["error"]["code"], -32602);
4408
4409        let attached = service
4410            .handle_async(request(
4411                2,
4412                "harness.v1.runtimes.attach_existing",
4413                json!({"harness": "codex", "runtime_id": "thread-1"}),
4414            ))
4415            .await;
4416        assert_eq!(attached["error"]["code"], -32000);
4417        assert!(attached["error"]["message"]
4418            .as_str()
4419            .unwrap()
4420            .contains("runtimes.resume"));
4421    }
4422
4423    #[test]
4424    fn invalid_params_and_unknown_methods_use_json_rpc_errors() {
4425        let mut service = HarnessSessionService::new();
4426        let invalid = service.handle(request(1, "harness.v1.sessions.load", json!({})));
4427        assert_eq!(invalid["error"]["code"], -32602);
4428        let unknown = service.handle(request(2, "harness.v1.unknown", json!({})));
4429        assert_eq!(unknown["error"]["code"], -32601);
4430    }
4431
4432    #[cfg(unix)]
4433    #[tokio::test]
4434    // The test mutates process-wide harness environment and deliberately
4435    // holds the global test lock until every async runtime operation ends.
4436    #[allow(clippy::await_holding_lock)]
4437    async fn async_service_drives_a_generic_acp_runtime() {
4438        let _environment_guard = crate::live_runtime::test_environment_lock();
4439        let script = r#"
4440            i=0
4441            while IFS= read -r line; do
4442              i=$((i + 1))
4443              case "$i" in
4444                1) printf '%s\n' '{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentCapabilities":{},"authMethods":[]}}' ;;
4445                2) printf '%s\n' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"svc_acp"}}' ;;
4446                3)
4447                  printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"ok"}}}}'
4448                  printf '%s\n' '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
4449                  ;;
4450                4)
4451                  printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"from terminal"}}}}'
4452                  printf '%s\n' '{"jsonrpc":"2.0","id":4,"result":{"stopReason":"end_turn"}}'
4453                  ;;
4454              esac
4455            done
4456        "#;
4457        let mut service = HarnessSessionService::new();
4458        let started = service
4459            .handle_async(request(
4460                1,
4461                "harness.v1.runtimes.start",
4462                json!({
4463                    "harness": "codex",
4464                    "protocol": "acp",
4465                    "cwd": std::env::current_dir().unwrap(),
4466                    "launch": {"program": "/bin/sh", "arguments": ["-c", script], "env": {}},
4467                }),
4468            ))
4469            .await;
4470        assert_eq!(started["result"]["connection"], "runtime-1");
4471        assert_eq!(started["result"]["handle"]["runtime_id"], "svc_acp");
4472
4473        let terminal = service
4474            .handle_async(request(
4475                9,
4476                "harness.v1.runtimes.terminal_instructions",
4477                json!({"connection":"runtime-1"}),
4478            ))
4479            .await;
4480        let arguments = terminal["result"]["launch"]["arguments"]
4481            .as_array()
4482            .expect("hosted runtime should return terminal arguments");
4483        let endpoint_index = arguments
4484            .iter()
4485            .position(|value| value == "--endpoint")
4486            .expect("terminal command should use an opaque endpoint");
4487        let endpoint = LiveRuntimeEndpoint::parse(
4488            arguments[endpoint_index + 1]
4489                .as_str()
4490                .expect("endpoint argument should be text"),
4491        )
4492        .unwrap();
4493        assert!(!terminal.to_string().contains("Bearer"));
4494        let workspace = std::env::current_dir().unwrap();
4495        let receipt = resolve_live_runtime(
4496            &endpoint,
4497            &LiveRuntimeSource {
4498                harness: "codex".into(),
4499                session_id: "svc_acp".into(),
4500                workspace,
4501            },
4502        )
4503        .unwrap();
4504        let remote = crate::HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
4505            .await
4506            .unwrap();
4507        let mut attachment = crate::FrontendRuntime::attach(remote.as_ref(), 100)
4508            .await
4509            .unwrap();
4510
4511        let sent = service
4512            .handle_async(request(
4513                2,
4514                "harness.v1.runtimes.send_input",
4515                json!({"connection": "runtime-1", "text": "hi"}),
4516            ))
4517            .await;
4518        assert_eq!(sent["result"]["turn_id"], "3");
4519
4520        let mut events = Vec::new();
4521        for _ in 0..20 {
4522            events.extend(service.poll_runtimes().await);
4523            if events.len() >= 2 {
4524                break;
4525            }
4526            tokio::time::sleep(Duration::from_millis(2)).await;
4527        }
4528        assert!(events
4529            .iter()
4530            .any(|event| { event["params"]["event"]["kind"] == "session/update" }));
4531        assert!(events.iter().any(|event| {
4532            event["params"]["event"]["kind"] == "supercode/acp_request_completed"
4533        }));
4534
4535        let saw_editor_reply = tokio::time::timeout(Duration::from_secs(2), async {
4536            loop {
4537                let event = attachment.next_event().await.unwrap();
4538                if event.kind == "text_delta" && event.payload["text"] == "ok" {
4539                    break;
4540                }
4541            }
4542        })
4543        .await;
4544        assert!(
4545            saw_editor_reply.is_ok(),
4546            "terminal should observe the editor-driven turn"
4547        );
4548
4549        crate::FrontendRuntime::submit(remote.as_ref(), "DRIVE FROM TERMINAL".into())
4550            .await
4551            .unwrap();
4552        let saw_terminal_reply = tokio::time::timeout(Duration::from_secs(2), async {
4553            loop {
4554                let event = attachment.next_event().await.unwrap();
4555                if event.kind == "text_delta" && event.payload["text"] == "from terminal" {
4556                    break;
4557                }
4558            }
4559        })
4560        .await;
4561        assert!(
4562            saw_terminal_reply.is_ok(),
4563            "terminal should drive the same runtime"
4564        );
4565
4566        let closed = service
4567            .handle_async(request(
4568                3,
4569                "harness.v1.runtimes.close",
4570                json!({"connection": "runtime-1"}),
4571            ))
4572            .await;
4573        assert_eq!(closed["result"]["closed"], true);
4574    }
4575}