1use std::collections::{BTreeMap, VecDeque};
4use std::net::TcpListener;
5use std::path::{Path, PathBuf};
6use std::process::Stdio;
7use std::sync::Arc;
8use std::time::{Duration, SystemTime, UNIX_EPOCH};
9
10use async_trait::async_trait;
11use futures::StreamExt;
12use serde_json::{json, Value};
13use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
14use tokio::process::{Child, ChildStdin, Command};
15use tokio::sync::{mpsc, Mutex};
16
17use super::{
18 BearerToken, HarnessEvent, JsonLineClient, McpServerLaunch, RuntimeAttachRequest,
19 RuntimeBackend, RuntimeCapabilities, RuntimeConnection, RuntimeEndpoint, RuntimeHandle,
20 RuntimeInput, RuntimeLaunch, RuntimeStartRequest,
21};
22use crate::{Error, HarnessId, Result};
23
24#[derive(Debug, Clone)]
26pub struct PiRuntimeBackend {
27 launch: RuntimeLaunch,
28}
29
30impl Default for PiRuntimeBackend {
31 fn default() -> Self {
32 Self::new()
33 }
34}
35
36impl PiRuntimeBackend {
37 pub fn new() -> Self {
39 Self {
40 launch: RuntimeLaunch {
41 program: "pi".into(),
42 arguments: vec!["--mode".into(), "rpc".into()],
43 env: BTreeMap::new(),
44 },
45 }
46 }
47
48 pub fn with_launch(launch: RuntimeLaunch) -> Self {
50 Self { launch }
51 }
52
53 async fn open(
54 &self,
55 cwd: &Path,
56 runtime_id: String,
57 launch: Option<RuntimeLaunch>,
58 mcp_servers: &[McpServerLaunch],
59 resume: bool,
60 ) -> Result<Box<dyn RuntimeConnection>> {
61 let mut launch = launch.unwrap_or_else(|| self.launch.clone());
62 if !mcp_servers.is_empty() {
67 let file = mcp_config_file("pi", &runtime_id, mcp_servers).await?;
68 launch.arguments.extend([
69 "--extension".into(),
70 std::env::var(PI_MCP_ADAPTER_ENV)
71 .ok()
72 .filter(|value| !value.is_empty())
73 .unwrap_or_else(|| PI_MCP_ADAPTER.into()),
74 "--mcp-config".into(),
75 file.to_string_lossy().into_owned(),
76 ]);
77 }
78 if resume {
79 launch
80 .arguments
81 .extend(["--session".into(), runtime_id.clone()]);
82 } else {
83 launch
84 .arguments
85 .extend(["--session-id".into(), runtime_id.clone()]);
86 }
87 let transport = RawLineTransport::spawn(&launch, Some(cwd), "pi-rpc-jsonl").await?;
88 let handle = RuntimeHandle {
89 harness: HarnessId::from(HarnessId::PI),
90 runtime_id,
91 endpoint: transport.endpoint.clone(),
92 };
93 Ok(Box::new(PiRuntimeConnection {
94 handle,
95 transport,
96 next_request: 1,
97 }))
98 }
99}
100
101#[async_trait]
102impl RuntimeBackend for PiRuntimeBackend {
103 fn harness(&self) -> HarnessId {
104 HarnessId::from(HarnessId::PI)
105 }
106
107 fn capabilities(&self) -> RuntimeCapabilities {
108 RuntimeCapabilities {
109 start_session: true,
110 resume_session: true,
111 attach_existing_process: false,
112 send_input: true,
113 stream_events: true,
114 interrupt: true,
115 steer: false,
116 respond_to_requests: true,
117 }
118 }
119
120 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
121 self.open(
122 &request.cwd,
123 generated_session_id(),
124 request.launch,
125 &request.mcp_servers,
126 false,
127 )
128 .await
129 }
130
131 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
132 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
133 self.open(
134 &cwd,
135 request.runtime_id,
136 request.launch,
137 &request.mcp_servers,
138 true,
139 )
140 .await
141 }
142}
143
144const PI_MCP_ADAPTER: &str = "npm:pi-mcp-adapter@2.37.0";
147const PI_MCP_ADAPTER_ENV: &str = "SUPERCODE_PI_MCP_ADAPTER";
151
152struct PiRuntimeConnection {
153 handle: RuntimeHandle,
154 transport: RawLineTransport,
155 next_request: u64,
156}
157
158#[async_trait]
159impl RuntimeConnection for PiRuntimeConnection {
160 fn handle(&self) -> &RuntimeHandle {
161 &self.handle
162 }
163
164 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
165 if !input.image_urls.is_empty() {
166 return Err(Error::Other(
167 "Pi RPC image input is not verified by the installed protocol contract".into(),
168 ));
169 }
170 let id = format!("supercode-{}", self.next_request);
171 self.next_request += 1;
172 self.transport
173 .write(json!({"id": id, "type": "prompt", "message": input.text}))
174 .await?;
175 Ok(Some(id))
176 }
177
178 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
179 raw_next_event(&mut self.transport.receiver).await
180 }
181
182 async fn interrupt(&mut self) -> Result<()> {
183 self.transport.write(json!({"type": "abort"})).await
184 }
185
186 async fn respond(&mut self, request_id: Value, mut response: Value) -> Result<()> {
187 if let Value::Object(object) = &mut response {
188 object.entry("id").or_insert(request_id);
189 self.transport.write(response).await
190 } else {
191 self.transport
192 .write(json!({"id": request_id, "response": response}))
193 .await
194 }
195 }
196
197 async fn close(&mut self) -> Result<()> {
198 self.transport.close().await
199 }
200}
201
202#[derive(Debug, Clone)]
209pub struct ClaudeCodeRuntimeBackend {
210 launch: RuntimeLaunch,
211 permission_timeout: Duration,
212}
213
214impl Default for ClaudeCodeRuntimeBackend {
215 fn default() -> Self {
216 Self::new()
217 }
218}
219
220impl ClaudeCodeRuntimeBackend {
221 pub fn new() -> Self {
224 Self {
225 launch: RuntimeLaunch {
226 program: "claude".into(),
227 arguments: vec![
228 "--print".into(),
229 "--input-format".into(),
230 "stream-json".into(),
231 "--output-format".into(),
232 "stream-json".into(),
233 "--verbose".into(),
234 "--permission-prompt-tool".into(),
243 "stdio".into(),
244 ],
245 env: BTreeMap::new(),
246 },
247 permission_timeout: CLAUDE_PERMISSION_RESPONSE_TIMEOUT,
248 }
249 }
250
251 pub fn launch(&self) -> &RuntimeLaunch {
256 &self.launch
257 }
258
259 pub fn with_launch(launch: RuntimeLaunch) -> Self {
261 Self {
262 launch,
263 permission_timeout: CLAUDE_PERMISSION_RESPONSE_TIMEOUT,
264 }
265 }
266
267 pub fn with_permission_timeout(mut self, timeout: Duration) -> Self {
273 self.permission_timeout = timeout;
274 self
275 }
276
277 async fn open(
278 &self,
279 cwd: &Path,
280 runtime_id: String,
281 launch: Option<RuntimeLaunch>,
282 mcp_servers: &[McpServerLaunch],
283 resume: bool,
284 ) -> Result<Box<dyn RuntimeConnection>> {
285 let mut prefix = launch.unwrap_or_else(|| self.launch.clone());
289 if !mcp_servers.is_empty() {
293 let file = mcp_config_file("claude", &runtime_id, mcp_servers).await?;
294 prefix
295 .arguments
296 .extend(["--mcp-config".into(), file.to_string_lossy().into_owned()]);
297 }
298 let mut launch = prefix.clone();
299 launch.arguments.extend(if resume {
300 vec!["--resume".into(), runtime_id.clone()]
301 } else {
302 vec!["--session-id".into(), runtime_id.clone()]
303 });
304 let transport = RawLineTransport::spawn(&launch, Some(cwd), "claude-stream-json").await?;
305 Ok(Box::new(ClaudeRuntimeConnection {
306 handle: RuntimeHandle {
307 harness: HarnessId::from(HarnessId::CLAUDE_CODE),
308 runtime_id,
309 endpoint: transport.endpoint.clone(),
310 },
311 transport,
312 prefix,
313 cwd: cwd.to_path_buf(),
314 spoke: false,
315 buffered_events: VecDeque::new(),
316 next_control_request: 1,
317 control_timeout: CLAUDE_CONTROL_RESPONSE_TIMEOUT,
318 pending_permissions: Vec::new(),
319 permission_timeout: self.permission_timeout,
320 mounted_mcp_servers: mcp_servers
321 .iter()
322 .map(|server| server.name.clone())
323 .collect(),
324 }))
325 }
326}
327
328const CLAUDE_CONTROL_RESPONSE_TIMEOUT: Duration = Duration::from_secs(10);
338
339pub const CLAUDE_PERMISSION_RESPONSE_TIMEOUT: Duration = Duration::from_secs(300);
352
353const CLAUDE_PERMISSION_BEHAVIORS: [&str; 2] = ["allow", "deny"];
360
361const CLAUDE_PERMISSION_TIMEOUT_MESSAGE: &str =
363 "Volter Harness denied this permission request: no answer arrived before the adapter's \
364 permission timeout elapsed";
365
366#[async_trait]
367impl RuntimeBackend for ClaudeCodeRuntimeBackend {
368 fn harness(&self) -> HarnessId {
369 HarnessId::from(HarnessId::CLAUDE_CODE)
370 }
371
372 fn capabilities(&self) -> RuntimeCapabilities {
373 RuntimeCapabilities {
374 start_session: true,
375 resume_session: true,
376 attach_existing_process: false,
377 send_input: true,
378 stream_events: true,
379 interrupt: true,
380 steer: true,
381 respond_to_requests: true,
382 }
383 }
384
385 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
386 let session_id = generated_session_id();
387 let env = request
390 .launch
391 .as_ref()
392 .map(|launch| launch.env.clone())
393 .unwrap_or_default();
394 let home = env.get("CLAUDE_CONFIG_DIR").map(PathBuf::from).or_else(|| {
395 crate::HarnessHomes::default()
396 .claude_code
397 .parent()
398 .map(Path::to_path_buf)
399 });
400 if let Some(home) = home {
401 if let Err(error) = crate::launch_agent::record(&home, &session_id, &env) {
402 tracing::warn!(
403 "could not record the launch agent of Claude Code session {session_id}: {error}"
404 );
405 }
406 }
407 let launch = match &request.model {
410 Some(model) => {
411 let mut launch = request.launch.unwrap_or_else(|| self.launch.clone());
412 launch.arguments.extend(["--model".into(), model.clone()]);
413 Some(launch)
414 }
415 None => request.launch,
416 };
417 self.open(
418 &request.cwd,
419 session_id,
420 launch,
421 &request.mcp_servers,
422 false,
423 )
424 .await
425 }
426
427 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
428 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
429 self.open(
430 &cwd,
431 request.runtime_id,
432 request.launch,
433 &request.mcp_servers,
434 true,
435 )
436 .await
437 }
438}
439
440async fn mcp_config_file(
444 harness: &str,
445 runtime_id: &str,
446 servers: &[McpServerLaunch],
447) -> Result<PathBuf> {
448 let mut entries = serde_json::Map::new();
449 for server in servers {
450 entries.insert(
451 server.name.clone(),
452 json!({
453 "type": "stdio",
454 "command": server.command,
455 "args": server.arguments,
456 "env": server.env,
457 }),
458 );
459 }
460 let safe: String = runtime_id
461 .chars()
462 .filter(|c| c.is_ascii_alphanumeric() || *c == '-' || *c == '_')
463 .collect();
464 let path = std::env::temp_dir().join(format!("supercode-{harness}-mcp-{safe}.json"));
465 let body = serde_json::to_vec_pretty(&json!({ "mcpServers": entries })).map_err(|error| {
466 Error::Other(format!(
467 "{harness} mcp config could not be encoded: {error}"
468 ))
469 })?;
470 tokio::fs::write(&path, body).await.map_err(|error| {
471 Error::Other(format!(
472 "{harness} mcp config could not be written: {error}"
473 ))
474 })?;
475 #[cfg(unix)]
476 {
477 use std::os::unix::fs::PermissionsExt;
478 tokio::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600))
479 .await
480 .map_err(|error| {
481 Error::Other(format!(
482 "{harness} mcp config could not be protected: {error}"
483 ))
484 })?;
485 }
486 Ok(path)
487}
488
489struct ClaudeRuntimeConnection {
490 handle: RuntimeHandle,
491 transport: RawLineTransport,
492 prefix: RuntimeLaunch,
496 cwd: PathBuf,
500 spoke: bool,
505 buffered_events: VecDeque<Value>,
509 next_control_request: u64,
510 control_timeout: Duration,
511 pending_permissions: Vec<PendingPermission>,
515 permission_timeout: Duration,
516 mounted_mcp_servers: Vec<String>,
520}
521
522struct PendingPermission {
524 request_id: String,
526 deadline: tokio::time::Instant,
528}
529
530impl ClaudeRuntimeConnection {
531 fn is_relay(&self) -> bool {
532 self.prefix
533 .env
534 .get("SUPERCODE_CLAUDE_RELAY")
535 .is_some_and(|value| value == "1")
536 }
537
538 async fn process_ended(&self) -> bool {
545 matches!(self.transport.child.lock().await.try_wait(), Ok(Some(_)))
546 }
547
548 async fn reopen(&mut self) -> Result<()> {
558 let mut launch = self.prefix.clone();
559 launch
560 .arguments
561 .extend(["--resume".into(), self.handle.runtime_id.clone()]);
562 let transport = RawLineTransport::spawn(&launch, Some(&self.cwd), "claude-stream-json")
563 .await
564 .map_err(|error| {
565 Error::Other(format!(
566 "could not resume Claude Code session `{}` after its process exited: {error}",
567 self.handle.runtime_id
568 ))
569 })?;
570 self.handle.endpoint = transport.endpoint.clone();
571 self.transport = transport;
572 self.pending_permissions.clear();
575 self.spoke = false;
576 Ok(())
577 }
578
579 async fn write_turn(&mut self, frame: Value) -> Result<()> {
587 if self.process_ended().await {
588 if self.is_relay() {
589 return Err(Error::Other("Claude relay process exited".into()));
590 }
591 self.reopen().await?;
592 }
593 match self.transport.write(frame.clone()).await {
594 Ok(()) => Ok(()),
595 Err(error) if broken_pipe(&error) && !self.is_relay() => {
596 self.reopen().await?;
597 self.transport.write(frame).await
598 }
599 Err(error) => Err(error),
600 }
601 }
602
603 fn is_control_response(value: &Value) -> bool {
608 value.get("type").and_then(Value::as_str) == Some("control_response")
609 }
610
611 fn control_result(value: &Value, request_id: &str) -> Option<Result<()>> {
619 let response = value.get("response")?;
620 if response.get("request_id").and_then(Value::as_str) != Some(request_id) {
621 return None;
622 }
623 match response.get("subtype").and_then(Value::as_str) {
624 Some("success") => Some(Ok(())),
625 other => Some(Err(Error::Other(format!(
626 "Claude Code rejected the interrupt control request: {}",
627 response
628 .get("error")
629 .and_then(Value::as_str)
630 .map(str::to_string)
631 .unwrap_or_else(|| format!(
632 "control_response subtype {}",
633 other.unwrap_or("(missing)")
634 ))
635 )))),
636 }
637 }
638
639 fn permission_request_id(value: &Value) -> Option<&str> {
646 if value.get("type").and_then(Value::as_str)? != "control_request" {
647 return None;
648 }
649 let request = value.get("request")?;
650 if request.get("subtype").and_then(Value::as_str)? != "can_use_tool" {
651 return None;
652 }
653 value.get("request_id").and_then(Value::as_str)
654 }
655
656 fn mounted_tool_permission_request(&self, payload: &Value) -> Option<String> {
660 let request_id = Self::permission_request_id(payload)?;
661 let tool = payload
662 .get("request")?
663 .get("tool_name")
664 .and_then(Value::as_str)?;
665 let mounted = self.mounted_mcp_servers.iter().any(|name| {
666 tool.strip_prefix("mcp__")
667 .and_then(|rest| rest.strip_prefix(name.as_str()))
668 .is_some_and(|rest| rest.starts_with("__"))
669 });
670 mounted.then(|| request_id.to_string())
671 }
672
673 fn note_permission_request(&mut self, payload: &Value) {
675 let Some(request_id) = Self::permission_request_id(payload) else {
676 return;
677 };
678 if self
679 .pending_permissions
680 .iter()
681 .any(|pending| pending.request_id == request_id)
682 {
683 return;
684 }
685 self.pending_permissions.push(PendingPermission {
686 request_id: request_id.to_string(),
687 deadline: tokio::time::Instant::now() + self.permission_timeout,
688 });
689 }
690
691 async fn write_permission_response(&mut self, request_id: &str, body: Value) -> Result<()> {
693 self.transport
694 .write(json!({
695 "type": "control_response",
696 "response": {
697 "subtype": "success",
698 "request_id": request_id,
699 "response": body,
700 },
701 }))
702 .await
703 }
704
705 async fn deny_expired_permissions(&mut self) -> Result<()> {
710 let now = tokio::time::Instant::now();
711 let expired = self
712 .pending_permissions
713 .iter()
714 .filter(|pending| pending.deadline <= now)
715 .map(|pending| pending.request_id.clone())
716 .collect::<Vec<_>>();
717 self.pending_permissions
718 .retain(|pending| pending.deadline > now);
719 for request_id in expired {
720 self.write_permission_response(
721 &request_id,
722 json!({"behavior": "deny", "message": CLAUDE_PERMISSION_TIMEOUT_MESSAGE}),
723 )
724 .await?;
725 }
726 Ok(())
727 }
728
729 async fn transport_ended(&mut self) -> Result<Option<HarnessEvent>> {
746 if self.is_relay() || !self.spoke {
749 return Ok(None);
750 }
751 std::future::pending().await
752 }
753
754 fn next_permission_deadline(&self) -> Option<Duration> {
756 let now = tokio::time::Instant::now();
757 self.pending_permissions
758 .iter()
759 .map(|pending| pending.deadline.saturating_duration_since(now))
760 .min()
761 }
762}
763
764fn claude_permission_result(response: Value) -> Result<Value> {
774 let Value::Object(mut body) = response else {
775 return Err(claude_permission_shape_error(&response));
776 };
777 match body.get("behavior").and_then(Value::as_str) {
778 Some("allow") => {}
779 Some("deny") => {
780 let empty = body
782 .get("message")
783 .and_then(Value::as_str)
784 .is_none_or(str::is_empty);
785 if empty {
786 body.insert(
787 "message".into(),
788 Value::String("Volter Harness denied this permission request".into()),
789 );
790 }
791 }
792 _ => return Err(claude_permission_shape_error(&Value::Object(body))),
793 }
794 Ok(Value::Object(body))
795}
796
797fn claude_permission_shape_error(response: &Value) -> Error {
798 Error::Other(format!(
799 "Claude Code permission answers must carry a `behavior` of {}; got {response}",
800 CLAUDE_PERMISSION_BEHAVIORS
801 .map(|behavior| format!("`{behavior}`"))
802 .join(" or "),
803 ))
804}
805
806#[async_trait]
807impl RuntimeConnection for ClaudeRuntimeConnection {
808 fn handle(&self) -> &RuntimeHandle {
809 &self.handle
810 }
811
812 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
813 let content = if input.image_urls.is_empty() {
814 Value::String(input.text)
815 } else {
816 let mut parts = Vec::new();
817 if !input.text.is_empty() {
818 parts.push(json!({"type":"text", "text":input.text}));
819 }
820 for url in input.image_urls {
821 parts.push(claude_image_part(&url)?);
822 }
823 Value::Array(parts)
824 };
825 self.write_turn(json!({
826 "type": "user",
827 "session_id": self.handle.runtime_id,
828 "message": {"role": "user", "content": content},
829 }))
830 .await?;
831 Ok(None)
832 }
833
834 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
835 if let Some(payload) = self.buffered_events.pop_front() {
836 self.spoke = true;
837 return Ok(Some(harness_event(payload)));
838 }
839 loop {
840 self.deny_expired_permissions().await?;
844 if self.is_relay() && self.process_ended().await {
845 return Ok(None);
846 }
847 let payload = match self
848 .next_permission_deadline()
849 .map(|remaining| remaining.min(Duration::from_millis(200)))
850 {
851 Some(remaining) => {
852 match tokio::time::timeout(remaining, self.transport.receiver.recv()).await {
853 Err(_) => continue,
854 Ok(None) => return self.transport_ended().await,
855 Ok(Some(payload)) => payload,
856 }
857 }
858 None => match tokio::time::timeout(
859 Duration::from_millis(200),
860 self.transport.receiver.recv(),
861 )
862 .await
863 {
864 Err(_) => continue,
865 Ok(None) => return self.transport_ended().await,
866 Ok(Some(payload)) => payload,
867 },
868 };
869 self.spoke = true;
870 if Self::is_control_response(&payload) {
871 continue;
872 }
873 if let Some(request_id) = self.mounted_tool_permission_request(&payload) {
874 self.write_permission_response(&request_id, json!({"behavior": "allow"}))
875 .await?;
876 continue;
877 }
878 self.note_permission_request(&payload);
879 return Ok(Some(harness_event(payload)));
880 }
881 }
882
883 async fn interrupt(&mut self) -> Result<()> {
892 if self.process_ended().await {
896 return Ok(());
897 }
898 let request_id = format!(
899 "supercode-{}-interrupt-{}",
900 self.handle.runtime_id, self.next_control_request
901 );
902 self.next_control_request += 1;
903 self.transport
904 .write(json!({
905 "type": "control_request",
906 "request_id": request_id,
907 "request": {"subtype": "interrupt"},
908 }))
909 .await?;
910
911 let deadline = tokio::time::Instant::now() + self.control_timeout;
912 loop {
913 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
914 if remaining.is_zero() {
915 return Err(claude_interrupt_timeout(self.control_timeout));
916 }
917 match tokio::time::timeout(remaining, self.transport.receiver.recv()).await {
918 Err(_) => return Err(claude_interrupt_timeout(self.control_timeout)),
919 Ok(None) => return Err(Error::Other(
920 "Claude Code stream-json transport closed before acknowledging the interrupt"
921 .into(),
922 )),
923 Ok(Some(payload)) => {
924 if Self::is_control_response(&payload) {
925 if let Some(result) = Self::control_result(&payload, &request_id) {
926 return result;
927 }
928 continue;
929 }
930 if let Some(request_id) = self.mounted_tool_permission_request(&payload) {
934 self.write_permission_response(&request_id, json!({"behavior": "allow"}))
935 .await?;
936 continue;
937 }
938 self.note_permission_request(&payload);
939 self.buffered_events.push_back(payload);
940 }
941 }
942 }
943 }
944
945 async fn steer(&mut self, text: String) -> Result<()> {
946 self.send_input(RuntimeInput {
947 text,
948 image_urls: Vec::new(),
949 })
950 .await
951 .map(|_| ())
952 }
953
954 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
963 let Some(request_id) = request_id.as_str().map(str::to_string) else {
964 return Err(Error::Other(format!(
965 "Claude Code control requests are identified by a string `request_id`; got \
966 {request_id}"
967 )));
968 };
969 let Some(index) = self
970 .pending_permissions
971 .iter()
972 .position(|pending| pending.request_id == request_id)
973 else {
974 return Err(Error::Other(format!(
975 "no Claude Code permission request `{request_id}` is waiting on this connection — \
976 a `can_use_tool` request is answerable only while its turn is blocked on it, and \
977 only until it is answered or denied on timeout"
978 )));
979 };
980 let body = claude_permission_result(response)?;
981 self.pending_permissions.remove(index);
982 self.write_permission_response(&request_id, body).await
983 }
984
985 async fn close(&mut self) -> Result<()> {
986 self.transport.close().await
987 }
988}
989
990#[derive(Debug, Clone)]
992pub struct AcpRuntimeBackend {
993 harness: HarnessId,
994 launch: RuntimeLaunch,
995 resume_session: bool,
996}
997
998impl AcpRuntimeBackend {
999 pub fn new(harness: HarnessId, launch: RuntimeLaunch) -> Self {
1001 Self {
1002 harness,
1003 launch,
1004 resume_session: false,
1005 }
1006 }
1007
1008 pub fn with_resume_support(mut self, supported: bool) -> Self {
1012 self.resume_session = supported;
1013 self
1014 }
1015
1016 async fn connect(
1017 &self,
1018 cwd: &Path,
1019 launch: Option<RuntimeLaunch>,
1020 ) -> Result<(
1021 Arc<JsonLineClient>,
1022 mpsc::UnboundedReceiver<Value>,
1023 RuntimeEndpoint,
1024 Value,
1025 )> {
1026 let launch = launch.unwrap_or_else(|| self.launch.clone());
1027 let (client, receiver, endpoint) =
1028 JsonLineClient::spawn(&launch, Some(cwd), true, "acp-v1-jsonrpc").await?;
1029 let initialized = client
1030 .request(
1031 "initialize",
1032 json!({
1033 "protocolVersion": 1,
1034 "clientCapabilities": {},
1035 "clientInfo": {
1036 "name": "supercode",
1037 "title": "Volter Harness",
1038 "version": env!("CARGO_PKG_VERSION"),
1039 },
1040 }),
1041 )
1042 .await?;
1043 if initialized.get("protocolVersion").and_then(Value::as_u64) != Some(1) {
1044 return Err(Error::Other(format!(
1045 "ACP agent negotiated unsupported protocol version: {}",
1046 initialized
1047 .get("protocolVersion")
1048 .cloned()
1049 .unwrap_or(Value::Null)
1050 )));
1051 }
1052 Ok((client, receiver, endpoint, initialized))
1053 }
1054
1055 async fn session_request(
1056 &self,
1057 client: &JsonLineClient,
1058 initialized: &Value,
1059 method: &str,
1060 params: Value,
1061 ) -> Result<Value> {
1062 match client.request(method, params.clone()).await {
1063 Ok(response) => Ok(response),
1064 Err(error) if acp_auth_required(&error.to_string()) => {
1065 let cached = initialized
1066 .get("authMethods")
1067 .and_then(Value::as_array)
1068 .and_then(|methods| {
1069 methods.iter().find_map(|candidate| {
1070 (candidate.get("id").and_then(Value::as_str) == Some("cached_token"))
1071 .then_some("cached_token")
1072 })
1073 });
1074 let Some(method_id) = cached else {
1075 return Err(Error::Other(
1076 "ACP agent requires authentication but did not advertise the non-interactive `cached_token` method"
1077 .into(),
1078 ));
1079 };
1080 client
1081 .request(
1082 "authenticate",
1083 json!({"methodId": method_id, "_meta": {"headless": true}}),
1084 )
1085 .await?;
1086 client.request(method, params).await
1087 }
1088 Err(error) => Err(error),
1089 }
1090 }
1091
1092 async fn connection(
1093 &self,
1094 cwd: &Path,
1095 runtime_id: Option<String>,
1096 launch: Option<RuntimeLaunch>,
1097 mcp_servers: Vec<McpServerLaunch>,
1098 ) -> Result<Box<dyn RuntimeConnection>> {
1099 let mcp_servers = acp_mcp_servers(&mcp_servers);
1100 let (launch, runtime_id) = match runtime_id {
1107 Some(session_id) if self.harness.as_str() == HarnessId::SUPERCODE => {
1108 let mut continuation = launch.unwrap_or_else(|| self.launch.clone());
1109 let flags = continuation
1110 .arguments
1111 .iter()
1112 .skip_while(|argument| argument.as_str() != "acp")
1113 .skip(1)
1114 .cloned()
1115 .collect::<Vec<_>>();
1116 continuation.arguments = ["resume", &session_id, "--harness", "supercode", "--acp"]
1117 .into_iter()
1118 .map(String::from)
1119 .chain(flags)
1120 .collect();
1121 (Some(continuation), None)
1122 }
1123 other => (launch, other),
1124 };
1125 let (client, mut receiver, endpoint, initialized) = self.connect(cwd, launch).await?;
1126 let session_id = if let Some(session_id) = runtime_id {
1127 let resume = initialized
1128 .pointer("/agentCapabilities/sessionCapabilities/resume")
1129 .is_some();
1130 let load = initialized
1131 .pointer("/agentCapabilities/loadSession")
1132 .and_then(Value::as_bool)
1133 .unwrap_or(false);
1134 let method = if resume {
1135 "session/resume"
1136 } else if load {
1137 "session/load"
1138 } else {
1139 return Err(Error::Other(
1140 "ACP agent did not advertise session resume or load".into(),
1141 ));
1142 };
1143 self.session_request(
1144 client.as_ref(),
1145 &initialized,
1146 method,
1147 json!({"sessionId": session_id, "cwd": cwd, "mcpServers": mcp_servers}),
1148 )
1149 .await?;
1150 session_id
1151 } else {
1152 self.session_request(
1153 client.as_ref(),
1154 &initialized,
1155 "session/new",
1156 json!({"cwd": cwd, "mcpServers": mcp_servers}),
1157 )
1158 .await?
1159 .get("sessionId")
1160 .and_then(Value::as_str)
1161 .ok_or_else(|| Error::Other("ACP session/new omitted sessionId".into()))?
1162 .to_string()
1163 };
1164 while receiver.try_recv().is_ok() {}
1173 Ok(Box::new(AcpRuntimeConnection {
1174 handle: RuntimeHandle {
1175 harness: self.harness.clone(),
1176 runtime_id: session_id,
1177 endpoint,
1178 },
1179 client,
1180 receiver,
1181 active_prompt: None,
1182 }))
1183 }
1184}
1185
1186fn acp_mcp_servers(servers: &[McpServerLaunch]) -> Value {
1191 Value::Array(
1192 servers
1193 .iter()
1194 .map(|server| {
1195 json!({
1196 "name": server.name,
1197 "command": server.command,
1198 "args": server.arguments,
1199 "env": server
1200 .env
1201 .iter()
1202 .map(|(name, value)| json!({"name": name, "value": value}))
1203 .collect::<Vec<_>>(),
1204 })
1205 })
1206 .collect::<Vec<_>>(),
1207 )
1208}
1209
1210fn acp_auth_required(message: &str) -> bool {
1211 let message = message.to_ascii_lowercase();
1212 [
1213 "auth",
1214 "login",
1215 "sign in",
1216 "sign-in",
1217 "unauthorized",
1218 "forbidden",
1219 "credential",
1220 ]
1221 .iter()
1222 .any(|needle| message.contains(needle))
1223}
1224
1225#[async_trait]
1226impl RuntimeBackend for AcpRuntimeBackend {
1227 fn harness(&self) -> HarnessId {
1228 self.harness.clone()
1229 }
1230
1231 fn capabilities(&self) -> RuntimeCapabilities {
1232 RuntimeCapabilities {
1233 start_session: true,
1234 resume_session: self.resume_session,
1237 attach_existing_process: false,
1238 send_input: true,
1239 stream_events: true,
1240 interrupt: true,
1241 steer: false,
1242 respond_to_requests: true,
1243 }
1244 }
1245
1246 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
1247 self.connection(&request.cwd, None, request.launch, request.mcp_servers)
1248 .await
1249 }
1250
1251 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
1252 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1253 self.connection(
1254 &cwd,
1255 Some(request.runtime_id),
1256 request.launch,
1257 request.mcp_servers,
1258 )
1259 .await
1260 }
1261}
1262
1263struct AcpRuntimeConnection {
1264 handle: RuntimeHandle,
1265 client: Arc<JsonLineClient>,
1266 receiver: mpsc::UnboundedReceiver<Value>,
1267 active_prompt: Option<u64>,
1268}
1269
1270#[async_trait]
1271impl RuntimeConnection for AcpRuntimeConnection {
1272 fn handle(&self) -> &RuntimeHandle {
1273 &self.handle
1274 }
1275
1276 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
1277 let mut prompt = Vec::new();
1278 if !input.text.is_empty() {
1279 prompt.push(json!({"type": "text", "text": input.text}));
1280 }
1281 for url in input.image_urls {
1282 let (mime_type, data) = data_image_parts(&url).ok_or_else(|| {
1283 Error::Other("ACP image prompts require base64 image data URLs".into())
1284 })?;
1285 prompt.push(json!({"type":"image", "mimeType":mime_type, "data":data}));
1286 }
1287 let (id, response) = self
1288 .client
1289 .begin_request(
1290 "session/prompt",
1291 json!({
1292 "sessionId": self.handle.runtime_id,
1293 "prompt": prompt,
1294 }),
1295 )
1296 .await?;
1297 self.active_prompt = Some(id);
1298 let client = self.client.clone();
1299 tokio::spawn(async move {
1300 let result = match response.await {
1301 Ok(Ok(result)) => json!({"id": id, "result": result}),
1302 Ok(Err(error)) => json!({"id": id, "error": error}),
1303 Err(_) => json!({"id": id, "error": "response channel closed"}),
1304 };
1305 client.emit(json!({
1306 "jsonrpc": "2.0",
1307 "method": "supercode/acp_request_completed",
1308 "params": result,
1309 }));
1310 });
1311 Ok(Some(id.to_string()))
1312 }
1313
1314 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
1315 let Some(payload) = self.receiver.recv().await else {
1316 return Ok(None);
1317 };
1318 let kind = payload
1319 .get("method")
1320 .and_then(Value::as_str)
1321 .or_else(|| payload.get("type").and_then(Value::as_str))
1322 .unwrap_or("protocol")
1323 .to_string();
1324 if kind == "supercode/acp_request_completed" {
1325 self.active_prompt = None;
1326 }
1327 Ok(Some(HarnessEvent {
1328 sequence: None,
1329 kind,
1330 payload,
1331 }))
1332 }
1333
1334 async fn interrupt(&mut self) -> Result<()> {
1335 self.client
1336 .notify(
1337 "session/cancel",
1338 json!({"sessionId": self.handle.runtime_id}),
1339 )
1340 .await
1341 }
1342
1343 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
1344 self.client.respond(request_id, response).await
1345 }
1346
1347 async fn close(&mut self) -> Result<()> {
1348 self.client.close().await
1349 }
1350}
1351
1352#[derive(Debug, Clone)]
1356pub struct OpenCodeRuntimeBackend {
1357 launch: RuntimeLaunch,
1358 base_url: Option<String>,
1359 bearer: Option<BearerToken>,
1360}
1361
1362impl Default for OpenCodeRuntimeBackend {
1363 fn default() -> Self {
1364 Self::new()
1365 }
1366}
1367
1368impl OpenCodeRuntimeBackend {
1369 pub fn new() -> Self {
1371 Self {
1372 launch: RuntimeLaunch {
1373 program: "opencode".into(),
1374 arguments: vec!["serve".into()],
1375 env: BTreeMap::new(),
1376 },
1377 base_url: None,
1378 bearer: None,
1379 }
1380 }
1381
1382 pub fn connect(base_url: impl Into<String>) -> Self {
1385 Self {
1386 base_url: Some(base_url.into().trim_end_matches('/').to_string()),
1387 ..Self::new()
1388 }
1389 }
1390
1391 pub fn with_launch(mut self, launch: RuntimeLaunch) -> Self {
1393 self.launch = launch;
1394 self
1395 }
1396
1397 pub fn with_bearer(mut self, token: BearerToken) -> Self {
1400 self.bearer = Some(token);
1401 self
1402 }
1403
1404 fn http_client(&self) -> Result<reqwest::Client> {
1405 let Some(token) = &self.bearer else {
1406 return Ok(reqwest::Client::new());
1407 };
1408 let mut headers = reqwest::header::HeaderMap::new();
1409 let mut value =
1410 reqwest::header::HeaderValue::from_str(&format!("Bearer {}", token.secret())).map_err(
1411 |_| Error::Other("connect-mode bearer token is not a valid header value".into()),
1412 )?;
1413 value.set_sensitive(true);
1414 headers.insert(reqwest::header::AUTHORIZATION, value);
1415 reqwest::Client::builder()
1416 .default_headers(headers)
1417 .build()
1418 .map_err(|error| Error::Other(format!("could not build HTTP client: {error}")))
1419 }
1420
1421 async fn service(
1422 &self,
1423 client: &reqwest::Client,
1424 launch: Option<RuntimeLaunch>,
1425 ) -> Result<(String, Option<super::GroupLeader>)> {
1426 if let Some(base_url) = &self.base_url {
1427 wait_for_health(client, base_url).await?;
1428 return Ok((base_url.clone(), None));
1429 }
1430 let port = TcpListener::bind(("127.0.0.1", 0))?.local_addr()?.port();
1431 let mut launch = launch.unwrap_or_else(|| self.launch.clone());
1432 launch.arguments.extend([
1433 "--hostname".into(),
1434 "127.0.0.1".into(),
1435 "--port".into(),
1436 port.to_string(),
1437 ]);
1438 let mut command = Command::new(&launch.program);
1439 command
1440 .args(&launch.arguments)
1441 .envs(&launch.env)
1442 .stdin(Stdio::null())
1443 .stdout(Stdio::null())
1444 .stderr(Stdio::inherit())
1445 .kill_on_drop(true);
1446 #[cfg(unix)]
1450 command.process_group(0);
1451 let mut child = super::GroupLeader(command.spawn().map_err(|error| {
1456 Error::Other(format!("could not launch {}: {error}", launch.program))
1457 })?);
1458 let base_url = format!("http://127.0.0.1:{port}");
1459 if let Err(error) = wait_for_health(client, &base_url).await {
1460 let _ = terminate_opencode_server(&mut child).await;
1461 return Err(error);
1462 }
1463 Ok((base_url, Some(child)))
1464 }
1465
1466 async fn open(
1467 &self,
1468 cwd: &Path,
1469 runtime_id: Option<String>,
1470 launch: Option<RuntimeLaunch>,
1471 ) -> Result<Box<dyn RuntimeConnection>> {
1472 let client = self.http_client()?;
1473 let (base_url, child) = self.service(&client, launch).await?;
1474 let cwd_string = cwd.to_string_lossy().to_string();
1475 let runtime_id = match runtime_id {
1476 Some(id) => {
1477 http_ok(
1478 client
1479 .get(format!("{base_url}/session/{id}"))
1480 .query(&[("directory", &cwd_string)])
1481 .send()
1482 .await,
1483 )
1484 .await?;
1485 id
1486 }
1487 None => {
1488 let response = http_ok(
1489 client
1490 .post(format!("{base_url}/session"))
1491 .query(&[("directory", &cwd_string)])
1492 .json(&json!({}))
1493 .send()
1494 .await,
1495 )
1496 .await?;
1497 response
1498 .json::<Value>()
1499 .await
1500 .map_err(http_error)?
1501 .get("id")
1502 .and_then(Value::as_str)
1503 .ok_or_else(|| Error::Other("OpenCode create session omitted id".into()))?
1504 .to_string()
1505 }
1506 };
1507 let receiver = spawn_sse(
1508 client.clone(),
1509 format!("{base_url}/event"),
1510 cwd_string.clone(),
1511 );
1512 Ok(Box::new(OpenCodeRuntimeConnection {
1513 handle: RuntimeHandle {
1514 harness: HarnessId::from(HarnessId::OPENCODE),
1515 runtime_id,
1516 endpoint: RuntimeEndpoint::Http {
1517 base_url: base_url.clone(),
1518 protocol: "opencode-http-sse".into(),
1519 },
1520 },
1521 base_url,
1522 cwd: cwd_string,
1523 client,
1524 receiver,
1525 child,
1526 }))
1527 }
1528}
1529
1530#[async_trait]
1531impl RuntimeBackend for OpenCodeRuntimeBackend {
1532 fn harness(&self) -> HarnessId {
1533 HarnessId::from(HarnessId::OPENCODE)
1534 }
1535
1536 fn capabilities(&self) -> RuntimeCapabilities {
1537 RuntimeCapabilities {
1538 start_session: true,
1539 resume_session: true,
1540 attach_existing_process: self.base_url.is_some(),
1541 send_input: true,
1542 stream_events: true,
1543 interrupt: true,
1544 steer: false,
1545 respond_to_requests: true,
1546 }
1547 }
1548
1549 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
1550 self.open(&request.cwd, None, request.launch).await
1551 }
1552
1553 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
1554 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1555 self.open(&cwd, Some(request.runtime_id), request.launch)
1556 .await
1557 }
1558
1559 async fn attach_existing(
1560 &self,
1561 request: RuntimeAttachRequest,
1562 ) -> Result<Box<dyn RuntimeConnection>> {
1563 if self.base_url.is_none() {
1564 return Err(Error::Other(
1565 "OpenCode live attach requires the existing server's `base_url`".into(),
1566 ));
1567 }
1568 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1569 self.open(&cwd, Some(request.runtime_id), request.launch)
1570 .await
1571 }
1572}
1573
1574struct OpenCodeRuntimeConnection {
1575 handle: RuntimeHandle,
1576 base_url: String,
1577 cwd: String,
1578 client: reqwest::Client,
1579 receiver: mpsc::UnboundedReceiver<Value>,
1580 child: Option<super::GroupLeader>,
1581}
1582
1583#[async_trait]
1584impl RuntimeConnection for OpenCodeRuntimeConnection {
1585 fn handle(&self) -> &RuntimeHandle {
1586 &self.handle
1587 }
1588
1589 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
1590 let mut parts = Vec::new();
1591 if !input.text.is_empty() {
1592 parts.push(json!({"type": "text", "text": input.text}));
1593 }
1594 for url in input.image_urls {
1595 let mime = image_mime_type(&url).ok_or_else(|| {
1596 Error::Other("OpenCode image prompts require a recognizable image MIME type".into())
1597 })?;
1598 parts.push(json!({"type":"file", "mime":mime, "url":url}));
1599 }
1600 http_ok(
1601 self.client
1602 .post(format!(
1603 "{}/session/{}/prompt_async",
1604 self.base_url, self.handle.runtime_id
1605 ))
1606 .query(&[("directory", &self.cwd)])
1607 .json(&json!({"parts": parts}))
1608 .send()
1609 .await,
1610 )
1611 .await?;
1612 Ok(None)
1613 }
1614
1615 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
1616 loop {
1617 let Some(payload) = self.receiver.recv().await else {
1618 return Ok(None);
1619 };
1620 if opencode_event_session_id(&payload)
1621 .is_some_and(|session_id| session_id != self.handle.runtime_id)
1622 {
1623 continue;
1624 }
1625 let kind = payload
1626 .get("type")
1627 .and_then(Value::as_str)
1628 .unwrap_or("event")
1629 .to_string();
1630 return Ok(Some(HarnessEvent {
1631 sequence: None,
1632 kind,
1633 payload,
1634 }));
1635 }
1636 }
1637
1638 async fn interrupt(&mut self) -> Result<()> {
1639 http_ok(
1640 self.client
1641 .post(format!(
1642 "{}/session/{}/abort",
1643 self.base_url, self.handle.runtime_id
1644 ))
1645 .query(&[("directory", &self.cwd)])
1646 .send()
1647 .await,
1648 )
1649 .await?;
1650 Ok(())
1651 }
1652
1653 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
1654 let permission = request_id.as_str().ok_or_else(|| {
1655 Error::Other("OpenCode permission request id must be a string".into())
1656 })?;
1657 http_ok(
1658 self.client
1659 .post(format!(
1660 "{}/session/{}/permissions/{permission}",
1661 self.base_url, self.handle.runtime_id
1662 ))
1663 .query(&[("directory", &self.cwd)])
1664 .json(&response)
1665 .send()
1666 .await,
1667 )
1668 .await?;
1669 Ok(())
1670 }
1671
1672 async fn close(&mut self) -> Result<()> {
1673 if let Some(child) = &mut self.child {
1674 terminate_opencode_server(child).await?;
1675 }
1676 Ok(())
1677 }
1678}
1679
1680fn data_image_parts(url: &str) -> Option<(&str, &str)> {
1681 let rest = url.strip_prefix("data:")?;
1682 let (mime_type, data) = rest.split_once(";base64,")?;
1683 mime_type.starts_with("image/").then_some((mime_type, data))
1684}
1685
1686fn image_mime_type(url: &str) -> Option<&str> {
1687 if let Some((mime_type, _)) = data_image_parts(url) {
1688 return Some(mime_type);
1689 }
1690 let path = url.split(['?', '#']).next()?.to_ascii_lowercase();
1691 if path.ends_with(".png") {
1692 Some("image/png")
1693 } else if path.ends_with(".jpg") || path.ends_with(".jpeg") {
1694 Some("image/jpeg")
1695 } else if path.ends_with(".gif") {
1696 Some("image/gif")
1697 } else if path.ends_with(".webp") {
1698 Some("image/webp")
1699 } else {
1700 None
1701 }
1702}
1703
1704fn claude_image_part(url: &str) -> Result<Value> {
1705 if let Some((media_type, data)) = data_image_parts(url) {
1706 return Ok(json!({
1707 "type":"image",
1708 "source":{"type":"base64", "media_type":media_type, "data":data}
1709 }));
1710 }
1711 if url.starts_with("https://") || url.starts_with("http://") {
1712 return Ok(json!({"type":"image", "source":{"type":"url", "url":url}}));
1713 }
1714 Err(Error::Other(
1715 "Claude image prompts require image data URLs or HTTP(S) URLs".into(),
1716 ))
1717}
1718
1719fn opencode_event_session_id(payload: &Value) -> Option<&str> {
1720 let properties = payload.get("properties").unwrap_or(payload);
1721 properties
1722 .get("sessionID")
1723 .and_then(Value::as_str)
1724 .or_else(|| {
1725 properties
1726 .get("part")
1727 .and_then(|part| part.get("sessionID"))
1728 .and_then(Value::as_str)
1729 })
1730 .or_else(|| {
1731 properties
1732 .get("info")
1733 .and_then(|info| info.get("sessionID"))
1734 .and_then(Value::as_str)
1735 })
1736}
1737
1738async fn terminate_opencode_server(child: &mut Child) -> Result<()> {
1739 #[cfg(unix)]
1740 let process_group = child.id();
1741 let leader_exited = child.try_wait()?.is_some();
1742 if leader_exited {
1743 #[cfg(unix)]
1744 if let Some(pid) = process_group.filter(|pid| process_group_exists(*pid)) {
1745 crate::lsp::kill_process_group(pid);
1746 wait_for_process_group_exit(pid, Duration::from_secs(3)).await?;
1747 }
1748 return Ok(());
1749 }
1750 #[cfg(unix)]
1756 if let Some(pid) = process_group {
1757 unsafe {
1758 libc::kill(-(pid as libc::pid_t), libc::SIGTERM);
1759 }
1760 let mut leader_reaped = false;
1761 if let Ok(status) = tokio::time::timeout(Duration::from_millis(500), child.wait()).await {
1762 status?;
1763 leader_reaped = true;
1764 if !process_group_exists(pid) {
1765 return Ok(());
1766 }
1767 }
1768 crate::lsp::kill_process_group(pid);
1771 if leader_reaped {
1772 return wait_for_process_group_exit(pid, Duration::from_secs(3)).await;
1773 }
1774 }
1775 #[cfg(not(unix))]
1776 child.start_kill()?;
1777 tokio::time::timeout(Duration::from_secs(3), child.wait())
1778 .await
1779 .map_err(|_| Error::Other("timed out reaping the OpenCode server".into()))??;
1780 #[cfg(unix)]
1781 if let Some(pid) = process_group {
1782 wait_for_process_group_exit(pid, Duration::from_secs(3)).await?;
1783 }
1784 Ok(())
1785}
1786
1787#[cfg(unix)]
1788fn process_group_exists(pid: u32) -> bool {
1789 let result = unsafe { libc::kill(-(pid as libc::pid_t), 0) };
1790 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
1791}
1792
1793#[cfg(unix)]
1794async fn wait_for_process_group_exit(pid: u32, timeout: Duration) -> Result<()> {
1795 let deadline = tokio::time::Instant::now() + timeout;
1796 while process_group_exists(pid) {
1797 if tokio::time::Instant::now() >= deadline {
1798 return Err(Error::Other(format!(
1799 "timed out stopping OpenCode process group {pid}"
1800 )));
1801 }
1802 tokio::time::sleep(Duration::from_millis(10)).await;
1803 }
1804 Ok(())
1805}
1806
1807struct RawLineTransport {
1808 stdin: Mutex<ChildStdin>,
1809 child: Mutex<super::GroupLeader>,
1810 receiver: mpsc::UnboundedReceiver<Value>,
1811 endpoint: RuntimeEndpoint,
1812}
1813
1814impl RawLineTransport {
1815 async fn spawn(launch: &RuntimeLaunch, cwd: Option<&Path>, protocol: &str) -> Result<Self> {
1816 let mut command = Command::new(&launch.program);
1817 command
1818 .args(&launch.arguments)
1819 .envs(&launch.env)
1820 .stdin(Stdio::piped())
1821 .stdout(Stdio::piped())
1822 .stderr(Stdio::inherit())
1823 .kill_on_drop(true);
1824 #[cfg(unix)]
1828 command.process_group(0);
1829 if let Some(cwd) = cwd {
1830 command.current_dir(cwd);
1831 }
1832 let mut child = command.spawn().map_err(|error| {
1833 Error::Other(format!("could not launch {}: {error}", launch.program))
1834 })?;
1835 let pid = child.id();
1836 let stdin = child
1837 .stdin
1838 .take()
1839 .ok_or_else(|| Error::Other("runtime child has no stdin".into()))?;
1840 let stdout = child
1841 .stdout
1842 .take()
1843 .ok_or_else(|| Error::Other("runtime child has no stdout".into()))?;
1844 let (sender, receiver) = mpsc::unbounded_channel();
1845 tokio::spawn(async move {
1846 let mut lines = BufReader::new(stdout).lines();
1847 while let Ok(Some(line)) = lines.next_line().await {
1848 let value = serde_json::from_str(&line)
1849 .unwrap_or_else(|_| json!({"type": "malformed_output", "line": line}));
1850 let _ = sender.send(value);
1851 }
1852 });
1853 Ok(Self {
1854 stdin: Mutex::new(stdin),
1855 child: Mutex::new(super::GroupLeader(child)),
1856 receiver,
1857 endpoint: RuntimeEndpoint::LocalProcess {
1858 pid,
1859 command: std::iter::once(launch.program.clone())
1860 .chain(launch.arguments.iter().cloned())
1861 .collect(),
1862 protocol: protocol.into(),
1863 },
1864 })
1865 }
1866
1867 async fn write(&self, value: Value) -> Result<()> {
1868 let mut stdin = self.stdin.lock().await;
1869 stdin.write_all(value.to_string().as_bytes()).await?;
1870 stdin.write_all(b"\n").await?;
1871 stdin.flush().await?;
1872 Ok(())
1873 }
1874
1875 async fn close(&self) -> Result<()> {
1876 let mut child = self.child.lock().await;
1877 if child.try_wait()?.is_some() {
1878 return Ok(());
1879 }
1880 #[cfg(unix)]
1883 if let Some(pid) = child.id() {
1884 crate::lsp::kill_process_group(pid);
1885 tokio::time::timeout(Duration::from_secs(3), child.wait())
1886 .await
1887 .map_err(|_| Error::Other("timed out reaping runtime process group".into()))??;
1888 return Ok(());
1889 }
1890 #[cfg(not(unix))]
1891 child.kill().await?;
1892 Ok(())
1893 }
1894}
1895
1896async fn raw_next_event(
1897 receiver: &mut mpsc::UnboundedReceiver<Value>,
1898) -> Result<Option<HarnessEvent>> {
1899 let Some(payload) = receiver.recv().await else {
1900 return Ok(None);
1901 };
1902 Ok(Some(harness_event(payload)))
1903}
1904
1905fn harness_event(payload: Value) -> HarnessEvent {
1906 let kind = payload
1907 .get("type")
1908 .and_then(Value::as_str)
1909 .unwrap_or("event")
1910 .to_string();
1911 HarnessEvent {
1912 sequence: None,
1913 kind,
1914 payload,
1915 }
1916}
1917
1918fn broken_pipe(error: &Error) -> bool {
1925 matches!(error, Error::Io(io) if io.kind() == std::io::ErrorKind::BrokenPipe)
1926}
1927
1928fn claude_interrupt_timeout(bound: Duration) -> Error {
1929 Error::Other(format!(
1930 "Claude Code did not acknowledge the interrupt control request within {}s",
1931 bound.as_secs_f32()
1932 ))
1933}
1934
1935pub(crate) fn generated_session_id() -> String {
1936 let mut bytes = [0_u8; 16];
1937 if getrandom::getrandom(&mut bytes).is_err() {
1938 let nanos = SystemTime::now()
1939 .duration_since(UNIX_EPOCH)
1940 .unwrap_or_default()
1941 .as_nanos()
1942 .to_le_bytes();
1943 bytes.copy_from_slice(&nanos);
1944 }
1945 bytes[6] = (bytes[6] & 0x0f) | 0x40;
1946 bytes[8] = (bytes[8] & 0x3f) | 0x80;
1947 format!(
1948 "{:02x}{:02x}{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}{:02x}{:02x}{:02x}{:02x}",
1949 bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5], bytes[6], bytes[7],
1950 bytes[8], bytes[9], bytes[10], bytes[11], bytes[12], bytes[13], bytes[14], bytes[15]
1951 )
1952}
1953
1954async fn wait_for_health(client: &reqwest::Client, base_url: &str) -> Result<()> {
1955 wait_for_health_for(client, base_url, Duration::from_secs(10)).await
1956}
1957
1958async fn wait_for_health_for(
1959 client: &reqwest::Client,
1960 base_url: &str,
1961 total_timeout: Duration,
1962) -> Result<()> {
1963 let url = format!("{base_url}/global/health");
1964 let mut last = None;
1965 let deadline = tokio::time::Instant::now() + total_timeout;
1966 loop {
1971 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
1972 if remaining.is_zero() {
1973 break;
1974 }
1975 let request_timeout = remaining.min(Duration::from_millis(500));
1976 match tokio::time::timeout(request_timeout, client.get(&url).send()).await {
1977 Ok(Ok(response)) if response.status().is_success() => return Ok(()),
1978 Ok(Ok(response)) => last = Some(format!("HTTP {}", response.status())),
1979 Ok(Err(error)) => last = Some(error.to_string()),
1980 Err(_) => last = Some("health request timed out".into()),
1981 }
1982 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
1983 if !remaining.is_zero() {
1984 tokio::time::sleep(remaining.min(Duration::from_millis(100))).await;
1985 }
1986 }
1987 Err(Error::Other(format!(
1988 "OpenCode server at {base_url} did not become healthy: {}",
1989 last.unwrap_or_else(|| "no response".into())
1990 )))
1991}
1992
1993async fn http_ok(
1994 response: std::result::Result<reqwest::Response, reqwest::Error>,
1995) -> Result<reqwest::Response> {
1996 response
1997 .map_err(http_error)?
1998 .error_for_status()
1999 .map_err(http_error)
2000}
2001
2002fn http_error(error: reqwest::Error) -> Error {
2003 Error::Other(format!("runtime HTTP request failed: {error}"))
2004}
2005
2006fn spawn_sse(
2007 client: reqwest::Client,
2008 url: String,
2009 directory: String,
2010) -> mpsc::UnboundedReceiver<Value> {
2011 let (sender, receiver) = mpsc::unbounded_channel();
2012 tokio::spawn(async move {
2013 let response = client
2014 .get(url)
2015 .query(&[("directory", directory)])
2016 .send()
2017 .await;
2018 let Ok(response) = response.and_then(reqwest::Response::error_for_status) else {
2019 let _ = sender.send(
2020 json!({"type": "stream_error", "message": "could not open OpenCode SSE stream"}),
2021 );
2022 return;
2023 };
2024 let mut stream = response.bytes_stream();
2025 let mut buffer = String::new();
2026 while let Some(chunk) = stream.next().await {
2027 let Ok(chunk) = chunk else {
2028 break;
2029 };
2030 buffer.push_str(&String::from_utf8_lossy(&chunk));
2031 while let Some(newline) = buffer.find('\n') {
2032 let line = buffer[..newline].trim_end_matches('\r').to_string();
2033 buffer.drain(..=newline);
2034 if let Some(data) = line.strip_prefix("data:") {
2035 let data = data.trim();
2036 if let Ok(value) = serde_json::from_str(data) {
2037 let _ = sender.send(value);
2038 }
2039 }
2040 }
2041 }
2042 });
2043 receiver
2044}