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