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