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_sessions, load_session, load_session_with_fidelity, SdkCapabilities, SdkError,
19 SdkErrorCode, SdkEvent, SdkOperation, SdkRequest, SdkRuntimeEvent, SdkService,
20};
21use crate::watch::normalized_session_json;
22use crate::Fidelity;
23#[cfg(feature = "adapter-api")]
24use crate::SupercodeHttpRuntimeBackend;
25use crate::{
26 discover_live_runtime, harness_support_registry, AcpRuntimeBackend, ClaudeCodeRuntimeBackend,
27 CodexRuntimeBackend, DiscoveryQuery, HarnessCatalog, HarnessId, ImplementationKind,
28 LiveRuntimeEndpoint, LiveRuntimeSource, OpenCodeRuntimeBackend, PiRuntimeBackend,
29 RuntimeAttachRequest, RuntimeBackend, RuntimeConnection, RuntimeInput, RuntimeLaunch,
30 RuntimeStartRequest, Session, SessionFollower, SessionFormat, SessionLocator, SessionSource,
31};
32#[cfg(feature = "adapter-api")]
33use crate::{register_live_runtime, resolve_live_runtime, LiveRuntimeRegistration};
34
35pub const HARNESS_SERVICE_VERSION: &str = "harness.v1";
37pub const SESSION_EVENT_METHOD: &str = "harness.v1.sessions.event";
39pub const RUNTIME_EVENT_METHOD: &str = "harness.v1.runtimes.event";
41
42pub struct HarnessSessionService {
45 catalog: HarnessCatalog,
46 followers: BTreeMap<String, SessionFollower>,
47 followed_sources: BTreeMap<String, FollowedSource>,
48 next_subscription: u64,
49 runtimes: BTreeMap<String, Box<dyn RuntimeConnection>>,
50 terminal_launches: BTreeMap<String, StructuredLaunch>,
51 runtime_sequences: BTreeMap<String, u64>,
52 next_runtime: u64,
53}
54
55impl Default for HarnessSessionService {
56 fn default() -> Self {
57 Self::new()
58 }
59}
60
61impl HarnessSessionService {
62 pub fn new() -> Self {
64 Self {
65 catalog: HarnessCatalog::new(),
66 followers: BTreeMap::new(),
67 followed_sources: BTreeMap::new(),
68 next_subscription: 1,
69 runtimes: BTreeMap::new(),
70 terminal_launches: BTreeMap::new(),
71 runtime_sequences: BTreeMap::new(),
72 next_runtime: 1,
73 }
74 }
75
76 #[cfg(feature = "adapter-api")]
78 pub fn handle(&mut self, request: Value) -> Value {
79 let id = request.get("id").cloned().unwrap_or(Value::Null);
80 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
81 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
82 }
83 let Some(method) = request.get("method").and_then(Value::as_str) else {
84 return rpc_error(id, -32600, "request is missing `method`");
85 };
86 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
87 match self.call(method, params) {
88 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
89 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
90 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
91 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
92 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
93 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
94 }
95 }
96
97 #[cfg(feature = "adapter-api")]
100 pub async fn handle_async(&mut self, request: Value) -> Value {
101 let method = request
102 .get("method")
103 .and_then(Value::as_str)
104 .unwrap_or_default();
105 if matches!(
106 method,
107 "harness.v1.harnesses.list" | "harness.v1.harnesses.probe"
108 ) {
109 let id = request.get("id").cloned().unwrap_or(Value::Null);
110 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
111 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
112 }
113 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
114 return match self.inventory_call(method, params).await {
115 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
116 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
117 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
118 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
119 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
120 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
121 };
122 }
123 if method == "harness.v1.sessions.message" {
124 let id = request.get("id").cloned().unwrap_or(Value::Null);
125 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
126 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
127 }
128 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
129 return match self.message_call(params).await {
130 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
131 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
132 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
133 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
134 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
135 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
136 };
137 }
138 if let Some(operation) = SdkOperation::from_method(method) {
139 let id = request.get("id").cloned().unwrap_or(Value::Null);
140 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
141 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
142 }
143 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
144 return match self.execute(SdkRequest { operation, params }).await {
145 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
146 Err(error) => sdk_rpc_error(id, &error),
147 };
148 }
149 if !method.starts_with("harness.v1.runtimes.") {
150 return self.handle(request);
151 }
152 let id = request.get("id").cloned().unwrap_or(Value::Null);
153 if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
154 return rpc_error(id, -32600, "expected a JSON-RPC 2.0 request");
155 }
156 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
157 match self.runtime_call(method, params).await {
158 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
159 Err(ServiceError::InvalidParams(message)) => rpc_error(id, -32602, &message),
160 Err(ServiceError::MethodNotFound) => rpc_error(id, -32601, "method not found"),
161 Err(ServiceError::UnsupportedAction(message)) => rpc_error(id, -32020, &message),
162 Err(ServiceError::Operation(message)) => rpc_error(id, -32000, &message),
163 Err(ServiceError::Sdk(error)) => sdk_rpc_error(id, &error),
164 }
165 }
166
167 #[cfg(feature = "adapter-api")]
170 pub fn poll(&mut self) -> Vec<Value> {
171 let mut notifications = Vec::new();
172 for (subscription, follower) in &mut self.followers {
173 match follower.poll() {
174 Ok(Some(event)) => notifications.push(json!({
175 "jsonrpc": "2.0",
176 "method": SESSION_EVENT_METHOD,
177 "params": {
178 "subscription": subscription,
179 "event": event.to_json(),
180 }
181 })),
182 Ok(None) => {}
183 Err(error) => notifications.push(json!({
184 "jsonrpc": "2.0",
185 "method": SESSION_EVENT_METHOD,
186 "params": {
187 "subscription": subscription,
188 "event": {
189 "type": "watch_error",
190 "recoverable": true,
191 "message": error.to_string(),
192 },
193 }
194 })),
195 }
196 }
197 notifications
198 }
199
200 #[cfg(feature = "adapter-api")]
211 pub async fn poll_session_runtime_states(&mut self) -> Vec<Value> {
212 let registry = crate::LocalRuntimeRegistry::new();
213 let authorization = crate::RuntimeAuthorization::observer();
214 let mut notifications = Vec::new();
215 for (subscription, source) in &mut self.followed_sources {
216 let state = match registry
217 .source_state(&source.harness, &source.session_id, &authorization)
218 .await
219 {
220 Ok(Some(state)) => state,
221 Ok(None) => crate::RuntimeRegistryState::Persisted,
222 Err(_) => continue,
224 };
225 if source.reported.as_deref() == Some(state.as_str()) {
226 continue;
227 }
228 source.reported = Some(state.as_str().to_string());
229 notifications.push(json!({
230 "jsonrpc": "2.0",
231 "method": SESSION_EVENT_METHOD,
232 "params": {
233 "subscription": subscription,
234 "event": {"type": "runtime_state", "state": state.as_str()},
235 },
236 }));
237 }
238 notifications
239 }
240
241 #[cfg(feature = "adapter-api")]
243 pub async fn poll_runtimes(&mut self) -> Vec<Value> {
244 self.poll_sdk_events()
245 .await
246 .into_iter()
247 .map(|(connection, runtime_event)| {
248 json!({
249 "jsonrpc": "2.0",
250 "method": RUNTIME_EVENT_METHOD,
251 "params": {
252 "connection": connection,
253 "session_id": runtime_event.session_id,
254 "sequence": runtime_event.event.sequence,
255 "event": {
256 "kind": runtime_event.event.kind,
257 "payload": runtime_event.event.payload,
258 },
259 },
260 })
261 })
262 .collect()
263 }
264
265 async fn poll_sdk_events(&mut self) -> Vec<(String, SdkRuntimeEvent)> {
266 let mut events = Vec::new();
267 let mut closed = Vec::new();
268 for (connection, runtime) in &mut self.runtimes {
269 let session_id = runtime.handle().runtime_id.clone();
270 match tokio::time::timeout(Duration::from_millis(1), runtime.next_event()).await {
271 Ok(Ok(Some(event))) => {
272 let terminal = event.kind == "transport_closed";
273 let next_sequence = self
274 .runtime_sequences
275 .entry(session_id.clone())
276 .or_insert(0);
277 let sequence = event.sequence.unwrap_or_else(|| {
278 *next_sequence = next_sequence.saturating_add(1);
279 *next_sequence
280 });
281 *next_sequence = (*next_sequence).max(sequence);
282 events.push((
283 connection.clone(),
284 SdkRuntimeEvent {
285 session_id: session_id.clone(),
286 event: SdkEvent {
287 sequence,
288 kind: event.kind,
289 payload: event.payload,
290 },
291 },
292 ));
293 if terminal {
294 closed.push(connection.clone());
295 }
296 }
297 Ok(Ok(None)) => {
298 let sequence = self
299 .runtime_sequences
300 .entry(session_id.clone())
301 .or_insert(0);
302 *sequence = sequence.saturating_add(1);
303 events.push((
304 connection.clone(),
305 SdkRuntimeEvent {
306 session_id,
307 event: SdkEvent {
308 sequence: *sequence,
309 kind: "transport_closed".into(),
310 payload: json!({"message": "Harness runtime transport closed."}),
311 },
312 },
313 ));
314 closed.push(connection.clone());
315 }
316 Err(_) => {}
317 Ok(Err(error)) => {
318 let sequence = self
319 .runtime_sequences
320 .entry(session_id.clone())
321 .or_insert(0);
322 *sequence = sequence.saturating_add(1);
323 events.push((
324 connection.clone(),
325 SdkRuntimeEvent {
326 session_id,
327 event: SdkEvent {
328 sequence: *sequence,
329 kind: "transport_error".into(),
330 payload: json!({"message": error.to_string(), "terminal": true}),
331 },
332 },
333 ));
334 closed.push(connection.clone());
335 }
336 }
337 }
338 for connection in closed {
339 if let Some(runtime) = self.runtimes.remove(&connection) {
340 self.runtime_sequences.remove(&runtime.handle().runtime_id);
341 }
342 self.terminal_launches.remove(&connection);
343 }
344 events
345 }
346
347 fn call(&mut self, method: &str, params: Value) -> std::result::Result<Value, ServiceError> {
348 match method {
349 "harness.v1.capabilities" => Ok(json!({
350 "version": HARNESS_SERVICE_VERSION,
351 "sdk": self.capabilities(),
352 "methods": [
353 "harness.v1.support.report",
354 "harness.v1.harnesses.list",
355 "harness.v1.harnesses.probe",
356 "harness.v1.sessions.discover",
357 "harness.v1.sessions.load",
358 "harness.v1.sessions.follow",
359 "harness.v1.sessions.unfollow",
360 "harness.v1.sessions.message",
361 "harness.v1.sessions.import",
362 "harness.v1.sessions.export",
363 "harness.v1.sessions.translate",
364 "harness.v1.sessions.branch",
365 "harness.v1.sessions.handoff",
366 "harness.v1.sessions.resume_instructions",
367 "harness.v1.runtimes.capabilities",
368 "harness.v1.runtimes.start",
369 "harness.v1.runtimes.resume",
370 "harness.v1.runtimes.attach_existing",
371 "harness.v1.runtimes.attach",
372 "harness.v1.runtimes.send_input",
373 "harness.v1.runtimes.interrupt",
374 "harness.v1.runtimes.steer",
375 "harness.v1.runtimes.respond",
376 "harness.v1.runtimes.terminal_instructions",
377 "harness.v1.runtimes.close",
378 ],
379 "notifications": [SESSION_EVENT_METHOD, RUNTIME_EVENT_METHOD],
380 "harnesses": harness_support_registry()
381 .harnesses
382 .into_iter()
383 .map(|harness| harness.id)
384 .collect::<Vec<_>>(),
385 })),
386 "harness.v1.support.report" => serde_json::to_value(harness_support_registry())
387 .map_err(|error| ServiceError::Operation(error.to_string())),
388 "harness.v1.sessions.discover" => {
389 let query = decode::<DiscoveryQuery>(params)?;
390 let sessions = discover_sessions(&query).map_err(operation)?;
391 let peers = if sessions
396 .iter()
397 .any(|session| session.locator.harness.as_str() == HarnessId::CLAUDE_CODE)
398 {
399 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(
400 &query.homes,
401 ))
402 } else {
403 Vec::new()
404 };
405 let sessions = sessions
406 .into_iter()
407 .map(|session| {
408 let mut value = serde_json::to_value(&session)
409 .map_err(|error| ServiceError::Operation(error.to_string()))?;
410 if let Some(workspace) = &session.cwd {
411 let source = LiveRuntimeSource {
412 harness: session.locator.harness.as_str().to_string(),
413 session_id: session.locator.session_id.clone(),
414 workspace: workspace.clone(),
415 };
416 if let Some(endpoint) = discover_live_runtime(&source)
417 .map_err(|error| ServiceError::Operation(error.to_string()))?
418 {
419 value["live_endpoint"] = json!(endpoint.as_str());
420 }
421 }
422 if let Some(peer) = peers.iter().find(|peer| {
423 session.locator.harness.as_str() == HarnessId::CLAUDE_CODE
424 && peer.session_id == session.locator.session_id
425 }) {
426 if value.get("live_endpoint").is_none() {
430 value["live_endpoint"] = json!(peer.endpoint().as_str());
431 }
432 if let Some(status) = peer.status {
433 value["live_status"] = json!(status.as_str());
434 }
435 }
436 Ok(value)
437 })
438 .collect::<std::result::Result<Vec<_>, ServiceError>>()?;
439 Ok(json!({"sessions": sessions}))
440 }
441 "harness.v1.sessions.load" => {
442 let params = decode::<LocatorParams>(params)?;
443 load_session_with_fidelity(¶ms.locator, params.read_fidelity())
444 .map(|session| json!({"session": normalized_session_json(&session)}))
445 .map_err(operation)
446 }
447 "harness.v1.sessions.follow" => {
448 let params = decode::<LocatorParams>(params)?;
449 let mut follower = self
450 .catalog
451 .follow_with_fidelity(¶ms.locator, params.read_fidelity())
452 .map_err(operation)?;
453 let initial = follower
454 .poll()
455 .map_err(operation)?
456 .map(|event| event.to_json());
457 let subscription = format!("sub-{}", self.next_subscription);
458 self.next_subscription += 1;
459 self.followers.insert(subscription.clone(), follower);
460 self.followed_sources.insert(
461 subscription.clone(),
462 FollowedSource {
463 harness: params.locator.harness.as_str().to_string(),
464 session_id: params.locator.session_id.clone(),
465 reported: None,
466 },
467 );
468 Ok(json!({"subscription": subscription, "initial": initial}))
469 }
470 "harness.v1.sessions.unfollow" => {
471 let params = decode::<UnfollowParams>(params)?;
472 self.followed_sources.remove(¶ms.subscription);
473 Ok(json!({
474 "removed": self.followers.remove(¶ms.subscription).is_some()
475 }))
476 }
477 "harness.v1.sessions.import" => {
478 let params = decode::<ImportSessionParams>(params)?;
479 let session = Session::load_str(¶ms.content, params.source_harness.into())
480 .map_err(operation)?;
481 Ok(json!({"session": normalized_session_json(&session)}))
482 }
483 "harness.v1.sessions.export" | "harness.v1.sessions.translate" => {
484 let params = decode::<ExportSessionParams>(params)?;
485 let session = load_session(¶ms.locator).map_err(operation)?;
486 let artifact = session_artifact(¶ms.locator, &session, params.target_harness)?;
487 Ok(json!({"artifact": artifact}))
488 }
489 "harness.v1.sessions.branch" => {
490 let params = decode::<BranchSessionParams>(params)?;
491 let session = load_session(¶ms.locator).map_err(operation)?;
492 let storage = params.locator.storage.path().display().to_string();
493 let bootstrap_prompt = format!(
494 "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.",
495 params.locator.harness.as_str(), params.locator.session_id, storage
496 );
497 let artifact = params
498 .target_harness
499 .map(|target| session_artifact(¶ms.locator, &session, target))
500 .transpose()?;
501 Ok(json!({
502 "parent": params.locator,
503 "session": normalized_session_json(&session),
504 "bootstrap_prompt": bootstrap_prompt,
505 "artifact": artifact,
506 }))
507 }
508 "harness.v1.sessions.handoff" => {
509 let params = decode::<HandoffSessionParams>(params)?;
510 let session = load_session(¶ms.locator).map_err(operation)?;
511 let cwd = params
512 .cwd
513 .or_else(|| session.meta.cwd.clone())
514 .unwrap_or_else(|| PathBuf::from("."));
515 let artifact =
516 handoff_artifact(¶ms.locator, &session, params.target_harness, &cwd)?;
517 let target_session_id = artifact.session_id.as_deref().ok_or_else(|| {
518 ServiceError::Operation(
519 "handoff artifact omitted target session identity".into(),
520 )
521 })?;
522 let instructions =
523 handoff_instructions(params.target_harness, target_session_id, &cwd);
524 Ok(json!({
525 "artifact": artifact,
526 "launch": instructions.launch,
527 "materialize": instructions.materialize,
528 "requires_materialization": instructions.requires_materialization,
529 "note": instructions.note,
530 }))
531 }
532 "harness.v1.sessions.resume_instructions" => {
533 let params = decode::<ResumeInstructionsParams>(params)?;
534 let session = load_session(¶ms.locator).map_err(operation)?;
535 let cwd = params
536 .cwd
537 .or(session.meta.cwd)
538 .unwrap_or_else(|| PathBuf::from("."));
539 let launch = resume_launch(
540 params.locator.harness.as_str(),
541 ¶ms.locator.session_id,
542 &cwd,
543 params.policy,
544 )?;
545 Ok(json!({"launch": launch}))
546 }
547 _ => Err(ServiceError::MethodNotFound),
548 }
549 }
550
551 async fn runtime_call(
552 &mut self,
553 method: &str,
554 params: Value,
555 ) -> std::result::Result<Value, ServiceError> {
556 match method {
557 "harness.v1.runtimes.capabilities" => {
558 let params = decode::<RuntimeBackendParams>(params)?;
559 let backend = runtime_backend(¶ms)?;
560 Ok(json!({
561 "harness": backend.harness(),
562 "capabilities": backend.capabilities(),
563 }))
564 }
565 "harness.v1.runtimes.start" => {
566 let params = decode::<RuntimeStartParams>(params)?;
567 let backend = runtime_backend(¶ms.backend)?;
568 let capabilities = backend.capabilities();
569 let workspace = params.cwd.clone();
570 let runtime = backend
571 .start(RuntimeStartRequest {
572 cwd: params.cwd,
573 launch: runtime_launch(¶ms.backend),
574 })
575 .await
576 .map_err(operation)?;
577 self.insert_hosted_runtime(runtime, capabilities, workspace)
578 .await
579 }
580 "harness.v1.runtimes.resume" | "harness.v1.runtimes.attach" => {
581 let params = decode::<RuntimeAttachParams>(params)?;
582 let backend = runtime_backend(¶ms.backend)?;
583 let capabilities = backend.capabilities();
584 let workspace = params.cwd.clone().unwrap_or_else(|| {
585 std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."))
586 });
587 let runtime = backend
588 .attach(RuntimeAttachRequest {
589 runtime_id: params.runtime_id,
590 cwd: params.cwd,
591 launch: runtime_launch(¶ms.backend),
592 })
593 .await
594 .map_err(operation)?;
595 self.insert_hosted_runtime(runtime, capabilities, workspace)
596 .await
597 }
598 "harness.v1.runtimes.attach_existing" => {
599 let params = decode::<RuntimeAttachParams>(params)?;
600 let backend: Box<dyn RuntimeBackend> = match params
601 .backend
602 .base_url
603 .as_deref()
604 .and_then(|value| LiveRuntimeEndpoint::parse(value).ok())
605 {
606 Some(endpoint) => {
607 #[cfg(not(feature = "adapter-api"))]
608 {
609 let _ = endpoint;
610 return Err(ServiceError::UnsupportedAction(
611 "live HTTP attachment adapter is not compiled".into(),
612 ));
613 }
614 #[cfg(feature = "adapter-api")]
615 {
616 let workspace = params.cwd.clone().ok_or_else(|| {
617 ServiceError::InvalidParams(
618 "Supercode live attach requires the project cwd".into(),
619 )
620 })?;
621 let source = LiveRuntimeSource {
622 harness: params.backend.harness.as_str().to_string(),
623 session_id: params.runtime_id.clone(),
624 workspace,
625 };
626 let receipt = resolve_live_runtime(&endpoint, &source)
627 .map_err(|error| ServiceError::Operation(error.to_string()))?;
628 Box::new(SupercodeHttpRuntimeBackend::new(receipt))
629 }
630 }
631 None => runtime_backend(¶ms.backend)?,
632 };
633 if !backend.capabilities().attach_existing_process {
634 return Err(ServiceError::Operation(format!(
635 "{} cannot attach to an already-running process; use runtimes.resume for a persisted session",
636 backend.harness().as_str()
637 )));
638 }
639 let runtime = backend
640 .attach_existing(RuntimeAttachRequest {
641 runtime_id: params.runtime_id,
642 cwd: params.cwd,
643 launch: runtime_launch(¶ms.backend),
644 })
645 .await
646 .map_err(operation)?;
647 self.insert_runtime(runtime)
648 }
649 "harness.v1.runtimes.send_input" => {
650 let params = decode::<RuntimeInputParams>(params)?;
651 let runtime = self.runtime_mut(¶ms.connection)?;
652 let turn_id = runtime
653 .send_input(RuntimeInput { text: params.text })
654 .await
655 .map_err(operation)?;
656 Ok(json!({"turn_id": turn_id}))
657 }
658 "harness.v1.runtimes.interrupt" => {
659 let params = decode::<RuntimeConnectionParams>(params)?;
660 self.runtime_mut(¶ms.connection)?
661 .interrupt()
662 .await
663 .map_err(operation)?;
664 Ok(json!({}))
665 }
666 "harness.v1.runtimes.steer" => Err(ServiceError::UnsupportedAction(
667 "steer is not supported by this harness-native runtime adapter".into(),
668 )),
669 "harness.v1.runtimes.respond" => {
670 let params = decode::<RuntimeRespondParams>(params)?;
671 self.runtime_mut(¶ms.connection)?
672 .respond(params.request_id, params.response)
673 .await
674 .map_err(operation)?;
675 Ok(json!({}))
676 }
677 "harness.v1.runtimes.terminal_instructions" => {
678 let params = decode::<RuntimeConnectionParams>(params)?;
679 let launch = self
680 .terminal_launches
681 .get(¶ms.connection)
682 .ok_or_else(|| {
683 ServiceError::Operation(
684 "this runtime is not hosted for terminal attachment".into(),
685 )
686 })?;
687 Ok(json!({"launch":launch}))
688 }
689 "harness.v1.runtimes.close" => {
690 let params = decode::<RuntimeConnectionParams>(params)?;
691 let Some(mut runtime) = self.runtimes.remove(¶ms.connection) else {
692 return Err(ServiceError::InvalidParams(format!(
693 "unknown runtime connection `{}`",
694 params.connection
695 )));
696 };
697 self.terminal_launches.remove(¶ms.connection);
698 self.runtime_sequences.remove(&runtime.handle().runtime_id);
699 runtime.close().await.map_err(operation)?;
700 Ok(json!({"closed": true}))
701 }
702 _ => Err(ServiceError::MethodNotFound),
703 }
704 }
705
706 #[cfg(feature = "adapter-api")]
708 async fn message_call(&self, params: Value) -> std::result::Result<Value, ServiceError> {
709 let params = decode::<MessageSessionParams>(params)?;
710 Ok(message_live_session(¶ms, &crate::claude_peer::ProcessCourierRunner).await)
711 }
712
713 fn insert_runtime(
714 &mut self,
715 runtime: Box<dyn RuntimeConnection>,
716 ) -> std::result::Result<Value, ServiceError> {
717 let connection = format!("runtime-{}", self.next_runtime);
718 self.next_runtime += 1;
719 let handle = runtime.handle().clone();
720 self.runtime_sequences
721 .entry(handle.runtime_id.clone())
722 .or_insert(0);
723 self.runtimes.insert(connection.clone(), runtime);
724 Ok(json!({"connection": connection, "handle": handle}))
725 }
726
727 #[cfg(feature = "adapter-api")]
728 async fn insert_hosted_runtime(
729 &mut self,
730 runtime: Box<dyn RuntimeConnection>,
731 capabilities: crate::RuntimeCapabilities,
732 workspace: PathBuf,
733 ) -> std::result::Result<Value, ServiceError> {
734 let (host, connection) = HostedHarnessRuntime::spawn(runtime, capabilities);
735 let token: std::sync::Arc<str> = crate::server::generate_token().into();
736 let server = crate::server::run_frontend_http(
737 host.clone(),
738 host.frontend_sender(),
739 "127.0.0.1:0",
740 token.clone(),
741 )
742 .await
743 .map_err(|error| ServiceError::Operation(error.to_string()))?;
744 let source = LiveRuntimeSource {
745 harness: connection.handle().harness.as_str().to_string(),
746 session_id: connection.handle().runtime_id.clone(),
747 workspace: workspace.clone(),
748 };
749 let registration = register_live_runtime(
750 connection.handle().runtime_id.clone(),
751 source.clone(),
752 format!("http://{}", server.address()),
753 token.to_string(),
754 )
755 .map_err(|error| ServiceError::Operation(error.to_string()))?;
756 let endpoint = registration.endpoint().to_string();
757 let launch = StructuredLaunch {
758 cwd: workspace,
759 program: std::env::current_exe()
763 .ok()
764 .map(|path| path.to_string_lossy().into_owned())
765 .unwrap_or_else(|| "supercode".into()),
766 arguments: vec![
767 "harness".into(),
768 "attach".into(),
769 "--endpoint".into(),
770 endpoint,
771 "--harness".into(),
772 source.harness,
773 "--session".into(),
774 source.session_id,
775 ],
776 env: BTreeMap::new(),
777 };
778 let lease = HostedRuntimeLease {
779 connection,
780 _host: host,
781 _registration: registration,
782 _server: server,
783 };
784 let opened = self.insert_runtime(Box::new(lease))?;
785 let connection_id = opened["connection"]
786 .as_str()
787 .expect("insert_runtime returns a connection id")
788 .to_string();
789 self.terminal_launches.insert(connection_id, launch);
790 Ok(opened)
791 }
792
793 #[cfg(not(feature = "adapter-api"))]
794 async fn insert_hosted_runtime(
795 &mut self,
796 runtime: Box<dyn RuntimeConnection>,
797 _capabilities: crate::RuntimeCapabilities,
798 _workspace: PathBuf,
799 ) -> std::result::Result<Value, ServiceError> {
800 self.insert_runtime(runtime)
801 }
802
803 fn runtime_mut(
804 &mut self,
805 connection: &str,
806 ) -> std::result::Result<&mut Box<dyn RuntimeConnection>, ServiceError> {
807 self.runtimes.get_mut(connection).ok_or_else(|| {
808 ServiceError::InvalidParams(format!("unknown runtime connection `{connection}`"))
809 })
810 }
811
812 async fn inventory_call(
813 &self,
814 method: &str,
815 params: Value,
816 ) -> std::result::Result<Value, ServiceError> {
817 let mut params = decode::<HarnessInventoryParams>(params)?;
818 if method == "harness.v1.harnesses.probe" {
819 let harness = params.harness.take().ok_or_else(|| {
820 ServiceError::InvalidParams("harnesses.probe requires `harness`".into())
821 })?;
822 params.harnesses = vec![harness];
823 }
824 let selected = params
825 .harnesses
826 .iter()
827 .map(HarnessId::as_str)
828 .collect::<std::collections::BTreeSet<_>>();
829 let supported = harness_support_registry()
830 .harnesses
831 .into_iter()
832 .filter(|descriptor| selected.is_empty() || selected.contains(descriptor.id.as_str()))
833 .collect::<Vec<_>>();
834 if !params.harnesses.is_empty() && supported.len() != selected.len() {
835 let known = supported
836 .iter()
837 .map(|harness| harness.id.as_str())
838 .collect::<std::collections::BTreeSet<_>>();
839 let missing = params
840 .harnesses
841 .iter()
842 .filter(|id| !known.contains(id.as_str()))
843 .map(HarnessId::as_str)
844 .collect::<Vec<_>>();
845 return Err(ServiceError::InvalidParams(format!(
846 "unknown harness(es): {}",
847 missing.join(", ")
848 )));
849 }
850 let global_counts = params
851 .include_sessions
852 .then(|| self.session_counts(None, ¶ms.harnesses));
853 let workspace_counts = params.include_sessions.then(|| {
854 params
855 .workspace
856 .as_deref()
857 .map(|workspace| self.session_counts(Some(workspace), ¶ms.harnesses))
858 });
859 let probes = supported.into_iter().map(|descriptor| {
860 let global = global_counts
861 .as_ref()
862 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
863 let workspace = workspace_counts
864 .as_ref()
865 .and_then(Option::as_ref)
866 .map(|counts| counts.get(descriptor.id.as_str()).copied().unwrap_or(0));
867 self.probe_harness(descriptor, ¶ms, global, workspace)
868 });
869 let harnesses = futures::future::join_all(probes).await;
870 serde_json::to_value(HarnessInventoryReport {
871 probe: params.probe,
872 workspace: params.workspace,
873 harnesses,
874 })
875 .map_err(|error| ServiceError::Operation(error.to_string()))
876 }
877
878 async fn probe_harness(
879 &self,
880 descriptor: crate::HarnessSupportDescriptor,
881 params: &HarnessInventoryParams,
882 global: Option<usize>,
883 workspace: Option<usize>,
884 ) -> LocalHarness {
885 let launch = descriptor.runtime.default_launch.as_ref();
886 let executable = launch.and_then(|launch| find_executable(&launch.program));
887 let installed = executable.is_some();
888 let version = match executable.as_deref() {
889 Some(path) => executable_version(path).await,
890 None => None,
891 };
892 let configured = auth_evidence(descriptor.id.as_str());
893 let mut auth = if configured {
894 HarnessAuthState::Configured
895 } else {
896 HarnessAuthState::Unknown
897 };
898 let mut runtime = if installed {
899 HarnessRuntimeState::Degraded
900 } else {
901 HarnessRuntimeState::Unavailable
902 };
903 let mut reason = (!installed).then(|| {
904 format!(
905 "{} is supported but `{}` was not found on PATH",
906 descriptor.display_name,
907 launch
908 .map(|launch| launch.program.as_str())
909 .unwrap_or("executable")
910 )
911 });
912 let mut repair = (!installed).then(|| {
913 format!(
914 "Install {} and ensure `{}` is on PATH.",
915 descriptor.display_name,
916 launch
917 .map(|launch| launch.program.as_str())
918 .unwrap_or("its executable")
919 )
920 });
921
922 if installed && params.probe == HarnessProbeLevel::Handshake {
923 let backend_params = RuntimeBackendParams {
924 harness: descriptor.id.clone(),
925 protocol: None,
926 launch: None,
927 base_url: None,
928 policy: RuntimePolicy::Default,
929 };
930 match runtime_backend(&backend_params) {
931 Ok(backend) => {
932 let cwd = params
933 .workspace
934 .clone()
935 .or_else(|| std::env::current_dir().ok())
936 .unwrap_or_else(|| PathBuf::from("."));
937 match tokio::time::timeout(
938 Duration::from_secs(30),
939 backend.start(RuntimeStartRequest { cwd, launch: None }),
940 )
941 .await
942 {
943 Ok(Ok(mut connection)) => {
944 match stabilize_handshake(connection.as_mut()).await {
945 Ok(()) => {
946 auth = HarnessAuthState::Ready;
947 runtime = HarnessRuntimeState::Ready;
948 reason = Some(
949 "No-prompt runtime handshake remained healthy through the startup stabilization window; no model request was sent."
950 .into(),
951 );
952 repair = None;
953 }
954 Err(message) => {
955 auth = if looks_like_auth_error(&message) {
956 HarnessAuthState::Required
957 } else if configured {
958 HarnessAuthState::Configured
959 } else {
960 HarnessAuthState::Unknown
961 };
962 reason = Some(format!(
963 "No-prompt runtime handshake became unhealthy during startup: {message}"
964 ));
965 repair = Some(if auth == HarnessAuthState::Required {
966 format!(
967 "Run `{}` interactively once and complete sign-in, then probe again.",
968 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
969 )
970 } else {
971 "Run the harness directly to inspect its startup failure, then probe again."
972 .into()
973 });
974 }
975 }
976 let _ =
977 tokio::time::timeout(Duration::from_secs(3), connection.close())
978 .await;
979 }
980 Ok(Err(error)) => {
981 let message = truncate_text(&error.to_string(), 500);
982 auth = if looks_like_auth_error(&message) {
983 HarnessAuthState::Required
984 } else if configured {
985 HarnessAuthState::Configured
986 } else {
987 HarnessAuthState::Unknown
988 };
989 reason = Some(format!("No-prompt runtime handshake failed: {message}"));
990 repair = Some(if auth == HarnessAuthState::Required {
991 format!(
992 "Run `{}` interactively once and complete sign-in, then probe again.",
993 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
994 )
995 } else {
996 "Check the harness installation and run the handshake probe again."
997 .into()
998 });
999 }
1000 Err(_) => {
1001 reason = Some(
1002 "No-prompt runtime handshake timed out after 30 seconds.".into(),
1003 );
1004 repair = Some("Run the harness directly to check startup or authentication, then probe again.".into());
1005 }
1006 }
1007 }
1008 Err(error) => {
1009 reason = Some(error_message(error));
1010 }
1011 }
1012 } else if installed && configured {
1013 reason = Some("Executable and local authentication evidence found; use a handshake probe to verify readiness.".into());
1014 } else if installed {
1015 reason = Some("Executable found; authentication readiness is unknown until a no-prompt handshake succeeds.".into());
1016 repair =
1017 Some(format!(
1018 "Run `{}` interactively once if sign-in is required, or use `--probe handshake`.",
1019 launch.map(|launch| launch.program.as_str()).unwrap_or("the harness")
1020 ));
1021 }
1022
1023 let effective_capabilities = if installed {
1024 descriptor.runtime.capabilities.clone()
1025 } else {
1026 unavailable_capabilities()
1027 };
1028 LocalHarness {
1029 id: descriptor.id,
1030 display_name: descriptor.display_name,
1031 supported: true,
1032 installed,
1033 executable: executable.map(|path| path.to_string_lossy().into_owned()),
1034 version,
1035 auth,
1036 runtime,
1037 protocol: descriptor.runtime.protocol,
1038 capabilities: descriptor.runtime.capabilities,
1039 effective_capabilities,
1040 sessions: HarnessSessionCounts { global, workspace },
1041 reason,
1042 repair,
1043 }
1044 }
1045
1046 fn session_counts(
1047 &self,
1048 workspace: Option<&Path>,
1049 harnesses: &[HarnessId],
1050 ) -> BTreeMap<String, usize> {
1051 let mut counts = BTreeMap::new();
1052 for session in self
1053 .catalog
1054 .discover(&DiscoveryQuery {
1055 workspace: workspace.map(Path::to_path_buf),
1056 harnesses: harnesses.to_vec(),
1057 ..DiscoveryQuery::default()
1058 })
1059 .unwrap_or_default()
1060 {
1061 *counts
1062 .entry(session.locator.harness.as_str().to_string())
1063 .or_insert(0) += 1;
1064 }
1065 counts
1066 }
1067}
1068
1069#[async_trait::async_trait]
1070impl SdkService for HarnessSessionService {
1071 fn capabilities(&self) -> SdkCapabilities {
1072 SdkCapabilities::default()
1073 }
1074
1075 async fn execute(&mut self, request: SdkRequest) -> Result<Value, SdkError> {
1076 if request.operation == SdkOperation::Events {
1077 let events = self
1078 .poll_sdk_events()
1079 .await
1080 .into_iter()
1081 .map(|(_, event)| event)
1082 .collect::<Vec<_>>();
1083 return serde_json::to_value(events).map_err(|error| {
1084 SdkError::new(
1085 SdkErrorCode::Execution,
1086 request.operation,
1087 error.to_string(),
1088 )
1089 });
1090 }
1091 let method = request
1092 .operation
1093 .method()
1094 .ok_or_else(|| SdkError::unsupported(request.operation))?;
1095 let result = match request.operation {
1096 SdkOperation::Discover | SdkOperation::Load | SdkOperation::Export => {
1097 self.call(method, request.params)
1098 }
1099 SdkOperation::Start
1100 | SdkOperation::Resume
1101 | SdkOperation::Input
1102 | SdkOperation::Interrupt
1103 | SdkOperation::Steer
1104 | SdkOperation::Respond
1105 | SdkOperation::Close => self.runtime_call(method, request.params).await,
1106 SdkOperation::Events => unreachable!("handled before method dispatch"),
1107 };
1108 result.map_err(|error| sdk_error(request.operation, error))
1109 }
1110
1111 async fn events(&mut self) -> Result<Vec<SdkRuntimeEvent>, SdkError> {
1112 Ok(self
1113 .poll_sdk_events()
1114 .await
1115 .into_iter()
1116 .map(|(_, event)| event)
1117 .collect())
1118 }
1119}
1120
1121#[cfg(feature = "adapter-api")]
1122struct HostedRuntimeLease {
1123 connection: HostedHarnessConnection,
1124 _host: std::sync::Arc<HostedHarnessRuntime>,
1125 _registration: LiveRuntimeRegistration,
1126 _server: crate::server::FrontendHttpServer,
1127}
1128
1129#[async_trait::async_trait]
1130#[cfg(feature = "adapter-api")]
1131impl RuntimeConnection for HostedRuntimeLease {
1132 fn handle(&self) -> &crate::RuntimeHandle {
1133 self.connection.handle()
1134 }
1135
1136 async fn send_input(&mut self, input: RuntimeInput) -> crate::Result<Option<String>> {
1137 self.connection.send_input(input).await
1138 }
1139
1140 async fn next_event(&mut self) -> crate::Result<Option<crate::HarnessEvent>> {
1141 self.connection.next_event().await
1142 }
1143
1144 async fn interrupt(&mut self) -> crate::Result<()> {
1145 self.connection.interrupt().await
1146 }
1147
1148 async fn respond(&mut self, request_id: Value, response: Value) -> crate::Result<()> {
1149 self.connection.respond(request_id, response).await
1150 }
1151
1152 async fn close(&mut self) -> crate::Result<()> {
1153 self.connection.close().await
1154 }
1155}
1156
1157async fn stabilize_handshake(connection: &mut dyn RuntimeConnection) -> Result<(), String> {
1158 let deadline = tokio::time::Instant::now() + Duration::from_secs(3);
1159 loop {
1160 let now = tokio::time::Instant::now();
1161 if now >= deadline {
1162 return Ok(());
1163 }
1164 match tokio::time::timeout(deadline - now, connection.next_event()).await {
1165 Err(_) => return Ok(()),
1166 Ok(Ok(Some(event))) => {
1167 if let Some(message) = handshake_event_failure(&event) {
1168 return Err(truncate_text(&message, 500));
1169 }
1170 }
1171 Ok(Ok(None)) => return Err("runtime transport closed during startup".into()),
1172 Ok(Err(error)) => return Err(error.to_string()),
1173 }
1174 }
1175}
1176
1177fn handshake_event_failure(event: &crate::HarnessEvent) -> Option<String> {
1178 let detail = event
1179 .payload
1180 .get("message")
1181 .or_else(|| event.payload.get("line"))
1182 .and_then(Value::as_str)
1183 .unwrap_or(event.kind.as_str());
1184 match event.kind.as_str() {
1185 "transport_closed" => Some("runtime transport closed during startup".into()),
1186 "transport_error" => Some(format!("runtime transport error: {detail}")),
1187 "malformed_output" => Some(format!("runtime emitted non-protocol output: {detail}")),
1188 _ => None,
1193 }
1194}
1195
1196#[derive(Deserialize)]
1197struct LocatorParams {
1198 locator: SessionLocator,
1199 #[serde(default)]
1212 fidelity: Option<Fidelity>,
1213}
1214
1215impl LocatorParams {
1216 fn read_fidelity(&self) -> Fidelity {
1217 self.fidelity.unwrap_or(Fidelity::Semantic)
1218 }
1219}
1220
1221#[derive(Deserialize)]
1222struct UnfollowParams {
1223 subscription: String,
1224}
1225
1226#[derive(Deserialize)]
1227struct MessageSessionParams {
1228 locator: SessionLocator,
1229 text: String,
1230 #[serde(default)]
1233 homes: crate::HarnessHomes,
1234}
1235
1236#[cfg(feature = "adapter-api")]
1248async fn message_live_session(
1249 params: &MessageSessionParams,
1250 runner: &dyn crate::claude_peer::CourierRunner,
1251) -> Value {
1252 if params.locator.harness.as_str() != HarnessId::CLAUDE_CODE {
1253 return json!({
1254 "delivered_to_bus": false,
1255 "refusal": {
1256 "reason": crate::claude_peer::ClaudePeerRefusal::HarnessUnsupported.as_str(),
1257 "message": format!(
1258 "`{}` does not publish a live-session registry; only claude-code sessions can be messaged in place",
1259 params.locator.harness.as_str()
1260 ),
1261 },
1262 });
1263 }
1264 match crate::claude_peer::message_claude_peer(
1265 ¶ms.homes,
1266 ¶ms.locator.session_id,
1267 ¶ms.text,
1268 runner,
1269 )
1270 .await
1271 {
1272 Ok(delivery) => json!({
1273 "delivered_to_bus": true,
1274 "target": {
1275 "session_id": delivery.target.session_id,
1276 "name": delivery.target.name,
1277 "pid": delivery.target.pid,
1278 "cwd": delivery.target.cwd,
1279 "status": delivery.target.status.map(|status| status.as_str()),
1280 },
1281 "courier": {
1282 "model": crate::claude_peer::COURIER_MODEL,
1283 "report": delivery.courier_report,
1284 },
1285 }),
1286 Err(refusal) => json!({
1287 "delivered_to_bus": false,
1288 "refusal": {"reason": refusal.reason.as_str(), "message": refusal.message},
1289 }),
1290 }
1291}
1292
1293#[cfg_attr(not(feature = "adapter-api"), allow(dead_code))]
1298struct FollowedSource {
1299 harness: String,
1300 session_id: String,
1301 reported: Option<String>,
1302}
1303
1304#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize)]
1305#[serde(rename_all = "kebab-case")]
1306enum TransferFormat {
1307 ClaudeCode,
1308 Codex,
1309 #[serde(rename = "opencode", alias = "open-code")]
1310 OpenCode,
1311 Pi,
1312 Grok,
1313}
1314
1315impl TransferFormat {
1316 fn id(self) -> &'static str {
1317 match self {
1318 Self::ClaudeCode => HarnessId::CLAUDE_CODE,
1319 Self::Codex => HarnessId::CODEX,
1320 Self::OpenCode => HarnessId::OPENCODE,
1321 Self::Pi => HarnessId::PI,
1322 Self::Grok => HarnessId::GROK,
1323 }
1324 }
1325}
1326
1327impl From<TransferFormat> for SessionFormat {
1328 fn from(value: TransferFormat) -> Self {
1329 match value {
1330 TransferFormat::ClaudeCode => Self::ClaudeCode,
1331 TransferFormat::Codex => Self::Codex,
1332 TransferFormat::OpenCode => Self::OpenCode,
1333 TransferFormat::Pi => Self::Pi,
1334 TransferFormat::Grok => Self::Grok,
1335 }
1336 }
1337}
1338
1339#[derive(Deserialize)]
1340struct ImportSessionParams {
1341 source_harness: TransferFormat,
1342 content: String,
1343}
1344
1345#[derive(Deserialize)]
1346struct ExportSessionParams {
1347 locator: SessionLocator,
1348 target_harness: TransferFormat,
1349}
1350
1351#[derive(Deserialize)]
1352struct BranchSessionParams {
1353 locator: SessionLocator,
1354 #[serde(default)]
1355 target_harness: Option<TransferFormat>,
1356}
1357
1358#[derive(Deserialize)]
1359struct HandoffSessionParams {
1360 locator: SessionLocator,
1361 target_harness: TransferFormat,
1362 #[serde(default)]
1363 cwd: Option<PathBuf>,
1364}
1365
1366#[derive(Debug, Clone, Copy, Default, Deserialize)]
1367#[serde(rename_all = "snake_case")]
1368enum ResumePolicy {
1369 #[default]
1370 Default,
1371 Yolo,
1372}
1373
1374#[derive(Deserialize)]
1375struct ResumeInstructionsParams {
1376 locator: SessionLocator,
1377 #[serde(default)]
1378 cwd: Option<PathBuf>,
1379 #[serde(default)]
1380 policy: ResumePolicy,
1381}
1382
1383#[derive(Serialize)]
1384struct SessionArtifact {
1385 source_harness: HarnessId,
1386 target_harness: &'static str,
1387 session_id: Option<String>,
1388 content: String,
1389 suggested_filename: String,
1390 files: Vec<SessionArtifactFile>,
1391 fidelity: Fidelity,
1392 residue: Vec<String>,
1393}
1394
1395#[derive(Serialize)]
1396struct SessionArtifactFile {
1397 path: String,
1398 content: String,
1399 role: ArtifactFileRole,
1400}
1401
1402#[derive(Serialize)]
1403#[serde(rename_all = "snake_case")]
1404enum ArtifactFileRole {
1405 Primary,
1406 Subagent,
1407 Bundle,
1408 SourceRecovery,
1409}
1410
1411#[derive(Serialize)]
1412struct StructuredLaunch {
1413 cwd: PathBuf,
1414 program: String,
1415 arguments: Vec<String>,
1416 env: BTreeMap<String, String>,
1417}
1418
1419struct HandoffInstructions {
1420 launch: StructuredLaunch,
1421 materialize: Option<StructuredLaunch>,
1422 requires_materialization: bool,
1423 note: String,
1424}
1425
1426#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
1427#[serde(rename_all = "snake_case")]
1428enum HarnessProbeLevel {
1429 #[default]
1430 Passive,
1431 Handshake,
1432}
1433
1434#[derive(Default, Deserialize)]
1435#[serde(default)]
1436struct HarnessInventoryParams {
1437 harness: Option<HarnessId>,
1438 harnesses: Vec<HarnessId>,
1439 workspace: Option<PathBuf>,
1440 probe: HarnessProbeLevel,
1441 include_sessions: bool,
1442}
1443
1444#[derive(Serialize)]
1445struct HarnessInventoryReport {
1446 probe: HarnessProbeLevel,
1447 workspace: Option<PathBuf>,
1448 harnesses: Vec<LocalHarness>,
1449}
1450
1451#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1452#[serde(rename_all = "snake_case")]
1453enum HarnessAuthState {
1454 Ready,
1455 Configured,
1456 Required,
1457 Unknown,
1458}
1459
1460#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
1461#[serde(rename_all = "snake_case")]
1462enum HarnessRuntimeState {
1463 Ready,
1464 Degraded,
1465 Unavailable,
1466}
1467
1468#[derive(Serialize)]
1469struct HarnessSessionCounts {
1470 global: Option<usize>,
1471 workspace: Option<usize>,
1472}
1473
1474#[derive(Serialize)]
1475struct LocalHarness {
1476 id: HarnessId,
1477 display_name: String,
1478 supported: bool,
1479 installed: bool,
1480 executable: Option<String>,
1481 version: Option<String>,
1482 auth: HarnessAuthState,
1483 runtime: HarnessRuntimeState,
1484 protocol: String,
1485 capabilities: crate::RuntimeCapabilities,
1486 effective_capabilities: crate::RuntimeCapabilities,
1487 sessions: HarnessSessionCounts,
1488 reason: Option<String>,
1489 repair: Option<String>,
1490}
1491
1492#[derive(Clone, Deserialize)]
1493struct RuntimeBackendParams {
1494 harness: HarnessId,
1495 #[serde(default)]
1496 protocol: Option<String>,
1497 #[serde(default)]
1498 launch: Option<RuntimeLaunch>,
1499 #[serde(default)]
1500 base_url: Option<String>,
1501 #[serde(default)]
1502 policy: RuntimePolicy,
1503}
1504
1505#[derive(Debug, Clone, Copy, Default, Deserialize)]
1506#[serde(rename_all = "snake_case")]
1507enum RuntimePolicy {
1508 #[default]
1509 Default,
1510 Yolo,
1511}
1512
1513#[derive(Deserialize)]
1514struct RuntimeStartParams {
1515 #[serde(flatten)]
1516 backend: RuntimeBackendParams,
1517 cwd: PathBuf,
1518}
1519
1520#[derive(Deserialize)]
1521struct RuntimeAttachParams {
1522 #[serde(flatten)]
1523 backend: RuntimeBackendParams,
1524 runtime_id: String,
1525 #[serde(default)]
1526 cwd: Option<PathBuf>,
1527}
1528
1529#[derive(Deserialize)]
1530struct RuntimeConnectionParams {
1531 connection: String,
1532}
1533
1534#[derive(Deserialize)]
1535struct RuntimeInputParams {
1536 connection: String,
1537 text: String,
1538}
1539
1540#[derive(Deserialize)]
1541struct RuntimeRespondParams {
1542 connection: String,
1543 request_id: Value,
1544 response: Value,
1545}
1546
1547fn session_artifact(
1548 locator: &SessionLocator,
1549 session: &Session,
1550 target: TransferFormat,
1551) -> std::result::Result<SessionArtifact, ServiceError> {
1552 session_artifact_with_id(locator, session, target, None)
1553}
1554
1555fn session_artifact_with_id(
1556 locator: &SessionLocator,
1557 session: &Session,
1558 target: TransferFormat,
1559 target_session_id: Option<&str>,
1560) -> std::result::Result<SessionArtifact, ServiceError> {
1561 let format: SessionFormat = target.into();
1562 let diagonal = format.source() == session.meta.source;
1563 let has_appended_turns = session
1564 .imported_message_count
1565 .is_some_and(|imported| imported < session.messages.len());
1566 let content = if let Some(id) = target_session_id {
1567 if diagonal && format != SessionFormat::OpenCode {
1568 session
1569 .to_jsonl_spliced(format, Some(id))
1570 .map_err(operation)?
1571 } else {
1572 let mut rewritten = session.clone();
1573 rewritten.meta.session_id = Some(id.to_string());
1574 rewritten.to_jsonl(format).map_err(operation)?
1575 }
1576 } else if diagonal && session.raw_is_verbatim && !has_appended_turns {
1577 session.raw_verbatim()
1578 } else if diagonal {
1579 session.to_jsonl_spliced(format, None).map_err(operation)?
1580 } else {
1581 session.to_jsonl(format).map_err(operation)?
1582 };
1583 let stem = sanitize_filename(
1584 target_session_id
1585 .or(session.meta.session_id.as_deref())
1586 .unwrap_or(&locator.session_id),
1587 );
1588 let suggested_filename = if diagonal && target == TransferFormat::Grok {
1589 "chat_history.jsonl".to_string()
1590 } else {
1591 format!("{stem}.{}.jsonl", target.id())
1592 };
1593 let mut files = vec![SessionArtifactFile {
1594 path: suggested_filename.clone(),
1595 content: content.clone(),
1596 role: ArtifactFileRole::Primary,
1597 }];
1598 if target == TransferFormat::ClaudeCode {
1599 let bundle_stem = Path::new(&suggested_filename)
1600 .file_stem()
1601 .and_then(|stem| stem.to_str())
1602 .unwrap_or(&stem);
1603 let mut child_paths = BTreeSet::new();
1604 for (index, subagent) in session.subagents.iter().enumerate() {
1605 let agent_id = subagent
1606 .meta
1607 .agent_id
1608 .as_deref()
1609 .map(|id| id.strip_prefix("agent-").unwrap_or(id))
1610 .map(sanitize_filename)
1611 .filter(|id| !id.is_empty())
1612 .unwrap_or_else(|| format!("subagent-{}", index + 1));
1613 let child_has_appended_turns = subagent
1614 .imported_message_count
1615 .is_some_and(|imported| imported < subagent.messages.len());
1616 let child_content = if target_session_id.is_none()
1617 && subagent.meta.source == SessionSource::ClaudeCode
1618 && subagent.raw_is_verbatim
1619 && !child_has_appended_turns
1620 {
1621 subagent.raw_verbatim()
1622 } else if subagent.meta.source == SessionSource::ClaudeCode {
1623 subagent
1624 .to_jsonl_spliced(SessionFormat::ClaudeCode, target_session_id)
1625 .map_err(operation)?
1626 } else {
1627 let mut child = subagent.clone();
1628 if let Some(id) = target_session_id {
1629 child.meta.session_id = Some(id.to_string());
1630 }
1631 child
1632 .to_jsonl(SessionFormat::ClaudeCode)
1633 .map_err(operation)?
1634 };
1635 let path = format!("{bundle_stem}/subagents/agent-{agent_id}.jsonl");
1636 if !child_paths.insert(path.clone()) {
1637 return Err(ServiceError::Operation(format!(
1638 "Claude subagent ids collide at artifact path `{path}`"
1639 )));
1640 }
1641 files.push(SessionArtifactFile {
1642 path,
1643 content: child_content,
1644 role: ArtifactFileRole::Subagent,
1645 });
1646 }
1647 }
1648 if diagonal && target == TransferFormat::Grok {
1649 append_grok_bundle_files(locator, "", ArtifactFileRole::Bundle, &mut files)?;
1650 }
1651 if !diagonal || !session.raw_is_verbatim {
1652 files.push(SessionArtifactFile {
1653 path: "recovery/source.supercode.jsonl".into(),
1654 content: session.to_native_jsonl(),
1655 role: ArtifactFileRole::SourceRecovery,
1656 });
1657 for (index, subagent) in session.subagents.iter().enumerate() {
1658 let id = subagent
1659 .meta
1660 .agent_id
1661 .as_deref()
1662 .map(sanitize_filename)
1663 .unwrap_or_else(|| format!("subagent-{}", index + 1));
1664 files.push(SessionArtifactFile {
1665 path: format!("recovery/subagents/{id}.supercode.jsonl"),
1666 content: subagent.to_native_jsonl(),
1667 role: ArtifactFileRole::SourceRecovery,
1668 });
1669 }
1670 }
1671 if !diagonal && session.meta.source == SessionSource::Grok {
1672 append_grok_bundle_files(
1673 locator,
1674 "recovery/grok/",
1675 ArtifactFileRole::SourceRecovery,
1676 &mut files,
1677 )?;
1678 }
1679 let (fidelity, residue) = if diagonal
1680 && target_session_id.is_none()
1681 && session.raw_is_verbatim
1682 && !has_appended_turns
1683 {
1684 (Fidelity::ByteLossless, Vec::new())
1685 } else if diagonal && !(target_session_id.is_some() && target == TransferFormat::OpenCode) {
1686 (
1687 Fidelity::ValueLossless,
1688 vec![if target_session_id.is_some() {
1689 "target identity was rewritten, so the artifact intentionally differs from source bytes".into()
1690 } else {
1691 "source storage was reconstructed as a native-value-equivalent export; original container bytes were not captured".into()
1692 }],
1693 )
1694 } else {
1695 (
1696 Fidelity::Semantic,
1697 vec!["target schema has no portable slot for every source-native record and metadata field".into()],
1698 )
1699 };
1700 Ok(SessionArtifact {
1701 source_harness: locator.harness.clone(),
1702 target_harness: target.id(),
1703 session_id: target_session_id
1704 .map(str::to_string)
1705 .or_else(|| session.meta.session_id.clone()),
1706 content,
1707 suggested_filename,
1708 files,
1709 fidelity,
1710 residue,
1711 })
1712}
1713
1714fn append_grok_bundle_files(
1715 locator: &SessionLocator,
1716 prefix: &str,
1717 role: ArtifactFileRole,
1718 files: &mut Vec<SessionArtifactFile>,
1719) -> std::result::Result<(), ServiceError> {
1720 let primary = locator.storage.path();
1721 if primary.file_name().and_then(|name| name.to_str()) != Some("chat_history.jsonl") {
1722 return Err(ServiceError::Operation(format!(
1723 "Grok bundle locator must name chat_history.jsonl, got {}",
1724 primary.display()
1725 )));
1726 }
1727 let parent = primary.parent().ok_or_else(|| {
1728 ServiceError::Operation("Grok chat_history.jsonl has no session directory".into())
1729 })?;
1730 for name in ["summary.json", "updates.jsonl"] {
1731 let path = parent.join(name);
1732 let metadata = match std::fs::symlink_metadata(&path) {
1733 Ok(metadata) => metadata,
1734 Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
1735 Err(error) => return Err(ServiceError::Operation(error.to_string())),
1736 };
1737 if metadata.file_type().is_symlink() || !metadata.is_file() {
1738 return Err(ServiceError::Operation(format!(
1739 "refusing non-regular Grok bundle member {}",
1740 path.display()
1741 )));
1742 }
1743 let content = std::fs::read_to_string(&path).map_err(|error| {
1744 ServiceError::Operation(format!(
1745 "Grok bundle member {} is not representable as UTF-8: {error}",
1746 path.display()
1747 ))
1748 })?;
1749 files.push(SessionArtifactFile {
1750 path: format!("{prefix}{name}"),
1751 content,
1752 role: match role {
1753 ArtifactFileRole::Bundle => ArtifactFileRole::Bundle,
1754 _ => ArtifactFileRole::SourceRecovery,
1755 },
1756 });
1757 }
1758 Ok(())
1759}
1760
1761fn handoff_artifact(
1762 locator: &SessionLocator,
1763 session: &Session,
1764 target: TransferFormat,
1765 cwd: &Path,
1766) -> std::result::Result<SessionArtifact, ServiceError> {
1767 if target != TransferFormat::Grok {
1768 let target_session_id = target_session_id(target);
1769 return session_artifact_with_id(locator, session, target, Some(&target_session_id));
1770 }
1771
1772 let mut importable = session.clone();
1776 importable.meta.session_id = Some(target_session_id(TransferFormat::ClaudeCode));
1782 importable.meta.cwd = Some(if cwd.is_absolute() {
1783 cwd.to_path_buf()
1784 } else {
1785 std::env::current_dir()
1786 .map_err(|error| ServiceError::Operation(error.to_string()))?
1787 .join(cwd)
1788 });
1789 let content = importable
1790 .to_jsonl(SessionFormat::ClaudeCode)
1791 .map_err(operation)?;
1792 let stem = sanitize_filename(
1793 importable
1794 .meta
1795 .session_id
1796 .as_deref()
1797 .unwrap_or(&locator.session_id),
1798 );
1799 let suggested_filename = format!("{stem}.grok-import.claude-code.jsonl");
1800 Ok(SessionArtifact {
1801 source_harness: locator.harness.clone(),
1802 target_harness: TransferFormat::ClaudeCode.id(),
1805 session_id: importable.meta.session_id.clone(),
1806 content: content.clone(),
1807 suggested_filename: suggested_filename.clone(),
1808 files: vec![SessionArtifactFile {
1809 path: suggested_filename,
1810 content,
1811 role: ArtifactFileRole::Primary,
1812 }],
1813 fidelity: Fidelity::Semantic,
1814 residue: vec!["Grok's stock importer accepts a Claude Code transcript, not a complete Grok updates/session bundle".into()],
1815 })
1816}
1817
1818fn target_session_id(target: TransferFormat) -> String {
1819 let uuid = generated_session_id();
1820 match target {
1821 TransferFormat::OpenCode => format!("ses_{}", uuid.replace('-', "")),
1822 TransferFormat::ClaudeCode
1823 | TransferFormat::Codex
1824 | TransferFormat::Pi
1825 | TransferFormat::Grok => uuid,
1826 }
1827}
1828
1829fn sanitize_filename(value: &str) -> String {
1830 let value = value
1831 .chars()
1832 .map(|character| {
1833 if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
1834 character
1835 } else {
1836 '-'
1837 }
1838 })
1839 .collect::<String>();
1840 let value = value.trim_matches('-');
1841 if value.is_empty() {
1842 "session".into()
1843 } else {
1844 value.chars().take(100).collect()
1845 }
1846}
1847
1848fn handoff_instructions(
1849 target: TransferFormat,
1850 session_id: &str,
1851 cwd: &Path,
1852) -> HandoffInstructions {
1853 let launch = |program: &str, arguments: Vec<String>| StructuredLaunch {
1854 cwd: cwd.to_path_buf(),
1855 program: program.into(),
1856 arguments,
1857 env: BTreeMap::new(),
1858 };
1859 match target {
1860 TransferFormat::ClaudeCode => HandoffInstructions {
1861 launch: launch("claude", vec!["--resume".into(), session_id.into()]),
1862 materialize: None,
1863 requires_materialization: true,
1864 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(),
1865 },
1866 TransferFormat::Codex => HandoffInstructions {
1867 launch: launch("codex", vec!["resume".into(), session_id.into()]),
1868 materialize: None,
1869 requires_materialization: true,
1870 note: "Write the artifact into Codex's native rollout store before running the resume launch; Codex has no general transcript-import command.".into(),
1871 },
1872 TransferFormat::OpenCode => HandoffInstructions {
1873 launch: launch("opencode", vec!["--session".into(), session_id.into()]),
1874 materialize: Some(launch(
1875 "opencode",
1876 vec!["import".into(), "{artifact_path}".into()],
1877 )),
1878 requires_materialization: true,
1879 note: "Write the artifact to a file, run the materialize command with its path, then launch the imported session.".into(),
1880 },
1881 TransferFormat::Pi => HandoffInstructions {
1882 launch: launch("pi", vec!["--session".into(), "{artifact_path}".into()]),
1883 materialize: None,
1884 requires_materialization: true,
1885 note: "Write the artifact to a file and replace {artifact_path} in the launch arguments; Pi can resume that file directly.".into(),
1886 },
1887 TransferFormat::Grok => HandoffInstructions {
1888 launch: launch(
1889 "grok",
1890 vec![
1891 "--resume".into(),
1892 "{imported_session_id}".into(),
1893 "--fork-session".into(),
1894 ],
1895 ),
1896 materialize: Some(launch(
1897 "grok",
1898 vec!["import".into(), "--json".into(), "{artifact_path}".into()],
1899 )),
1900 requires_materialization: true,
1901 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(),
1902 },
1903 }
1904}
1905
1906fn resume_launch(
1907 harness: &str,
1908 session_id: &str,
1909 cwd: &Path,
1910 policy: ResumePolicy,
1911) -> std::result::Result<StructuredLaunch, ServiceError> {
1912 let mut arguments = Vec::new();
1913 let program = match harness {
1914 HarnessId::GROK => {
1915 if matches!(policy, ResumePolicy::Yolo) {
1916 arguments.extend([
1917 "--sandbox".into(),
1918 "workspace".into(),
1919 "--always-approve".into(),
1920 ]);
1921 }
1922 arguments.extend(["--resume".into(), session_id.into()]);
1923 "grok"
1924 }
1925 HarnessId::CODEX => {
1926 if matches!(policy, ResumePolicy::Yolo) {
1927 arguments.extend([
1928 "--dangerously-bypass-approvals-and-sandbox".into(),
1929 "--dangerously-bypass-hook-trust".into(),
1930 ]);
1931 }
1932 arguments.extend(["resume".into(), session_id.into()]);
1933 "codex"
1934 }
1935 HarnessId::CLAUDE_CODE => {
1936 if matches!(policy, ResumePolicy::Yolo) {
1937 arguments.push("--dangerously-skip-permissions".into());
1938 }
1939 arguments.extend(["--resume".into(), session_id.into()]);
1940 "claude"
1941 }
1942 HarnessId::PI => {
1943 if matches!(policy, ResumePolicy::Yolo) {
1944 arguments.push("--approve".into());
1945 }
1946 arguments.extend(["--session".into(), session_id.into()]);
1947 "pi"
1948 }
1949 HarnessId::OPENCODE => {
1950 arguments.extend(["--session".into(), session_id.into()]);
1951 "opencode"
1952 }
1953 other => {
1954 return Err(ServiceError::InvalidParams(format!(
1955 "no structured resume launch is registered for harness `{other}`"
1956 )))
1957 }
1958 };
1959 Ok(StructuredLaunch {
1960 cwd: cwd.to_path_buf(),
1961 program: program.into(),
1962 arguments,
1963 env: BTreeMap::new(),
1964 })
1965}
1966
1967fn runtime_backend(
1968 params: &RuntimeBackendParams,
1969) -> std::result::Result<Box<dyn RuntimeBackend>, ServiceError> {
1970 if params.protocol.as_deref() == Some("acp") {
1971 let launch = params
1972 .launch
1973 .clone()
1974 .or_else(|| {
1975 harness_support_registry()
1976 .harnesses
1977 .into_iter()
1978 .find(|harness| harness.id == params.harness)
1979 .filter(|harness| {
1980 harness.runtime.implementation == ImplementationKind::GenericProtocol
1981 && harness.runtime.protocol.starts_with("acp")
1982 })
1983 .and_then(|harness| harness.runtime.default_launch)
1984 })
1985 .ok_or_else(|| {
1986 ServiceError::InvalidParams(
1987 "an ACP runtime requires `launch` unless the harness has a registered default"
1988 .into(),
1989 )
1990 })?;
1991 let resume_session = params.harness.as_str() == HarnessId::GROK;
1992 return Ok(Box::new(
1993 AcpRuntimeBackend::new(params.harness.clone(), launch)
1994 .with_resume_support(resume_session),
1995 ));
1996 }
1997 let backend: Box<dyn RuntimeBackend> = match params.harness.as_str() {
1998 HarnessId::CODEX => Box::new(CodexRuntimeBackend::new()),
1999 HarnessId::CLAUDE_CODE => Box::new(ClaudeCodeRuntimeBackend::new()),
2000 HarnessId::PI => Box::new(PiRuntimeBackend::new()),
2001 HarnessId::OPENCODE => match ¶ms.base_url {
2002 Some(url) => Box::new(OpenCodeRuntimeBackend::connect(url)),
2003 None => Box::new(OpenCodeRuntimeBackend::new()),
2004 },
2005 HarnessId::GROK => {
2006 let descriptor = harness_support_registry()
2007 .harnesses
2008 .into_iter()
2009 .find(|harness| harness.id.as_str() == HarnessId::GROK)
2010 .expect("Grok descriptor is part of the canonical registry");
2011 Box::new(
2012 AcpRuntimeBackend::new(
2013 descriptor.id,
2014 descriptor
2015 .runtime
2016 .default_launch
2017 .expect("Grok registry includes its ACP launch"),
2018 )
2019 .with_resume_support(true),
2020 )
2021 }
2022 other => {
2023 return Err(ServiceError::InvalidParams(format!(
2024 "no runtime adapter for harness `{other}`; use protocol `acp` with a launch command"
2025 )))
2026 }
2027 };
2028 Ok(backend)
2029}
2030
2031fn runtime_launch(params: &RuntimeBackendParams) -> Option<RuntimeLaunch> {
2032 if let Some(launch) = ¶ms.launch {
2033 return Some(launch.clone());
2034 }
2035 if !matches!(params.policy, RuntimePolicy::Yolo) {
2036 return None;
2037 }
2038 let launch = match params.harness.as_str() {
2039 HarnessId::GROK => RuntimeLaunch {
2040 program: "grok".into(),
2041 arguments: vec![
2042 "--sandbox".into(),
2043 "workspace".into(),
2044 "--always-approve".into(),
2045 "agent".into(),
2046 "--no-leader".into(),
2047 "stdio".into(),
2048 ],
2049 env: BTreeMap::from([("GROK_AGENT_DASHBOARD".into(), "0".into())]),
2050 },
2051 HarnessId::CODEX => RuntimeLaunch {
2052 program: "codex".into(),
2053 arguments: vec![
2054 "--dangerously-bypass-approvals-and-sandbox".into(),
2055 "--dangerously-bypass-hook-trust".into(),
2056 "app-server".into(),
2057 ],
2058 env: BTreeMap::new(),
2059 },
2060 HarnessId::CLAUDE_CODE => RuntimeLaunch {
2061 program: "claude".into(),
2062 arguments: vec![
2063 "--dangerously-skip-permissions".into(),
2064 "--print".into(),
2065 "--input-format".into(),
2066 "stream-json".into(),
2067 "--output-format".into(),
2068 "stream-json".into(),
2069 "--verbose".into(),
2070 ],
2071 env: BTreeMap::new(),
2072 },
2073 HarnessId::PI => RuntimeLaunch {
2074 program: "pi".into(),
2075 arguments: vec!["--approve".into(), "--mode".into(), "rpc".into()],
2076 env: BTreeMap::new(),
2077 },
2078 HarnessId::OPENCODE => RuntimeLaunch {
2079 program: "opencode".into(),
2080 arguments: vec!["serve".into()],
2081 env: BTreeMap::new(),
2082 },
2083 _ => return None,
2084 };
2085 Some(launch)
2086}
2087
2088fn find_executable(program: &str) -> Option<PathBuf> {
2089 let candidate = PathBuf::from(program);
2090 if candidate.components().count() > 1 {
2091 return candidate.is_file().then_some(candidate);
2092 }
2093 let path = std::env::var_os("PATH")?;
2094 for directory in std::env::split_paths(&path) {
2095 let candidate = directory.join(program);
2096 if candidate.is_file() {
2097 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
2098 }
2099 #[cfg(windows)]
2100 {
2101 for extension in ["exe", "cmd", "bat"] {
2102 let candidate = directory.join(format!("{program}.{extension}"));
2103 if candidate.is_file() {
2104 return std::fs::canonicalize(&candidate).ok().or(Some(candidate));
2105 }
2106 }
2107 }
2108 }
2109 None
2110}
2111
2112async fn executable_version(executable: &Path) -> Option<String> {
2113 let mut command = tokio::process::Command::new(executable);
2114 command
2115 .arg("--version")
2116 .stdin(std::process::Stdio::null())
2117 .stdout(std::process::Stdio::piped())
2118 .stderr(std::process::Stdio::piped())
2119 .kill_on_drop(true);
2120 let output = tokio::time::timeout(Duration::from_secs(3), command.output())
2121 .await
2122 .ok()?
2123 .ok()?;
2124 let stdout = String::from_utf8_lossy(&output.stdout);
2125 let stderr = String::from_utf8_lossy(&output.stderr);
2126 stdout
2127 .lines()
2128 .chain(stderr.lines())
2129 .map(str::trim)
2130 .find(|line| !line.is_empty())
2131 .map(|line| truncate_text(line, 200))
2132}
2133
2134fn auth_evidence(harness: &str) -> bool {
2135 let env_names: &[&str] = match harness {
2136 HarnessId::CLAUDE_CODE => &["ANTHROPIC_API_KEY", "CLAUDE_CODE_OAUTH_TOKEN"],
2137 HarnessId::CODEX => &["OPENAI_API_KEY"],
2138 HarnessId::OPENCODE => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
2139 HarnessId::PI => &["ANTHROPIC_API_KEY", "OPENAI_API_KEY", "OPENROUTER_API_KEY"],
2140 HarnessId::GROK => &["XAI_API_KEY", "GROK_API_KEY"],
2141 _ => &[],
2142 };
2143 if env_names
2144 .iter()
2145 .any(|name| std::env::var_os(name).is_some_and(|value| !value.is_empty()))
2146 {
2147 return true;
2148 }
2149 let Some(home) = std::env::var_os("HOME").map(PathBuf::from) else {
2150 return false;
2151 };
2152 let files: Vec<PathBuf> = match harness {
2153 HarnessId::CLAUDE_CODE => vec![home.join(".claude/.credentials.json")],
2154 HarnessId::CODEX => vec![home.join(".codex/auth.json")],
2155 HarnessId::OPENCODE => vec![
2156 home.join(".local/share/opencode/auth.json"),
2157 home.join(".config/opencode/auth.json"),
2158 ],
2159 HarnessId::PI => vec![home.join(".pi/agent/auth.json")],
2160 HarnessId::GROK => vec![home.join(".grok/auth.json")],
2161 _ => Vec::new(),
2162 };
2163 if files.into_iter().any(|path| {
2164 std::fs::metadata(path)
2165 .map(|metadata| metadata.is_file() && metadata.len() > 2)
2166 .unwrap_or(false)
2167 }) {
2168 return true;
2169 }
2170 if harness == HarnessId::CLAUDE_CODE {
2177 return std::fs::read_to_string(home.join(".claude.json"))
2178 .map(|text| text.contains("\"oauthAccount\""))
2179 .unwrap_or(false);
2180 }
2181 false
2182}
2183
2184fn looks_like_auth_error(message: &str) -> bool {
2185 let message = message.to_ascii_lowercase();
2186 [
2187 "auth",
2188 "login",
2189 "sign in",
2190 "sign-in",
2191 "credential",
2192 "unauthorized",
2193 "forbidden",
2194 "token",
2195 ]
2196 .iter()
2197 .any(|needle| message.contains(needle))
2198}
2199
2200fn unavailable_capabilities() -> crate::RuntimeCapabilities {
2201 crate::RuntimeCapabilities {
2202 start_session: false,
2203 resume_session: false,
2204 attach_existing_process: false,
2205 send_input: false,
2206 stream_events: false,
2207 interrupt: false,
2208 respond_to_requests: false,
2209 }
2210}
2211
2212fn truncate_text(text: &str, max_chars: usize) -> String {
2213 let mut chars = text.chars();
2214 let truncated = chars.by_ref().take(max_chars).collect::<String>();
2215 if chars.next().is_some() {
2216 format!("{truncated}…")
2217 } else {
2218 truncated
2219 }
2220}
2221
2222fn error_message(error: ServiceError) -> String {
2223 match error {
2224 ServiceError::InvalidParams(message)
2225 | ServiceError::Operation(message)
2226 | ServiceError::UnsupportedAction(message) => message,
2227 ServiceError::MethodNotFound => "runtime adapter is not available".into(),
2228 ServiceError::Sdk(error) => error.to_string(),
2229 }
2230}
2231
2232enum ServiceError {
2233 InvalidParams(String),
2234 MethodNotFound,
2235 UnsupportedAction(String),
2236 Operation(String),
2237 Sdk(SdkError),
2238}
2239
2240fn sdk_error(operation: SdkOperation, error: ServiceError) -> SdkError {
2241 match error {
2242 ServiceError::InvalidParams(message) => {
2243 SdkError::new(SdkErrorCode::InvalidArgument, operation, message)
2244 }
2245 ServiceError::MethodNotFound | ServiceError::UnsupportedAction(_) => {
2246 SdkError::unsupported(operation)
2247 }
2248 ServiceError::Operation(message) => {
2249 let code = if message.contains("already in progress") {
2250 SdkErrorCode::Busy
2251 } else if message.contains("not supported by this runtime") {
2252 SdkErrorCode::UnsupportedAction
2253 } else if message.contains("unknown runtime connection") {
2254 SdkErrorCode::NotFound
2255 } else {
2256 SdkErrorCode::Execution
2257 };
2258 SdkError::new(code, operation, message)
2259 }
2260 ServiceError::Sdk(error) => error,
2261 }
2262}
2263
2264fn sdk_rpc_error(id: Value, error: &SdkError) -> Value {
2265 let error_code = error.code();
2266 let code = match error_code {
2267 SdkErrorCode::Unauthenticated => -32030,
2268 SdkErrorCode::Unauthorized => -32031,
2269 SdkErrorCode::ControllerRequired => -32032,
2270 SdkErrorCode::LeaseExpired => -32033,
2271 SdkErrorCode::InvalidArgument => -32602,
2272 SdkErrorCode::NotFound => -32004,
2273 SdkErrorCode::Busy => -32000,
2274 SdkErrorCode::UnsupportedAction => -32020,
2275 SdkErrorCode::Execution => -32002,
2276 SdkErrorCode::Transport => -32003,
2277 };
2278 json!({
2279 "jsonrpc": "2.0",
2280 "id": id,
2281 "error": {
2282 "code": code,
2283 "name": error_code,
2284 "operation": error.operation(),
2285 "message": error.to_string(),
2286 },
2287 })
2288}
2289
2290fn decode<T: for<'de> Deserialize<'de>>(value: Value) -> std::result::Result<T, ServiceError> {
2291 serde_json::from_value(value).map_err(|error| ServiceError::InvalidParams(error.to_string()))
2292}
2293
2294fn operation(error: crate::Error) -> ServiceError {
2295 match error {
2296 crate::Error::Sdk(error) => ServiceError::Sdk(error),
2297 error => ServiceError::Operation(error.to_string()),
2298 }
2299}
2300
2301fn rpc_error(id: Value, code: i64, message: &str) -> Value {
2302 json!({
2303 "jsonrpc": "2.0",
2304 "id": id,
2305 "error": {"code": code, "message": message},
2306 })
2307}
2308
2309#[cfg(test)]
2310mod tests {
2311 use super::*;
2312 use crate::{HarnessEvent, HarnessId, RuntimeEndpoint, RuntimeHandle, StorageLocator};
2313 use async_trait::async_trait;
2314 use std::path::PathBuf;
2315
2316 struct EndingRuntime {
2317 handle: RuntimeHandle,
2318 event: Option<HarnessEvent>,
2319 }
2320
2321 #[async_trait]
2322 impl RuntimeConnection for EndingRuntime {
2323 fn handle(&self) -> &RuntimeHandle {
2324 &self.handle
2325 }
2326
2327 async fn send_input(&mut self, _input: RuntimeInput) -> crate::Result<Option<String>> {
2328 unreachable!("ending runtime does not accept input")
2329 }
2330
2331 async fn next_event(&mut self) -> crate::Result<Option<HarnessEvent>> {
2332 Ok(self.event.take())
2333 }
2334
2335 async fn interrupt(&mut self) -> crate::Result<()> {
2336 Ok(())
2337 }
2338
2339 async fn respond(&mut self, _request_id: Value, _response: Value) -> crate::Result<()> {
2340 Ok(())
2341 }
2342
2343 async fn close(&mut self) -> crate::Result<()> {
2344 Ok(())
2345 }
2346 }
2347
2348 fn ending_runtime(event: Option<HarnessEvent>) -> Box<dyn RuntimeConnection> {
2349 Box::new(EndingRuntime {
2350 handle: RuntimeHandle {
2351 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2352 runtime_id: "ending-session".into(),
2353 endpoint: RuntimeEndpoint::LocalProcess {
2354 pid: None,
2355 command: vec!["ending-runtime".into()],
2356 protocol: "test".into(),
2357 },
2358 },
2359 event,
2360 })
2361 }
2362
2363 fn request(id: u64, method: &str, params: Value) -> Value {
2364 json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})
2365 }
2366
2367 fn pi_locator() -> SessionLocator {
2368 SessionLocator {
2369 harness: HarnessId::from(HarnessId::PI),
2370 session_id: "1e6f2a3b-0000-4000-8000-000000000001".into(),
2371 storage: StorageLocator::File {
2372 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2373 .join("tests/fixtures/pi_session.jsonl"),
2374 },
2375 }
2376 }
2377
2378 fn opencode_locator() -> SessionLocator {
2379 let session_id = "ses_fixtureAAAAAAAAAAAAAAA1";
2380 SessionLocator {
2381 harness: HarnessId::from(HarnessId::OPENCODE),
2382 session_id: session_id.into(),
2383 storage: StorageLocator::Sqlite {
2384 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2385 .join("tests/fixtures/opencode_fixture/opencode.db"),
2386 selector: session_id.into(),
2387 },
2388 }
2389 }
2390
2391 fn grok_locator() -> SessionLocator {
2392 SessionLocator {
2393 harness: HarnessId::from(HarnessId::GROK),
2394 session_id: "73c09283-4b33-41fa-90f1-0bcb0f7be523".into(),
2395 storage: StorageLocator::File {
2396 path: PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2397 .join("tests/fixtures/grok_session/chat_history.jsonl"),
2398 },
2399 }
2400 }
2401
2402 #[test]
2403 fn capabilities_are_explicit_and_versioned() {
2404 let mut service = HarnessSessionService::new();
2405 let response = service.handle(request(1, "harness.v1.capabilities", json!({})));
2406 assert_eq!(response["result"]["version"], HARNESS_SERVICE_VERSION);
2407 assert_eq!(
2408 response["result"]["sdk"]["schema_version"],
2409 crate::SDK_SCHEMA_VERSION
2410 );
2411 assert_eq!(
2412 response["result"]["sdk"]["operations"]
2413 .as_array()
2414 .unwrap()
2415 .len(),
2416 SdkOperation::ALL.len()
2417 );
2418 assert_eq!(response["result"]["harnesses"].as_array().unwrap().len(), 5);
2419 assert!(response["result"]["harnesses"]
2420 .as_array()
2421 .unwrap()
2422 .iter()
2423 .any(|harness| harness == HarnessId::GROK));
2424 }
2425
2426 #[test]
2427 fn handshake_health_uses_protocol_liveness_not_stderr_severity() {
2428 let noisy_stderr = crate::HarnessEvent {
2429 sequence: None,
2430 kind: "transport_stderr".into(),
2431 payload: json!({"line": "ERROR optional worker AuthorizationRequired"}),
2432 };
2433 assert_eq!(handshake_event_failure(&noisy_stderr), None);
2434
2435 let closed = crate::HarnessEvent {
2436 sequence: None,
2437 kind: "transport_closed".into(),
2438 payload: json!({}),
2439 };
2440 assert!(handshake_event_failure(&closed).is_some());
2441 }
2442
2443 #[tokio::test]
2444 async fn runtime_eof_is_notified_and_removed_for_raw_and_explicit_close() {
2445 let mut service = HarnessSessionService::new();
2446 service
2447 .runtimes
2448 .insert("raw-eof".into(), ending_runtime(None));
2449 service.runtimes.insert(
2450 "explicit-close".into(),
2451 ending_runtime(Some(HarnessEvent {
2452 sequence: None,
2453 kind: "transport_closed".into(),
2454 payload: json!({"message": "native transport exited"}),
2455 })),
2456 );
2457
2458 let notifications = service.poll_runtimes().await;
2459
2460 assert_eq!(notifications.len(), 2);
2461 assert!(notifications
2462 .iter()
2463 .all(|notification| { notification["params"]["event"]["kind"] == "transport_closed" }));
2464 assert!(notifications.iter().all(|notification| {
2465 notification["params"]["session_id"] == "ending-session"
2466 && notification["params"]["connection"].is_string()
2467 }));
2468 let mut sequences = notifications
2469 .iter()
2470 .filter_map(|notification| notification["params"]["sequence"].as_u64())
2471 .collect::<Vec<_>>();
2472 sequences.sort_unstable();
2473 assert_eq!(sequences, vec![1, 2]);
2474 assert!(service.runtimes.is_empty());
2475 }
2476
2477 #[tokio::test]
2478 async fn sdk_facade_returns_named_unsupported_actions() {
2479 let mut service = HarnessSessionService::new();
2480 let error = service
2481 .execute(SdkRequest {
2482 operation: SdkOperation::Steer,
2483 params: json!({"connection": "runtime-1", "text": "go left"}),
2484 })
2485 .await
2486 .unwrap_err();
2487 assert_eq!(error.code(), SdkErrorCode::UnsupportedAction);
2488 assert_eq!(error.operation(), Some(SdkOperation::Steer));
2489
2490 let response = service
2491 .handle_async(request(
2492 7,
2493 "harness.v1.runtimes.steer",
2494 json!({"connection": "runtime-1", "text": "go left"}),
2495 ))
2496 .await;
2497 assert_eq!(response["error"]["name"], "unsupported_action");
2498 assert_eq!(response["error"]["operation"], "steer");
2499 }
2500
2501 #[test]
2502 fn support_report_and_grok_default_binding_share_the_registry() {
2503 let mut service = HarnessSessionService::new();
2504 let response = service.handle(request(1, "harness.v1.support.report", json!({})));
2505 assert_eq!(response["result"]["schema"], crate::SUPPORT_REGISTRY_SCHEMA);
2506 let params = RuntimeBackendParams {
2507 harness: HarnessId::from(HarnessId::GROK),
2508 protocol: None,
2509 launch: None,
2510 base_url: None,
2511 policy: RuntimePolicy::Default,
2512 };
2513 let backend = match runtime_backend(¶ms) {
2514 Ok(backend) => backend,
2515 Err(_) => panic!("Grok should bind through its registered ACP launch"),
2516 };
2517 assert_eq!(backend.harness().as_str(), HarnessId::GROK);
2518 assert!(backend.capabilities().start_session);
2519 let registered = harness_support_registry()
2520 .harnesses
2521 .into_iter()
2522 .find(|harness| harness.id.as_str() == HarnessId::GROK)
2523 .and_then(|harness| harness.runtime.default_launch)
2524 .unwrap();
2525 assert!(!registered
2526 .arguments
2527 .iter()
2528 .any(|argument| argument == "--always-approve"));
2529 assert!(runtime_launch(¶ms).is_none());
2530
2531 let yolo = RuntimeBackendParams {
2532 policy: RuntimePolicy::Yolo,
2533 ..params
2534 };
2535 assert!(runtime_launch(&yolo)
2536 .unwrap()
2537 .arguments
2538 .iter()
2539 .any(|argument| argument == "--always-approve"));
2540
2541 let mismatched_protocol = RuntimeBackendParams {
2542 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2543 protocol: Some("acp".into()),
2544 launch: None,
2545 base_url: None,
2546 policy: RuntimePolicy::Default,
2547 };
2548 assert!(runtime_backend(&mismatched_protocol).is_err());
2549 }
2550
2551 #[test]
2552 fn load_follow_and_unfollow_share_the_same_locator() {
2553 let mut service = HarnessSessionService::new();
2554 let locator = pi_locator();
2555 let loaded = service.handle(request(
2556 1,
2557 "harness.v1.sessions.load",
2558 json!({"locator": locator}),
2559 ));
2560 assert_eq!(
2561 loaded["result"]["session"]["session_id"],
2562 locator.session_id
2563 );
2564
2565 let followed = service.handle(request(
2566 2,
2567 "harness.v1.sessions.follow",
2568 json!({"locator": locator}),
2569 ));
2570 assert_eq!(followed["result"]["subscription"], "sub-1");
2571 assert_eq!(followed["result"]["initial"]["type"], "session_snapshot");
2572 assert!(service.poll().is_empty());
2573
2574 let unfollowed = service.handle(request(
2575 3,
2576 "harness.v1.sessions.unfollow",
2577 json!({"subscription": "sub-1"}),
2578 ));
2579 assert_eq!(unfollowed["result"]["removed"], true);
2580 }
2581
2582 #[test]
2583 fn import_translate_branch_and_handoff_use_typed_artifacts() {
2584 let mut service = HarnessSessionService::new();
2585 let locator = pi_locator();
2586 let translated = service.handle(request(
2587 1,
2588 "harness.v1.sessions.translate",
2589 json!({"locator": locator, "target_harness": "grok"}),
2590 ));
2591 assert_eq!(translated["result"]["artifact"]["source_harness"], "pi");
2592 assert_eq!(translated["result"]["artifact"]["target_harness"], "grok");
2593 assert!(translated["result"]["artifact"]["content"]
2594 .as_str()
2595 .is_some_and(|content| !content.is_empty()));
2596
2597 for target in ["opencode", "open-code"] {
2598 let opencode = service.handle(request(
2599 6,
2600 "harness.v1.sessions.translate",
2601 json!({"locator": locator, "target_harness": target}),
2602 ));
2603 assert_eq!(opencode["result"]["artifact"]["target_harness"], "opencode");
2604 }
2605
2606 let imported = service.handle(request(
2607 2,
2608 "harness.v1.sessions.import",
2609 json!({
2610 "source_harness": "grok",
2611 "content": translated["result"]["artifact"]["content"],
2612 }),
2613 ));
2614 assert_eq!(imported["result"]["session"]["source"], "grok");
2615
2616 let branched = service.handle(request(
2617 3,
2618 "harness.v1.sessions.branch",
2619 json!({"locator": locator, "target_harness": "codex"}),
2620 ));
2621 assert_eq!(branched["result"]["parent"]["harness"], "pi");
2622 assert!(branched["result"]["bootstrap_prompt"]
2623 .as_str()
2624 .unwrap()
2625 .contains("frozen parent transcript"));
2626 assert_eq!(branched["result"]["artifact"]["target_harness"], "codex");
2627
2628 let handoff = service.handle(request(
2629 4,
2630 "harness.v1.sessions.handoff",
2631 json!({"locator": locator, "target_harness": "pi", "cwd": "/tmp/project"}),
2632 ));
2633 assert_eq!(handoff["result"]["launch"]["program"], "pi");
2634 assert_eq!(handoff["result"]["launch"]["cwd"], "/tmp/project");
2635 assert_eq!(handoff["result"]["requires_materialization"], true);
2636
2637 let resumed = service.handle(request(
2638 5,
2639 "harness.v1.sessions.resume_instructions",
2640 json!({"locator": locator, "cwd": "/tmp/project", "policy": "yolo"}),
2641 ));
2642 assert_eq!(resumed["result"]["launch"]["program"], "pi");
2643 assert_eq!(resumed["result"]["launch"]["arguments"][0], "--approve");
2644 }
2645
2646 #[test]
2647 fn read_surfaces_view_a_severed_claude_graph_while_transfer_still_refuses_it() {
2648 let temp = std::env::temp_dir().join(format!(
2649 "supercode-severed-view-{}-{}",
2650 std::process::id(),
2651 generated_session_id()
2652 ));
2653 std::fs::create_dir_all(&temp).unwrap();
2654 let path = temp.join("severed.jsonl");
2655 std::fs::write(
2658 &path,
2659 concat!(
2660 r#"{"type":"user","uuid":"orphan-u","parentUuid":null,"message":{"role":"user","content":"stranded prompt"}}"#,
2661 "\n",
2662 r#"{"type":"assistant","uuid":"live-a","parentUuid":"pruned","message":{"id":"m","role":"assistant","content":[{"type":"text","text":"live answer"}]}}"#,
2663 "\n",
2664 ),
2665 )
2666 .unwrap();
2667 let locator = SessionLocator {
2668 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2669 session_id: "severed".into(),
2670 storage: StorageLocator::File { path },
2671 };
2672 let mut service = HarnessSessionService::new();
2673
2674 let viewed = service.handle(request(
2675 1,
2676 "harness.v1.sessions.load",
2677 json!({"locator": locator}),
2678 ));
2679 let session = &viewed["result"]["session"];
2680 assert_eq!(session["fidelity"], "semantic");
2681 assert_eq!(session["messages"].as_array().unwrap().len(), 2);
2682 assert!(session["residue"].as_array().unwrap().iter().any(|entry| {
2683 entry
2684 .as_str()
2685 .is_some_and(|entry| entry.contains("live-a") && entry.contains("pruned"))
2686 }));
2687
2688 let strict = service.handle(request(
2691 2,
2692 "harness.v1.sessions.load",
2693 json!({"locator": locator, "fidelity": "byte_lossless"}),
2694 ));
2695 assert!(strict["error"]["message"]
2696 .as_str()
2697 .unwrap()
2698 .contains("cannot reconstruct lossless Claude continuation"));
2699
2700 let translated = service.handle(request(
2702 3,
2703 "harness.v1.sessions.translate",
2704 json!({"locator": locator, "target_harness": "codex"}),
2705 ));
2706 assert!(translated["error"]["message"]
2707 .as_str()
2708 .unwrap()
2709 .contains("cannot reconstruct lossless Claude continuation"));
2710 let resumed = service.handle(request(
2711 4,
2712 "harness.v1.sessions.resume_instructions",
2713 json!({"locator": locator}),
2714 ));
2715 assert!(resumed["error"]["message"]
2716 .as_str()
2717 .unwrap()
2718 .contains("cannot reconstruct lossless Claude continuation"));
2719
2720 let _ = std::fs::remove_dir_all(&temp);
2721 }
2722
2723 #[test]
2724 fn diagonal_artifacts_preserve_claude_subagents_and_grok_bundle_members() {
2725 let temp = std::env::temp_dir().join(format!(
2726 "supercode-harness-artifact-{}-{}",
2727 std::process::id(),
2728 generated_session_id()
2729 ));
2730 let main_path = temp.join("parent.jsonl");
2731 let subagent_path = temp.join("parent/subagents/agent-child.jsonl");
2732 std::fs::create_dir_all(subagent_path.parent().unwrap()).unwrap();
2733 let fixture = std::fs::read_to_string(
2734 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2735 .join("tests/fixtures/claude_code_session.jsonl"),
2736 )
2737 .unwrap();
2738 let parent = fixture.trim_end_matches('\n');
2739 let child = fixture.trim_end_matches('\n');
2740 std::fs::write(&main_path, parent).unwrap();
2741 std::fs::write(&subagent_path, child).unwrap();
2742 let locator = SessionLocator {
2743 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
2744 session_id: "213bb148-51ea-453f-9206-f8b4b1168547".into(),
2745 storage: StorageLocator::File {
2746 path: main_path.clone(),
2747 },
2748 };
2749 let mut service = HarnessSessionService::new();
2750 let claude = service.handle(request(
2751 1,
2752 "harness.v1.sessions.translate",
2753 json!({"locator": locator, "target_harness": "claude-code"}),
2754 ));
2755 let artifact = &claude["result"]["artifact"];
2756 assert_eq!(artifact["fidelity"], "byte_lossless");
2757 assert_eq!(artifact["content"], parent);
2758 let files = artifact["files"].as_array().unwrap();
2759 assert!(files.iter().any(|file| {
2760 file["role"] == "subagent"
2761 && file["path"]
2762 .as_str()
2763 .is_some_and(|path| path.ends_with("/subagents/agent-child.jsonl"))
2764 && file["content"] == child
2765 }));
2766 assert!(!artifact["content"].as_str().unwrap().ends_with('\n'));
2767
2768 let grok = service.handle(request(
2769 2,
2770 "harness.v1.sessions.translate",
2771 json!({"locator": grok_locator(), "target_harness": "grok"}),
2772 ));
2773 let files = grok["result"]["artifact"]["files"].as_array().unwrap();
2774 for name in ["summary.json", "updates.jsonl"] {
2775 let expected = std::fs::read_to_string(
2776 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2777 .join("tests/fixtures/grok_session")
2778 .join(name),
2779 )
2780 .unwrap();
2781 assert!(files.iter().any(|file| {
2782 file["path"] == name && file["role"] == "bundle" && file["content"] == expected
2783 }));
2784 }
2785 std::fs::remove_dir_all(temp).ok();
2786 }
2787
2788 #[test]
2789 fn every_non_grok_handoff_mints_and_uses_a_fresh_target_identity() {
2790 let mut service = HarnessSessionService::new();
2791 let source = pi_locator();
2792 for (target, format) in [
2793 ("claude-code", SessionFormat::ClaudeCode),
2794 ("codex", SessionFormat::Codex),
2795 ("opencode", SessionFormat::OpenCode),
2796 ("pi", SessionFormat::Pi),
2797 ] {
2798 let result = service.handle(request(
2799 1,
2800 "harness.v1.sessions.handoff",
2801 json!({"locator": source, "target_harness": target, "cwd": "/tmp/project"}),
2802 ));
2803 let artifact = &result["result"]["artifact"];
2804 let target_id = artifact["session_id"].as_str().unwrap();
2805 assert_ne!(target_id, source.session_id, "{target}");
2806 let parsed = Session::load_str(artifact["content"].as_str().unwrap(), format).unwrap();
2807 assert_eq!(
2808 parsed.meta.session_id.as_deref(),
2809 Some(target_id),
2810 "{target}"
2811 );
2812 if target != "pi" {
2813 assert!(result["result"]["launch"]["arguments"]
2814 .as_array()
2815 .unwrap()
2816 .iter()
2817 .any(|argument| argument == target_id));
2818 }
2819 if target == "opencode" {
2820 assert!(target_id.starts_with("ses_"));
2821 fn assert_session_ids(value: &Value, target_id: &str) {
2822 match value {
2823 Value::Object(fields) => {
2824 if let Some(session_id) = fields.get("sessionID") {
2825 assert_eq!(session_id, target_id);
2826 }
2827 for child in fields.values() {
2828 assert_session_ids(child, target_id);
2829 }
2830 }
2831 Value::Array(values) => {
2832 for child in values {
2833 assert_session_ids(child, target_id);
2834 }
2835 }
2836 _ => {}
2837 }
2838 }
2839 let document: Value =
2840 serde_json::from_str(artifact["content"].as_str().unwrap()).unwrap();
2841 assert_session_ids(&document, target_id);
2842 }
2843 }
2844
2845 let first = service.handle(request(
2846 2,
2847 "harness.v1.sessions.handoff",
2848 json!({"locator": source, "target_harness": "codex"}),
2849 ));
2850 let second = service.handle(request(
2851 3,
2852 "harness.v1.sessions.handoff",
2853 json!({"locator": source, "target_harness": "codex"}),
2854 ));
2855 assert_ne!(
2856 first["result"]["artifact"]["session_id"],
2857 second["result"]["artifact"]["session_id"]
2858 );
2859 }
2860
2861 #[test]
2862 fn grok_handoff_uses_the_official_importer_contract() {
2863 let mut service = HarnessSessionService::new();
2864 let source = opencode_locator();
2865 let response = service.handle(request(
2866 1,
2867 "harness.v1.sessions.handoff",
2868 json!({
2869 "locator": source,
2870 "target_harness": "grok",
2871 "cwd": "/tmp/grok-handoff-project",
2872 }),
2873 ));
2874 let result = &response["result"];
2875
2876 assert_eq!(result["artifact"]["target_harness"], "claude-code");
2880 assert!(result["artifact"]["suggested_filename"]
2881 .as_str()
2882 .unwrap()
2883 .ends_with(".grok-import.claude-code.jsonl"));
2884 let artifact = Session::load_str(
2885 result["artifact"]["content"].as_str().unwrap(),
2886 SessionFormat::ClaudeCode,
2887 )
2888 .unwrap();
2889 assert_eq!(
2890 artifact.meta.cwd.as_deref(),
2891 Some(Path::new("/tmp/grok-handoff-project"))
2892 );
2893 let target_session_id = artifact.meta.session_id.as_deref().unwrap();
2894 assert_eq!(target_session_id.len(), 36);
2895 assert_eq!(target_session_id.as_bytes()[14], b'4');
2896 assert_ne!(target_session_id, opencode_locator().session_id);
2897 assert_eq!(
2898 result["artifact"]["session_id"],
2899 artifact.meta.session_id.as_deref().unwrap()
2900 );
2901
2902 assert_eq!(
2903 result["materialize"]["arguments"],
2904 json!(["import", "--json", "{artifact_path}"])
2905 );
2906 assert_eq!(
2907 result["launch"]["arguments"],
2908 json!(["--resume", "{imported_session_id}", "--fork-session"])
2909 );
2910 assert!(result["note"]
2911 .as_str()
2912 .unwrap()
2913 .contains("outcome=imported"));
2914 assert!(!result["launch"]["arguments"]
2915 .as_array()
2916 .unwrap()
2917 .iter()
2918 .any(|argument| argument == &opencode_locator().session_id));
2919 }
2920
2921 #[tokio::test]
2922 async fn inventory_rejects_unknown_harnesses_and_runtime_attach_is_honest() {
2923 let mut service = HarnessSessionService::new();
2924 let inventory = service
2925 .handle_async(request(
2926 1,
2927 "harness.v1.harnesses.list",
2928 json!({"harnesses": ["missing"]}),
2929 ))
2930 .await;
2931 assert_eq!(inventory["error"]["code"], -32602);
2932
2933 let attached = service
2934 .handle_async(request(
2935 2,
2936 "harness.v1.runtimes.attach_existing",
2937 json!({"harness": "codex", "runtime_id": "thread-1"}),
2938 ))
2939 .await;
2940 assert_eq!(attached["error"]["code"], -32000);
2941 assert!(attached["error"]["message"]
2942 .as_str()
2943 .unwrap()
2944 .contains("runtimes.resume"));
2945 }
2946
2947 #[test]
2948 fn invalid_params_and_unknown_methods_use_json_rpc_errors() {
2949 let mut service = HarnessSessionService::new();
2950 let invalid = service.handle(request(1, "harness.v1.sessions.load", json!({})));
2951 assert_eq!(invalid["error"]["code"], -32602);
2952 let unknown = service.handle(request(2, "harness.v1.unknown", json!({})));
2953 assert_eq!(unknown["error"]["code"], -32601);
2954 }
2955
2956 #[cfg(unix)]
2957 #[tokio::test]
2958 #[allow(clippy::await_holding_lock)]
2961 async fn async_service_drives_a_generic_acp_runtime() {
2962 let _environment_guard = crate::live_runtime::test_environment_lock();
2963 let script = r#"
2964 i=0
2965 while IFS= read -r line; do
2966 i=$((i + 1))
2967 case "$i" in
2968 1) printf '%s\n' '{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":1,"agentCapabilities":{},"authMethods":[]}}' ;;
2969 2) printf '%s\n' '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"svc_acp"}}' ;;
2970 3)
2971 printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"ok"}}}}'
2972 printf '%s\n' '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
2973 ;;
2974 4)
2975 printf '%s\n' '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"svc_acp","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"from terminal"}}}}'
2976 printf '%s\n' '{"jsonrpc":"2.0","id":4,"result":{"stopReason":"end_turn"}}'
2977 ;;
2978 esac
2979 done
2980 "#;
2981 let mut service = HarnessSessionService::new();
2982 let started = service
2983 .handle_async(request(
2984 1,
2985 "harness.v1.runtimes.start",
2986 json!({
2987 "harness": "codex",
2988 "protocol": "acp",
2989 "cwd": std::env::current_dir().unwrap(),
2990 "launch": {"program": "/bin/sh", "arguments": ["-c", script], "env": {}},
2991 }),
2992 ))
2993 .await;
2994 assert_eq!(started["result"]["connection"], "runtime-1");
2995 assert_eq!(started["result"]["handle"]["runtime_id"], "svc_acp");
2996
2997 let terminal = service
2998 .handle_async(request(
2999 9,
3000 "harness.v1.runtimes.terminal_instructions",
3001 json!({"connection":"runtime-1"}),
3002 ))
3003 .await;
3004 let arguments = terminal["result"]["launch"]["arguments"]
3005 .as_array()
3006 .expect("hosted runtime should return terminal arguments");
3007 let endpoint_index = arguments
3008 .iter()
3009 .position(|value| value == "--endpoint")
3010 .expect("terminal command should use an opaque endpoint");
3011 let endpoint = LiveRuntimeEndpoint::parse(
3012 arguments[endpoint_index + 1]
3013 .as_str()
3014 .expect("endpoint argument should be text"),
3015 )
3016 .unwrap();
3017 assert!(!terminal.to_string().contains("Bearer"));
3018 let workspace = std::env::current_dir().unwrap();
3019 let receipt = resolve_live_runtime(
3020 &endpoint,
3021 &LiveRuntimeSource {
3022 harness: "codex".into(),
3023 session_id: "svc_acp".into(),
3024 workspace,
3025 },
3026 )
3027 .unwrap();
3028 let remote = crate::HttpFrontendRuntime::connect(receipt.base_url, receipt.token)
3029 .await
3030 .unwrap();
3031 let mut attachment = crate::FrontendRuntime::attach(remote.as_ref(), 100)
3032 .await
3033 .unwrap();
3034
3035 let sent = service
3036 .handle_async(request(
3037 2,
3038 "harness.v1.runtimes.send_input",
3039 json!({"connection": "runtime-1", "text": "hi"}),
3040 ))
3041 .await;
3042 assert_eq!(sent["result"]["turn_id"], "3");
3043
3044 let mut events = Vec::new();
3045 for _ in 0..20 {
3046 events.extend(service.poll_runtimes().await);
3047 if events.len() >= 2 {
3048 break;
3049 }
3050 tokio::time::sleep(Duration::from_millis(2)).await;
3051 }
3052 assert!(events
3053 .iter()
3054 .any(|event| { event["params"]["event"]["kind"] == "session/update" }));
3055 assert!(events.iter().any(|event| {
3056 event["params"]["event"]["kind"] == "supercode/acp_request_completed"
3057 }));
3058
3059 let saw_editor_reply = tokio::time::timeout(Duration::from_secs(2), async {
3060 loop {
3061 let event = attachment.next_event().await.unwrap();
3062 if event.kind == "text_delta" && event.payload["text"] == "ok" {
3063 break;
3064 }
3065 }
3066 })
3067 .await;
3068 assert!(
3069 saw_editor_reply.is_ok(),
3070 "terminal should observe the editor-driven turn"
3071 );
3072
3073 crate::FrontendRuntime::submit(remote.as_ref(), "DRIVE FROM TERMINAL".into())
3074 .await
3075 .unwrap();
3076 let saw_terminal_reply = tokio::time::timeout(Duration::from_secs(2), async {
3077 loop {
3078 let event = attachment.next_event().await.unwrap();
3079 if event.kind == "text_delta" && event.payload["text"] == "from terminal" {
3080 break;
3081 }
3082 }
3083 })
3084 .await;
3085 assert!(
3086 saw_terminal_reply.is_ok(),
3087 "terminal should drive the same runtime"
3088 );
3089
3090 let closed = service
3091 .handle_async(request(
3092 3,
3093 "harness.v1.runtimes.close",
3094 json!({"connection": "runtime-1"}),
3095 ))
3096 .await;
3097 assert_eq!(closed["result"]["closed"], true);
3098 }
3099}