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 parent_session_id: None,
3681 child_session_count: 0,
3682 };
3683 let peer = crate::claude_peer::ClaudePeerSession {
3684 pid: 42,
3685 session_id: "live-session".into(),
3686 cwd: Some(PathBuf::from("/project")),
3687 name: "peer".into(),
3688 socket_path: PathBuf::from("/tmp/peer.sock"),
3689 status: Some(crate::claude_peer::ClaudePeerStatus::Busy),
3690 updated_at_ms: Some(1),
3691 version: Some("test".into()),
3692 };
3693
3694 let value = live_descriptor_value(&descriptor, &[peer]).unwrap();
3695 assert!(value["live_endpoint"]
3696 .as_str()
3697 .is_some_and(|endpoint| endpoint.starts_with("cc-peer:v1:42:peer:")));
3698 }
3699
3700 struct EndingRuntime {
3701 handle: RuntimeHandle,
3702 event: Option<HarnessEvent>,
3703 }
3704
3705 #[async_trait]
3706 impl RuntimeConnection for EndingRuntime {
3707 fn handle(&self) -> &RuntimeHandle {
3708 &self.handle
3709 }
3710
3711 async fn send_input(&mut self, _input: RuntimeInput) -> crate::Result<Option<String>> {
3712 unreachable!("ending runtime does not accept input")
3713 }
3714
3715 async fn next_event(&mut self) -> crate::Result<Option<HarnessEvent>> {
3716 Ok(self.event.take())
3717 }
3718
3719 async fn interrupt(&mut self) -> crate::Result<()> {
3720 Ok(())
3721 }
3722
3723 async fn respond(&mut self, _request_id: Value, _response: Value) -> crate::Result<()> {
3724 Ok(())
3725 }
3726
3727 async fn close(&mut self) -> crate::Result<()> {
3728 Ok(())
3729 }
3730 }
3731
3732 fn ending_runtime(event: Option<HarnessEvent>) -> Box<dyn RuntimeConnection> {
3733 Box::new(EndingRuntime {
3734 handle: RuntimeHandle {
3735 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3736 runtime_id: "ending-session".into(),
3737 endpoint: RuntimeEndpoint::LocalProcess {
3738 pid: None,
3739 command: vec!["ending-runtime".into()],
3740 protocol: "test".into(),
3741 },
3742 },
3743 event,
3744 })
3745 }
3746
3747 fn request(id: u64, method: &str, params: Value) -> Value {
3748 json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})
3749 }
3750
3751 fn pi_locator() -> SessionLocator {
3752 SessionLocator {
3753 harness: HarnessId::from(HarnessId::PI),
3754 session_id: "1e6f2a3b-0000-4000-8000-000000000001".into(),
3755 storage: StorageLocator::File {
3756 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3757 .join("tests/fixtures/pi_session.jsonl"),
3758 },
3759 }
3760 }
3761
3762 fn opencode_locator() -> SessionLocator {
3763 let session_id = "ses_fixtureAAAAAAAAAAAAAAA1";
3764 SessionLocator {
3765 harness: HarnessId::from(HarnessId::OPENCODE),
3766 session_id: session_id.into(),
3767 storage: StorageLocator::Sqlite {
3768 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3769 .join("tests/fixtures/opencode_fixture/opencode.db"),
3770 selector: session_id.into(),
3771 },
3772 }
3773 }
3774
3775 fn grok_locator() -> SessionLocator {
3776 SessionLocator {
3777 harness: HarnessId::from(HarnessId::GROK),
3778 session_id: "73c09283-4b33-41fa-90f1-0bcb0f7be523".into(),
3779 storage: StorageLocator::File {
3780 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3781 .join("tests/fixtures/grok_session/chat_history.jsonl"),
3782 },
3783 }
3784 }
3785
3786 #[test]
3787 fn capabilities_are_explicit_and_versioned() {
3788 let mut service = HarnessSessionService::new();
3789 let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
3790 assert_eq!(response["result"]["version"], HARNESS_SERVICE_VERSION);
3791 assert_eq!(
3792 response["result"]["sdk"]["schema_version"],
3793 crate::SDK_SCHEMA_VERSION
3794 );
3795 assert_eq!(
3796 response["result"]["sdk"]["operations"]
3797 .as_array()
3798 .unwrap()
3799 .len(),
3800 SdkOperation::ALL.len()
3801 );
3802 assert_eq!(response["result"]["harnesses"].as_array().unwrap().len(), 8);
3803 assert!(response["result"]["harnesses"]
3804 .as_array()
3805 .unwrap()
3806 .iter()
3807 .any(|harness| harness == HarnessId::GROK));
3808 assert!(response["result"]["harnesses"]
3809 .as_array()
3810 .unwrap()
3811 .iter()
3812 .any(|harness| harness == HarnessId::GOOSE));
3813 }
3814
3815 #[test]
3816 fn handshake_health_uses_protocol_liveness_not_stderr_severity() {
3817 let noisy_stderr = crate::HarnessEvent {
3818 sequence: None,
3819 kind: "transport_stderr".into(),
3820 payload: json!({"line": "ERROR optional worker AuthorizationRequired"}),
3821 };
3822 assert_eq!(handshake_event_failure(&noisy_stderr), None);
3823
3824 let closed = crate::HarnessEvent {
3825 sequence: None,
3826 kind: "transport_closed".into(),
3827 payload: json!({}),
3828 };
3829 assert!(handshake_event_failure(&closed).is_some());
3830 }
3831
3832 #[tokio::test]
3833 async fn runtime_eof_is_notified_and_removed_for_raw_and_explicit_close() {
3834 let mut service = HarnessSessionService::new();
3835 service
3836 .runtimes
3837 .insert("raw-eof".into(), ending_runtime(None));
3838 service.runtimes.insert(
3839 "explicit-close".into(),
3840 ending_runtime(Some(HarnessEvent {
3841 sequence: None,
3842 kind: "transport_closed".into(),
3843 payload: json!({"message": "native transport exited"}),
3844 })),
3845 );
3846
3847 let notifications = service.poll_runtimes().await;
3848
3849 assert_eq!(notifications.len(), 2);
3850 assert!(notifications
3851 .iter()
3852 .all(|notification| { notification["params"]["event"]["kind"] == "transport_closed" }));
3853 assert!(notifications.iter().all(|notification| {
3854 notification["params"]["session_id"] == "ending-session"
3855 && notification["params"]["connection"].is_string()
3856 }));
3857 let mut sequences = notifications
3858 .iter()
3859 .filter_map(|notification| notification["params"]["sequence"].as_u64())
3860 .collect::<Vec<_>>();
3861 sequences.sort_unstable();
3862 assert_eq!(sequences, vec![1, 2]);
3863 assert!(service.runtimes.is_empty());
3864 }
3865
3866 #[test]
3867 fn support_report_and_grok_default_binding_share_the_registry() {
3868 let mut service = HarnessSessionService::new();
3869 let response = service.handle(request(1, "harness.v1.support.report", json!({})));
3870 assert_eq!(response["result"]["schema"], crate::SUPPORT_REGISTRY_SCHEMA);
3871 let params = RuntimeBackendParams {
3872 harness: HarnessId::from(HarnessId::GROK),
3873 protocol: None,
3874 launch: None,
3875 base_url: None,
3876 policy: RuntimePolicy::Default,
3877 };
3878 let backend = match runtime_backend(¶ms) {
3879 Ok(backend) => backend,
3880 Err(_) => panic!("Grok should bind through its registered ACP launch"),
3881 };
3882 assert_eq!(backend.harness().as_str(), HarnessId::GROK);
3883 assert!(backend.capabilities().start_session);
3884 let registered = harness_support_registry()
3885 .harnesses
3886 .into_iter()
3887 .find(|harness| harness.id.as_str() == HarnessId::GROK)
3888 .and_then(|harness| harness.runtime.default_launch)
3889 .unwrap();
3890 assert!(!registered
3891 .arguments
3892 .iter()
3893 .any(|argument| argument == "--always-approve"));
3894 assert!(runtime_launch(¶ms).is_none());
3895
3896 let yolo = RuntimeBackendParams {
3897 policy: RuntimePolicy::Yolo,
3898 ..params
3899 };
3900 assert!(runtime_launch(&yolo)
3901 .unwrap()
3902 .arguments
3903 .iter()
3904 .any(|argument| argument == "--always-approve"));
3905
3906 let mismatched_protocol = RuntimeBackendParams {
3907 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3908 protocol: Some("acp".into()),
3909 launch: None,
3910 base_url: None,
3911 policy: RuntimePolicy::Default,
3912 };
3913 assert!(runtime_backend(&mismatched_protocol).is_err());
3914 }
3915
3916 #[test]
3917 fn load_follow_and_unfollow_share_the_same_locator() {
3918 let mut service = HarnessSessionService::new();
3919 let locator = pi_locator();
3920 let loaded = service.handle(request(
3921 1,
3922 "harness.v1.sessions.load",
3923 json!({"locator": locator}),
3924 ));
3925 assert_eq!(
3926 loaded["result"]["session"]["session_id"],
3927 locator.session_id
3928 );
3929
3930 let followed = service.handle(request(
3931 2,
3932 "harness.v1.sessions.follow",
3933 json!({"locator": locator}),
3934 ));
3935 assert_eq!(followed["result"]["subscription"], "sub-1");
3936 assert_eq!(followed["result"]["initial"]["type"], "session_snapshot");
3937 assert!(service.poll().is_empty());
3938
3939 let unfollowed = service.handle(request(
3940 3,
3941 "harness.v1.sessions.unfollow",
3942 json!({"subscription": "sub-1"}),
3943 ));
3944 assert_eq!(unfollowed["result"]["removed"], true);
3945 }
3946
3947 #[test]
3948 fn bounded_read_view_excludes_subagents_and_keeps_only_the_tail() {
3949 let temp = std::env::temp_dir().join(format!(
3950 "supercode-bounded-view-{}-{}",
3951 std::process::id(),
3952 generated_session_id()
3953 ));
3954 let path = temp.join("parent.jsonl");
3955 let subagents = temp.join("parent/subagents");
3956 std::fs::create_dir_all(&subagents).unwrap();
3957 let long_last = "x".repeat(300);
3958 let parent_records = [
3959 json!({"type":"user","uuid":"u1","parentUuid":null,"message":{"role":"user","content":"first"}}),
3960 json!({"type":"assistant","uuid":"a1","parentUuid":"u1","message":{"role":"assistant","content":[{"type":"text","text":"middle"}]}}),
3961 json!({"type":"user","uuid":"u2","parentUuid":"a1","message":{"role":"user","content":long_last}}),
3962 ];
3963 std::fs::write(
3964 &path,
3965 format!(
3966 "{}\n",
3967 parent_records
3968 .iter()
3969 .map(Value::to_string)
3970 .collect::<Vec<_>>()
3971 .join("\n")
3972 ),
3973 )
3974 .unwrap();
3975 std::fs::write(
3976 subagents.join("agent-child.jsonl"),
3977 concat!(
3978 r#"{"type":"user","uuid":"cu","parentUuid":null,"agentId":"child","message":{"role":"user","content":"child work"}}"#,
3979 "\n",
3980 ),
3981 )
3982 .unwrap();
3983 let locator = SessionLocator {
3984 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3985 session_id: "parent".into(),
3986 storage: StorageLocator::File { path },
3987 };
3988 let mut service = HarnessSessionService::new();
3989
3990 let complete = service.handle(request(
3991 1,
3992 "harness.v1.sessions.load",
3993 json!({"locator": locator}),
3994 ));
3995 assert_eq!(
3996 complete["result"]["session"]["subagents"]
3997 .as_array()
3998 .unwrap()
3999 .len(),
4000 1
4001 );
4002
4003 let bounded = service.handle(request(
4004 2,
4005 "harness.v1.sessions.load",
4006 json!({
4007 "locator": locator,
4008 "view": {
4009 "tail_messages": 1,
4010 "max_message_chars": 256,
4011 "include_subagents": false
4012 },
4013 }),
4014 ));
4015 let session = &bounded["result"]["session"];
4016 assert!(session["subagents"].as_array().unwrap().is_empty());
4017 assert_eq!(session["messages"].as_array().unwrap().len(), 1);
4018 assert_eq!(
4019 session["messages"][0]["content"],
4020 format!("{}\n…", "x".repeat(256))
4021 );
4022
4023 let followed = service.handle(request(
4024 3,
4025 "harness.v1.sessions.follow",
4026 json!({
4027 "locator": locator,
4028 "view": {
4029 "tail_messages": 1,
4030 "max_message_chars": 256,
4031 "include_subagents": false
4032 },
4033 }),
4034 ));
4035 let initial = &followed["result"]["initial"]["session"];
4036 assert!(initial["subagents"].as_array().unwrap().is_empty());
4037 assert_eq!(initial["messages"].as_array().unwrap().len(), 1);
4038
4039 let _ = std::fs::remove_dir_all(&temp);
4040 }
4041
4042 #[test]
4043 fn forty_megabyte_display_load_is_bounded_and_prompt() {
4044 let temp = std::env::temp_dir().join(format!(
4045 "supercode-large-display-view-{}-{}",
4046 std::process::id(),
4047 generated_session_id()
4048 ));
4049 std::fs::create_dir_all(&temp).unwrap();
4050 let path = temp.join("rollout.jsonl");
4051 let mut file = std::io::BufWriter::new(std::fs::File::create(&path).unwrap());
4052 writeln!(
4053 file,
4054 r#"{{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{{"id":"large-display","cwd":"/tmp"}}}}"#
4055 )
4056 .unwrap();
4057 let padding = "x".repeat(80 * 1024);
4058 for index in 0..512 {
4059 let marker = if index == 0 {
4060 "OLDEST-SHOULD-NOT-LOAD"
4061 } else if index == 511 {
4062 "LATEST-MUST-LOAD"
4063 } else {
4064 "bulk"
4065 };
4066 writeln!(
4067 file,
4068 "{}",
4069 json!({
4070 "timestamp": "2026-01-01T00:00:01Z",
4071 "type": "response_item",
4072 "payload": {
4073 "type": "message",
4074 "role": "assistant",
4075 "content": [{"type": "output_text", "text": format!("{marker}:{padding}")}],
4076 },
4077 })
4078 )
4079 .unwrap();
4080 }
4081 file.flush().unwrap();
4082 drop(file);
4083 assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
4084
4085 let locator = SessionLocator {
4086 harness: HarnessId::from(HarnessId::CODEX),
4087 session_id: "large-display".into(),
4088 storage: StorageLocator::File { path },
4089 };
4090 let started = Instant::now();
4091 let response = HarnessSessionService::new().handle(request(
4092 1,
4093 "harness.v1.sessions.load",
4094 json!({
4095 "locator": locator,
4096 "view": {
4097 "tail_messages": 500,
4098 "max_message_chars": 1024,
4099 "include_subagents": false,
4100 "display_history": true,
4101 },
4102 }),
4103 ));
4104 let elapsed = started.elapsed();
4105 let wire = response.to_string();
4106 eprintln!(
4107 "bounded 40 MiB display load: {elapsed:?}, {} response bytes",
4108 wire.len()
4109 );
4110 assert!(response.get("error").is_none(), "{response:#}");
4111 assert!(wire.contains("LATEST-MUST-LOAD"));
4112 assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
4113 assert!(
4114 wire.len() < 2 * 1024 * 1024,
4115 "bounded wire was {} bytes",
4116 wire.len()
4117 );
4118 assert!(
4119 elapsed.as_secs_f64() < 3.0,
4120 "bounded 40 MiB load took {elapsed:?}"
4121 );
4122
4123 let _ = std::fs::remove_dir_all(&temp);
4124 }
4125
4126 #[test]
4127 fn forty_megabyte_goose_store_display_load_reads_only_the_tail() {
4128 let temp = std::env::temp_dir().join(format!(
4129 "supercode-large-goose-view-{}-{}",
4130 std::process::id(),
4131 generated_session_id()
4132 ));
4133 std::fs::create_dir_all(&temp).unwrap();
4134 let path = temp.join("sessions.db");
4135 let connection = rusqlite::Connection::open(&path).unwrap();
4136 connection
4137 .execute_batch(
4138 "CREATE TABLE sessions (
4139 id TEXT PRIMARY KEY, name TEXT NOT NULL, working_dir TEXT NOT NULL,
4140 created_at TEXT NOT NULL, updated_at TEXT NOT NULL,
4141 session_type TEXT NOT NULL, extension_data TEXT,
4142 goose_mode TEXT NOT NULL, provider_name TEXT, model_config_json TEXT,
4143 archived_at TEXT
4144 );
4145 CREATE TABLE messages (
4146 id INTEGER PRIMARY KEY, session_id TEXT NOT NULL, message_id TEXT,
4147 role TEXT NOT NULL, content_json TEXT NOT NULL,
4148 created_timestamp INTEGER NOT NULL, metadata_json TEXT
4149 );",
4150 )
4151 .unwrap();
4152 connection
4153 .execute(
4154 "INSERT INTO sessions VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL)",
4155 rusqlite::params![
4156 "goose-large",
4157 "Large Goose session",
4158 "/tmp",
4159 "2026-01-01 00:00:00",
4160 "2026-01-01 00:00:02",
4161 "user",
4162 "{}",
4163 "auto",
4164 "anthropic",
4165 r#"{"model_name":"claude-sonnet"}"#,
4166 ],
4167 )
4168 .unwrap();
4169 let old_content = serde_json::to_string(&vec![json!({
4170 "type": "text",
4171 "text": format!("OLDEST-SHOULD-NOT-LOAD:{}", "x".repeat(40 * 1024 * 1024)),
4172 })])
4173 .unwrap();
4174 connection
4175 .execute(
4176 "INSERT INTO messages VALUES (1, ?1, 'old', 'user', ?2, 1, '{}')",
4177 rusqlite::params!["goose-large", old_content],
4178 )
4179 .unwrap();
4180 connection
4181 .execute(
4182 "INSERT INTO messages VALUES (2, ?1, 'new', 'assistant', ?2, 2, '{}')",
4183 rusqlite::params![
4184 "goose-large",
4185 r#"[{"type":"text","text":"LATEST-MUST-LOAD"}]"#
4186 ],
4187 )
4188 .unwrap();
4189 drop(connection);
4190 assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
4191
4192 let locator = SessionLocator {
4193 harness: HarnessId::from(HarnessId::GOOSE),
4194 session_id: "goose-large".into(),
4195 storage: StorageLocator::Sqlite {
4196 path,
4197 selector: "goose-large".into(),
4198 },
4199 };
4200 let started = Instant::now();
4201 let response = HarnessSessionService::new().handle(request(
4202 1,
4203 "harness.v1.sessions.load",
4204 json!({
4205 "locator": locator,
4206 "view": {
4207 "tail_messages": 1,
4208 "max_message_chars": 1024,
4209 "include_subagents": false,
4210 "display_history": true,
4211 },
4212 }),
4213 ));
4214 let elapsed = started.elapsed();
4215 let wire = response.to_string();
4216 eprintln!(
4217 "bounded 40 MiB Goose display load: {elapsed:?}, {} response bytes",
4218 wire.len()
4219 );
4220 assert!(response.get("error").is_none(), "{response:#}");
4221 assert!(wire.contains("LATEST-MUST-LOAD"));
4222 assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
4223 assert!(
4224 wire.len() < 64 * 1024,
4225 "bounded wire was {} bytes",
4226 wire.len()
4227 );
4228 assert!(
4229 elapsed.as_secs_f64() < 1.0,
4230 "bounded Goose load took {elapsed:?}"
4231 );
4232
4233 let _ = std::fs::remove_dir_all(&temp);
4234 }
4235
4236 #[test]
4237 fn display_view_keeps_codex_assistant_history_across_compaction() {
4238 let temp = std::env::temp_dir().join(format!(
4239 "supercode-codex-display-view-{}-{}",
4240 std::process::id(),
4241 generated_session_id()
4242 ));
4243 std::fs::create_dir_all(&temp).unwrap();
4244 let path = temp.join("rollout.jsonl");
4245 std::fs::write(
4246 &path,
4247 concat!(
4248 r#"{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{"id":"codex-display","cwd":"/tmp"}}"#,
4249 "\n",
4250 r#"{"timestamp":"2026-01-01T00:00:01Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"old prompt"}]}}"#,
4251 "\n",
4252 r#"{"timestamp":"2026-01-01T00:00:02Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"old answer"}]}}"#,
4253 "\n",
4254 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"}]}}"#,
4255 "\n",
4256 r#"{"timestamp":"2026-01-01T00:00:04Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"new prompt"}]}}"#,
4257 "\n",
4258 r#"{"timestamp":"2026-01-01T00:00:05Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"new answer"}]}}"#,
4259 "\n",
4260 ),
4261 )
4262 .unwrap();
4263 let locator = SessionLocator {
4264 harness: HarnessId::from(HarnessId::CODEX),
4265 session_id: "codex-display".into(),
4266 storage: StorageLocator::File { path },
4267 };
4268 let mut service = HarnessSessionService::new();
4269
4270 let continuation = service.handle(request(
4271 1,
4272 "harness.v1.sessions.load",
4273 json!({"locator": locator}),
4274 ));
4275 let continuation_text = continuation["result"]["session"]["messages"].to_string();
4276 assert!(!continuation_text.contains("old answer"));
4277
4278 let display = service.handle(request(
4279 2,
4280 "harness.v1.sessions.load",
4281 json!({
4282 "locator": locator,
4283 "view": {
4284 "tail_messages": 10,
4285 "include_subagents": false,
4286 "display_history": true,
4287 },
4288 }),
4289 ));
4290 let display_text = display["result"]["session"]["messages"].to_string();
4291 assert!(display_text.contains("old prompt"));
4292 assert!(display_text.contains("old answer"));
4293 assert!(display_text.contains("new prompt"));
4294 assert!(display_text.contains("new answer"));
4295
4296 let _ = std::fs::remove_dir_all(&temp);
4297 }
4298
4299 #[test]
4300 fn load_supports_bounded_windows_and_media_metadata() {
4301 let mut service = HarnessSessionService::new();
4302 let locator = pi_locator();
4303 let bounded = service.handle(request(
4304 1,
4305 "harness.v1.sessions.load",
4306 json!({
4307 "locator": locator,
4308 "options": {
4309 "include_subagents": false,
4310 "message_limit": 2,
4311 "message_offset": 1
4312 }
4313 }),
4314 ));
4315 assert_eq!(bounded["result"]["window"]["offset"], 1);
4316 assert_eq!(bounded["result"]["window"]["returned"], 2);
4317 assert!(bounded["result"]["summary"]["first_message"].is_object());
4318 assert!(bounded["result"]["summary"]["last_message"].is_object());
4319 assert_eq!(
4320 bounded["result"]["session"]["messages"]
4321 .as_array()
4322 .unwrap()
4323 .len(),
4324 2
4325 );
4326 assert!(bounded["result"]["session"]["subagents"]
4327 .as_array()
4328 .unwrap()
4329 .is_empty());
4330
4331 let tail = service.handle(request(
4332 2,
4333 "harness.v1.sessions.load",
4334 json!({"locator": locator, "options": {"message_tail": 1}}),
4335 ));
4336 assert_eq!(tail["result"]["window"]["returned"], 1);
4337 assert_eq!(tail["result"]["window"]["has_more"], true);
4338 assert_eq!(tail["result"]["window"]["has_older"], true);
4339 assert!(tail["result"]["window"]["older_items"].as_u64().unwrap() > 0);
4340 assert!(tail["result"]["summary"]["first_message"].is_object());
4341
4342 let metadata_only = service.handle(request(
4343 3,
4344 "harness.v1.sessions.load",
4345 json!({"locator": locator, "options": {"inline_media": "metadata"}}),
4346 ));
4347 assert!(metadata_only["result"]["session"]
4348 .to_string()
4349 .contains("media_reference"));
4350 assert!(!metadata_only["result"]["session"]
4351 .to_string()
4352 .contains("data:image/"));
4353 }
4354
4355 #[test]
4356 fn import_translate_branch_and_handoff_use_typed_artifacts() {
4357 let mut service = HarnessSessionService::new();
4358 let locator = pi_locator();
4359 let translated = service.handle(request(
4360 1,
4361 "harness.v1.sessions.translate",
4362 json!({"locator": locator, "target_harness": "grok"}),
4363 ));
4364 assert_eq!(translated["result"]["artifact"]["source_harness"], "pi");
4365 assert_eq!(translated["result"]["artifact"]["target_harness"], "grok");
4366 assert!(translated["result"]["artifact"]["content"]
4367 .as_str()
4368 .is_some_and(|content| !content.is_empty()));
4369
4370 for target in ["opencode", "open-code"] {
4371 let opencode = service.handle(request(
4372 6,
4373 "harness.v1.sessions.translate",
4374 json!({"locator": locator, "target_harness": target}),
4375 ));
4376 assert_eq!(opencode["result"]["artifact"]["target_harness"], "opencode");
4377 }
4378 let goose = service.handle(request(
4379 7,
4380 "harness.v1.sessions.translate",
4381 json!({"locator": locator, "target_harness": "goose"}),
4382 ));
4383 assert_eq!(goose["result"]["artifact"]["target_harness"], "goose");
4384 assert!(serde_json::from_str::<Value>(
4385 goose["result"]["artifact"]["content"].as_str().unwrap()
4386 )
4387 .unwrap()["conversation"]
4388 .is_array());
4389
4390 let imported = service.handle(request(
4391 2,
4392 "harness.v1.sessions.import",
4393 json!({
4394 "source_harness": "grok",
4395 "content": translated["result"]["artifact"]["content"],
4396 }),
4397 ));
4398 assert_eq!(imported["result"]["session"]["source"], "grok");
4399
4400 let branched = service.handle(request(
4401 3,
4402 "harness.v1.sessions.branch",
4403 json!({"locator": locator, "target_harness": "codex"}),
4404 ));
4405 assert_eq!(branched["result"]["parent"]["harness"], "pi");
4406 assert!(branched["result"]["bootstrap_prompt"]
4407 .as_str()
4408 .unwrap()
4409 .contains("frozen parent transcript"));
4410 assert_eq!(branched["result"]["artifact"]["target_harness"], "codex");
4411
4412 let handoff = service.handle(request(
4413 4,
4414 "harness.v1.sessions.handoff",
4415 json!({"locator": locator, "target_harness": "pi", "cwd": "/tmp/project"}),
4416 ));
4417 assert_eq!(handoff["result"]["launch"]["program"], "pi");
4418 assert_eq!(handoff["result"]["launch"]["cwd"], "/tmp/project");
4419 assert_eq!(handoff["result"]["requires_materialization"], true);
4420
4421 let goose_handoff = service.handle(request(
4422 8,
4423 "harness.v1.sessions.handoff",
4424 json!({"locator": locator, "target_harness": "goose", "cwd": "/tmp/project"}),
4425 ));
4426 assert_eq!(goose_handoff["result"]["launch"]["program"], "goose");
4427 assert_eq!(
4428 goose_handoff["result"]["materialize"]["arguments"],
4429 json!(["session", "import", "{artifact_path}"])
4430 );
4431
4432 let resumed = service.handle(request(
4433 5,
4434 "harness.v1.sessions.resume_instructions",
4435 json!({"locator": locator, "cwd": "/tmp/project", "policy": "yolo"}),
4436 ));
4437 assert_eq!(resumed["result"]["launch"]["program"], "pi");
4438 assert_eq!(resumed["result"]["launch"]["arguments"][0], "--approve");
4439 }
4440
4441 #[test]
4442 fn reduce_persists_and_reloads_a_byte_exact_reversible_bundle() {
4443 let temp = std::env::temp_dir().join(format!(
4444 "supercode-service-reduce-{}-{}",
4445 std::process::id(),
4446 generated_session_id()
4447 ));
4448 let source_path = temp.join("source.jsonl");
4449 let store_root = temp.join("store");
4450 std::fs::create_dir_all(&temp).unwrap();
4451
4452 let mut records = vec![json!({
4453 "timestamp": "2026-01-01T00:00:00Z",
4454 "type": "session_meta",
4455 "payload": {"id": "codex-reduce", "cwd": "/tmp/project"},
4456 })];
4457 for turn in 0..16 {
4458 records.push(json!({
4459 "timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 1),
4460 "type": "response_item",
4461 "payload": {
4462 "type": "message",
4463 "role": "user",
4464 "content": [{
4465 "type": "input_text",
4466 "text": format!("request {turn}: {}", "context ".repeat(80)),
4467 }],
4468 },
4469 }));
4470 records.push(json!({
4471 "timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 2),
4472 "type": "response_item",
4473 "payload": {
4474 "type": "message",
4475 "role": "assistant",
4476 "content": [{
4477 "type": "output_text",
4478 "text": format!("answer {turn}: {}", "implementation detail ".repeat(80)),
4479 }],
4480 },
4481 }));
4482 }
4483 let source = format!(
4484 "{}\n",
4485 records
4486 .iter()
4487 .map(Value::to_string)
4488 .collect::<Vec<_>>()
4489 .join("\n")
4490 );
4491 std::fs::write(&source_path, &source).unwrap();
4492 let locator = SessionLocator {
4493 harness: HarnessId::from(HarnessId::CODEX),
4494 session_id: "codex-reduce".into(),
4495 storage: StorageLocator::File {
4496 path: source_path.clone(),
4497 },
4498 };
4499 let original = load_session(&locator).unwrap();
4500 let mut service =
4501 HarnessSessionService::new().with_reduction_store_root(store_root.clone());
4502
4503 let response = service.handle(request(
4504 1,
4505 "harness.v1.sessions.reduce",
4506 json!({
4507 "locator": locator,
4508 "target_harness": "claude-code",
4509 "keep_last": 4,
4510 }),
4511 ));
4512 assert!(response.get("error").is_none(), "{response:#}");
4513 let receipt = &response["result"]["receipt"];
4514 assert_eq!(receipt["source_harness"], "codex");
4515 assert_eq!(receipt["target_harness"], "claude-code");
4516 assert_eq!(receipt["verified"], true);
4517 assert_eq!(receipt["reversible"], true);
4518 assert!(receipt["reductions"].as_u64().unwrap() > 0);
4519 assert!(
4520 receipt["source_tokens"].as_u64().unwrap()
4521 > receipt["reduced_tokens"].as_u64().unwrap()
4522 );
4523 assert!(receipt["ratio"].as_f64().unwrap() > 1.0);
4524 assert!(response["result"]["bootstrap_prompt"]
4525 .as_str()
4526 .unwrap()
4527 .contains("Do not guess hidden content"));
4528
4529 let rescue_id = receipt["id"].as_str().unwrap();
4530 let store = crate::SessionStore::open(&store_root).unwrap();
4531 let sidecar =
4532 Session::from_sidecar_str(&store.load_sidecar(rescue_id).unwrap().unwrap()).unwrap();
4533 let log = store.load_reduction_log(rescue_id).unwrap().unwrap();
4534 let persisted_view = parse_messages_jsonl(&store.load(rescue_id).unwrap()).unwrap();
4535 let policy = reduce::ReductionPolicy {
4536 clear_turns_older_than: Some(4),
4537 ..Default::default()
4538 };
4539 let (restamped_view, reapplied_log) =
4540 reduce::project_messages(&sidecar.messages, &policy, &log);
4541 assert_eq!(
4542 messages_jsonl(&persisted_view).unwrap(),
4543 messages_jsonl(&restamped_view).unwrap()
4544 );
4545 assert_eq!(reapplied_log, log);
4546 reduce::verify_log(&log, &sidecar).unwrap();
4547 assert_eq!(
4548 reduce::invert(&restamped_view, &log, &sidecar).unwrap(),
4549 original.messages
4550 );
4551 assert_eq!(std::fs::read_to_string(&source_path).unwrap(), source);
4552
4553 std::fs::remove_dir_all(temp).ok();
4554 }
4555
4556 #[test]
4557 fn read_surfaces_view_a_severed_claude_graph_while_transfer_still_refuses_it() {
4558 let temp = std::env::temp_dir().join(format!(
4559 "supercode-severed-view-{}-{}",
4560 std::process::id(),
4561 generated_session_id()
4562 ));
4563 std::fs::create_dir_all(&temp).unwrap();
4564 let path = temp.join("severed.jsonl");
4565 std::fs::write(
4568 &path,
4569 concat!(
4570 r#"{"type":"user","uuid":"orphan-u","parentUuid":null,"message":{"role":"user","content":"stranded prompt"}}"#,
4571 "\n",
4572 r#"{"type":"assistant","uuid":"live-a","parentUuid":"pruned","message":{"id":"m","role":"assistant","content":[{"type":"text","text":"live answer"}]}}"#,
4573 "\n",
4574 ),
4575 )
4576 .unwrap();
4577 let locator = SessionLocator {
4578 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
4579 session_id: "severed".into(),
4580 storage: StorageLocator::File { path },
4581 };
4582 let mut service = HarnessSessionService::new();
4583
4584 let viewed = service.handle(request(
4585 1,
4586 "harness.v1.sessions.load",
4587 json!({"locator": locator}),
4588 ));
4589 let session = &viewed["result"]["session"];
4590 assert_eq!(session["fidelity"], "semantic");
4591 assert_eq!(session["messages"].as_array().unwrap().len(), 2);
4592 assert!(session["residue"].as_array().unwrap().iter().any(|entry| {
4593 entry
4594 .as_str()
4595 .is_some_and(|entry| entry.contains("live-a") && entry.contains("pruned"))
4596 }));
4597
4598 let strict = service.handle(request(
4601 2,
4602 "harness.v1.sessions.load",
4603 json!({"locator": locator, "fidelity": "byte_lossless"}),
4604 ));
4605 assert!(strict["error"]["message"]
4606 .as_str()
4607 .unwrap()
4608 .contains("cannot reconstruct lossless Claude continuation"));
4609
4610 let translated = service.handle(request(
4612 3,
4613 "harness.v1.sessions.translate",
4614 json!({"locator": locator, "target_harness": "codex"}),
4615 ));
4616 assert!(translated["error"]["message"]
4617 .as_str()
4618 .unwrap()
4619 .contains("cannot reconstruct lossless Claude continuation"));
4620 let resumed = service.handle(request(
4621 4,
4622 "harness.v1.sessions.resume_instructions",
4623 json!({"locator": locator}),
4624 ));
4625 assert!(resumed["error"]["message"]
4626 .as_str()
4627 .unwrap()
4628 .contains("cannot reconstruct lossless Claude continuation"));
4629
4630 let _ = std::fs::remove_dir_all(&temp);
4631 }
4632
4633 #[test]
4634 fn structured_resume_launches_cover_gemini_goose_and_supercode() {
4635 let codex = resume_launch(
4636 HarnessId::CODEX,
4637 "codex-session",
4638 Path::new("/tmp/project"),
4639 ResumePolicy::Yolo,
4640 )
4641 .unwrap_or_else(|_| panic!("Codex resume launch must be registered"));
4642 assert_eq!(codex.program, "codex");
4643 assert_eq!(
4644 codex.arguments,
4645 [
4646 "-c",
4647 "check_for_update_on_startup=false",
4648 "-c",
4649 "projects.\"/tmp/project\".trust_level=\"trusted\"",
4650 "--dangerously-bypass-approvals-and-sandbox",
4651 "--dangerously-bypass-hook-trust",
4652 "resume",
4653 "codex-session",
4654 ]
4655 );
4656
4657 let gemini = resume_launch(
4658 HarnessId::GEMINI,
4659 "gemini-session",
4660 Path::new("/tmp/project"),
4661 ResumePolicy::Yolo,
4662 )
4663 .unwrap_or_else(|_| panic!("Gemini resume launch must be registered"));
4664 assert_eq!(gemini.program, "gemini");
4665 assert_eq!(gemini.arguments, ["--yolo", "--resume", "gemini-session"]);
4666
4667 let goose = resume_launch(
4668 HarnessId::GOOSE,
4669 "goose-session",
4670 Path::new("/tmp/project"),
4671 ResumePolicy::Yolo,
4672 )
4673 .unwrap_or_else(|_| panic!("Goose resume launch must be registered"));
4674 assert_eq!(goose.program, "goose");
4675 assert_eq!(
4676 goose.arguments,
4677 ["session", "--resume", "--session-id", "goose-session"]
4678 );
4679
4680 let supercode = resume_launch(
4681 HarnessId::SUPERCODE,
4682 "supercode-session",
4683 Path::new("/tmp/project"),
4684 ResumePolicy::Yolo,
4685 )
4686 .unwrap_or_else(|_| panic!("Supercode resume launch must be registered"));
4687 assert_eq!(supercode.program, "supercode");
4688 assert_eq!(
4689 supercode.arguments,
4690 ["--dangerous", "resume", "supercode-session"]
4691 );
4692 }
4693
4694 #[test]
4695 fn diagonal_artifacts_preserve_claude_subagents_and_grok_bundle_members() {
4696 let temp = std::env::temp_dir().join(format!(
4697 "supercode-harness-artifact-{}-{}",
4698 std::process::id(),
4699 generated_session_id()
4700 ));
4701 let main_path = temp.join("parent.jsonl");
4702 let subagent_path = temp.join("parent/subagents/agent-child.jsonl");
4703 std::fs::create_dir_all(subagent_path.parent().unwrap()).unwrap();
4704 let fixture = std::fs::read_to_string(
4705 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4706 .join("tests/fixtures/claude_code_session.jsonl"),
4707 )
4708 .unwrap();
4709 let parent = fixture.trim_end_matches('\n');
4710 let child = fixture.trim_end_matches('\n');
4711 std::fs::write(&main_path, parent).unwrap();
4712 std::fs::write(&subagent_path, child).unwrap();
4713 let locator = SessionLocator {
4714 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
4715 session_id: "213bb148-51ea-453f-9206-f8b4b1168547".into(),
4716 storage: StorageLocator::File {
4717 path: main_path.clone(),
4718 },
4719 };
4720 let mut service = HarnessSessionService::new();
4721 let claude = service.handle(request(
4722 1,
4723 "harness.v1.sessions.translate",
4724 json!({"locator": locator, "target_harness": "claude-code"}),
4725 ));
4726 let artifact = &claude["result"]["artifact"];
4727 assert_eq!(artifact["fidelity"], "byte_lossless");
4728 assert_eq!(artifact["content"], parent);
4729 let files = artifact["files"].as_array().unwrap();
4730 assert!(files.iter().any(|file| {
4731 file["role"] == "subagent"
4732 && file["path"]
4733 .as_str()
4734 .is_some_and(|path| path.ends_with("/subagents/agent-child.jsonl"))
4735 && file["content"] == child
4736 }));
4737 assert!(!artifact["content"].as_str().unwrap().ends_with('\n'));
4738
4739 let grok = service.handle(request(
4740 2,
4741 "harness.v1.sessions.translate",
4742 json!({"locator": grok_locator(), "target_harness": "grok"}),
4743 ));
4744 let files = grok["result"]["artifact"]["files"].as_array().unwrap();
4745 for name in ["summary.json", "updates.jsonl"] {
4746 let expected = std::fs::read_to_string(
4747 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4748 .join("tests/fixtures/grok_session")
4749 .join(name),
4750 )
4751 .unwrap();
4752 assert!(files.iter().any(|file| {
4753 file["path"] == name && file["role"] == "bundle" && file["content"] == expected
4754 }));
4755 }
4756 std::fs::remove_dir_all(temp).ok();
4757 }
4758
4759 #[test]
4760 fn every_non_grok_handoff_mints_and_uses_a_fresh_target_identity() {
4761 let mut service = HarnessSessionService::new();
4762 let source = pi_locator();
4763 for (target, format) in [
4764 ("claude-code", SessionFormat::ClaudeCode),
4765 ("codex", SessionFormat::Codex),
4766 ("opencode", SessionFormat::OpenCode),
4767 ("pi", SessionFormat::Pi),
4768 ] {
4769 let result = service.handle(request(
4770 1,
4771 "harness.v1.sessions.handoff",
4772 json!({"locator": source, "target_harness": target, "cwd": "/tmp/project"}),
4773 ));
4774 let artifact = &result["result"]["artifact"];
4775 let target_id = artifact["session_id"].as_str().unwrap();
4776 assert_ne!(target_id, source.session_id, "{target}");
4777 let parsed = Session::load_str(artifact["content"].as_str().unwrap(), format).unwrap();
4778 assert_eq!(
4779 parsed.meta.session_id.as_deref(),
4780 Some(target_id),
4781 "{target}"
4782 );
4783 if target != "pi" {
4784 assert!(result["result"]["launch"]["arguments"]
4785 .as_array()
4786 .unwrap()
4787 .iter()
4788 .any(|argument| argument == target_id));
4789 }
4790 if target == "opencode" {
4791 assert!(target_id.starts_with("ses_"));
4792 fn assert_session_ids(value: &Value, target_id: &str) {
4793 match value {
4794 Value::Object(fields) => {
4795 if let Some(session_id) = fields.get("sessionID") {
4796 assert_eq!(session_id, target_id);
4797 }
4798 for child in fields.values() {
4799 assert_session_ids(child, target_id);
4800 }
4801 }
4802 Value::Array(values) => {
4803 for child in values {
4804 assert_session_ids(child, target_id);
4805 }
4806 }
4807 _ => {}
4808 }
4809 }
4810 let document: Value =
4811 serde_json::from_str(artifact["content"].as_str().unwrap()).unwrap();
4812 assert_session_ids(&document, target_id);
4813 }
4814 }
4815
4816 let first = service.handle(request(
4817 2,
4818 "harness.v1.sessions.handoff",
4819 json!({"locator": source, "target_harness": "codex"}),
4820 ));
4821 let second = service.handle(request(
4822 3,
4823 "harness.v1.sessions.handoff",
4824 json!({"locator": source, "target_harness": "codex"}),
4825 ));
4826 assert_ne!(
4827 first["result"]["artifact"]["session_id"],
4828 second["result"]["artifact"]["session_id"]
4829 );
4830 }
4831
4832 #[test]
4833 fn grok_handoff_uses_the_official_importer_contract() {
4834 let mut service = HarnessSessionService::new();
4835 let source = opencode_locator();
4836 let response = service.handle(request(
4837 1,
4838 "harness.v1.sessions.handoff",
4839 json!({
4840 "locator": source,
4841 "target_harness": "grok",
4842 "cwd": "/tmp/grok-handoff-project",
4843 }),
4844 ));
4845 let result = &response["result"];
4846
4847 assert_eq!(result["artifact"]["target_harness"], "claude-code");
4851 assert!(result["artifact"]["suggested_filename"]
4852 .as_str()
4853 .unwrap()
4854 .ends_with(".grok-import.claude-code.jsonl"));
4855 let artifact = Session::load_str(
4856 result["artifact"]["content"].as_str().unwrap(),
4857 SessionFormat::ClaudeCode,
4858 )
4859 .unwrap();
4860 assert_eq!(
4861 artifact.meta.cwd.as_deref(),
4862 Some(Path::new("/tmp/grok-handoff-project"))
4863 );
4864 let target_session_id = artifact.meta.session_id.as_deref().unwrap();
4865 assert_eq!(target_session_id.len(), 36);
4866 assert_eq!(target_session_id.as_bytes()[14], b'4');
4867 assert_ne!(target_session_id, opencode_locator().session_id);
4868 assert_eq!(
4869 result["artifact"]["session_id"],
4870 artifact.meta.session_id.as_deref().unwrap()
4871 );
4872
4873 assert_eq!(
4874 result["materialize"]["arguments"],
4875 json!(["import", "--json", "{artifact_path}"])
4876 );
4877 assert_eq!(
4878 result["launch"]["arguments"],
4879 json!(["--resume", "{imported_session_id}", "--fork-session"])
4880 );
4881 assert!(result["note"]
4882 .as_str()
4883 .unwrap()
4884 .contains("outcome=imported"));
4885 assert!(!result["launch"]["arguments"]
4886 .as_array()
4887 .unwrap()
4888 .iter()
4889 .any(|argument| argument == &opencode_locator().session_id));
4890 }
4891
4892 #[tokio::test]
4893 async fn inventory_rejects_unknown_harnesses_and_runtime_attach_is_honest() {
4894 let mut service = HarnessSessionService::new();
4895 let inventory = service
4896 .handle_async(request(
4897 1,
4898 "harness.v1.harnesses.list",
4899 json!({"harnesses": ["missing"]}),
4900 ))
4901 .await;
4902 assert_eq!(inventory["error"]["code"], -32602);
4903
4904 let attached = service
4905 .handle_async(request(
4906 2,
4907 "harness.v1.runtimes.attach_existing",
4908 json!({"harness": "codex", "runtime_id": "thread-1"}),
4909 ))
4910 .await;
4911 assert_eq!(attached["error"]["code"], -32000);
4912 assert!(attached["error"]["message"]
4913 .as_str()
4914 .unwrap()
4915 .contains("runtimes.resume"));
4916 }
4917
4918 #[test]
4919 fn invalid_params_and_unknown_methods_use_json_rpc_errors() {
4920 let mut service = HarnessSessionService::new();
4921 let invalid = service.handle(request(1, "harness.v1.sessions.load", json!({})));
4922 assert_eq!(invalid["error"]["code"], -32602);
4923 let unknown = service.handle(request(2, "harness.v1.unknown", json!({})));
4924 assert_eq!(unknown["error"]["code"], -32601);
4925 }
4926
4927 #[cfg(unix)]
4928 #[tokio::test]
4929 #[allow(clippy::await_holding_lock)]
4932 async fn async_service_drives_a_generic_acp_runtime() {
4933 let _environment_guard = crate::live_runtime::test_environment_lock();
4934 let script = r#"
4935 i=0
4936 while IFS= read -r line; do
4937 i=$((i + 1))
4938 case "$i" in
4939 1) printf '%s\n' '{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentCapabilities":{},"authMethods":[]}}' ;;
4940 2) printf '%s\n' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"svc_acp"}}' ;;
4941 3)
4942 printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"ok"}}}}'
4943 printf '%s\n' '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
4944 ;;
4945 4)
4946 printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"from terminal"}}}}'
4947 printf '%s\n' '{"jsonrpc":"2.0","id":4,"result":{"stopReason":"end_turn"}}'
4948 ;;
4949 esac
4950 done
4951 "#;
4952 let mut service = HarnessSessionService::new();
4953 let started = service
4954 .handle_async(request(
4955 1,
4956 "harness.v1.runtimes.start",
4957 json!({
4958 "harness": "codex",
4959 "protocol": "acp",
4960 "cwd": std::env::current_dir().unwrap(),
4961 "launch": {"program": "/bin/sh", "arguments": ["-c", script], "env": {}},
4962 }),
4963 ))
4964 .await;
4965 assert_eq!(started["result"]["connection"], "runtime-1");
4966 assert_eq!(started["result"]["handle"]["runtime_id"], "svc_acp");
4967
4968 let terminal = service
4969 .handle_async(request(
4970 9,
4971 "harness.v1.runtimes.terminal_instructions",
4972 json!({"connection":"runtime-1"}),
4973 ))
4974 .await;
4975 let arguments = terminal["result"]["launch"]["arguments"]
4976 .as_array()
4977 .expect("hosted runtime should return terminal arguments");
4978 let endpoint_index = arguments
4979 .iter()
4980 .position(|value| value == "--endpoint")
4981 .expect("terminal command should use an opaque endpoint");
4982 let endpoint = LiveRuntimeEndpoint::parse(
4983 arguments[endpoint_index + 1]
4984 .as_str()
4985 .expect("endpoint argument should be text"),
4986 )
4987 .unwrap();
4988 assert!(!terminal.to_string().contains("Bearer"));
4989 let workspace = std::env::current_dir().unwrap();
4990 let receipt = resolve_live_runtime(
4991 &endpoint,
4992 &LiveRuntimeSource {
4993 harness: "codex".into(),
4994 session_id: "svc_acp".into(),
4995 workspace,
4996 },
4997 )
4998 .unwrap();
4999 let remote = crate::HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
5000 .await
5001 .unwrap();
5002 let mut attachment = crate::FrontendRuntime::attach(remote.as_ref(), 100)
5003 .await
5004 .unwrap();
5005
5006 let sent = service
5007 .handle_async(request(
5008 2,
5009 "harness.v1.runtimes.send_input",
5010 json!({"connection": "runtime-1", "text": "hi"}),
5011 ))
5012 .await;
5013 assert_eq!(sent["result"]["turn_id"], "3");
5014
5015 let mut events = Vec::new();
5016 for _ in 0..20 {
5017 events.extend(service.poll_runtimes().await);
5018 if events.len() >= 2 {
5019 break;
5020 }
5021 tokio::time::sleep(Duration::from_millis(2)).await;
5022 }
5023 assert!(events
5024 .iter()
5025 .any(|event| { event["params"]["event"]["kind"] == "session/update" }));
5026 assert!(events.iter().any(|event| {
5027 event["params"]["event"]["kind"] == "supercode/acp_request_completed"
5028 }));
5029
5030 let saw_editor_reply = tokio::time::timeout(Duration::from_secs(2), async {
5031 loop {
5032 let event = attachment.next_event().await.unwrap();
5033 if event.kind == "text_delta" && event.payload["text"] == "ok" {
5034 break;
5035 }
5036 }
5037 })
5038 .await;
5039 assert!(
5040 saw_editor_reply.is_ok(),
5041 "terminal should observe the editor-driven turn"
5042 );
5043
5044 crate::FrontendRuntime::submit(remote.as_ref(), "DRIVE FROM TERMINAL".into())
5045 .await
5046 .unwrap();
5047 let saw_terminal_reply = tokio::time::timeout(Duration::from_secs(2), async {
5048 loop {
5049 let event = attachment.next_event().await.unwrap();
5050 if event.kind == "text_delta" && event.payload["text"] == "from terminal" {
5051 break;
5052 }
5053 }
5054 })
5055 .await;
5056 assert!(
5057 saw_terminal_reply.is_ok(),
5058 "terminal should drive the same runtime"
5059 );
5060
5061 let closed = service
5062 .handle_async(request(
5063 3,
5064 "harness.v1.runtimes.close",
5065 json!({"connection": "runtime-1"}),
5066 ))
5067 .await;
5068 assert_eq!(closed["result"]["closed"], true);
5069 }
5070}