1use 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
36pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
38pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
40pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
42
43pub 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 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 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 #[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 #[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 #[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 #[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 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 #[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 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 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) = ¶ms.options {
459 options.validate()?;
460 return load_session(¶ms.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 ¶ms.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(¶ms.read.locator, params.read.read_fidelity())
474 } else {
475 self.catalog
476 .load_parent_with_fidelity(
477 ¶ms.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 ¶ms.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(¶ms.subscription);
519 Ok(json!({
520 "removed": self.followers.remove(¶ms.subscription).is_some()
521 }))
522 }
523 "harness.v1.sessions.import" => {
524 let params = decode::<ImportSessionParams>(params)?;
525 let session = Session::load_str(¶ms.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(¶ms.locator).map_err(operation)?;
532 let artifact = session_artifact(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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 ¶ms.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(¶ms.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 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 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 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 ¶ms.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(¶ms)?;
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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.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(¶ms.connection) else {
898 return Err(ServiceError::InvalidParams(format!(
899 "unknown runtime connection `{}`",
900 params.connection
901 )));
902 };
903 self.terminal_launches.remove(¶ms.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 #[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(¶ms, &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 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, ¶ms.harnesses));
1059 let workspace_counts = params.include_sessions.then(|| {
1060 params
1061 .workspace
1062 .as_deref()
1063 .map(|workspace| self.session_counts(Some(workspace), ¶ms.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, ¶ms, 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 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 _ => 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 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 #[serde(default)]
1678 fidelity: Option<Fidelity>,
1679 #[serde(default)]
1682 view: Option<SessionReadView>,
1683}
1684
1685#[derive(Deserialize)]
1686struct SessionReadView {
1687 #[serde(default)]
1690 tail_messages: Option<usize>,
1691 #[serde(default)]
1694 include_subagents: bool,
1695 #[serde(default)]
1697 display_history: bool,
1698 #[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 #[serde(default)]
1792 homes: crate::HarnessHomes,
1793}
1794
1795#[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 ¶ms.homes,
1825 ¶ms.locator.session_id,
1826 ¶ms.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#[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 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 let mut importable = session.clone();
2422 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 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 ¶ms.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) = ¶ms.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
2808struct 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 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(¶ms) {
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(¶ms).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 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 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 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 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 #[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}