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" => Err(ServiceError::UnsupportedAction(
1096 "steer is not supported by this harness-native runtime adapter".into(),
1097 )),
1098 "harness.v1.runtimes.respond" => {
1099 let params = decode::<RuntimeRespondParams>(params)?;
1100 self.runtime_mut(¶ms.connection)?
1101 .respond(params.request_id, params.response)
1102 .await
1103 .map_err(operation)?;
1104 Ok(json!({}))
1105 }
1106 "harness.v1.runtimes.terminal_instructions" => {
1107 let params = decode::<RuntimeConnectionParams>(params)?;
1108 let launch = self
1109 .terminal_launches
1110 .get(¶ms.connection)
1111 .ok_or_else(|| {
1112 ServiceError::Operation(
1113 "this runtime is not hosted for terminal attachment".into(),
1114 )
1115 })?;
1116 Ok(json!({"launch":launch}))
1117 }
1118 "harness.v1.runtimes.close" => {
1119 let params = decode::<RuntimeConnectionParams>(params)?;
1120 let Some(mut runtime) = self.runtimes.remove(¶ms.connection) else {
1121 return Err(ServiceError::InvalidParams(format!(
1122 "unknown runtime connection `{}`",
1123 params.connection
1124 )));
1125 };
1126 self.terminal_launches.remove(¶ms.connection);
1127 self.runtime_sequences.remove(&runtime.handle().runtime_id);
1128 runtime.close().await.map_err(operation)?;
1129 Ok(json!({"closed": true}))
1130 }
1131 _ => Err(ServiceError::MethodNotFound),
1132 }
1133 }
1134
1135 #[cfg(feature = "adapter-api")]
1137 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
1138 let params = decode::<MessageSessionParams>(params)?;
1139 Ok(message_live_session(¶ms, &crate::claude_peer::ProcessCourierRunner).await)
1140 }
1141
1142 #[cfg(feature = "adapter-api")]
1143 fn harness_settings_call(
1144 &self,
1145 method: &str,
1146 params: Value,
1147 ) -> std::result::Result<Value, ServiceError> {
1148 let homes = crate::HarnessHomes::default();
1149 match method {
1150 "harness.v1.harnesses.settings" => {
1151 let params = decode::<HarnessSettingsParams>(params)?;
1152 let report = crate::inspect_harness_interop_settings(&homes, ¶ms.harness)
1153 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1154 serde_json::to_value(report)
1155 .map_err(|error| ServiceError::Operation(error.to_string()))
1156 }
1157 "harness.v1.harnesses.configure" => {
1158 let params = decode::<ConfigureHarnessParams>(params)?;
1159 let report = crate::configure_harness_interop_settings(
1160 &homes,
1161 ¶ms.harness,
1162 ¶ms.changes,
1163 params.expected_revision.as_deref(),
1164 )
1165 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1166 serde_json::to_value(report)
1167 .map_err(|error| ServiceError::Operation(error.to_string()))
1168 }
1169 _ => Err(ServiceError::MethodNotFound),
1170 }
1171 }
1172
1173 fn insert_runtime(
1174 &mut self,
1175 runtime: Box<dyn RuntimeConnection>,
1176 ) -> std::result::Result<Value, ServiceError> {
1177 let connection = format!("runtime-{}", self.next_runtime);
1178 self.next_runtime += 1;
1179 let handle = runtime.handle().clone();
1180 self.runtime_sequences
1181 .entry(handle.runtime_id.clone())
1182 .or_insert(0);
1183 self.runtimes.insert(connection.clone(), runtime);
1184 Ok(json!({"connection": connection, "handle": handle}))
1185 }
1186
1187 #[cfg(feature = "adapter-api")]
1188 async fn insert_hosted_runtime(
1189 &mut self,
1190 runtime: Box<dyn RuntimeConnection>,
1191 capabilities: crate::RuntimeCapabilities,
1192 workspace: PathBuf,
1193 ) -> std::result::Result<Value, ServiceError> {
1194 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities);
1195 let token: std::sync::Arc<str> = crate::server::generate_token().into();
1196 let server = crate::server::run_frontend_http(
1197 host.clone(),
1198 host.frontend_sender(),
1199 "127.0.0.1:0",
1200 token.clone(),
1201 )
1202 .await
1203 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1204 let source = LiveRuntimeSource {
1205 harness: connection.handle().harness.as_str().to_string(),
1206 session_id: connection.handle().runtime_id.clone(),
1207 workspace: workspace.clone(),
1208 };
1209 let registration = register_live_runtime(
1210 connection.handle().runtime_id.clone(),
1211 source.clone(),
1212 format!("http://{}", server.address()),
1213 token.to_string(),
1214 )
1215 .map_err(|error| ServiceError::Operation(error.to_string()))?;
1216 let endpoint = registration.endpoint().to_string();
1217 let launch = StructuredLaunch {
1218 cwd: workspace,
1219 program: std::env::current_exe()
1223 .ok()
1224 .map(|path| path.to_string_lossy().into_owned())
1225 .unwrap_or_else(|| "supercode".into()),
1226 arguments: vec![
1227 "harness".into(),
1228 "attach".into(),
1229 "--endpoint".into(),
1230 endpoint,
1231 "--harness".into(),
1232 source.harness,
1233 "--session".into(),
1234 source.session_id,
1235 ],
1236 env: BTreeMap::new(),
1237 };
1238 let lease = HostedRuntimeLease {
1239 connection,
1240 _host: host,
1241 _registration: registration,
1242 _server: server,
1243 };
1244 let opened = self.insert_runtime(Box::new(lease))?;
1245 let connection_id = opened["connection"]
1246 .as_str()
1247 .expect("insert_runtime returns a connection id")
1248 .to_string();
1249 self.terminal_launches.insert(connection_id, launch);
1250 Ok(opened)
1251 }
1252
1253 #[cfg(not(feature = "adapter-api"))]
1254 async fn insert_hosted_runtime(
1255 &mut self,
1256 runtime: Box<dyn RuntimeConnection>,
1257 _capabilities: crate::RuntimeCapabilities,
1258 _workspace: PathBuf,
1259 ) -> std::result::Result<Value, ServiceError> {
1260 self.insert_runtime(runtime)
1261 }
1262
1263 fn runtime_mut(
1264 &mut self,
1265 connection: &str,
1266 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
1267 self.runtimes.get_mut(connection).ok_or_else(|| {
1268 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
1269 })
1270 }
1271
1272 async fn inventory_call(
1273 &self,
1274 method: &str,
1275 params: Value,
1276 ) -> std::result::Result<Value, ServiceError> {
1277 let mut params = decode::<HarnessInventoryParams>(params)?;
1278 if method == "harness.v1.harnesses.probe" {
1279 let harness = params.harness.take().ok_or_else(|| {
1280 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
1281 })?;
1282 params.harnesses = vec![harness];
1283 }
1284 let selected = params
1285 .harnesses
1286 .iter()
1287 .map(HarnessId::as_str)
1288 .collect::<std::collections::BTreeSet<_>>();
1289 let supported = harness_support_registry()
1290 .harnesses
1291 .into_iter()
1292 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
1293 .collect::<Vec<_>>();
1294 if !params.harnesses.is_empty() && supported.len() != selected.len() {
1295 let known = supported
1296 .iter()
1297 .map(|harness| harness.id.as_str())
1298 .collect::<std::collections::BTreeSet<_>>();
1299 let missing = params
1300 .harnesses
1301 .iter()
1302 .filter(|id| !known.contains(id.as_str()))
1303 .map(HarnessId::as_str)
1304 .collect::<Vec<_>>();
1305 return Err(ServiceError::InvalidParams(format!(
1306 "unknown harness(es): {}",
1307 missing.join(", ")
1308 )));
1309 }
1310 let global_counts = params
1311 .include_sessions
1312 .then(|| self.session_counts(None, ¶ms.harnesses));
1313 let workspace_counts = params.include_sessions.then(|| {
1314 params
1315 .workspace
1316 .as_deref()
1317 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
1318 });
1319 let probes = supported.into_iter().map(|descriptor| {
1320 let global = global_counts
1321 .as_ref()
1322 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
1323 let workspace = workspace_counts
1324 .as_ref()
1325 .and_then(Option::as_ref)
1326 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
1327 self.probe_harness(descriptor, ¶ms, global, workspace)
1328 });
1329 let harnesses = futures::future::join_all(probes).await;
1330 serde_json::to_value(HarnessInventoryReport {
1331 probe: params.probe,
1332 workspace: params.workspace,
1333 harnesses,
1334 })
1335 .map_err(|error| ServiceError::Operation(error.to_string()))
1336 }
1337
1338 async fn probe_harness(
1339 &self,
1340 descriptor: crate::HarnessSupportDescriptor,
1341 params: &HarnessInventoryParams,
1342 global: Option<usize>,
1343 workspace: Option<usize>,
1344 ) -> LocalHarness {
1345 let launch = descriptor.runtime.default_launch.as_ref();
1346 let executable = launch.and_then(|launch| find_executable(&launch.program));
1347 let installed = executable.is_some();
1348 let version = if params.skip_versions {
1349 None
1350 } else {
1351 match executable.as_deref() {
1352 Some(path) => executable_version(path).await,
1353 None => None,
1354 }
1355 };
1356 let configured = auth_evidence(descriptor.id.as_str());
1357 let mut auth = if configured {
1358 HarnessAuthState::Configured
1359 } else {
1360 HarnessAuthState::Unknown
1361 };
1362 let mut runtime = if installed {
1363 HarnessRuntimeState::Degraded
1364 } else {
1365 HarnessRuntimeState::Unavailable
1366 };
1367 let mut reason = (!installed).then(|| {
1368 format!(
1369 "{} is supported but `{}` was not found on PATH",
1370 descriptor.display_name,
1371 launch
1372 .map(|launch| launch.program.as_str())
1373 .unwrap_or("executable")
1374 )
1375 });
1376 let mut repair = (!installed).then(|| {
1377 format!(
1378 "Install {} and ensure `{}` is on PATH.",
1379 descriptor.display_name,
1380 launch
1381 .map(|launch| launch.program.as_str())
1382 .unwrap_or("its executable")
1383 )
1384 });
1385
1386 if installed && params.probe == HarnessProbeLevel::Handshake {
1387 let backend_params = RuntimeBackendParams {
1388 harness: descriptor.id.clone(),
1389 protocol: None,
1390 launch: None,
1391 base_url: None,
1392 policy: RuntimePolicy::Default,
1393 };
1394 match runtime_backend(&backend_params) {
1395 Ok(backend) => {
1396 let cwd = params
1397 .workspace
1398 .clone()
1399 .or_else(|| std::env::current_dir().ok())
1400 .unwrap_or_else(|| PathBuf::from("."));
1401 let isolated = descriptor
1402 .runtime
1403 .default_launch
1404 .clone()
1405 .and_then(|launch| {
1406 IsolatedProbeHome::new(descriptor.id.as_str(), launch).ok()
1407 });
1408 let Some(isolated) = isolated else {
1409 reason = Some(
1410 "No-prompt runtime handshake could not create its isolated harness home."
1411 .into(),
1412 );
1413 repair = Some(
1414 "Check temporary-directory permissions, then run the handshake probe again."
1415 .into(),
1416 );
1417 return LocalHarness {
1418 id: descriptor.id,
1419 display_name: descriptor.display_name,
1420 supported: true,
1421 installed,
1422 executable: executable.map(|path| path.to_string_lossy().into_owned()),
1423 version,
1424 auth,
1425 runtime,
1426 protocol: descriptor.runtime.protocol,
1427 capabilities: descriptor.runtime.capabilities.clone(),
1428 effective_capabilities: descriptor.runtime.capabilities,
1429 sessions: HarnessSessionCounts { global, workspace },
1430 reason,
1431 repair,
1432 };
1433 };
1434 match tokio::time::timeout(
1435 Duration::from_secs(30),
1436 backend.start(RuntimeStartRequest {
1437 cwd,
1438 launch: Some(isolated.launch.clone()),
1439 }),
1440 )
1441 .await
1442 {
1443 Ok(Ok(mut connection)) => {
1444 match stabilize_handshake(connection.as_mut()).await {
1445 Ok(()) => {
1446 auth = HarnessAuthState::Ready;
1447 runtime = HarnessRuntimeState::Ready;
1448 reason = Some(
1449 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
1450 .into(),
1451 );
1452 repair = None;
1453 }
1454 Err(message) => {
1455 auth = if looks_like_auth_error(&message) {
1456 HarnessAuthState::Required
1457 } else if configured {
1458 HarnessAuthState::Configured
1459 } else {
1460 HarnessAuthState::Unknown
1461 };
1462 reason = Some(format!(
1463 "No-prompt runtime handshake became unhealthy during startup: {message}"
1464 ));
1465 repair = Some(if auth == HarnessAuthState::Required {
1466 format!(
1467 "Run `{}` interactively once and complete sign-in, then probe again.",
1468 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1469 )
1470 } else {
1471 "Run the harness directly to inspect its startup failure, then probe again."
1472 .into()
1473 });
1474 }
1475 }
1476 let _ =
1477 tokio::time::timeout(Duration::from_secs(3), connection.close())
1478 .await;
1479 }
1480 Ok(Err(error)) => {
1481 let message = truncate_text(&error.to_string(), 500);
1482 auth = if looks_like_auth_error(&message) {
1483 HarnessAuthState::Required
1484 } else if configured {
1485 HarnessAuthState::Configured
1486 } else {
1487 HarnessAuthState::Unknown
1488 };
1489 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
1490 repair = Some(if auth == HarnessAuthState::Required {
1491 format!(
1492 "Run `{}` interactively once and complete sign-in, then probe again.",
1493 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1494 )
1495 } else {
1496 "Check the harness installation and run the handshake probe again."
1497 .into()
1498 });
1499 }
1500 Err(_) => {
1501 reason = Some(
1502 "No-prompt runtime handshake timed out after 30 seconds.".into(),
1503 );
1504 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
1505 }
1506 }
1507 let _ = isolated.cleanup();
1516 tokio::time::sleep(Duration::from_millis(250)).await;
1517 if let Err(error) = isolated.cleanup() {
1518 auth = if configured {
1519 HarnessAuthState::Configured
1520 } else {
1521 HarnessAuthState::Unknown
1522 };
1523 runtime = HarnessRuntimeState::Degraded;
1524 reason = Some(format!(
1525 "No-prompt runtime handshake could not remove its isolated harness home: {error}"
1526 ));
1527 repair = Some(
1528 "Check temporary-directory permissions, remove the reported disposable probe home, then run the handshake again."
1529 .into(),
1530 );
1531 }
1532 }
1533 Err(error) => {
1534 reason = Some(error_message(error));
1535 }
1536 }
1537 } else if installed && configured {
1538 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
1539 } else if installed {
1540 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
1541 repair =
1542 Some(format!(
1543 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
1544 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1545 ));
1546 }
1547
1548 let effective_capabilities = if installed {
1549 descriptor.runtime.capabilities.clone()
1550 } else {
1551 unavailable_capabilities()
1552 };
1553 LocalHarness {
1554 id: descriptor.id,
1555 display_name: descriptor.display_name,
1556 supported: true,
1557 installed,
1558 executable: executable.map(|path| path.to_string_lossy().into_owned()),
1559 version,
1560 auth,
1561 runtime,
1562 protocol: descriptor.runtime.protocol,
1563 capabilities: descriptor.runtime.capabilities,
1564 effective_capabilities,
1565 sessions: HarnessSessionCounts { global, workspace },
1566 reason,
1567 repair,
1568 }
1569 }
1570
1571 fn session_counts(
1572 &self,
1573 workspace: Option<&Path>,
1574 harnesses: &[HarnessId],
1575 ) -> BTreeMap<String, usize> {
1576 let mut counts = BTreeMap::new();
1577 for session in self
1578 .catalog
1579 .discover(&DiscoveryQuery {
1580 workspace: workspace.map(Path::to_path_buf),
1581 harnesses: harnesses.to_vec(),
1582 ..DiscoveryQuery::default()
1583 })
1584 .unwrap_or_default()
1585 {
1586 *counts
1587 .entry(session.locator.harness.as_str().to_string())
1588 .or_insert(0) += 1;
1589 }
1590 counts
1591 }
1592}
1593
1594#[async_trait::async_trait]
1595impl SdkService for HarnessSessionService {
1596 fn capabilities(&self) -> SdkCapabilities {
1597 SdkCapabilities::default()
1598 }
1599
1600 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
1601 if request.operation == SdkOperation::Events {
1602 let events = self
1603 .poll_sdk_events()
1604 .await
1605 .into_iter()
1606 .map(|(_, event)| event)
1607 .collect::<Vec<_>>();
1608 return serde_json::to_value(events).map_err(|error| {
1609 SdkError::new(
1610 SdkErrorCode::Execution,
1611 request.operation,
1612 error.to_string(),
1613 )
1614 });
1615 }
1616 let method = request
1617 .operation
1618 .method()
1619 .ok_or_else(|| SdkError::unsupported(request.operation))?;
1620 let result = match request.operation {
1621 SdkOperation::Discover | SdkOperation::Load | SdkOperation::Export => {
1622 self.call(method, request.params)
1623 }
1624 SdkOperation::Start
1625 | SdkOperation::Resume
1626 | SdkOperation::Input
1627 | SdkOperation::Interrupt
1628 | SdkOperation::Steer
1629 | SdkOperation::Respond
1630 | SdkOperation::Close => self.runtime_call(method, request.params).await,
1631 SdkOperation::Events => unreachable!("handled before method dispatch"),
1632 };
1633 result.map_err(|error| sdk_error(request.operation, error))
1634 }
1635
1636 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
1637 Ok(self
1638 .poll_sdk_events()
1639 .await
1640 .into_iter()
1641 .map(|(_, event)| event)
1642 .collect())
1643 }
1644}
1645
1646#[cfg(feature = "adapter-api")]
1647struct HostedRuntimeLease {
1648 connection: HostedHarnessConnection,
1649 _host: std::sync::Arc<HostedHarnessRuntime>,
1650 _registration: LiveRuntimeRegistration,
1651 _server: crate::server::FrontendHttpServer,
1652}
1653
1654#[async_trait::async_trait]
1655#[cfg(feature = "adapter-api")]
1656impl RuntimeConnection for HostedRuntimeLease {
1657 fn handle(&self) -> &crate::RuntimeHandle {
1658 self.connection.handle()
1659 }
1660
1661 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
1662 self.connection.send_input(input).await
1663 }
1664
1665 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
1666 self.connection.next_event().await
1667 }
1668
1669 async fn interrupt(&mut self) -> crate::Result<()> {
1670 self.connection.interrupt().await
1671 }
1672
1673 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
1674 self.connection.respond(request_id, response).await
1675 }
1676
1677 async fn close(&mut self) -> crate::Result<()> {
1678 self.connection.close().await
1679 }
1680}
1681
1682async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
1683 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
1684 loop {
1685 let now = tokio::time::Instant::now();
1686 if now >= deadline {
1687 return Ok(());
1688 }
1689 match tokio::time::timeout(deadline - now, connection.next_event()).await {
1690 Err(_) => return Ok(()),
1691 Ok(Ok(Some(event))) => {
1692 if let Some(message) = handshake_event_failure(&event) {
1693 return Err(truncate_text(&message, 500));
1694 }
1695 }
1696 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
1697 Ok(Err(error)) => return Err(error.to_string()),
1698 }
1699 }
1700}
1701
1702fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
1703 let detail = event
1704 .payload
1705 .get("message")
1706 .or_else(|| event.payload.get("line"))
1707 .and_then(Value::as_str)
1708 .unwrap_or(event.kind.as_str());
1709 match event.kind.as_str() {
1710 "transport_closed" => Some("runtime transport closed during startup".into()),
1711 "transport_error" => Some(format!("runtime transport error: {detail}")),
1712 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
1713 _ => None,
1718 }
1719}
1720
1721fn projected_session_result(session: &Session, options: &SessionLoadOptions) -> Value {
1722 let total_messages = session.messages.len();
1723 let (offset, end) = projected_message_window(total_messages, options);
1724 json!({
1725 "session": projected_session_json(session, options),
1726 "summary": projected_session_summary(session, options),
1727 "window": {
1728 "has_more": offset > 0 || end < total_messages,
1729 "has_newer": end < total_messages,
1730 "has_older": offset > 0,
1731 "newer_items": normalized_item_count(&session.messages[end..]),
1732 "offset": offset,
1733 "older_items": normalized_item_count(&session.messages[..offset]),
1734 "returned": end.saturating_sub(offset),
1735 "total_messages": total_messages,
1736 }
1737 })
1738}
1739
1740fn normalized_item_count(messages: &[crate::ChatMessage]) -> usize {
1741 messages
1742 .iter()
1743 .map(|message| {
1744 let conversation = usize::from(
1745 matches!(message.role, Role::Assistant | Role::User)
1746 && message_has_content(message),
1747 );
1748 let tool_result =
1749 usize::from(message.role == Role::Tool && message_has_content(message));
1750 conversation + tool_result + message.tool_calls().len()
1751 })
1752 .sum()
1753}
1754
1755fn projected_session_summary(session: &Session, options: &SessionLoadOptions) -> Value {
1756 let mut conversational = session.messages.iter().filter(|message| {
1757 matches!(message.role, Role::Assistant | Role::User) && message_has_content(message)
1758 });
1759 let first_message = conversational.clone().next();
1760 let last_message = conversational.next_back();
1761 let mut assistant = session
1762 .messages
1763 .iter()
1764 .filter(|message| message.role == Role::Assistant && message_has_content(message));
1765 let first_assistant_message = assistant.clone().next();
1766 let last_assistant_message = assistant.next_back();
1767 let end_of_turn = session
1768 .messages
1769 .iter()
1770 .rev()
1771 .find(|message| message.role != Role::System)
1772 .is_some_and(|message| {
1773 message.role == Role::Assistant
1774 && message_has_content(message)
1775 && message.tool_calls().is_empty()
1776 });
1777 let project = |message: Option<&crate::ChatMessage>| {
1778 message.map(|message| project_inline_media(message_json(message), options))
1779 };
1780 json!({
1781 "end_of_turn": end_of_turn,
1782 "first_assistant_message": project(first_assistant_message),
1783 "first_message": project(first_message),
1784 "last_assistant_message": project(last_assistant_message),
1785 "last_assistant_text": last_assistant_message.map(message_text).unwrap_or_default(),
1786 "last_message": project(last_message),
1787 })
1788}
1789
1790fn message_has_content(message: &crate::ChatMessage) -> bool {
1791 message
1792 .content
1793 .as_deref()
1794 .is_some_and(|content| !content.trim().is_empty())
1795 || message
1796 .content_parts
1797 .as_ref()
1798 .is_some_and(|parts| !parts.is_empty())
1799}
1800
1801fn message_text(message: &crate::ChatMessage) -> String {
1802 if let Some(content) = &message.content {
1803 return content.clone();
1804 }
1805 message
1806 .content_parts
1807 .as_ref()
1808 .into_iter()
1809 .flatten()
1810 .filter_map(|part| part.get("text").and_then(Value::as_str))
1811 .collect::<Vec<_>>()
1812 .join("\n")
1813}
1814
1815fn projected_session_json(session: &Session, options: &SessionLoadOptions) -> Value {
1816 let (offset, end) = projected_message_window(session.messages.len(), options);
1817 let messages = session.messages[offset..end]
1818 .iter()
1819 .map(|message| project_inline_media(message_json(message), options))
1820 .collect::<Vec<_>>();
1821 let subagents = if options.include_subagents.unwrap_or(true) {
1822 let subagent_options = SessionLoadOptions {
1827 message_limit: None,
1828 message_offset: None,
1829 message_tail: None,
1830 ..options.clone()
1831 };
1832 session
1833 .subagents
1834 .iter()
1835 .map(|subagent| projected_session_json(subagent, &subagent_options))
1836 .collect::<Vec<_>>()
1837 } else {
1838 Vec::new()
1839 };
1840 json!({
1841 "source": match session.meta.source {
1842 SessionSource::ClaudeCode => "claude_code",
1843 SessionSource::Codex => "codex",
1844 SessionSource::Gemini => "gemini",
1845 SessionSource::Goose => "goose",
1846 SessionSource::Grok => "grok",
1847 SessionSource::Native => "native",
1848 SessionSource::OpenCode => "opencode",
1849 SessionSource::Pi => "pi",
1850 },
1851 "session_id": session.meta.session_id,
1852 "model": session.meta.model,
1853 "cwd": session.meta.cwd,
1854 "system_prompt": session.meta.system_prompt,
1855 "agent_id": session.meta.agent_id,
1856 "parent_tool_use_id": session.meta.parent_tool_use_id,
1857 "lineage": session.meta.lineage,
1858 "messages": messages,
1859 "subagents": subagents,
1860 "raw_record_count": session.raw.len(),
1861 "parse_error_lines": session.parse_error_lines,
1862 })
1863}
1864
1865fn projected_message_window(total: usize, options: &SessionLoadOptions) -> (usize, usize) {
1866 if let Some(tail) = options.message_tail {
1867 return (total.saturating_sub(tail), total);
1868 }
1869 let offset = options.message_offset.unwrap_or(0).min(total);
1870 let end = options
1871 .message_limit
1872 .map(|limit| offset.saturating_add(limit).min(total))
1873 .unwrap_or(total);
1874 (offset, end)
1875}
1876
1877fn project_inline_media(mut message: Value, options: &SessionLoadOptions) -> Value {
1878 let Some(parts) = message.get_mut("content").and_then(Value::as_array_mut) else {
1879 return message;
1880 };
1881 for part in parts {
1882 let Some(url) = part
1883 .get("image_url")
1884 .and_then(|image| image.get("url"))
1885 .and_then(Value::as_str)
1886 else {
1887 continue;
1888 };
1889 let Some(rest) = url.strip_prefix("data:") else {
1890 continue;
1891 };
1892 let Some((media_type, encoded)) = rest.split_once(";base64,") else {
1893 continue;
1894 };
1895 let padding = usize::from(encoded.ends_with('=')) + usize::from(encoded.ends_with("=="));
1896 let decoded_bytes = encoded.len().saturating_mul(3) / 4;
1897 let decoded_bytes = decoded_bytes.saturating_sub(padding);
1898 let should_elide = matches!(options.inline_media, InlineMediaMode::Metadata)
1899 || options
1900 .max_inline_media_bytes
1901 .is_some_and(|limit| decoded_bytes > limit);
1902 if should_elide {
1903 *part = json!({
1904 "type": "media_reference",
1905 "media_type": media_type,
1906 "encoding": "base64",
1907 "encoded_bytes": encoded.len(),
1908 "decoded_bytes": decoded_bytes,
1909 "omitted": true,
1910 });
1911 }
1912 }
1913 message
1914}
1915
1916#[derive(Deserialize)]
1917struct LocatorParams {
1918 locator: SessionLocator,
1919 #[serde(default)]
1932 fidelity: Option<Fidelity>,
1933 #[serde(default)]
1936 view: Option<SessionReadView>,
1937}
1938
1939#[derive(Deserialize)]
1940struct SessionReadView {
1941 #[serde(default)]
1944 tail_messages: Option<usize>,
1945 #[serde(default)]
1948 include_subagents: bool,
1949 #[serde(default)]
1951 display_history: bool,
1952 #[serde(default)]
1955 max_message_chars: Option<usize>,
1956}
1957
1958impl LocatorParams {
1959 fn read_fidelity(&self) -> Fidelity {
1960 self.fidelity.unwrap_or(Fidelity::Semantic)
1961 }
1962
1963 fn include_subagents(&self) -> bool {
1964 self.view
1965 .as_ref()
1966 .map(|view| view.include_subagents)
1967 .unwrap_or(true)
1968 }
1969
1970 fn tail_messages(&self) -> Option<usize> {
1971 self.view
1972 .as_ref()
1973 .and_then(|view| view.tail_messages)
1974 .map(|limit| limit.clamp(1, 5_000))
1975 }
1976
1977 fn display_history(&self) -> bool {
1978 self.view.as_ref().is_some_and(|view| view.display_history)
1979 }
1980
1981 fn max_message_chars(&self) -> Option<usize> {
1982 self.view
1983 .as_ref()
1984 .and_then(|view| view.max_message_chars)
1985 .map(|limit| limit.clamp(256, 64_000))
1986 }
1987
1988 fn bound_session(&self, session: &mut Session) {
1989 bound_session_view(session, self.tail_messages(), self.max_message_chars());
1990 }
1991}
1992
1993#[derive(Debug, Clone, Copy, Default, Deserialize)]
1994#[serde(rename_all = "snake_case")]
1995enum InlineMediaMode {
1996 #[default]
1997 Full,
1998 Metadata,
1999}
2000
2001#[derive(Debug, Clone, Default, Deserialize)]
2002#[serde(default)]
2003struct SessionLoadOptions {
2004 include_subagents: Option<bool>,
2005 inline_media: InlineMediaMode,
2006 max_inline_media_bytes: Option<usize>,
2007 message_limit: Option<usize>,
2008 message_offset: Option<usize>,
2009 message_tail: Option<usize>,
2010}
2011
2012impl SessionLoadOptions {
2013 fn validate(&self) -> std::result::Result<(), ServiceError> {
2014 if self.message_tail.is_some()
2015 && (self.message_limit.is_some() || self.message_offset.is_some())
2016 {
2017 return Err(ServiceError::InvalidParams(
2018 "sessions.load options.message_tail cannot be combined with message_limit or message_offset"
2019 .into(),
2020 ));
2021 }
2022 Ok(())
2023 }
2024}
2025
2026#[derive(Deserialize)]
2027struct LoadSessionParams {
2028 #[serde(flatten)]
2029 read: LocatorParams,
2030 #[serde(default)]
2031 options: Option<SessionLoadOptions>,
2032}
2033
2034#[derive(Deserialize)]
2035struct UnfollowParams {
2036 subscription: String,
2037}
2038
2039#[derive(Deserialize)]
2040struct ActivitySubscribeParams {
2041 locators: Vec<SessionLocator>,
2042 #[serde(default)]
2043 homes: crate::HarnessHomes,
2044}
2045
2046#[derive(Deserialize)]
2047struct MessageSessionParams {
2048 locator: SessionLocator,
2049 text: String,
2050 #[serde(default)]
2053 homes: crate::HarnessHomes,
2054}
2055
2056#[derive(Deserialize)]
2057#[serde(deny_unknown_fields)]
2058struct HarnessSettingsParams {
2059 harness: String,
2060}
2061
2062#[derive(Deserialize)]
2063#[serde(deny_unknown_fields)]
2064struct ConfigureHarnessParams {
2065 harness: String,
2066 #[serde(default)]
2067 changes: Vec<crate::HarnessSettingChange>,
2068 #[serde(default)]
2069 expected_revision: Option<String>,
2070}
2071
2072fn claude_inbound_controls_or_error(homes: &crate::HarnessHomes) -> (Value, Value) {
2073 match crate::inspect_harness_interop_settings(homes, HarnessId::CLAUDE_CODE) {
2074 Ok(report) => (
2075 serde_json::to_value(report).unwrap_or(Value::Null),
2076 Value::Null,
2077 ),
2078 Err(error) => (
2079 Value::Null,
2080 Value::String(format!(
2081 "Supercode could not inspect Claude Code inbound controls: {error}"
2082 )),
2083 ),
2084 }
2085}
2086
2087#[cfg(feature = "adapter-api")]
2099async fn message_live_session(
2100 params: &MessageSessionParams,
2101 runner: &dyn crate::claude_peer::CourierRunner,
2102) -> Value {
2103 if params.locator.harness.as_str() != HarnessId::CLAUDE_CODE {
2104 return json!({
2105 "delivered_to_bus": false,
2106 "refusal": {
2107 "reason": crate::claude_peer::ClaudePeerRefusal::HarnessUnsupported.as_str(),
2108 "message": format!(
2109 "`{}` does not publish a live-session registry; only claude-code sessions can be messaged in place",
2110 params.locator.harness.as_str()
2111 ),
2112 },
2113 });
2114 }
2115 let (inbound_controls, inbound_controls_error) =
2116 claude_inbound_controls_or_error(¶ms.homes);
2117 match crate::claude_peer::message_claude_peer(
2118 ¶ms.homes,
2119 ¶ms.locator.session_id,
2120 ¶ms.text,
2121 runner,
2122 )
2123 .await
2124 {
2125 Ok(delivery) => json!({
2126 "delivered_to_bus": true,
2127 "target": {
2128 "session_id": delivery.target.session_id,
2129 "name": delivery.target.name,
2130 "pid": delivery.target.pid,
2131 "cwd": delivery.target.cwd,
2132 "status": delivery.target.status.map(|status| status.as_str()),
2133 },
2134 "courier": {
2135 "model": crate::claude_peer::COURIER_MODEL,
2136 "report": delivery.courier_report,
2137 },
2138 "inbound_controls": inbound_controls,
2139 "inbound_controls_error": inbound_controls_error,
2140 }),
2141 Err(refusal) => json!({
2142 "delivered_to_bus": false,
2143 "refusal": {"reason": refusal.reason.as_str(), "message": refusal.message},
2144 "inbound_controls": inbound_controls,
2145 "inbound_controls_error": inbound_controls_error,
2146 }),
2147 }
2148}
2149
2150#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
2155struct FollowedSource {
2156 harness: String,
2157 session_id: String,
2158 reported: Option<String>,
2159}
2160
2161#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
2162struct ActivitySubscription {
2163 locators: Vec<SessionLocator>,
2164 homes: crate::HarnessHomes,
2165 reported: BTreeMap<(String, String), crate::SessionActivity>,
2166}
2167
2168fn peers_for_descriptors(
2169 descriptors: &[SessionDescriptor],
2170 homes: &HarnessHomes,
2171) -> Vec<crate::claude_peer::ClaudePeerSession> {
2172 if descriptors
2173 .iter()
2174 .any(|session| session.locator.harness.as_str() == HarnessId::CLAUDE_CODE)
2175 {
2176 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
2177 } else {
2178 Vec::new()
2179 }
2180}
2181
2182fn live_descriptor_value(
2188 session: &SessionDescriptor,
2189 peers: &[crate::claude_peer::ClaudePeerSession],
2190) -> std::result::Result<Value, ServiceError> {
2191 let mut value = serde_json::to_value(session)
2192 .map_err(|error| ServiceError::Operation(error.to_string()))?;
2193 if let Some(workspace) = &session.cwd {
2194 let source = LiveRuntimeSource {
2195 harness: session.locator.harness.as_str().to_string(),
2196 session_id: session.locator.session_id.clone(),
2197 workspace: workspace.clone(),
2198 };
2199 if let Some(endpoint) = discover_live_runtime(&source)
2200 .map_err(|error| ServiceError::Operation(error.to_string()))?
2201 {
2202 value["live_endpoint"] = json!(endpoint.as_str());
2203 }
2204 }
2205 if value.get("live_endpoint").is_none() {
2206 if let Some(peer) = peers.iter().find(|peer| {
2207 session.locator.harness.as_str() == HarnessId::CLAUDE_CODE
2208 && peer.session_id == session.locator.session_id
2209 }) {
2210 value["live_endpoint"] = json!(peer.endpoint().as_str());
2211 }
2212 }
2213 Ok(value)
2214}
2215
2216fn live_index_changes(
2217 changes: Vec<crate::session_index::SessionIndexChange>,
2218 homes: &HarnessHomes,
2219) -> std::result::Result<Vec<Value>, ServiceError> {
2220 use crate::session_index::SessionIndexChange;
2221 let has_claude = changes.iter().any(|change| match change {
2222 SessionIndexChange::Added { descriptor } | SessionIndexChange::Updated { descriptor } => {
2223 descriptor.locator.harness.as_str() == HarnessId::CLAUDE_CODE
2224 }
2225 SessionIndexChange::Removed { .. } => false,
2226 });
2227 let peers = if has_claude {
2228 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
2229 } else {
2230 Vec::new()
2231 };
2232 changes
2233 .into_iter()
2234 .map(|change| match change {
2235 SessionIndexChange::Added { descriptor } => Ok(json!({
2236 "kind": "added",
2237 "descriptor": live_descriptor_value(&descriptor, &peers)?,
2238 })),
2239 SessionIndexChange::Updated { descriptor } => Ok(json!({
2240 "kind": "updated",
2241 "descriptor": live_descriptor_value(&descriptor, &peers)?,
2242 })),
2243 SessionIndexChange::Removed { key } => Ok(json!({
2244 "kind": "removed",
2245 "key": key,
2246 })),
2247 })
2248 .collect()
2249}
2250
2251fn legacy_live_status(activity: &crate::SessionActivity) -> Option<&'static str> {
2252 use crate::{SessionPresence, SessionTurnState};
2253 match (activity.presence, activity.turn) {
2254 (SessionPresence::Persisted, _) => None,
2255 (SessionPresence::Running, SessionTurnState::Working) => Some("busy"),
2256 (SessionPresence::Running, SessionTurnState::Idle) => Some("idle"),
2257 (SessionPresence::Running, _) | (SessionPresence::ShuttingDown, _) => Some("running"),
2258 }
2259}
2260
2261#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
2262#[serde(rename_all = "kebab-case")]
2263enum TransferFormat {
2264 ClaudeCode,
2265 Codex,
2266 #[serde(rename = "opencode", alias = "open-code")]
2267 OpenCode,
2268 Pi,
2269 Grok,
2270 Gemini,
2271 Goose,
2272}
2273
2274impl TransferFormat {
2275 fn id(self) -> &'static str {
2276 match self {
2277 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
2278 Self::Codex => HarnessId::CODEX,
2279 Self::OpenCode => HarnessId::OPENCODE,
2280 Self::Pi => HarnessId::PI,
2281 Self::Grok => HarnessId::GROK,
2282 Self::Gemini => HarnessId::GEMINI,
2283 Self::Goose => HarnessId::GOOSE,
2284 }
2285 }
2286}
2287
2288impl From<TransferFormat> for SessionFormat {
2289 fn from(value: TransferFormat) -> Self {
2290 match value {
2291 TransferFormat::ClaudeCode => Self::ClaudeCode,
2292 TransferFormat::Codex => Self::Codex,
2293 TransferFormat::OpenCode => Self::OpenCode,
2294 TransferFormat::Pi => Self::Pi,
2295 TransferFormat::Grok => Self::Grok,
2296 TransferFormat::Gemini => Self::Gemini,
2297 TransferFormat::Goose => Self::Goose,
2298 }
2299 }
2300}
2301
2302#[derive(Deserialize)]
2303struct ImportSessionParams {
2304 source_harness: TransferFormat,
2305 content: String,
2306}
2307
2308#[derive(Deserialize)]
2309struct ExportSessionParams {
2310 locator: SessionLocator,
2311 target_harness: TransferFormat,
2312}
2313
2314#[derive(Deserialize)]
2315struct ReduceSessionParams {
2316 locator: SessionLocator,
2317 target_harness: TransferFormat,
2318 #[serde(default = "default_keep_last")]
2319 keep_last: usize,
2320}
2321
2322fn default_keep_last() -> usize {
2323 6
2324}
2325
2326#[derive(Deserialize)]
2327struct BranchSessionParams {
2328 locator: SessionLocator,
2329 #[serde(default)]
2330 target_harness: Option<TransferFormat>,
2331}
2332
2333#[derive(Deserialize)]
2334struct HandoffSessionParams {
2335 locator: SessionLocator,
2336 target_harness: TransferFormat,
2337 #[serde(default)]
2338 cwd: Option<PathBuf>,
2339}
2340
2341#[derive(Debug, Clone, Copy, Default, Deserialize)]
2342#[serde(rename_all = "snake_case")]
2343enum ResumePolicy {
2344 #[default]
2345 Default,
2346 Yolo,
2347}
2348
2349#[derive(Deserialize)]
2350struct ResumeInstructionsParams {
2351 locator: SessionLocator,
2352 #[serde(default)]
2353 cwd: Option<PathBuf>,
2354 #[serde(default)]
2355 policy: ResumePolicy,
2356}
2357
2358#[derive(Serialize)]
2359struct SessionArtifact {
2360 source_harness: HarnessId,
2361 target_harness: &'static str,
2362 session_id: Option<String>,
2363 content: String,
2364 suggested_filename: String,
2365 files: Vec<SessionArtifactFile>,
2366 fidelity: Fidelity,
2367 residue: Vec<String>,
2368}
2369
2370#[derive(Serialize)]
2371struct SessionArtifactFile {
2372 path: String,
2373 content: String,
2374 role: ArtifactFileRole,
2375}
2376
2377#[derive(Serialize)]
2378#[serde(rename_all = "snake_case")]
2379enum ArtifactFileRole {
2380 Primary,
2381 Subagent,
2382 Bundle,
2383 SourceRecovery,
2384}
2385
2386#[derive(Serialize)]
2387struct StructuredLaunch {
2388 cwd: PathBuf,
2389 program: String,
2390 arguments: Vec<String>,
2391 env: BTreeMap<String, String>,
2392}
2393
2394struct HandoffInstructions {
2395 launch: StructuredLaunch,
2396 materialize: Option<StructuredLaunch>,
2397 requires_materialization: bool,
2398 note: String,
2399}
2400
2401#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
2402#[serde(rename_all = "snake_case")]
2403enum HarnessProbeLevel {
2404 #[default]
2405 Passive,
2406 Handshake,
2407}
2408
2409#[derive(Default, Deserialize)]
2410#[serde(default)]
2411struct HarnessInventoryParams {
2412 harness: Option<HarnessId>,
2413 harnesses: Vec<HarnessId>,
2414 workspace: Option<PathBuf>,
2415 probe: HarnessProbeLevel,
2416 include_sessions: bool,
2417 skip_versions: bool,
2419}
2420
2421#[derive(Serialize)]
2422struct HarnessInventoryReport {
2423 probe: HarnessProbeLevel,
2424 workspace: Option<PathBuf>,
2425 harnesses: Vec<LocalHarness>,
2426}
2427
2428#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
2429#[serde(rename_all = "snake_case")]
2430enum HarnessAuthState {
2431 Ready,
2432 Configured,
2433 Required,
2434 Unknown,
2435}
2436
2437#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
2438#[serde(rename_all = "snake_case")]
2439enum HarnessRuntimeState {
2440 Ready,
2441 Degraded,
2442 Unavailable,
2443}
2444
2445#[derive(Serialize)]
2446struct HarnessSessionCounts {
2447 global: Option<usize>,
2448 workspace: Option<usize>,
2449}
2450
2451#[derive(Serialize)]
2452struct LocalHarness {
2453 id: HarnessId,
2454 display_name: String,
2455 supported: bool,
2456 installed: bool,
2457 executable: Option<String>,
2458 version: Option<String>,
2459 auth: HarnessAuthState,
2460 runtime: HarnessRuntimeState,
2461 protocol: String,
2462 capabilities: crate::RuntimeCapabilities,
2463 effective_capabilities: crate::RuntimeCapabilities,
2464 sessions: HarnessSessionCounts,
2465 reason: Option<String>,
2466 repair: Option<String>,
2467}
2468
2469#[derive(Clone, Deserialize)]
2470struct RuntimeBackendParams {
2471 harness: HarnessId,
2472 #[serde(default)]
2473 protocol: Option<String>,
2474 #[serde(default)]
2475 launch: Option<RuntimeLaunch>,
2476 #[serde(default)]
2477 base_url: Option<String>,
2478 #[serde(default)]
2479 policy: RuntimePolicy,
2480}
2481
2482#[derive(Debug, Clone, Copy, Default, Deserialize)]
2483#[serde(rename_all = "snake_case")]
2484enum RuntimePolicy {
2485 #[default]
2486 Default,
2487 Yolo,
2488}
2489
2490#[derive(Deserialize)]
2491struct RuntimeStartParams {
2492 #[serde(flatten)]
2493 backend: RuntimeBackendParams,
2494 cwd: PathBuf,
2495}
2496
2497#[derive(Deserialize)]
2498struct RuntimeAttachParams {
2499 #[serde(flatten)]
2500 backend: RuntimeBackendParams,
2501 runtime_id: String,
2502 #[serde(default)]
2503 cwd: Option<PathBuf>,
2504}
2505
2506#[derive(Deserialize)]
2507struct RuntimeConnectionParams {
2508 connection: String,
2509}
2510
2511#[derive(Deserialize)]
2512struct RuntimeInputParams {
2513 connection: String,
2514 text: String,
2515 #[serde(default)]
2516 image_urls: Vec<String>,
2517}
2518
2519const MAX_RUNTIME_IMAGES: usize = 4;
2520const MAX_RUNTIME_IMAGE_URL_BYTES: usize = 12 * 1024 * 1024;
2521const MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL: usize = 32 * 1024 * 1024;
2522
2523fn validate_runtime_image_urls(image_urls: Vec<String>) -> Result<Vec<String>, ServiceError> {
2524 if image_urls.len() > MAX_RUNTIME_IMAGES {
2525 return Err(ServiceError::InvalidParams(format!(
2526 "a runtime prompt accepts at most {MAX_RUNTIME_IMAGES} images"
2527 )));
2528 }
2529 let mut total = 0usize;
2530 for url in &image_urls {
2531 if !(url.starts_with("data:image/")
2532 || url.starts_with("https://")
2533 || url.starts_with("http://"))
2534 {
2535 return Err(ServiceError::InvalidParams(
2536 "runtime images must be image data URLs or HTTP(S) URLs".into(),
2537 ));
2538 }
2539 if url.len() > MAX_RUNTIME_IMAGE_URL_BYTES {
2540 return Err(ServiceError::InvalidParams(format!(
2541 "one runtime image exceeds the {MAX_RUNTIME_IMAGE_URL_BYTES}-byte encoded limit"
2542 )));
2543 }
2544 total = total.saturating_add(url.len());
2545 }
2546 if total > MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL {
2547 return Err(ServiceError::InvalidParams(format!(
2548 "runtime images exceed the {MAX_RUNTIME_IMAGE_URL_BYTES_TOTAL}-byte encoded total limit"
2549 )));
2550 }
2551 Ok(image_urls)
2552}
2553
2554#[derive(Deserialize)]
2555struct RuntimeRespondParams {
2556 connection: String,
2557 request_id: Value,
2558 response: Value,
2559}
2560
2561fn default_reduction_store_root() -> PathBuf {
2562 if let Some(root) = std::env::var_os("SUPERCODE_HOME") {
2563 return PathBuf::from(root).join("sessions");
2564 }
2565 if let Some(home) = std::env::var_os("HOME") {
2566 return PathBuf::from(home).join(".supercode").join("sessions");
2567 }
2568 PathBuf::from(".supercode").join("sessions")
2569}
2570
2571fn messages_jsonl(messages: &[crate::ChatMessage]) -> std::result::Result<String, ServiceError> {
2572 let mut output = String::new();
2573 for message in messages {
2574 output.push_str(
2575 &serde_json::to_string(message)
2576 .map_err(|error| ServiceError::Operation(error.to_string()))?,
2577 );
2578 output.push('\n');
2579 }
2580 Ok(output)
2581}
2582
2583fn parse_messages_jsonl(
2584 content: &str,
2585) -> std::result::Result<Vec<crate::ChatMessage>, ServiceError> {
2586 content
2587 .lines()
2588 .enumerate()
2589 .filter(|(_, line)| !line.trim().is_empty())
2590 .map(|(index, line)| {
2591 serde_json::from_str::<crate::ChatMessage>(line).map_err(|error| {
2592 ServiceError::Operation(format!(
2593 "reduced transcript line {} is invalid: {error}",
2594 index + 1
2595 ))
2596 })
2597 })
2598 .collect()
2599}
2600
2601fn reduced_bootstrap_prompt(
2602 source: &SessionLocator,
2603 target: TransferFormat,
2604 view_jsonl: &str,
2605 sidecar_path: &Path,
2606 reduction_log_path: &Path,
2607) -> String {
2608 format!(
2609 "Continue the work from this losslessly reduced {source_harness} session in {target_harness}.\n\
2610 \n\
2611 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\
2612 \n\
2613 <supercode-reduced-session source-session=\"{source_id}\">\n\
2614 {view_jsonl}\
2615 </supercode-reduced-session>\n\
2616 \n\
2617 Resume from the latest unresolved user request and preserve the source session's decisions and constraints.",
2618 source_harness = source.harness.as_str(),
2619 target_harness = target.id(),
2620 sidecar = sidecar_path.display(),
2621 log = reduction_log_path.display(),
2622 source_id = source.session_id,
2623 )
2624}
2625
2626fn session_artifact(
2627 locator: &SessionLocator,
2628 session: &Session,
2629 target: TransferFormat,
2630) -> std::result::Result<SessionArtifact, ServiceError> {
2631 session_artifact_with_id(locator, session, target, None)
2632}
2633
2634fn session_artifact_with_id(
2635 locator: &SessionLocator,
2636 session: &Session,
2637 target: TransferFormat,
2638 target_session_id: Option<&str>,
2639) -> std::result::Result<SessionArtifact, ServiceError> {
2640 let format: SessionFormat = target.into();
2641 let diagonal = format.source() == session.meta.source;
2642 let has_appended_turns = session
2643 .imported_message_count
2644 .is_some_and(|imported| imported < session.messages.len());
2645 let content = if let Some(id) = target_session_id {
2646 if diagonal && format != SessionFormat::OpenCode {
2647 session
2648 .to_jsonl_spliced(format, Some(id))
2649 .map_err(operation)?
2650 } else {
2651 let mut rewritten = session.clone();
2652 rewritten.meta.session_id = Some(id.to_string());
2653 rewritten.to_jsonl(format).map_err(operation)?
2654 }
2655 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
2656 session.raw_verbatim()
2657 } else if diagonal {
2658 session.to_jsonl_spliced(format, None).map_err(operation)?
2659 } else {
2660 session.to_jsonl(format).map_err(operation)?
2661 };
2662 let stem = sanitize_filename(
2663 target_session_id
2664 .or(session.meta.session_id.as_deref())
2665 .unwrap_or(&locator.session_id),
2666 );
2667 let suggested_filename = if diagonal && target == TransferFormat::Grok {
2668 "chat_history.jsonl".to_string()
2669 } else if target == TransferFormat::Goose {
2670 format!("{stem}.goose.json")
2671 } else {
2672 format!("{stem}.{}.jsonl", target.id())
2673 };
2674 let mut files = vec![SessionArtifactFile {
2675 path: suggested_filename.clone(),
2676 content: content.clone(),
2677 role: ArtifactFileRole::Primary,
2678 }];
2679 if target == TransferFormat::ClaudeCode {
2680 let bundle_stem = Path::new(&suggested_filename)
2681 .file_stem()
2682 .and_then(|stem| stem.to_str())
2683 .unwrap_or(&stem);
2684 let mut child_paths = BTreeSet::new();
2685 for (index, subagent) in session.subagents.iter().enumerate() {
2686 let agent_id = subagent
2687 .meta
2688 .agent_id
2689 .as_deref()
2690 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
2691 .map(sanitize_filename)
2692 .filter(|id| !id.is_empty())
2693 .unwrap_or_else(|| format!("subagent-{}", index + 1));
2694 let child_has_appended_turns = subagent
2695 .imported_message_count
2696 .is_some_and(|imported| imported < subagent.messages.len());
2697 let child_content = if target_session_id.is_none()
2698 && subagent.meta.source == SessionSource::ClaudeCode
2699 && subagent.raw_is_verbatim
2700 && !child_has_appended_turns
2701 {
2702 subagent.raw_verbatim()
2703 } else if subagent.meta.source == SessionSource::ClaudeCode {
2704 subagent
2705 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
2706 .map_err(operation)?
2707 } else {
2708 let mut child = subagent.clone();
2709 if let Some(id) = target_session_id {
2710 child.meta.session_id = Some(id.to_string());
2711 }
2712 child
2713 .to_jsonl(SessionFormat::ClaudeCode)
2714 .map_err(operation)?
2715 };
2716 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
2717 if !child_paths.insert(path.clone()) {
2718 return Err(ServiceError::Operation(format!(
2719 "Claude subagent ids collide at artifact path `{path}`"
2720 )));
2721 }
2722 files.push(SessionArtifactFile {
2723 path,
2724 content: child_content,
2725 role: ArtifactFileRole::Subagent,
2726 });
2727 }
2728 }
2729 if diagonal && target == TransferFormat::Grok {
2730 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
2731 }
2732 if !diagonal || !session.raw_is_verbatim {
2733 files.push(SessionArtifactFile {
2734 path: "recovery/source.supercode.jsonl".into(),
2735 content: session.to_native_jsonl(),
2736 role: ArtifactFileRole::SourceRecovery,
2737 });
2738 for (index, subagent) in session.subagents.iter().enumerate() {
2739 let id = subagent
2740 .meta
2741 .agent_id
2742 .as_deref()
2743 .map(sanitize_filename)
2744 .unwrap_or_else(|| format!("subagent-{}", index + 1));
2745 files.push(SessionArtifactFile {
2746 path: format!("recovery/subagents/{id}.supercode.jsonl"),
2747 content: subagent.to_native_jsonl(),
2748 role: ArtifactFileRole::SourceRecovery,
2749 });
2750 }
2751 }
2752 if !diagonal && session.meta.source == SessionSource::Grok {
2753 append_grok_bundle_files(
2754 locator,
2755 "recovery/grok/",
2756 ArtifactFileRole::SourceRecovery,
2757 &mut files,
2758 )?;
2759 }
2760 let (fidelity, residue) = if diagonal
2761 && target_session_id.is_none()
2762 && session.raw_is_verbatim
2763 && !has_appended_turns
2764 {
2765 (Fidelity::ByteLossless, Vec::new())
2766 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
2767 (
2768 Fidelity::ValueLossless,
2769 vec![if target_session_id.is_some() {
2770 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
2771 } else {
2772 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
2773 }],
2774 )
2775 } else {
2776 (
2777 Fidelity::Semantic,
2778 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
2779 )
2780 };
2781 Ok(SessionArtifact {
2782 source_harness: locator.harness.clone(),
2783 target_harness: target.id(),
2784 session_id: target_session_id
2785 .map(str::to_string)
2786 .or_else(|| session.meta.session_id.clone()),
2787 content,
2788 suggested_filename,
2789 files,
2790 fidelity,
2791 residue,
2792 })
2793}
2794
2795fn append_grok_bundle_files(
2796 locator: &SessionLocator,
2797 prefix: &str,
2798 role: ArtifactFileRole,
2799 files: &mut Vec<SessionArtifactFile>,
2800) -> std::result::Result<(), ServiceError> {
2801 let primary = locator.storage.path();
2802 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
2803 return Err(ServiceError::Operation(format!(
2804 "Grok bundle locator must name chat_history.jsonl, got {}",
2805 primary.display()
2806 )));
2807 }
2808 let parent = primary.parent().ok_or_else(|| {
2809 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
2810 })?;
2811 for name in ["summary.json", "updates.jsonl"] {
2812 let path = parent.join(name);
2813 let metadata = match std::fs::symlink_metadata(&path) {
2814 Ok(metadata) => metadata,
2815 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
2816 Err(error) => return Err(ServiceError::Operation(error.to_string())),
2817 };
2818 if metadata.file_type().is_symlink() || !metadata.is_file() {
2819 return Err(ServiceError::Operation(format!(
2820 "refusing non-regular Grok bundle member {}",
2821 path.display()
2822 )));
2823 }
2824 let content = std::fs::read_to_string(&path).map_err(|error| {
2825 ServiceError::Operation(format!(
2826 "Grok bundle member {} is not representable as UTF-8: {error}",
2827 path.display()
2828 ))
2829 })?;
2830 files.push(SessionArtifactFile {
2831 path: format!("{prefix}{name}"),
2832 content,
2833 role: match role {
2834 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
2835 _ => ArtifactFileRole::SourceRecovery,
2836 },
2837 });
2838 }
2839 Ok(())
2840}
2841
2842fn handoff_artifact(
2843 locator: &SessionLocator,
2844 session: &Session,
2845 target: TransferFormat,
2846 cwd: &Path,
2847) -> std::result::Result<SessionArtifact, ServiceError> {
2848 if target != TransferFormat::Grok {
2849 let target_session_id = target_session_id(target);
2850 return session_artifact_with_id(locator, session, target, Some(&target_session_id));
2851 }
2852
2853 let mut importable = session.clone();
2857 importable.meta.session_id = Some(target_session_id(TransferFormat::ClaudeCode));
2863 importable.meta.cwd = Some(if cwd.is_absolute() {
2864 cwd.to_path_buf()
2865 } else {
2866 std::env::current_dir()
2867 .map_err(|error| ServiceError::Operation(error.to_string()))?
2868 .join(cwd)
2869 });
2870 let content = importable
2871 .to_jsonl(SessionFormat::ClaudeCode)
2872 .map_err(operation)?;
2873 let stem = sanitize_filename(
2874 importable
2875 .meta
2876 .session_id
2877 .as_deref()
2878 .unwrap_or(&locator.session_id),
2879 );
2880 let suggested_filename = format!("{stem}.grok-import.claude-code.jsonl");
2881 Ok(SessionArtifact {
2882 source_harness: locator.harness.clone(),
2883 target_harness: TransferFormat::ClaudeCode.id(),
2886 session_id: importable.meta.session_id.clone(),
2887 content: content.clone(),
2888 suggested_filename: suggested_filename.clone(),
2889 files: vec![SessionArtifactFile {
2890 path: suggested_filename,
2891 content,
2892 role: ArtifactFileRole::Primary,
2893 }],
2894 fidelity: Fidelity::Semantic,
2895 residue: vec!["Grok's stock importer accepts a Claude Code transcript, not a complete Grok updates/session bundle".into()],
2896 })
2897}
2898
2899fn target_session_id(target: TransferFormat) -> String {
2900 let uuid = generated_session_id();
2901 match target {
2902 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
2903 TransferFormat::ClaudeCode
2904 | TransferFormat::Codex
2905 | TransferFormat::Pi
2906 | TransferFormat::Grok
2907 | TransferFormat::Gemini
2908 | TransferFormat::Goose => uuid,
2909 }
2910}
2911
2912fn sanitize_filename(value: &str) -> String {
2913 let value = value
2914 .chars()
2915 .map(|character| {
2916 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
2917 character
2918 } else {
2919 '-'
2920 }
2921 })
2922 .collect::<String>();
2923 let value = value.trim_matches('-');
2924 if value.is_empty() {
2925 "session".into()
2926 } else {
2927 value.chars().take(100).collect()
2928 }
2929}
2930
2931fn handoff_instructions(
2932 target: TransferFormat,
2933 session_id: &str,
2934 cwd: &Path,
2935) -> HandoffInstructions {
2936 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
2937 cwd: cwd.to_path_buf(),
2938 program: program.into(),
2939 arguments,
2940 env: BTreeMap::new(),
2941 };
2942 match target {
2943 TransferFormat::ClaudeCode => HandoffInstructions {
2944 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
2945 materialize: None,
2946 requires_materialization: true,
2947 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(),
2948 },
2949 TransferFormat::Codex => HandoffInstructions {
2950 launch: launch("codex", vec!["resume".into(), session_id.into()]),
2951 materialize: None,
2952 requires_materialization: true,
2953 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
2954 },
2955 TransferFormat::OpenCode => HandoffInstructions {
2956 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
2957 materialize: Some(launch(
2958 "opencode",
2959 vec!["import".into(), "{artifact_path}".into()],
2960 )),
2961 requires_materialization: true,
2962 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
2963 },
2964 TransferFormat::Pi => HandoffInstructions {
2965 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
2966 materialize: None,
2967 requires_materialization: true,
2968 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
2969 },
2970 TransferFormat::Grok => HandoffInstructions {
2971 launch: launch(
2972 "grok",
2973 vec![
2974 "--resume".into(),
2975 "{imported_session_id}".into(),
2976 "--fork-session".into(),
2977 ],
2978 ),
2979 materialize: Some(launch(
2980 "grok",
2981 vec!["import".into(), "--json".into(), "{artifact_path}".into()],
2982 )),
2983 requires_materialization: true,
2984 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(),
2985 },
2986 TransferFormat::Gemini => HandoffInstructions {
2987 launch: launch(
2988 "gemini",
2989 vec!["--session-file".into(), "{artifact_path}".into()],
2990 ),
2991 materialize: None,
2992 requires_materialization: true,
2993 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(),
2994 },
2995 TransferFormat::Goose => HandoffInstructions {
2996 launch: launch(
2997 "goose",
2998 vec![
2999 "session".into(),
3000 "--resume".into(),
3001 "--session-id".into(),
3002 "{imported_session_id}".into(),
3003 ],
3004 ),
3005 materialize: Some(launch(
3006 "goose",
3007 vec!["session".into(), "import".into(), "{artifact_path}".into()],
3008 )),
3009 requires_materialization: true,
3010 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(),
3011 },
3012 }
3013}
3014
3015fn resume_launch(
3016 harness: &str,
3017 session_id: &str,
3018 cwd: &Path,
3019 policy: ResumePolicy,
3020) -> std::result::Result<StructuredLaunch, ServiceError> {
3021 let mut arguments = Vec::new();
3022 let program = match harness {
3023 HarnessId::GROK => {
3024 if matches!(policy, ResumePolicy::Yolo) {
3025 arguments.extend([
3026 "--sandbox".into(),
3027 "workspace".into(),
3028 "--always-approve".into(),
3029 ]);
3030 }
3031 arguments.extend(["--resume".into(), session_id.into()]);
3032 "grok"
3033 }
3034 HarnessId::CODEX => {
3035 if matches!(policy, ResumePolicy::Yolo) {
3036 arguments.extend([
3037 "--dangerously-bypass-approvals-and-sandbox".into(),
3038 "--dangerously-bypass-hook-trust".into(),
3039 ]);
3040 }
3041 arguments.extend(["resume".into(), session_id.into()]);
3042 "codex"
3043 }
3044 HarnessId::CLAUDE_CODE => {
3045 if matches!(policy, ResumePolicy::Yolo) {
3046 arguments.push("--dangerously-skip-permissions".into());
3047 }
3048 arguments.extend(["--resume".into(), session_id.into()]);
3049 "claude"
3050 }
3051 HarnessId::GEMINI => {
3052 if matches!(policy, ResumePolicy::Yolo) {
3053 arguments.push("--yolo".into());
3054 }
3055 arguments.extend(["--resume".into(), session_id.into()]);
3056 "gemini"
3057 }
3058 HarnessId::GOOSE => {
3059 arguments.extend([
3060 "session".into(),
3061 "--resume".into(),
3062 "--session-id".into(),
3063 session_id.into(),
3064 ]);
3065 "goose"
3066 }
3067 HarnessId::PI => {
3068 if matches!(policy, ResumePolicy::Yolo) {
3069 arguments.push("--approve".into());
3070 }
3071 arguments.extend(["--session".into(), session_id.into()]);
3072 "pi"
3073 }
3074 HarnessId::OPENCODE => {
3075 arguments.extend(["--session".into(), session_id.into()]);
3076 "opencode"
3077 }
3078 HarnessId::SUPERCODE => {
3079 if matches!(policy, ResumePolicy::Yolo) {
3080 arguments.push("--dangerous".into());
3081 }
3082 arguments.extend(["resume".into(), session_id.into()]);
3083 "supercode"
3084 }
3085 other => {
3086 return Err(ServiceError::InvalidParams(format!(
3087 "no structured resume launch is registered for harness `{other}`"
3088 )))
3089 }
3090 };
3091 Ok(StructuredLaunch {
3092 cwd: cwd.to_path_buf(),
3093 program: program.into(),
3094 arguments,
3095 env: BTreeMap::new(),
3096 })
3097}
3098
3099fn runtime_backend(
3100 params: &RuntimeBackendParams,
3101) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
3102 if params.protocol.as_deref() == Some("acp") {
3103 let launch = params
3104 .launch
3105 .clone()
3106 .or_else(|| {
3107 harness_support_registry()
3108 .harnesses
3109 .into_iter()
3110 .find(|harness| harness.id == params.harness)
3111 .filter(|harness| {
3112 harness.runtime.implementation == ImplementationKind::GenericProtocol
3113 && harness.runtime.protocol.starts_with("acp")
3114 })
3115 .and_then(|harness| harness.runtime.default_launch)
3116 })
3117 .ok_or_else(|| {
3118 ServiceError::InvalidParams(
3119 "an ACP runtime requires `launch` unless the harness has a registered default"
3120 .into(),
3121 )
3122 })?;
3123 let resume_session = harness_support_registry()
3124 .harnesses
3125 .into_iter()
3126 .find(|harness| harness.id == params.harness)
3127 .is_some_and(|harness| harness.runtime.capabilities.resume_session);
3128 return Ok(Box::new(
3129 AcpRuntimeBackend::new(params.harness.clone(), launch)
3130 .with_resume_support(resume_session),
3131 ));
3132 }
3133 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
3134 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
3135 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
3136 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
3137 HarnessId::OPENCODE => match ¶ms.base_url {
3138 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
3139 None => Box::new(OpenCodeRuntimeBackend::new()),
3140 },
3141 harness => {
3142 let descriptor = harness_support_registry()
3143 .harnesses
3144 .into_iter()
3145 .find(|descriptor| descriptor.id.as_str() == harness)
3146 .filter(|descriptor| {
3147 descriptor.runtime.implementation == ImplementationKind::GenericProtocol
3148 && descriptor.runtime.protocol.starts_with("acp")
3149 });
3150 let Some(descriptor) = descriptor else {
3151 return Err(ServiceError::InvalidParams(format!(
3152 "no runtime adapter for harness `{harness}`; use protocol `acp` with a launch command"
3153 )));
3154 };
3155 let resume = descriptor.runtime.capabilities.resume_session;
3156 Box::new(
3157 AcpRuntimeBackend::new(
3158 descriptor.id,
3159 descriptor
3160 .runtime
3161 .default_launch
3162 .expect("generic ACP registry entry includes its launch"),
3163 )
3164 .with_resume_support(resume),
3165 )
3166 }
3167 };
3168 Ok(backend)
3169}
3170
3171fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
3172 if let Some(launch) = ¶ms.launch {
3173 return Some(launch.clone());
3174 }
3175 if !matches!(params.policy, RuntimePolicy::Yolo) {
3176 return None;
3177 }
3178 let launch = match params.harness.as_str() {
3179 HarnessId::GROK => RuntimeLaunch {
3180 program: "grok".into(),
3181 arguments: vec![
3182 "--sandbox".into(),
3183 "workspace".into(),
3184 "--always-approve".into(),
3185 "agent".into(),
3186 "--no-leader".into(),
3187 "stdio".into(),
3188 ],
3189 env: BTreeMap::from([("GROK_AGENT_DASHBOARD".into(), "0".into())]),
3190 },
3191 HarnessId::CODEX => RuntimeLaunch {
3192 program: "codex".into(),
3193 arguments: vec![
3194 "--dangerously-bypass-approvals-and-sandbox".into(),
3195 "--dangerously-bypass-hook-trust".into(),
3196 "app-server".into(),
3197 ],
3198 env: BTreeMap::new(),
3199 },
3200 HarnessId::CLAUDE_CODE => RuntimeLaunch {
3201 program: "claude".into(),
3202 arguments: vec![
3203 "--dangerously-skip-permissions".into(),
3204 "--print".into(),
3205 "--input-format".into(),
3206 "stream-json".into(),
3207 "--output-format".into(),
3208 "stream-json".into(),
3209 "--verbose".into(),
3210 ],
3211 env: BTreeMap::new(),
3212 },
3213 HarnessId::PI => RuntimeLaunch {
3214 program: "pi".into(),
3215 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
3216 env: BTreeMap::new(),
3217 },
3218 HarnessId::OPENCODE => RuntimeLaunch {
3219 program: "opencode".into(),
3220 arguments: vec!["serve".into()],
3221 env: BTreeMap::new(),
3222 },
3223 HarnessId::GEMINI => RuntimeLaunch {
3224 program: "gemini".into(),
3225 arguments: vec!["--acp".into(), "--yolo".into()],
3226 env: BTreeMap::new(),
3227 },
3228 HarnessId::GOOSE => RuntimeLaunch {
3229 program: "goose".into(),
3230 arguments: vec!["acp".into()],
3231 env: BTreeMap::new(),
3232 },
3233 HarnessId::SUPERCODE => RuntimeLaunch {
3234 program: "supercode".into(),
3235 arguments: vec!["acp".into(), "--dangerous".into()],
3236 env: BTreeMap::new(),
3237 },
3238 _ => return None,
3239 };
3240 Some(launch)
3241}
3242
3243struct IsolatedProbeHome {
3249 launch: RuntimeLaunch,
3250 root: PathBuf,
3251}
3252
3253impl IsolatedProbeHome {
3254 fn new(harness: &str, mut launch: RuntimeLaunch) -> std::io::Result<Self> {
3255 let root = std::env::temp_dir().join(format!(
3256 "supercode-harness-probe-{harness}-{}",
3257 generated_session_id()
3258 ));
3259 std::fs::create_dir_all(&root)?;
3260 set_private_dir_permissions(&root)?;
3261
3262 if let Some(source_home) = std::env::var_os("HOME").map(PathBuf::from) {
3263 for relative in probe_auth_files(harness) {
3264 copy_probe_file(&source_home, &root, relative)?;
3265 }
3266 }
3267 configure_isolated_probe_auth(harness, &root)?;
3268
3269 let root_text = root.to_string_lossy().into_owned();
3270 for (key, value) in [
3271 ("HOME", root_text.clone()),
3272 (
3273 "XDG_CACHE_HOME",
3274 root.join(".cache").to_string_lossy().into_owned(),
3275 ),
3276 (
3277 "XDG_CONFIG_HOME",
3278 root.join(".config").to_string_lossy().into_owned(),
3279 ),
3280 (
3281 "XDG_DATA_HOME",
3282 root.join(".local/share").to_string_lossy().into_owned(),
3283 ),
3284 ] {
3285 launch.env.insert(key.into(), value);
3286 }
3287 let scoped = match harness {
3288 HarnessId::CLAUDE_CODE => Some(("CLAUDE_CONFIG_DIR", root.join(".claude"))),
3289 HarnessId::CODEX => Some(("CODEX_HOME", root.join(".codex"))),
3290 HarnessId::GEMINI => Some(("GEMINI_CLI_HOME", root.clone())),
3291 HarnessId::GROK => Some(("GROK_HOME", root.join(".grok"))),
3292 HarnessId::PI => Some(("PI_CODING_AGENT_DIR", root.join(".pi/agent"))),
3293 HarnessId::SUPERCODE => Some(("SUPERCODE_HOME", root.join(".config/supercode"))),
3294 _ => None,
3295 };
3296 if let Some((key, value)) = scoped {
3297 launch
3298 .env
3299 .insert(key.into(), value.to_string_lossy().into_owned());
3300 }
3301 Ok(Self { launch, root })
3302 }
3303
3304 fn cleanup(&self) -> std::io::Result<()> {
3305 match std::fs::remove_dir_all(&self.root) {
3306 Ok(()) => Ok(()),
3307 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
3308 Err(error) => Err(error),
3309 }
3310 }
3311}
3312
3313impl Drop for IsolatedProbeHome {
3314 fn drop(&mut self) {
3315 let _ = self.cleanup();
3316 }
3317}
3318
3319fn probe_auth_files(harness: &str) -> &'static [&'static str] {
3320 match harness {
3321 HarnessId::CLAUDE_CODE => &[".claude/.credentials.json", ".claude.json"],
3322 HarnessId::CODEX => &[".codex/auth.json"],
3323 HarnessId::GEMINI => &[
3324 ".gemini/google_accounts.json",
3325 ".gemini/oauth_creds.json",
3326 ".gemini/settings.json",
3327 ],
3328 HarnessId::GROK => &[".grok/auth.json", ".grok/config.toml"],
3329 HarnessId::OPENCODE => &[
3330 ".config/opencode/auth.json",
3331 ".local/share/opencode/auth.json",
3332 ],
3333 HarnessId::PI => &[".pi/agent/auth.json"],
3334 HarnessId::SUPERCODE => &[
3335 ".config/supercode/config.toml",
3336 ".config/supercode/credentials.toml",
3337 ],
3338 _ => &[],
3339 }
3340}
3341
3342fn copy_probe_file(source_home: &Path, probe_home: &Path, relative: &str) -> std::io::Result<()> {
3343 let source = source_home.join(relative);
3344 if !source.is_file() {
3345 return Ok(());
3346 }
3347 let destination = probe_home.join(relative);
3348 if let Some(parent) = destination.parent() {
3349 std::fs::create_dir_all(parent)?;
3350 set_private_dir_permissions(parent)?;
3351 }
3352 std::fs::copy(source, &destination)?;
3353 set_private_file_permissions(&destination)
3354}
3355
3356fn configure_isolated_probe_auth(harness: &str, probe_home: &Path) -> std::io::Result<()> {
3357 if harness != HarnessId::GEMINI {
3358 return Ok(());
3359 }
3360 let oauth = probe_home.join(".gemini/oauth_creds.json");
3361 if !oauth.is_file() {
3362 return Ok(());
3363 }
3364 let settings_path = probe_home.join(".gemini/settings.json");
3365 let mut settings = std::fs::read_to_string(&settings_path)
3366 .ok()
3367 .and_then(|raw| serde_json::from_str::<Value>(&raw).ok())
3368 .unwrap_or_else(|| json!({}));
3369 settings["security"]["auth"]["selectedType"] = Value::String("oauth-personal".into());
3370 std::fs::write(
3371 &settings_path,
3372 serde_json::to_vec_pretty(&settings).map_err(std::io::Error::other)?,
3373 )?;
3374 set_private_file_permissions(&settings_path)
3375}
3376
3377#[cfg(unix)]
3378fn set_private_dir_permissions(path: &Path) -> std::io::Result<()> {
3379 use std::os::unix::fs::PermissionsExt;
3380 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))
3381}
3382
3383#[cfg(not(unix))]
3384fn set_private_dir_permissions(_path: &Path) -> std::io::Result<()> {
3385 Ok(())
3386}
3387
3388#[cfg(unix)]
3389fn set_private_file_permissions(path: &Path) -> std::io::Result<()> {
3390 use std::os::unix::fs::PermissionsExt;
3391 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
3392}
3393
3394#[cfg(not(unix))]
3395fn set_private_file_permissions(_path: &Path) -> std::io::Result<()> {
3396 Ok(())
3397}
3398
3399fn find_executable(program: &str) -> Option<PathBuf> {
3400 let candidate = PathBuf::from(program);
3401 if candidate.components().count() > 1 {
3402 return candidate.is_file().then_some(candidate);
3403 }
3404 let path = std::env::var_os("PATH")?;
3405 for directory in std::env::split_paths(&path) {
3406 let candidate = directory.join(program);
3407 if candidate.is_file() {
3408 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
3409 }
3410 #[cfg(windows)]
3411 {
3412 for extension in ["exe", "cmd", "bat"] {
3413 let candidate = directory.join(format!("{program}.{extension}"));
3414 if candidate.is_file() {
3415 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
3416 }
3417 }
3418 }
3419 }
3420 None
3421}
3422
3423async fn executable_version(executable: &Path) -> Option<String> {
3424 let mut command = tokio::process::Command::new(executable);
3425 command
3426 .arg("--version")
3427 .stdin(std::process::Stdio::null())
3428 .stdout(std::process::Stdio::piped())
3429 .stderr(std::process::Stdio::piped())
3430 .kill_on_drop(true);
3431 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
3432 .await
3433 .ok()?
3434 .ok()?;
3435 let stdout = String::from_utf8_lossy(&output.stdout);
3436 let stderr = String::from_utf8_lossy(&output.stderr);
3437 stdout
3438 .lines()
3439 .chain(stderr.lines())
3440 .map(str::trim)
3441 .find(|line| !line.is_empty())
3442 .map(|line| truncate_text(line, 200))
3443}
3444
3445fn auth_evidence(harness: &str) -> bool {
3446 let env_names: &[&str] = match harness {
3447 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
3448 HarnessId::CODEX => &["OPENAI_API_KEY"],
3449 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
3450 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
3451 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
3452 HarnessId::GEMINI => &["GEMINI_API_KEY", "GOOGLE_API_KEY"],
3453 HarnessId::SUPERCODE => &["OPENROUTER_API_KEY"],
3454 _ => &[],
3455 };
3456 if env_names
3457 .iter()
3458 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
3459 {
3460 return true;
3461 }
3462 let Some(home) = std::env::var_os("HOME").map(PathBuf::from) else {
3463 return false;
3464 };
3465 let files: Vec<PathBuf> = match harness {
3466 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
3467 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
3468 HarnessId::OPENCODE => vec![
3469 home.join(".local/share/opencode/auth.json"),
3470 home.join(".config/opencode/auth.json"),
3471 ],
3472 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
3473 HarnessId::GROK => vec![home.join(".grok/auth.json")],
3474 HarnessId::GEMINI => vec![
3475 home.join(".gemini/oauth_creds.json"),
3476 home.join(".gemini/google_accounts.json"),
3477 ],
3478 HarnessId::SUPERCODE => vec![home.join(".config/supercode/credentials.toml")],
3479 _ => Vec::new(),
3480 };
3481 if files.into_iter().any(|path| {
3482 std::fs::metadata(path)
3483 .map(|metadata| metadata.is_file() && metadata.len() > 2)
3484 .unwrap_or(false)
3485 }) {
3486 return true;
3487 }
3488 if harness == HarnessId::CLAUDE_CODE {
3495 return std::fs::read_to_string(home.join(".claude.json"))
3496 .map(|text| text.contains("\"oauthAccount\""))
3497 .unwrap_or(false);
3498 }
3499 false
3500}
3501
3502fn looks_like_auth_error(message: &str) -> bool {
3503 let message = message.to_ascii_lowercase();
3504 [
3505 "auth",
3506 "login",
3507 "sign in",
3508 "sign-in",
3509 "credential",
3510 "unauthorized",
3511 "forbidden",
3512 "token",
3513 ]
3514 .iter()
3515 .any(|needle| message.contains(needle))
3516}
3517
3518fn unavailable_capabilities() -> crate::RuntimeCapabilities {
3519 crate::RuntimeCapabilities {
3520 start_session: false,
3521 resume_session: false,
3522 attach_existing_process: false,
3523 send_input: false,
3524 stream_events: false,
3525 interrupt: false,
3526 respond_to_requests: false,
3527 }
3528}
3529
3530fn truncate_text(text: &str, max_chars: usize) -> String {
3531 let mut chars = text.chars();
3532 let truncated = chars.by_ref().take(max_chars).collect::<String>();
3533 if chars.next().is_some() {
3534 format!("{truncated}…")
3535 } else {
3536 truncated
3537 }
3538}
3539
3540fn error_message(error: ServiceError) -> String {
3541 match error {
3542 ServiceError::InvalidParams(message)
3543 | ServiceError::Operation(message)
3544 | ServiceError::UnsupportedAction(message) => message,
3545 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
3546 ServiceError::Sdk(error) => error.to_string(),
3547 }
3548}
3549
3550#[derive(Debug)]
3551enum ServiceError {
3552 InvalidParams(String),
3553 MethodNotFound,
3554 UnsupportedAction(String),
3555 Operation(String),
3556 Sdk(SdkError),
3557}
3558
3559fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
3560 match error {
3561 ServiceError::InvalidParams(message) => {
3562 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
3563 }
3564 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
3565 SdkError::unsupported(operation)
3566 }
3567 ServiceError::Operation(message) => {
3568 let code = if message.contains("already in progress") {
3569 SdkErrorCode::Busy
3570 } else if message.contains("not supported by this runtime") {
3571 SdkErrorCode::UnsupportedAction
3572 } else if message.contains("unknown runtime connection") {
3573 SdkErrorCode::NotFound
3574 } else {
3575 SdkErrorCode::Execution
3576 };
3577 SdkError::new(code, operation, message)
3578 }
3579 ServiceError::Sdk(error) => error,
3580 }
3581}
3582
3583fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
3584 let error_code = error.code();
3585 let code = match error_code {
3586 SdkErrorCode::Unauthenticated => -32030,
3587 SdkErrorCode::Unauthorized => -32031,
3588 SdkErrorCode::ControllerRequired => -32032,
3589 SdkErrorCode::LeaseExpired => -32033,
3590 SdkErrorCode::InvalidArgument => -32602,
3591 SdkErrorCode::NotFound => -32004,
3592 SdkErrorCode::Busy => -32000,
3593 SdkErrorCode::UnsupportedAction => -32020,
3594 SdkErrorCode::Execution => -32002,
3595 SdkErrorCode::Transport => -32003,
3596 };
3597 json!({
3598 "jsonrpc": "2.0",
3599 "id": id,
3600 "error": {
3601 "code": code,
3602 "name": error_code,
3603 "operation": error.operation(),
3604 "message": error.to_string(),
3605 },
3606 })
3607}
3608
3609fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
3610 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
3611}
3612
3613fn operation(error: impl Into<crate::Error>) -> ServiceError {
3614 let error = error.into();
3615 match error {
3616 crate::Error::Sdk(error) => ServiceError::Sdk(error),
3617 error => ServiceError::Operation(error.to_string()),
3618 }
3619}
3620
3621fn rpc_error(id: Value, code: i64, message: &str) -> Value {
3622 json!({
3623 "jsonrpc": "2.0",
3624 "id": id,
3625 "error": {"code": code, "message": message},
3626 })
3627}
3628
3629#[cfg(test)]
3630mod tests {
3631 use super::*;
3632 use crate::{HarnessEvent, HarnessId, RuntimeEndpoint, RuntimeHandle, StorageLocator};
3633 use async_trait::async_trait;
3634 use std::io::Write;
3635 use std::path::PathBuf;
3636 use std::time::Instant;
3637
3638 #[test]
3639 fn indexed_claude_descriptor_keeps_the_live_peer_address() {
3640 let descriptor = SessionDescriptor {
3641 locator: SessionLocator {
3642 harness: HarnessId::new(HarnessId::CLAUDE_CODE),
3643 session_id: "live-session".into(),
3644 storage: StorageLocator::File {
3645 path: PathBuf::from("/tmp/live-session.jsonl"),
3646 },
3647 },
3648 cwd: Some(PathBuf::from("/project")),
3649 title: None,
3650 preview_candidates: Vec::new(),
3651 latest_message_candidates: Vec::new(),
3652 updated_at_ms: Some(1),
3653 message_count: None,
3654 model: None,
3655 };
3656 let peer = crate::claude_peer::ClaudePeerSession {
3657 pid: 42,
3658 session_id: "live-session".into(),
3659 cwd: Some(PathBuf::from("/project")),
3660 name: "peer".into(),
3661 socket_path: PathBuf::from("/tmp/peer.sock"),
3662 status: Some(crate::claude_peer::ClaudePeerStatus::Busy),
3663 updated_at_ms: Some(1),
3664 version: Some("test".into()),
3665 };
3666
3667 let value = live_descriptor_value(&descriptor, &[peer]).unwrap();
3668 assert!(value["live_endpoint"]
3669 .as_str()
3670 .is_some_and(|endpoint| endpoint.starts_with("cc-peer:v1:42:peer:")));
3671 }
3672
3673 struct EndingRuntime {
3674 handle: RuntimeHandle,
3675 event: Option<HarnessEvent>,
3676 }
3677
3678 #[async_trait]
3679 impl RuntimeConnection for EndingRuntime {
3680 fn handle(&self) -> &RuntimeHandle {
3681 &self.handle
3682 }
3683
3684 async fn send_input(&mut self, _input: RuntimeInput) -> crate::Result<Option<String>> {
3685 unreachable!("ending runtime does not accept input")
3686 }
3687
3688 async fn next_event(&mut self) -> crate::Result<Option<HarnessEvent>> {
3689 Ok(self.event.take())
3690 }
3691
3692 async fn interrupt(&mut self) -> crate::Result<()> {
3693 Ok(())
3694 }
3695
3696 async fn respond(&mut self, _request_id: Value, _response: Value) -> crate::Result<()> {
3697 Ok(())
3698 }
3699
3700 async fn close(&mut self) -> crate::Result<()> {
3701 Ok(())
3702 }
3703 }
3704
3705 fn ending_runtime(event: Option<HarnessEvent>) -> Box<dyn RuntimeConnection> {
3706 Box::new(EndingRuntime {
3707 handle: RuntimeHandle {
3708 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3709 runtime_id: "ending-session".into(),
3710 endpoint: RuntimeEndpoint::LocalProcess {
3711 pid: None,
3712 command: vec!["ending-runtime".into()],
3713 protocol: "test".into(),
3714 },
3715 },
3716 event,
3717 })
3718 }
3719
3720 fn request(id: u64, method: &str, params: Value) -> Value {
3721 json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})
3722 }
3723
3724 fn pi_locator() -> SessionLocator {
3725 SessionLocator {
3726 harness: HarnessId::from(HarnessId::PI),
3727 session_id: "1e6f2a3b-0000-4000-8000-000000000001".into(),
3728 storage: StorageLocator::File {
3729 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3730 .join("tests/fixtures/pi_session.jsonl"),
3731 },
3732 }
3733 }
3734
3735 fn opencode_locator() -> SessionLocator {
3736 let session_id = "ses_fixtureAAAAAAAAAAAAAAA1";
3737 SessionLocator {
3738 harness: HarnessId::from(HarnessId::OPENCODE),
3739 session_id: session_id.into(),
3740 storage: StorageLocator::Sqlite {
3741 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3742 .join("tests/fixtures/opencode_fixture/opencode.db"),
3743 selector: session_id.into(),
3744 },
3745 }
3746 }
3747
3748 fn grok_locator() -> SessionLocator {
3749 SessionLocator {
3750 harness: HarnessId::from(HarnessId::GROK),
3751 session_id: "73c09283-4b33-41fa-90f1-0bcb0f7be523".into(),
3752 storage: StorageLocator::File {
3753 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
3754 .join("tests/fixtures/grok_session/chat_history.jsonl"),
3755 },
3756 }
3757 }
3758
3759 #[test]
3760 fn capabilities_are_explicit_and_versioned() {
3761 let mut service = HarnessSessionService::new();
3762 let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
3763 assert_eq!(response["result"]["version"], HARNESS_SERVICE_VERSION);
3764 assert_eq!(
3765 response["result"]["sdk"]["schema_version"],
3766 crate::SDK_SCHEMA_VERSION
3767 );
3768 assert_eq!(
3769 response["result"]["sdk"]["operations"]
3770 .as_array()
3771 .unwrap()
3772 .len(),
3773 SdkOperation::ALL.len()
3774 );
3775 assert_eq!(response["result"]["harnesses"].as_array().unwrap().len(), 8);
3776 assert!(response["result"]["harnesses"]
3777 .as_array()
3778 .unwrap()
3779 .iter()
3780 .any(|harness| harness == HarnessId::GROK));
3781 assert!(response["result"]["harnesses"]
3782 .as_array()
3783 .unwrap()
3784 .iter()
3785 .any(|harness| harness == HarnessId::GOOSE));
3786 }
3787
3788 #[test]
3789 fn handshake_health_uses_protocol_liveness_not_stderr_severity() {
3790 let noisy_stderr = crate::HarnessEvent {
3791 sequence: None,
3792 kind: "transport_stderr".into(),
3793 payload: json!({"line": "ERROR optional worker AuthorizationRequired"}),
3794 };
3795 assert_eq!(handshake_event_failure(&noisy_stderr), None);
3796
3797 let closed = crate::HarnessEvent {
3798 sequence: None,
3799 kind: "transport_closed".into(),
3800 payload: json!({}),
3801 };
3802 assert!(handshake_event_failure(&closed).is_some());
3803 }
3804
3805 #[tokio::test]
3806 async fn runtime_eof_is_notified_and_removed_for_raw_and_explicit_close() {
3807 let mut service = HarnessSessionService::new();
3808 service
3809 .runtimes
3810 .insert("raw-eof".into(), ending_runtime(None));
3811 service.runtimes.insert(
3812 "explicit-close".into(),
3813 ending_runtime(Some(HarnessEvent {
3814 sequence: None,
3815 kind: "transport_closed".into(),
3816 payload: json!({"message": "native transport exited"}),
3817 })),
3818 );
3819
3820 let notifications = service.poll_runtimes().await;
3821
3822 assert_eq!(notifications.len(), 2);
3823 assert!(notifications
3824 .iter()
3825 .all(|notification| { notification["params"]["event"]["kind"] == "transport_closed" }));
3826 assert!(notifications.iter().all(|notification| {
3827 notification["params"]["session_id"] == "ending-session"
3828 && notification["params"]["connection"].is_string()
3829 }));
3830 let mut sequences = notifications
3831 .iter()
3832 .filter_map(|notification| notification["params"]["sequence"].as_u64())
3833 .collect::<Vec<_>>();
3834 sequences.sort_unstable();
3835 assert_eq!(sequences, vec![1, 2]);
3836 assert!(service.runtimes.is_empty());
3837 }
3838
3839 #[tokio::test]
3840 async fn sdk_facade_returns_named_unsupported_actions() {
3841 let mut service = HarnessSessionService::new();
3842 let error = service
3843 .execute(SdkRequest {
3844 operation: SdkOperation::Steer,
3845 params: json!({"connection": "runtime-1", "text": "go left"}),
3846 })
3847 .await
3848 .unwrap_err();
3849 assert_eq!(error.code(), SdkErrorCode::UnsupportedAction);
3850 assert_eq!(error.operation(), Some(SdkOperation::Steer));
3851
3852 let response = service
3853 .handle_async(request(
3854 7,
3855 "harness.v1.runtimes.steer",
3856 json!({"connection": "runtime-1", "text": "go left"}),
3857 ))
3858 .await;
3859 assert_eq!(response["error"]["name"], "unsupported_action");
3860 assert_eq!(response["error"]["operation"], "steer");
3861 }
3862
3863 #[test]
3864 fn support_report_and_grok_default_binding_share_the_registry() {
3865 let mut service = HarnessSessionService::new();
3866 let response = service.handle(request(1, "harness.v1.support.report", json!({})));
3867 assert_eq!(response["result"]["schema"], crate::SUPPORT_REGISTRY_SCHEMA);
3868 let params = RuntimeBackendParams {
3869 harness: HarnessId::from(HarnessId::GROK),
3870 protocol: None,
3871 launch: None,
3872 base_url: None,
3873 policy: RuntimePolicy::Default,
3874 };
3875 let backend = match runtime_backend(¶ms) {
3876 Ok(backend) => backend,
3877 Err(_) => panic!("Grok should bind through its registered ACP launch"),
3878 };
3879 assert_eq!(backend.harness().as_str(), HarnessId::GROK);
3880 assert!(backend.capabilities().start_session);
3881 let registered = harness_support_registry()
3882 .harnesses
3883 .into_iter()
3884 .find(|harness| harness.id.as_str() == HarnessId::GROK)
3885 .and_then(|harness| harness.runtime.default_launch)
3886 .unwrap();
3887 assert!(!registered
3888 .arguments
3889 .iter()
3890 .any(|argument| argument == "--always-approve"));
3891 assert!(runtime_launch(¶ms).is_none());
3892
3893 let yolo = RuntimeBackendParams {
3894 policy: RuntimePolicy::Yolo,
3895 ..params
3896 };
3897 assert!(runtime_launch(&yolo)
3898 .unwrap()
3899 .arguments
3900 .iter()
3901 .any(|argument| argument == "--always-approve"));
3902
3903 let mismatched_protocol = RuntimeBackendParams {
3904 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3905 protocol: Some("acp".into()),
3906 launch: None,
3907 base_url: None,
3908 policy: RuntimePolicy::Default,
3909 };
3910 assert!(runtime_backend(&mismatched_protocol).is_err());
3911 }
3912
3913 #[test]
3914 fn load_follow_and_unfollow_share_the_same_locator() {
3915 let mut service = HarnessSessionService::new();
3916 let locator = pi_locator();
3917 let loaded = service.handle(request(
3918 1,
3919 "harness.v1.sessions.load",
3920 json!({"locator": locator}),
3921 ));
3922 assert_eq!(
3923 loaded["result"]["session"]["session_id"],
3924 locator.session_id
3925 );
3926
3927 let followed = service.handle(request(
3928 2,
3929 "harness.v1.sessions.follow",
3930 json!({"locator": locator}),
3931 ));
3932 assert_eq!(followed["result"]["subscription"], "sub-1");
3933 assert_eq!(followed["result"]["initial"]["type"], "session_snapshot");
3934 assert!(service.poll().is_empty());
3935
3936 let unfollowed = service.handle(request(
3937 3,
3938 "harness.v1.sessions.unfollow",
3939 json!({"subscription": "sub-1"}),
3940 ));
3941 assert_eq!(unfollowed["result"]["removed"], true);
3942 }
3943
3944 #[test]
3945 fn bounded_read_view_excludes_subagents_and_keeps_only_the_tail() {
3946 let temp = std::env::temp_dir().join(format!(
3947 "supercode-bounded-view-{}-{}",
3948 std::process::id(),
3949 generated_session_id()
3950 ));
3951 let path = temp.join("parent.jsonl");
3952 let subagents = temp.join("parent/subagents");
3953 std::fs::create_dir_all(&subagents).unwrap();
3954 let long_last = "x".repeat(300);
3955 let parent_records = [
3956 json!({"type":"user","uuid":"u1","parentUuid":null,"message":{"role":"user","content":"first"}}),
3957 json!({"type":"assistant","uuid":"a1","parentUuid":"u1","message":{"role":"assistant","content":[{"type":"text","text":"middle"}]}}),
3958 json!({"type":"user","uuid":"u2","parentUuid":"a1","message":{"role":"user","content":long_last}}),
3959 ];
3960 std::fs::write(
3961 &path,
3962 format!(
3963 "{}\n",
3964 parent_records
3965 .iter()
3966 .map(Value::to_string)
3967 .collect::<Vec<_>>()
3968 .join("\n")
3969 ),
3970 )
3971 .unwrap();
3972 std::fs::write(
3973 subagents.join("agent-child.jsonl"),
3974 concat!(
3975 r#"{"type":"user","uuid":"cu","parentUuid":null,"agentId":"child","message":{"role":"user","content":"child work"}}"#,
3976 "\n",
3977 ),
3978 )
3979 .unwrap();
3980 let locator = SessionLocator {
3981 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
3982 session_id: "parent".into(),
3983 storage: StorageLocator::File { path },
3984 };
3985 let mut service = HarnessSessionService::new();
3986
3987 let complete = service.handle(request(
3988 1,
3989 "harness.v1.sessions.load",
3990 json!({"locator": locator}),
3991 ));
3992 assert_eq!(
3993 complete["result"]["session"]["subagents"]
3994 .as_array()
3995 .unwrap()
3996 .len(),
3997 1
3998 );
3999
4000 let bounded = service.handle(request(
4001 2,
4002 "harness.v1.sessions.load",
4003 json!({
4004 "locator": locator,
4005 "view": {
4006 "tail_messages": 1,
4007 "max_message_chars": 256,
4008 "include_subagents": false
4009 },
4010 }),
4011 ));
4012 let session = &bounded["result"]["session"];
4013 assert!(session["subagents"].as_array().unwrap().is_empty());
4014 assert_eq!(session["messages"].as_array().unwrap().len(), 1);
4015 assert_eq!(
4016 session["messages"][0]["content"],
4017 format!("{}\n…", "x".repeat(256))
4018 );
4019
4020 let followed = service.handle(request(
4021 3,
4022 "harness.v1.sessions.follow",
4023 json!({
4024 "locator": locator,
4025 "view": {
4026 "tail_messages": 1,
4027 "max_message_chars": 256,
4028 "include_subagents": false
4029 },
4030 }),
4031 ));
4032 let initial = &followed["result"]["initial"]["session"];
4033 assert!(initial["subagents"].as_array().unwrap().is_empty());
4034 assert_eq!(initial["messages"].as_array().unwrap().len(), 1);
4035
4036 let _ = std::fs::remove_dir_all(&temp);
4037 }
4038
4039 #[test]
4040 fn forty_megabyte_display_load_is_bounded_and_prompt() {
4041 let temp = std::env::temp_dir().join(format!(
4042 "supercode-large-display-view-{}-{}",
4043 std::process::id(),
4044 generated_session_id()
4045 ));
4046 std::fs::create_dir_all(&temp).unwrap();
4047 let path = temp.join("rollout.jsonl");
4048 let mut file = std::io::BufWriter::new(std::fs::File::create(&path).unwrap());
4049 writeln!(
4050 file,
4051 r#"{{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{{"id":"large-display","cwd":"/tmp"}}}}"#
4052 )
4053 .unwrap();
4054 let padding = "x".repeat(80 * 1024);
4055 for index in 0..512 {
4056 let marker = if index == 0 {
4057 "OLDEST-SHOULD-NOT-LOAD"
4058 } else if index == 511 {
4059 "LATEST-MUST-LOAD"
4060 } else {
4061 "bulk"
4062 };
4063 writeln!(
4064 file,
4065 "{}",
4066 json!({
4067 "timestamp": "2026-01-01T00:00:01Z",
4068 "type": "response_item",
4069 "payload": {
4070 "type": "message",
4071 "role": "assistant",
4072 "content": [{"type": "output_text", "text": format!("{marker}:{padding}")}],
4073 },
4074 })
4075 )
4076 .unwrap();
4077 }
4078 file.flush().unwrap();
4079 drop(file);
4080 assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
4081
4082 let locator = SessionLocator {
4083 harness: HarnessId::from(HarnessId::CODEX),
4084 session_id: "large-display".into(),
4085 storage: StorageLocator::File { path },
4086 };
4087 let started = Instant::now();
4088 let response = HarnessSessionService::new().handle(request(
4089 1,
4090 "harness.v1.sessions.load",
4091 json!({
4092 "locator": locator,
4093 "view": {
4094 "tail_messages": 500,
4095 "max_message_chars": 1024,
4096 "include_subagents": false,
4097 "display_history": true,
4098 },
4099 }),
4100 ));
4101 let elapsed = started.elapsed();
4102 let wire = response.to_string();
4103 eprintln!(
4104 "bounded 40 MiB display load: {elapsed:?}, {} response bytes",
4105 wire.len()
4106 );
4107 assert!(response.get("error").is_none(), "{response:#}");
4108 assert!(wire.contains("LATEST-MUST-LOAD"));
4109 assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
4110 assert!(
4111 wire.len() < 2 * 1024 * 1024,
4112 "bounded wire was {} bytes",
4113 wire.len()
4114 );
4115 assert!(
4116 elapsed.as_secs_f64() < 3.0,
4117 "bounded 40 MiB load took {elapsed:?}"
4118 );
4119
4120 let _ = std::fs::remove_dir_all(&temp);
4121 }
4122
4123 #[test]
4124 fn forty_megabyte_goose_store_display_load_reads_only_the_tail() {
4125 let temp = std::env::temp_dir().join(format!(
4126 "supercode-large-goose-view-{}-{}",
4127 std::process::id(),
4128 generated_session_id()
4129 ));
4130 std::fs::create_dir_all(&temp).unwrap();
4131 let path = temp.join("sessions.db");
4132 let connection = rusqlite::Connection::open(&path).unwrap();
4133 connection
4134 .execute_batch(
4135 "CREATE TABLE sessions (
4136 id TEXT PRIMARY KEY, name TEXT NOT NULL, working_dir TEXT NOT NULL,
4137 created_at TEXT NOT NULL, updated_at TEXT NOT NULL,
4138 session_type TEXT NOT NULL, extension_data TEXT,
4139 goose_mode TEXT NOT NULL, provider_name TEXT, model_config_json TEXT,
4140 archived_at TEXT
4141 );
4142 CREATE TABLE messages (
4143 id INTEGER PRIMARY KEY, session_id TEXT NOT NULL, message_id TEXT,
4144 role TEXT NOT NULL, content_json TEXT NOT NULL,
4145 created_timestamp INTEGER NOT NULL, metadata_json TEXT
4146 );",
4147 )
4148 .unwrap();
4149 connection
4150 .execute(
4151 "INSERT INTO sessions VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, NULL)",
4152 rusqlite::params![
4153 "goose-large",
4154 "Large Goose session",
4155 "/tmp",
4156 "2026-01-01 00:00:00",
4157 "2026-01-01 00:00:02",
4158 "user",
4159 "{}",
4160 "auto",
4161 "anthropic",
4162 r#"{"model_name":"claude-sonnet"}"#,
4163 ],
4164 )
4165 .unwrap();
4166 let old_content = serde_json::to_string(&vec![json!({
4167 "type": "text",
4168 "text": format!("OLDEST-SHOULD-NOT-LOAD:{}", "x".repeat(40 * 1024 * 1024)),
4169 })])
4170 .unwrap();
4171 connection
4172 .execute(
4173 "INSERT INTO messages VALUES (1, ?1, 'old', 'user', ?2, 1, '{}')",
4174 rusqlite::params!["goose-large", old_content],
4175 )
4176 .unwrap();
4177 connection
4178 .execute(
4179 "INSERT INTO messages VALUES (2, ?1, 'new', 'assistant', ?2, 2, '{}')",
4180 rusqlite::params![
4181 "goose-large",
4182 r#"[{"type":"text","text":"LATEST-MUST-LOAD"}]"#
4183 ],
4184 )
4185 .unwrap();
4186 drop(connection);
4187 assert!(std::fs::metadata(&path).unwrap().len() >= 40 * 1024 * 1024);
4188
4189 let locator = SessionLocator {
4190 harness: HarnessId::from(HarnessId::GOOSE),
4191 session_id: "goose-large".into(),
4192 storage: StorageLocator::Sqlite {
4193 path,
4194 selector: "goose-large".into(),
4195 },
4196 };
4197 let started = Instant::now();
4198 let response = HarnessSessionService::new().handle(request(
4199 1,
4200 "harness.v1.sessions.load",
4201 json!({
4202 "locator": locator,
4203 "view": {
4204 "tail_messages": 1,
4205 "max_message_chars": 1024,
4206 "include_subagents": false,
4207 "display_history": true,
4208 },
4209 }),
4210 ));
4211 let elapsed = started.elapsed();
4212 let wire = response.to_string();
4213 eprintln!(
4214 "bounded 40 MiB Goose display load: {elapsed:?}, {} response bytes",
4215 wire.len()
4216 );
4217 assert!(response.get("error").is_none(), "{response:#}");
4218 assert!(wire.contains("LATEST-MUST-LOAD"));
4219 assert!(!wire.contains("OLDEST-SHOULD-NOT-LOAD"));
4220 assert!(
4221 wire.len() < 64 * 1024,
4222 "bounded wire was {} bytes",
4223 wire.len()
4224 );
4225 assert!(
4226 elapsed.as_secs_f64() < 1.0,
4227 "bounded Goose load took {elapsed:?}"
4228 );
4229
4230 let _ = std::fs::remove_dir_all(&temp);
4231 }
4232
4233 #[test]
4234 fn display_view_keeps_codex_assistant_history_across_compaction() {
4235 let temp = std::env::temp_dir().join(format!(
4236 "supercode-codex-display-view-{}-{}",
4237 std::process::id(),
4238 generated_session_id()
4239 ));
4240 std::fs::create_dir_all(&temp).unwrap();
4241 let path = temp.join("rollout.jsonl");
4242 std::fs::write(
4243 &path,
4244 concat!(
4245 r#"{"timestamp":"2026-01-01T00:00:00Z","type":"session_meta","payload":{"id":"codex-display","cwd":"/tmp"}}"#,
4246 "\n",
4247 r#"{"timestamp":"2026-01-01T00:00:01Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"old prompt"}]}}"#,
4248 "\n",
4249 r#"{"timestamp":"2026-01-01T00:00:02Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"old answer"}]}}"#,
4250 "\n",
4251 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"}]}}"#,
4252 "\n",
4253 r#"{"timestamp":"2026-01-01T00:00:04Z","type":"response_item","payload":{"type":"message","role":"user","content":[{"type":"input_text","text":"new prompt"}]}}"#,
4254 "\n",
4255 r#"{"timestamp":"2026-01-01T00:00:05Z","type":"response_item","payload":{"type":"message","role":"assistant","content":[{"type":"output_text","text":"new answer"}]}}"#,
4256 "\n",
4257 ),
4258 )
4259 .unwrap();
4260 let locator = SessionLocator {
4261 harness: HarnessId::from(HarnessId::CODEX),
4262 session_id: "codex-display".into(),
4263 storage: StorageLocator::File { path },
4264 };
4265 let mut service = HarnessSessionService::new();
4266
4267 let continuation = service.handle(request(
4268 1,
4269 "harness.v1.sessions.load",
4270 json!({"locator": locator}),
4271 ));
4272 let continuation_text = continuation["result"]["session"]["messages"].to_string();
4273 assert!(!continuation_text.contains("old answer"));
4274
4275 let display = service.handle(request(
4276 2,
4277 "harness.v1.sessions.load",
4278 json!({
4279 "locator": locator,
4280 "view": {
4281 "tail_messages": 10,
4282 "include_subagents": false,
4283 "display_history": true,
4284 },
4285 }),
4286 ));
4287 let display_text = display["result"]["session"]["messages"].to_string();
4288 assert!(display_text.contains("old prompt"));
4289 assert!(display_text.contains("old answer"));
4290 assert!(display_text.contains("new prompt"));
4291 assert!(display_text.contains("new answer"));
4292
4293 let _ = std::fs::remove_dir_all(&temp);
4294 }
4295
4296 #[test]
4297 fn load_supports_bounded_windows_and_media_metadata() {
4298 let mut service = HarnessSessionService::new();
4299 let locator = pi_locator();
4300 let bounded = service.handle(request(
4301 1,
4302 "harness.v1.sessions.load",
4303 json!({
4304 "locator": locator,
4305 "options": {
4306 "include_subagents": false,
4307 "message_limit": 2,
4308 "message_offset": 1
4309 }
4310 }),
4311 ));
4312 assert_eq!(bounded["result"]["window"]["offset"], 1);
4313 assert_eq!(bounded["result"]["window"]["returned"], 2);
4314 assert!(bounded["result"]["summary"]["first_message"].is_object());
4315 assert!(bounded["result"]["summary"]["last_message"].is_object());
4316 assert_eq!(
4317 bounded["result"]["session"]["messages"]
4318 .as_array()
4319 .unwrap()
4320 .len(),
4321 2
4322 );
4323 assert!(bounded["result"]["session"]["subagents"]
4324 .as_array()
4325 .unwrap()
4326 .is_empty());
4327
4328 let tail = service.handle(request(
4329 2,
4330 "harness.v1.sessions.load",
4331 json!({"locator": locator, "options": {"message_tail": 1}}),
4332 ));
4333 assert_eq!(tail["result"]["window"]["returned"], 1);
4334 assert_eq!(tail["result"]["window"]["has_more"], true);
4335 assert_eq!(tail["result"]["window"]["has_older"], true);
4336 assert!(tail["result"]["window"]["older_items"].as_u64().unwrap() > 0);
4337 assert!(tail["result"]["summary"]["first_message"].is_object());
4338
4339 let metadata_only = service.handle(request(
4340 3,
4341 "harness.v1.sessions.load",
4342 json!({"locator": locator, "options": {"inline_media": "metadata"}}),
4343 ));
4344 assert!(metadata_only["result"]["session"]
4345 .to_string()
4346 .contains("media_reference"));
4347 assert!(!metadata_only["result"]["session"]
4348 .to_string()
4349 .contains("data:image/"));
4350 }
4351
4352 #[test]
4353 fn import_translate_branch_and_handoff_use_typed_artifacts() {
4354 let mut service = HarnessSessionService::new();
4355 let locator = pi_locator();
4356 let translated = service.handle(request(
4357 1,
4358 "harness.v1.sessions.translate",
4359 json!({"locator": locator, "target_harness": "grok"}),
4360 ));
4361 assert_eq!(translated["result"]["artifact"]["source_harness"], "pi");
4362 assert_eq!(translated["result"]["artifact"]["target_harness"], "grok");
4363 assert!(translated["result"]["artifact"]["content"]
4364 .as_str()
4365 .is_some_and(|content| !content.is_empty()));
4366
4367 for target in ["opencode", "open-code"] {
4368 let opencode = service.handle(request(
4369 6,
4370 "harness.v1.sessions.translate",
4371 json!({"locator": locator, "target_harness": target}),
4372 ));
4373 assert_eq!(opencode["result"]["artifact"]["target_harness"], "opencode");
4374 }
4375 let goose = service.handle(request(
4376 7,
4377 "harness.v1.sessions.translate",
4378 json!({"locator": locator, "target_harness": "goose"}),
4379 ));
4380 assert_eq!(goose["result"]["artifact"]["target_harness"], "goose");
4381 assert!(serde_json::from_str::<Value>(
4382 goose["result"]["artifact"]["content"].as_str().unwrap()
4383 )
4384 .unwrap()["conversation"]
4385 .is_array());
4386
4387 let imported = service.handle(request(
4388 2,
4389 "harness.v1.sessions.import",
4390 json!({
4391 "source_harness": "grok",
4392 "content": translated["result"]["artifact"]["content"],
4393 }),
4394 ));
4395 assert_eq!(imported["result"]["session"]["source"], "grok");
4396
4397 let branched = service.handle(request(
4398 3,
4399 "harness.v1.sessions.branch",
4400 json!({"locator": locator, "target_harness": "codex"}),
4401 ));
4402 assert_eq!(branched["result"]["parent"]["harness"], "pi");
4403 assert!(branched["result"]["bootstrap_prompt"]
4404 .as_str()
4405 .unwrap()
4406 .contains("frozen parent transcript"));
4407 assert_eq!(branched["result"]["artifact"]["target_harness"], "codex");
4408
4409 let handoff = service.handle(request(
4410 4,
4411 "harness.v1.sessions.handoff",
4412 json!({"locator": locator, "target_harness": "pi", "cwd": "/tmp/project"}),
4413 ));
4414 assert_eq!(handoff["result"]["launch"]["program"], "pi");
4415 assert_eq!(handoff["result"]["launch"]["cwd"], "/tmp/project");
4416 assert_eq!(handoff["result"]["requires_materialization"], true);
4417
4418 let goose_handoff = service.handle(request(
4419 8,
4420 "harness.v1.sessions.handoff",
4421 json!({"locator": locator, "target_harness": "goose", "cwd": "/tmp/project"}),
4422 ));
4423 assert_eq!(goose_handoff["result"]["launch"]["program"], "goose");
4424 assert_eq!(
4425 goose_handoff["result"]["materialize"]["arguments"],
4426 json!(["session", "import", "{artifact_path}"])
4427 );
4428
4429 let resumed = service.handle(request(
4430 5,
4431 "harness.v1.sessions.resume_instructions",
4432 json!({"locator": locator, "cwd": "/tmp/project", "policy": "yolo"}),
4433 ));
4434 assert_eq!(resumed["result"]["launch"]["program"], "pi");
4435 assert_eq!(resumed["result"]["launch"]["arguments"][0], "--approve");
4436 }
4437
4438 #[test]
4439 fn reduce_persists_and_reloads_a_byte_exact_reversible_bundle() {
4440 let temp = std::env::temp_dir().join(format!(
4441 "supercode-service-reduce-{}-{}",
4442 std::process::id(),
4443 generated_session_id()
4444 ));
4445 let source_path = temp.join("source.jsonl");
4446 let store_root = temp.join("store");
4447 std::fs::create_dir_all(&temp).unwrap();
4448
4449 let mut records = vec![json!({
4450 "timestamp": "2026-01-01T00:00:00Z",
4451 "type": "session_meta",
4452 "payload": {"id": "codex-reduce", "cwd": "/tmp/project"},
4453 })];
4454 for turn in 0..16 {
4455 records.push(json!({
4456 "timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 1),
4457 "type": "response_item",
4458 "payload": {
4459 "type": "message",
4460 "role": "user",
4461 "content": [{
4462 "type": "input_text",
4463 "text": format!("request {turn}: {}", "context ".repeat(80)),
4464 }],
4465 },
4466 }));
4467 records.push(json!({
4468 "timestamp": format!("2026-01-01T00:00:{:02}Z", turn * 2 + 2),
4469 "type": "response_item",
4470 "payload": {
4471 "type": "message",
4472 "role": "assistant",
4473 "content": [{
4474 "type": "output_text",
4475 "text": format!("answer {turn}: {}", "implementation detail ".repeat(80)),
4476 }],
4477 },
4478 }));
4479 }
4480 let source = format!(
4481 "{}\n",
4482 records
4483 .iter()
4484 .map(Value::to_string)
4485 .collect::<Vec<_>>()
4486 .join("\n")
4487 );
4488 std::fs::write(&source_path, &source).unwrap();
4489 let locator = SessionLocator {
4490 harness: HarnessId::from(HarnessId::CODEX),
4491 session_id: "codex-reduce".into(),
4492 storage: StorageLocator::File {
4493 path: source_path.clone(),
4494 },
4495 };
4496 let original = load_session(&locator).unwrap();
4497 let mut service =
4498 HarnessSessionService::new().with_reduction_store_root(store_root.clone());
4499
4500 let response = service.handle(request(
4501 1,
4502 "harness.v1.sessions.reduce",
4503 json!({
4504 "locator": locator,
4505 "target_harness": "claude-code",
4506 "keep_last": 4,
4507 }),
4508 ));
4509 assert!(response.get("error").is_none(), "{response:#}");
4510 let receipt = &response["result"]["receipt"];
4511 assert_eq!(receipt["source_harness"], "codex");
4512 assert_eq!(receipt["target_harness"], "claude-code");
4513 assert_eq!(receipt["verified"], true);
4514 assert_eq!(receipt["reversible"], true);
4515 assert!(receipt["reductions"].as_u64().unwrap() > 0);
4516 assert!(
4517 receipt["source_tokens"].as_u64().unwrap()
4518 > receipt["reduced_tokens"].as_u64().unwrap()
4519 );
4520 assert!(receipt["ratio"].as_f64().unwrap() > 1.0);
4521 assert!(response["result"]["bootstrap_prompt"]
4522 .as_str()
4523 .unwrap()
4524 .contains("Do not guess hidden content"));
4525
4526 let rescue_id = receipt["id"].as_str().unwrap();
4527 let store = crate::SessionStore::open(&store_root).unwrap();
4528 let sidecar =
4529 Session::from_sidecar_str(&store.load_sidecar(rescue_id).unwrap().unwrap()).unwrap();
4530 let log = store.load_reduction_log(rescue_id).unwrap().unwrap();
4531 let persisted_view = parse_messages_jsonl(&store.load(rescue_id).unwrap()).unwrap();
4532 let policy = reduce::ReductionPolicy {
4533 clear_turns_older_than: Some(4),
4534 ..Default::default()
4535 };
4536 let (restamped_view, reapplied_log) =
4537 reduce::project_messages(&sidecar.messages, &policy, &log);
4538 assert_eq!(
4539 messages_jsonl(&persisted_view).unwrap(),
4540 messages_jsonl(&restamped_view).unwrap()
4541 );
4542 assert_eq!(reapplied_log, log);
4543 reduce::verify_log(&log, &sidecar).unwrap();
4544 assert_eq!(
4545 reduce::invert(&restamped_view, &log, &sidecar).unwrap(),
4546 original.messages
4547 );
4548 assert_eq!(std::fs::read_to_string(&source_path).unwrap(), source);
4549
4550 std::fs::remove_dir_all(temp).ok();
4551 }
4552
4553 #[test]
4554 fn read_surfaces_view_a_severed_claude_graph_while_transfer_still_refuses_it() {
4555 let temp = std::env::temp_dir().join(format!(
4556 "supercode-severed-view-{}-{}",
4557 std::process::id(),
4558 generated_session_id()
4559 ));
4560 std::fs::create_dir_all(&temp).unwrap();
4561 let path = temp.join("severed.jsonl");
4562 std::fs::write(
4565 &path,
4566 concat!(
4567 r#"{"type":"user","uuid":"orphan-u","parentUuid":null,"message":{"role":"user","content":"stranded prompt"}}"#,
4568 "\n",
4569 r#"{"type":"assistant","uuid":"live-a","parentUuid":"pruned","message":{"id":"m","role":"assistant","content":[{"type":"text","text":"live answer"}]}}"#,
4570 "\n",
4571 ),
4572 )
4573 .unwrap();
4574 let locator = SessionLocator {
4575 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
4576 session_id: "severed".into(),
4577 storage: StorageLocator::File { path },
4578 };
4579 let mut service = HarnessSessionService::new();
4580
4581 let viewed = service.handle(request(
4582 1,
4583 "harness.v1.sessions.load",
4584 json!({"locator": locator}),
4585 ));
4586 let session = &viewed["result"]["session"];
4587 assert_eq!(session["fidelity"], "semantic");
4588 assert_eq!(session["messages"].as_array().unwrap().len(), 2);
4589 assert!(session["residue"].as_array().unwrap().iter().any(|entry| {
4590 entry
4591 .as_str()
4592 .is_some_and(|entry| entry.contains("live-a") && entry.contains("pruned"))
4593 }));
4594
4595 let strict = service.handle(request(
4598 2,
4599 "harness.v1.sessions.load",
4600 json!({"locator": locator, "fidelity": "byte_lossless"}),
4601 ));
4602 assert!(strict["error"]["message"]
4603 .as_str()
4604 .unwrap()
4605 .contains("cannot reconstruct lossless Claude continuation"));
4606
4607 let translated = service.handle(request(
4609 3,
4610 "harness.v1.sessions.translate",
4611 json!({"locator": locator, "target_harness": "codex"}),
4612 ));
4613 assert!(translated["error"]["message"]
4614 .as_str()
4615 .unwrap()
4616 .contains("cannot reconstruct lossless Claude continuation"));
4617 let resumed = service.handle(request(
4618 4,
4619 "harness.v1.sessions.resume_instructions",
4620 json!({"locator": locator}),
4621 ));
4622 assert!(resumed["error"]["message"]
4623 .as_str()
4624 .unwrap()
4625 .contains("cannot reconstruct lossless Claude continuation"));
4626
4627 let _ = std::fs::remove_dir_all(&temp);
4628 }
4629
4630 #[test]
4631 fn structured_resume_launches_cover_gemini_goose_and_supercode() {
4632 let gemini = resume_launch(
4633 HarnessId::GEMINI,
4634 "gemini-session",
4635 Path::new("/tmp/project"),
4636 ResumePolicy::Yolo,
4637 )
4638 .unwrap_or_else(|_| panic!("Gemini resume launch must be registered"));
4639 assert_eq!(gemini.program, "gemini");
4640 assert_eq!(gemini.arguments, ["--yolo", "--resume", "gemini-session"]);
4641
4642 let goose = resume_launch(
4643 HarnessId::GOOSE,
4644 "goose-session",
4645 Path::new("/tmp/project"),
4646 ResumePolicy::Yolo,
4647 )
4648 .unwrap_or_else(|_| panic!("Goose resume launch must be registered"));
4649 assert_eq!(goose.program, "goose");
4650 assert_eq!(
4651 goose.arguments,
4652 ["session", "--resume", "--session-id", "goose-session"]
4653 );
4654
4655 let supercode = resume_launch(
4656 HarnessId::SUPERCODE,
4657 "supercode-session",
4658 Path::new("/tmp/project"),
4659 ResumePolicy::Yolo,
4660 )
4661 .unwrap_or_else(|_| panic!("Supercode resume launch must be registered"));
4662 assert_eq!(supercode.program, "supercode");
4663 assert_eq!(
4664 supercode.arguments,
4665 ["--dangerous", "resume", "supercode-session"]
4666 );
4667 }
4668
4669 #[test]
4670 fn diagonal_artifacts_preserve_claude_subagents_and_grok_bundle_members() {
4671 let temp = std::env::temp_dir().join(format!(
4672 "supercode-harness-artifact-{}-{}",
4673 std::process::id(),
4674 generated_session_id()
4675 ));
4676 let main_path = temp.join("parent.jsonl");
4677 let subagent_path = temp.join("parent/subagents/agent-child.jsonl");
4678 std::fs::create_dir_all(subagent_path.parent().unwrap()).unwrap();
4679 let fixture = std::fs::read_to_string(
4680 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4681 .join("tests/fixtures/claude_code_session.jsonl"),
4682 )
4683 .unwrap();
4684 let parent = fixture.trim_end_matches('\n');
4685 let child = fixture.trim_end_matches('\n');
4686 std::fs::write(&main_path, parent).unwrap();
4687 std::fs::write(&subagent_path, child).unwrap();
4688 let locator = SessionLocator {
4689 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
4690 session_id: "213bb148-51ea-453f-9206-f8b4b1168547".into(),
4691 storage: StorageLocator::File {
4692 path: main_path.clone(),
4693 },
4694 };
4695 let mut service = HarnessSessionService::new();
4696 let claude = service.handle(request(
4697 1,
4698 "harness.v1.sessions.translate",
4699 json!({"locator": locator, "target_harness": "claude-code"}),
4700 ));
4701 let artifact = &claude["result"]["artifact"];
4702 assert_eq!(artifact["fidelity"], "byte_lossless");
4703 assert_eq!(artifact["content"], parent);
4704 let files = artifact["files"].as_array().unwrap();
4705 assert!(files.iter().any(|file| {
4706 file["role"] == "subagent"
4707 && file["path"]
4708 .as_str()
4709 .is_some_and(|path| path.ends_with("/subagents/agent-child.jsonl"))
4710 && file["content"] == child
4711 }));
4712 assert!(!artifact["content"].as_str().unwrap().ends_with('\n'));
4713
4714 let grok = service.handle(request(
4715 2,
4716 "harness.v1.sessions.translate",
4717 json!({"locator": grok_locator(), "target_harness": "grok"}),
4718 ));
4719 let files = grok["result"]["artifact"]["files"].as_array().unwrap();
4720 for name in ["summary.json", "updates.jsonl"] {
4721 let expected = std::fs::read_to_string(
4722 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
4723 .join("tests/fixtures/grok_session")
4724 .join(name),
4725 )
4726 .unwrap();
4727 assert!(files.iter().any(|file| {
4728 file["path"] == name && file["role"] == "bundle" && file["content"] == expected
4729 }));
4730 }
4731 std::fs::remove_dir_all(temp).ok();
4732 }
4733
4734 #[test]
4735 fn every_non_grok_handoff_mints_and_uses_a_fresh_target_identity() {
4736 let mut service = HarnessSessionService::new();
4737 let source = pi_locator();
4738 for (target, format) in [
4739 ("claude-code", SessionFormat::ClaudeCode),
4740 ("codex", SessionFormat::Codex),
4741 ("opencode", SessionFormat::OpenCode),
4742 ("pi", SessionFormat::Pi),
4743 ] {
4744 let result = service.handle(request(
4745 1,
4746 "harness.v1.sessions.handoff",
4747 json!({"locator": source, "target_harness": target, "cwd": "/tmp/project"}),
4748 ));
4749 let artifact = &result["result"]["artifact"];
4750 let target_id = artifact["session_id"].as_str().unwrap();
4751 assert_ne!(target_id, source.session_id, "{target}");
4752 let parsed = Session::load_str(artifact["content"].as_str().unwrap(), format).unwrap();
4753 assert_eq!(
4754 parsed.meta.session_id.as_deref(),
4755 Some(target_id),
4756 "{target}"
4757 );
4758 if target != "pi" {
4759 assert!(result["result"]["launch"]["arguments"]
4760 .as_array()
4761 .unwrap()
4762 .iter()
4763 .any(|argument| argument == target_id));
4764 }
4765 if target == "opencode" {
4766 assert!(target_id.starts_with("ses_"));
4767 fn assert_session_ids(value: &Value, target_id: &str) {
4768 match value {
4769 Value::Object(fields) => {
4770 if let Some(session_id) = fields.get("sessionID") {
4771 assert_eq!(session_id, target_id);
4772 }
4773 for child in fields.values() {
4774 assert_session_ids(child, target_id);
4775 }
4776 }
4777 Value::Array(values) => {
4778 for child in values {
4779 assert_session_ids(child, target_id);
4780 }
4781 }
4782 _ => {}
4783 }
4784 }
4785 let document: Value =
4786 serde_json::from_str(artifact["content"].as_str().unwrap()).unwrap();
4787 assert_session_ids(&document, target_id);
4788 }
4789 }
4790
4791 let first = service.handle(request(
4792 2,
4793 "harness.v1.sessions.handoff",
4794 json!({"locator": source, "target_harness": "codex"}),
4795 ));
4796 let second = service.handle(request(
4797 3,
4798 "harness.v1.sessions.handoff",
4799 json!({"locator": source, "target_harness": "codex"}),
4800 ));
4801 assert_ne!(
4802 first["result"]["artifact"]["session_id"],
4803 second["result"]["artifact"]["session_id"]
4804 );
4805 }
4806
4807 #[test]
4808 fn grok_handoff_uses_the_official_importer_contract() {
4809 let mut service = HarnessSessionService::new();
4810 let source = opencode_locator();
4811 let response = service.handle(request(
4812 1,
4813 "harness.v1.sessions.handoff",
4814 json!({
4815 "locator": source,
4816 "target_harness": "grok",
4817 "cwd": "/tmp/grok-handoff-project",
4818 }),
4819 ));
4820 let result = &response["result"];
4821
4822 assert_eq!(result["artifact"]["target_harness"], "claude-code");
4826 assert!(result["artifact"]["suggested_filename"]
4827 .as_str()
4828 .unwrap()
4829 .ends_with(".grok-import.claude-code.jsonl"));
4830 let artifact = Session::load_str(
4831 result["artifact"]["content"].as_str().unwrap(),
4832 SessionFormat::ClaudeCode,
4833 )
4834 .unwrap();
4835 assert_eq!(
4836 artifact.meta.cwd.as_deref(),
4837 Some(Path::new("/tmp/grok-handoff-project"))
4838 );
4839 let target_session_id = artifact.meta.session_id.as_deref().unwrap();
4840 assert_eq!(target_session_id.len(), 36);
4841 assert_eq!(target_session_id.as_bytes()[14], b'4');
4842 assert_ne!(target_session_id, opencode_locator().session_id);
4843 assert_eq!(
4844 result["artifact"]["session_id"],
4845 artifact.meta.session_id.as_deref().unwrap()
4846 );
4847
4848 assert_eq!(
4849 result["materialize"]["arguments"],
4850 json!(["import", "--json", "{artifact_path}"])
4851 );
4852 assert_eq!(
4853 result["launch"]["arguments"],
4854 json!(["--resume", "{imported_session_id}", "--fork-session"])
4855 );
4856 assert!(result["note"]
4857 .as_str()
4858 .unwrap()
4859 .contains("outcome=imported"));
4860 assert!(!result["launch"]["arguments"]
4861 .as_array()
4862 .unwrap()
4863 .iter()
4864 .any(|argument| argument == &opencode_locator().session_id));
4865 }
4866
4867 #[tokio::test]
4868 async fn inventory_rejects_unknown_harnesses_and_runtime_attach_is_honest() {
4869 let mut service = HarnessSessionService::new();
4870 let inventory = service
4871 .handle_async(request(
4872 1,
4873 "harness.v1.harnesses.list",
4874 json!({"harnesses": ["missing"]}),
4875 ))
4876 .await;
4877 assert_eq!(inventory["error"]["code"], -32602);
4878
4879 let attached = service
4880 .handle_async(request(
4881 2,
4882 "harness.v1.runtimes.attach_existing",
4883 json!({"harness": "codex", "runtime_id": "thread-1"}),
4884 ))
4885 .await;
4886 assert_eq!(attached["error"]["code"], -32000);
4887 assert!(attached["error"]["message"]
4888 .as_str()
4889 .unwrap()
4890 .contains("runtimes.resume"));
4891 }
4892
4893 #[test]
4894 fn invalid_params_and_unknown_methods_use_json_rpc_errors() {
4895 let mut service = HarnessSessionService::new();
4896 let invalid = service.handle(request(1, "harness.v1.sessions.load", json!({})));
4897 assert_eq!(invalid["error"]["code"], -32602);
4898 let unknown = service.handle(request(2, "harness.v1.unknown", json!({})));
4899 assert_eq!(unknown["error"]["code"], -32601);
4900 }
4901
4902 #[cfg(unix)]
4903 #[tokio::test]
4904 #[allow(clippy::await_holding_lock)]
4907 async fn async_service_drives_a_generic_acp_runtime() {
4908 let _environment_guard = crate::live_runtime::test_environment_lock();
4909 let script = r#"
4910 i=0
4911 while IFS= read -r line; do
4912 i=$((i + 1))
4913 case "$i" in
4914 1) printf '%s\n' '{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentCapabilities":{},"authMethods":[]}}' ;;
4915 2) printf '%s\n' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"svc_acp"}}' ;;
4916 3)
4917 printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"ok"}}}}'
4918 printf '%s\n' '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
4919 ;;
4920 4)
4921 printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"from terminal"}}}}'
4922 printf '%s\n' '{"jsonrpc":"2.0","id":4,"result":{"stopReason":"end_turn"}}'
4923 ;;
4924 esac
4925 done
4926 "#;
4927 let mut service = HarnessSessionService::new();
4928 let started = service
4929 .handle_async(request(
4930 1,
4931 "harness.v1.runtimes.start",
4932 json!({
4933 "harness": "codex",
4934 "protocol": "acp",
4935 "cwd": std::env::current_dir().unwrap(),
4936 "launch": {"program": "/bin/sh", "arguments": ["-c", script], "env": {}},
4937 }),
4938 ))
4939 .await;
4940 assert_eq!(started["result"]["connection"], "runtime-1");
4941 assert_eq!(started["result"]["handle"]["runtime_id"], "svc_acp");
4942
4943 let terminal = service
4944 .handle_async(request(
4945 9,
4946 "harness.v1.runtimes.terminal_instructions",
4947 json!({"connection":"runtime-1"}),
4948 ))
4949 .await;
4950 let arguments = terminal["result"]["launch"]["arguments"]
4951 .as_array()
4952 .expect("hosted runtime should return terminal arguments");
4953 let endpoint_index = arguments
4954 .iter()
4955 .position(|value| value == "--endpoint")
4956 .expect("terminal command should use an opaque endpoint");
4957 let endpoint = LiveRuntimeEndpoint::parse(
4958 arguments[endpoint_index + 1]
4959 .as_str()
4960 .expect("endpoint argument should be text"),
4961 )
4962 .unwrap();
4963 assert!(!terminal.to_string().contains("Bearer"));
4964 let workspace = std::env::current_dir().unwrap();
4965 let receipt = resolve_live_runtime(
4966 &endpoint,
4967 &LiveRuntimeSource {
4968 harness: "codex".into(),
4969 session_id: "svc_acp".into(),
4970 workspace,
4971 },
4972 )
4973 .unwrap();
4974 let remote = crate::HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
4975 .await
4976 .unwrap();
4977 let mut attachment = crate::FrontendRuntime::attach(remote.as_ref(), 100)
4978 .await
4979 .unwrap();
4980
4981 let sent = service
4982 .handle_async(request(
4983 2,
4984 "harness.v1.runtimes.send_input",
4985 json!({"connection": "runtime-1", "text": "hi"}),
4986 ))
4987 .await;
4988 assert_eq!(sent["result"]["turn_id"], "3");
4989
4990 let mut events = Vec::new();
4991 for _ in 0..20 {
4992 events.extend(service.poll_runtimes().await);
4993 if events.len() >= 2 {
4994 break;
4995 }
4996 tokio::time::sleep(Duration::from_millis(2)).await;
4997 }
4998 assert!(events
4999 .iter()
5000 .any(|event| { event["params"]["event"]["kind"] == "session/update" }));
5001 assert!(events.iter().any(|event| {
5002 event["params"]["event"]["kind"] == "supercode/acp_request_completed"
5003 }));
5004
5005 let saw_editor_reply = tokio::time::timeout(Duration::from_secs(2), async {
5006 loop {
5007 let event = attachment.next_event().await.unwrap();
5008 if event.kind == "text_delta" && event.payload["text"] == "ok" {
5009 break;
5010 }
5011 }
5012 })
5013 .await;
5014 assert!(
5015 saw_editor_reply.is_ok(),
5016 "terminal should observe the editor-driven turn"
5017 );
5018
5019 crate::FrontendRuntime::submit(remote.as_ref(), "DRIVE FROM TERMINAL".into())
5020 .await
5021 .unwrap();
5022 let saw_terminal_reply = tokio::time::timeout(Duration::from_secs(2), async {
5023 loop {
5024 let event = attachment.next_event().await.unwrap();
5025 if event.kind == "text_delta" && event.payload["text"] == "from terminal" {
5026 break;
5027 }
5028 }
5029 })
5030 .await;
5031 assert!(
5032 saw_terminal_reply.is_ok(),
5033 "terminal should drive the same runtime"
5034 );
5035
5036 let closed = service
5037 .handle_async(request(
5038 3,
5039 "harness.v1.runtimes.close",
5040 json!({"connection": "runtime-1"}),
5041 ))
5042 .await;
5043 assert_eq!(closed["result"]["closed"], true);
5044 }
5045}