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 self.open(
408 &request.cwd,
409 session_id,
410 request.launch,
411 &request.mcp_servers,
412 false,
413 )
414 .await
415 }
416
417 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
418 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
419 self.open(
420 &cwd,
421 request.runtime_id,
422 request.launch,
423 &request.mcp_servers,
424 true,
425 )
426 .await
427 }
428}
429
430async fn mcp_config_file(
434 harness: &str,
435 runtime_id: &str,
436 servers: &[McpServerLaunch],
437) -> Result<PathBuf> {
438 let mut entries = serde_json::Map::new();
439 for server in servers {
440 entries.insert(
441 server.name.clone(),
442 json!({
443 "type": "stdio",
444 "command": server.command,
445 "args": server.arguments,
446 "env": server.env,
447 }),
448 );
449 }
450 let safe: String = runtime_id
451 .chars()
452 .filter(|c| c.is_ascii_alphanumeric() || *c == '-' || *c == '_')
453 .collect();
454 let path = std::env::temp_dir().join(format!("supercode-{harness}-mcp-{safe}.json"));
455 let body = serde_json::to_vec_pretty(&json!({ "mcpServers": entries })).map_err(|error| {
456 Error::Other(format!(
457 "{harness} mcp config could not be encoded: {error}"
458 ))
459 })?;
460 tokio::fs::write(&path, body).await.map_err(|error| {
461 Error::Other(format!(
462 "{harness} mcp config could not be written: {error}"
463 ))
464 })?;
465 #[cfg(unix)]
466 {
467 use std::os::unix::fs::PermissionsExt;
468 tokio::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600))
469 .await
470 .map_err(|error| {
471 Error::Other(format!(
472 "{harness} mcp config could not be protected: {error}"
473 ))
474 })?;
475 }
476 Ok(path)
477}
478
479struct ClaudeRuntimeConnection {
480 handle: RuntimeHandle,
481 transport: RawLineTransport,
482 prefix: RuntimeLaunch,
486 cwd: PathBuf,
490 spoke: bool,
495 buffered_events: VecDeque<Value>,
499 next_control_request: u64,
500 control_timeout: Duration,
501 pending_permissions: Vec<PendingPermission>,
505 permission_timeout: Duration,
506 mounted_mcp_servers: Vec<String>,
510}
511
512struct PendingPermission {
514 request_id: String,
516 deadline: tokio::time::Instant,
518}
519
520impl ClaudeRuntimeConnection {
521 fn is_relay(&self) -> bool {
522 self.prefix
523 .env
524 .get("SUPERCODE_CLAUDE_RELAY")
525 .is_some_and(|value| value == "1")
526 }
527
528 async fn process_ended(&self) -> bool {
535 matches!(self.transport.child.lock().await.try_wait(), Ok(Some(_)))
536 }
537
538 async fn reopen(&mut self) -> Result<()> {
548 let mut launch = self.prefix.clone();
549 launch
550 .arguments
551 .extend(["--resume".into(), self.handle.runtime_id.clone()]);
552 let transport = RawLineTransport::spawn(&launch, Some(&self.cwd), "claude-stream-json")
553 .await
554 .map_err(|error| {
555 Error::Other(format!(
556 "could not resume Claude Code session `{}` after its process exited: {error}",
557 self.handle.runtime_id
558 ))
559 })?;
560 self.handle.endpoint = transport.endpoint.clone();
561 self.transport = transport;
562 self.pending_permissions.clear();
565 self.spoke = false;
566 Ok(())
567 }
568
569 async fn write_turn(&mut self, frame: Value) -> Result<()> {
577 if self.process_ended().await {
578 if self.is_relay() {
579 return Err(Error::Other("Claude relay process exited".into()));
580 }
581 self.reopen().await?;
582 }
583 match self.transport.write(frame.clone()).await {
584 Ok(()) => Ok(()),
585 Err(error) if broken_pipe(&error) && !self.is_relay() => {
586 self.reopen().await?;
587 self.transport.write(frame).await
588 }
589 Err(error) => Err(error),
590 }
591 }
592
593 fn is_control_response(value: &Value) -> bool {
598 value.get("type").and_then(Value::as_str) == Some("control_response")
599 }
600
601 fn control_result(value: &Value, request_id: &str) -> Option<Result<()>> {
609 let response = value.get("response")?;
610 if response.get("request_id").and_then(Value::as_str) != Some(request_id) {
611 return None;
612 }
613 match response.get("subtype").and_then(Value::as_str) {
614 Some("success") => Some(Ok(())),
615 other => Some(Err(Error::Other(format!(
616 "Claude Code rejected the interrupt control request: {}",
617 response
618 .get("error")
619 .and_then(Value::as_str)
620 .map(str::to_string)
621 .unwrap_or_else(|| format!(
622 "control_response subtype {}",
623 other.unwrap_or("(missing)")
624 ))
625 )))),
626 }
627 }
628
629 fn permission_request_id(value: &Value) -> Option<&str> {
636 if value.get("type").and_then(Value::as_str)? != "control_request" {
637 return None;
638 }
639 let request = value.get("request")?;
640 if request.get("subtype").and_then(Value::as_str)? != "can_use_tool" {
641 return None;
642 }
643 value.get("request_id").and_then(Value::as_str)
644 }
645
646 fn mounted_tool_permission_request(&self, payload: &Value) -> Option<String> {
650 let request_id = Self::permission_request_id(payload)?;
651 let tool = payload
652 .get("request")?
653 .get("tool_name")
654 .and_then(Value::as_str)?;
655 let mounted = self.mounted_mcp_servers.iter().any(|name| {
656 tool.strip_prefix("mcp__")
657 .and_then(|rest| rest.strip_prefix(name.as_str()))
658 .is_some_and(|rest| rest.starts_with("__"))
659 });
660 mounted.then(|| request_id.to_string())
661 }
662
663 fn note_permission_request(&mut self, payload: &Value) {
665 let Some(request_id) = Self::permission_request_id(payload) else {
666 return;
667 };
668 if self
669 .pending_permissions
670 .iter()
671 .any(|pending| pending.request_id == request_id)
672 {
673 return;
674 }
675 self.pending_permissions.push(PendingPermission {
676 request_id: request_id.to_string(),
677 deadline: tokio::time::Instant::now() + self.permission_timeout,
678 });
679 }
680
681 async fn write_permission_response(&mut self, request_id: &str, body: Value) -> Result<()> {
683 self.transport
684 .write(json!({
685 "type": "control_response",
686 "response": {
687 "subtype": "success",
688 "request_id": request_id,
689 "response": body,
690 },
691 }))
692 .await
693 }
694
695 async fn deny_expired_permissions(&mut self) -> Result<()> {
700 let now = tokio::time::Instant::now();
701 let expired = self
702 .pending_permissions
703 .iter()
704 .filter(|pending| pending.deadline <= now)
705 .map(|pending| pending.request_id.clone())
706 .collect::<Vec<_>>();
707 self.pending_permissions
708 .retain(|pending| pending.deadline > now);
709 for request_id in expired {
710 self.write_permission_response(
711 &request_id,
712 json!({"behavior": "deny", "message": CLAUDE_PERMISSION_TIMEOUT_MESSAGE}),
713 )
714 .await?;
715 }
716 Ok(())
717 }
718
719 async fn transport_ended(&mut self) -> Result<Option<HarnessEvent>> {
736 if self.is_relay() || !self.spoke {
739 return Ok(None);
740 }
741 std::future::pending().await
742 }
743
744 fn next_permission_deadline(&self) -> Option<Duration> {
746 let now = tokio::time::Instant::now();
747 self.pending_permissions
748 .iter()
749 .map(|pending| pending.deadline.saturating_duration_since(now))
750 .min()
751 }
752}
753
754fn claude_permission_result(response: Value) -> Result<Value> {
764 let Value::Object(mut body) = response else {
765 return Err(claude_permission_shape_error(&response));
766 };
767 match body.get("behavior").and_then(Value::as_str) {
768 Some("allow") => {}
769 Some("deny") => {
770 let empty = body
772 .get("message")
773 .and_then(Value::as_str)
774 .is_none_or(str::is_empty);
775 if empty {
776 body.insert(
777 "message".into(),
778 Value::String("Volter Harness denied this permission request".into()),
779 );
780 }
781 }
782 _ => return Err(claude_permission_shape_error(&Value::Object(body))),
783 }
784 Ok(Value::Object(body))
785}
786
787fn claude_permission_shape_error(response: &Value) -> Error {
788 Error::Other(format!(
789 "Claude Code permission answers must carry a `behavior` of {}; got {response}",
790 CLAUDE_PERMISSION_BEHAVIORS
791 .map(|behavior| format!("`{behavior}`"))
792 .join(" or "),
793 ))
794}
795
796#[async_trait]
797impl RuntimeConnection for ClaudeRuntimeConnection {
798 fn handle(&self) -> &RuntimeHandle {
799 &self.handle
800 }
801
802 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
803 let content = if input.image_urls.is_empty() {
804 Value::String(input.text)
805 } else {
806 let mut parts = Vec::new();
807 if !input.text.is_empty() {
808 parts.push(json!({"type":"text", "text":input.text}));
809 }
810 for url in input.image_urls {
811 parts.push(claude_image_part(&url)?);
812 }
813 Value::Array(parts)
814 };
815 self.write_turn(json!({
816 "type": "user",
817 "session_id": self.handle.runtime_id,
818 "message": {"role": "user", "content": content},
819 }))
820 .await?;
821 Ok(None)
822 }
823
824 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
825 if let Some(payload) = self.buffered_events.pop_front() {
826 self.spoke = true;
827 return Ok(Some(harness_event(payload)));
828 }
829 loop {
830 self.deny_expired_permissions().await?;
834 if self.is_relay() && self.process_ended().await {
835 return Ok(None);
836 }
837 let payload = match self
838 .next_permission_deadline()
839 .map(|remaining| remaining.min(Duration::from_millis(200)))
840 {
841 Some(remaining) => {
842 match tokio::time::timeout(remaining, self.transport.receiver.recv()).await {
843 Err(_) => continue,
844 Ok(None) => return self.transport_ended().await,
845 Ok(Some(payload)) => payload,
846 }
847 }
848 None => match tokio::time::timeout(
849 Duration::from_millis(200),
850 self.transport.receiver.recv(),
851 )
852 .await
853 {
854 Err(_) => continue,
855 Ok(None) => return self.transport_ended().await,
856 Ok(Some(payload)) => payload,
857 },
858 };
859 self.spoke = true;
860 if Self::is_control_response(&payload) {
861 continue;
862 }
863 if let Some(request_id) = self.mounted_tool_permission_request(&payload) {
864 self.write_permission_response(&request_id, json!({"behavior": "allow"}))
865 .await?;
866 continue;
867 }
868 self.note_permission_request(&payload);
869 return Ok(Some(harness_event(payload)));
870 }
871 }
872
873 async fn interrupt(&mut self) -> Result<()> {
882 if self.process_ended().await {
886 return Ok(());
887 }
888 let request_id = format!(
889 "supercode-{}-interrupt-{}",
890 self.handle.runtime_id, self.next_control_request
891 );
892 self.next_control_request += 1;
893 self.transport
894 .write(json!({
895 "type": "control_request",
896 "request_id": request_id,
897 "request": {"subtype": "interrupt"},
898 }))
899 .await?;
900
901 let deadline = tokio::time::Instant::now() + self.control_timeout;
902 loop {
903 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
904 if remaining.is_zero() {
905 return Err(claude_interrupt_timeout(self.control_timeout));
906 }
907 match tokio::time::timeout(remaining, self.transport.receiver.recv()).await {
908 Err(_) => return Err(claude_interrupt_timeout(self.control_timeout)),
909 Ok(None) => return Err(Error::Other(
910 "Claude Code stream-json transport closed before acknowledging the interrupt"
911 .into(),
912 )),
913 Ok(Some(payload)) => {
914 if Self::is_control_response(&payload) {
915 if let Some(result) = Self::control_result(&payload, &request_id) {
916 return result;
917 }
918 continue;
919 }
920 if let Some(request_id) = self.mounted_tool_permission_request(&payload) {
924 self.write_permission_response(&request_id, json!({"behavior": "allow"}))
925 .await?;
926 continue;
927 }
928 self.note_permission_request(&payload);
929 self.buffered_events.push_back(payload);
930 }
931 }
932 }
933 }
934
935 async fn steer(&mut self, text: String) -> Result<()> {
936 self.send_input(RuntimeInput {
937 text,
938 image_urls: Vec::new(),
939 })
940 .await
941 .map(|_| ())
942 }
943
944 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
953 let Some(request_id) = request_id.as_str().map(str::to_string) else {
954 return Err(Error::Other(format!(
955 "Claude Code control requests are identified by a string `request_id`; got \
956 {request_id}"
957 )));
958 };
959 let Some(index) = self
960 .pending_permissions
961 .iter()
962 .position(|pending| pending.request_id == request_id)
963 else {
964 return Err(Error::Other(format!(
965 "no Claude Code permission request `{request_id}` is waiting on this connection — \
966 a `can_use_tool` request is answerable only while its turn is blocked on it, and \
967 only until it is answered or denied on timeout"
968 )));
969 };
970 let body = claude_permission_result(response)?;
971 self.pending_permissions.remove(index);
972 self.write_permission_response(&request_id, body).await
973 }
974
975 async fn close(&mut self) -> Result<()> {
976 self.transport.close().await
977 }
978}
979
980#[derive(Debug, Clone)]
982pub struct AcpRuntimeBackend {
983 harness: HarnessId,
984 launch: RuntimeLaunch,
985 resume_session: bool,
986}
987
988impl AcpRuntimeBackend {
989 pub fn new(harness: HarnessId, launch: RuntimeLaunch) -> Self {
991 Self {
992 harness,
993 launch,
994 resume_session: false,
995 }
996 }
997
998 pub fn with_resume_support(mut self, supported: bool) -> Self {
1002 self.resume_session = supported;
1003 self
1004 }
1005
1006 async fn connect(
1007 &self,
1008 cwd: &Path,
1009 launch: Option<RuntimeLaunch>,
1010 ) -> Result<(
1011 Arc<JsonLineClient>,
1012 mpsc::UnboundedReceiver<Value>,
1013 RuntimeEndpoint,
1014 Value,
1015 )> {
1016 let launch = launch.unwrap_or_else(|| self.launch.clone());
1017 let (client, receiver, endpoint) =
1018 JsonLineClient::spawn(&launch, Some(cwd), true, "acp-v1-jsonrpc").await?;
1019 let initialized = client
1020 .request(
1021 "initialize",
1022 json!({
1023 "protocolVersion": 1,
1024 "clientCapabilities": {},
1025 "clientInfo": {
1026 "name": "supercode",
1027 "title": "Volter Harness",
1028 "version": env!("CARGO_PKG_VERSION"),
1029 },
1030 }),
1031 )
1032 .await?;
1033 if initialized.get("protocolVersion").and_then(Value::as_u64) != Some(1) {
1034 return Err(Error::Other(format!(
1035 "ACP agent negotiated unsupported protocol version: {}",
1036 initialized
1037 .get("protocolVersion")
1038 .cloned()
1039 .unwrap_or(Value::Null)
1040 )));
1041 }
1042 Ok((client, receiver, endpoint, initialized))
1043 }
1044
1045 async fn session_request(
1046 &self,
1047 client: &JsonLineClient,
1048 initialized: &Value,
1049 method: &str,
1050 params: Value,
1051 ) -> Result<Value> {
1052 match client.request(method, params.clone()).await {
1053 Ok(response) => Ok(response),
1054 Err(error) if acp_auth_required(&error.to_string()) => {
1055 let cached = initialized
1056 .get("authMethods")
1057 .and_then(Value::as_array)
1058 .and_then(|methods| {
1059 methods.iter().find_map(|candidate| {
1060 (candidate.get("id").and_then(Value::as_str) == Some("cached_token"))
1061 .then_some("cached_token")
1062 })
1063 });
1064 let Some(method_id) = cached else {
1065 return Err(Error::Other(
1066 "ACP agent requires authentication but did not advertise the non-interactive `cached_token` method"
1067 .into(),
1068 ));
1069 };
1070 client
1071 .request(
1072 "authenticate",
1073 json!({"methodId": method_id, "_meta": {"headless": true}}),
1074 )
1075 .await?;
1076 client.request(method, params).await
1077 }
1078 Err(error) => Err(error),
1079 }
1080 }
1081
1082 async fn connection(
1083 &self,
1084 cwd: &Path,
1085 runtime_id: Option<String>,
1086 launch: Option<RuntimeLaunch>,
1087 mcp_servers: Vec<McpServerLaunch>,
1088 ) -> Result<Box<dyn RuntimeConnection>> {
1089 let mcp_servers = acp_mcp_servers(&mcp_servers);
1090 let (launch, runtime_id) = match runtime_id {
1097 Some(session_id) if self.harness.as_str() == HarnessId::SUPERCODE => {
1098 let mut continuation = launch.unwrap_or_else(|| self.launch.clone());
1099 let flags = continuation
1100 .arguments
1101 .iter()
1102 .skip_while(|argument| argument.as_str() != "acp")
1103 .skip(1)
1104 .cloned()
1105 .collect::<Vec<_>>();
1106 continuation.arguments = ["resume", &session_id, "--harness", "supercode", "--acp"]
1107 .into_iter()
1108 .map(String::from)
1109 .chain(flags)
1110 .collect();
1111 (Some(continuation), None)
1112 }
1113 other => (launch, other),
1114 };
1115 let (client, mut receiver, endpoint, initialized) = self.connect(cwd, launch).await?;
1116 let session_id = if let Some(session_id) = runtime_id {
1117 let resume = initialized
1118 .pointer("/agentCapabilities/sessionCapabilities/resume")
1119 .is_some();
1120 let load = initialized
1121 .pointer("/agentCapabilities/loadSession")
1122 .and_then(Value::as_bool)
1123 .unwrap_or(false);
1124 let method = if resume {
1125 "session/resume"
1126 } else if load {
1127 "session/load"
1128 } else {
1129 return Err(Error::Other(
1130 "ACP agent did not advertise session resume or load".into(),
1131 ));
1132 };
1133 self.session_request(
1134 client.as_ref(),
1135 &initialized,
1136 method,
1137 json!({"sessionId": session_id, "cwd": cwd, "mcpServers": mcp_servers}),
1138 )
1139 .await?;
1140 session_id
1141 } else {
1142 self.session_request(
1143 client.as_ref(),
1144 &initialized,
1145 "session/new",
1146 json!({"cwd": cwd, "mcpServers": mcp_servers}),
1147 )
1148 .await?
1149 .get("sessionId")
1150 .and_then(Value::as_str)
1151 .ok_or_else(|| Error::Other("ACP session/new omitted sessionId".into()))?
1152 .to_string()
1153 };
1154 while receiver.try_recv().is_ok() {}
1163 Ok(Box::new(AcpRuntimeConnection {
1164 handle: RuntimeHandle {
1165 harness: self.harness.clone(),
1166 runtime_id: session_id,
1167 endpoint,
1168 },
1169 client,
1170 receiver,
1171 active_prompt: None,
1172 }))
1173 }
1174}
1175
1176fn acp_mcp_servers(servers: &[McpServerLaunch]) -> Value {
1181 Value::Array(
1182 servers
1183 .iter()
1184 .map(|server| {
1185 json!({
1186 "name": server.name,
1187 "command": server.command,
1188 "args": server.arguments,
1189 "env": server
1190 .env
1191 .iter()
1192 .map(|(name, value)| json!({"name": name, "value": value}))
1193 .collect::<Vec<_>>(),
1194 })
1195 })
1196 .collect::<Vec<_>>(),
1197 )
1198}
1199
1200fn acp_auth_required(message: &str) -> bool {
1201 let message = message.to_ascii_lowercase();
1202 [
1203 "auth",
1204 "login",
1205 "sign in",
1206 "sign-in",
1207 "unauthorized",
1208 "forbidden",
1209 "credential",
1210 ]
1211 .iter()
1212 .any(|needle| message.contains(needle))
1213}
1214
1215#[async_trait]
1216impl RuntimeBackend for AcpRuntimeBackend {
1217 fn harness(&self) -> HarnessId {
1218 self.harness.clone()
1219 }
1220
1221 fn capabilities(&self) -> RuntimeCapabilities {
1222 RuntimeCapabilities {
1223 start_session: true,
1224 resume_session: self.resume_session,
1227 attach_existing_process: false,
1228 send_input: true,
1229 stream_events: true,
1230 interrupt: true,
1231 steer: false,
1232 respond_to_requests: true,
1233 }
1234 }
1235
1236 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
1237 self.connection(&request.cwd, None, request.launch, request.mcp_servers)
1238 .await
1239 }
1240
1241 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
1242 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1243 self.connection(
1244 &cwd,
1245 Some(request.runtime_id),
1246 request.launch,
1247 request.mcp_servers,
1248 )
1249 .await
1250 }
1251}
1252
1253struct AcpRuntimeConnection {
1254 handle: RuntimeHandle,
1255 client: Arc<JsonLineClient>,
1256 receiver: mpsc::UnboundedReceiver<Value>,
1257 active_prompt: Option<u64>,
1258}
1259
1260#[async_trait]
1261impl RuntimeConnection for AcpRuntimeConnection {
1262 fn handle(&self) -> &RuntimeHandle {
1263 &self.handle
1264 }
1265
1266 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
1267 let mut prompt = Vec::new();
1268 if !input.text.is_empty() {
1269 prompt.push(json!({"type": "text", "text": input.text}));
1270 }
1271 for url in input.image_urls {
1272 let (mime_type, data) = data_image_parts(&url).ok_or_else(|| {
1273 Error::Other("ACP image prompts require base64 image data URLs".into())
1274 })?;
1275 prompt.push(json!({"type":"image", "mimeType":mime_type, "data":data}));
1276 }
1277 let (id, response) = self
1278 .client
1279 .begin_request(
1280 "session/prompt",
1281 json!({
1282 "sessionId": self.handle.runtime_id,
1283 "prompt": prompt,
1284 }),
1285 )
1286 .await?;
1287 self.active_prompt = Some(id);
1288 let client = self.client.clone();
1289 tokio::spawn(async move {
1290 let result = match response.await {
1291 Ok(Ok(result)) => json!({"id": id, "result": result}),
1292 Ok(Err(error)) => json!({"id": id, "error": error}),
1293 Err(_) => json!({"id": id, "error": "response channel closed"}),
1294 };
1295 client.emit(json!({
1296 "jsonrpc": "2.0",
1297 "method": "supercode/acp_request_completed",
1298 "params": result,
1299 }));
1300 });
1301 Ok(Some(id.to_string()))
1302 }
1303
1304 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
1305 let Some(payload) = self.receiver.recv().await else {
1306 return Ok(None);
1307 };
1308 let kind = payload
1309 .get("method")
1310 .and_then(Value::as_str)
1311 .or_else(|| payload.get("type").and_then(Value::as_str))
1312 .unwrap_or("protocol")
1313 .to_string();
1314 if kind == "supercode/acp_request_completed" {
1315 self.active_prompt = None;
1316 }
1317 Ok(Some(HarnessEvent {
1318 sequence: None,
1319 kind,
1320 payload,
1321 }))
1322 }
1323
1324 async fn interrupt(&mut self) -> Result<()> {
1325 self.client
1326 .notify(
1327 "session/cancel",
1328 json!({"sessionId": self.handle.runtime_id}),
1329 )
1330 .await
1331 }
1332
1333 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
1334 self.client.respond(request_id, response).await
1335 }
1336
1337 async fn close(&mut self) -> Result<()> {
1338 self.client.close().await
1339 }
1340}
1341
1342#[derive(Debug, Clone)]
1346pub struct OpenCodeRuntimeBackend {
1347 launch: RuntimeLaunch,
1348 base_url: Option<String>,
1349 bearer: Option<BearerToken>,
1350}
1351
1352impl Default for OpenCodeRuntimeBackend {
1353 fn default() -> Self {
1354 Self::new()
1355 }
1356}
1357
1358impl OpenCodeRuntimeBackend {
1359 pub fn new() -> Self {
1361 Self {
1362 launch: RuntimeLaunch {
1363 program: "opencode".into(),
1364 arguments: vec!["serve".into()],
1365 env: BTreeMap::new(),
1366 },
1367 base_url: None,
1368 bearer: None,
1369 }
1370 }
1371
1372 pub fn connect(base_url: impl Into<String>) -> Self {
1375 Self {
1376 base_url: Some(base_url.into().trim_end_matches('/').to_string()),
1377 ..Self::new()
1378 }
1379 }
1380
1381 pub fn with_launch(mut self, launch: RuntimeLaunch) -> Self {
1383 self.launch = launch;
1384 self
1385 }
1386
1387 pub fn with_bearer(mut self, token: BearerToken) -> Self {
1390 self.bearer = Some(token);
1391 self
1392 }
1393
1394 fn http_client(&self) -> Result<reqwest::Client> {
1395 let Some(token) = &self.bearer else {
1396 return Ok(reqwest::Client::new());
1397 };
1398 let mut headers = reqwest::header::HeaderMap::new();
1399 let mut value =
1400 reqwest::header::HeaderValue::from_str(&format!("Bearer {}", token.secret())).map_err(
1401 |_| Error::Other("connect-mode bearer token is not a valid header value".into()),
1402 )?;
1403 value.set_sensitive(true);
1404 headers.insert(reqwest::header::AUTHORIZATION, value);
1405 reqwest::Client::builder()
1406 .default_headers(headers)
1407 .build()
1408 .map_err(|error| Error::Other(format!("could not build HTTP client: {error}")))
1409 }
1410
1411 async fn service(
1412 &self,
1413 client: &reqwest::Client,
1414 launch: Option<RuntimeLaunch>,
1415 ) -> Result<(String, Option<super::GroupLeader>)> {
1416 if let Some(base_url) = &self.base_url {
1417 wait_for_health(client, base_url).await?;
1418 return Ok((base_url.clone(), None));
1419 }
1420 let port = TcpListener::bind(("127.0.0.1", 0))?.local_addr()?.port();
1421 let mut launch = launch.unwrap_or_else(|| self.launch.clone());
1422 launch.arguments.extend([
1423 "--hostname".into(),
1424 "127.0.0.1".into(),
1425 "--port".into(),
1426 port.to_string(),
1427 ]);
1428 let mut command = Command::new(&launch.program);
1429 command
1430 .args(&launch.arguments)
1431 .envs(&launch.env)
1432 .stdin(Stdio::null())
1433 .stdout(Stdio::null())
1434 .stderr(Stdio::inherit())
1435 .kill_on_drop(true);
1436 #[cfg(unix)]
1440 command.process_group(0);
1441 let mut child = super::GroupLeader(command.spawn().map_err(|error| {
1446 Error::Other(format!("could not launch {}: {error}", launch.program))
1447 })?);
1448 let base_url = format!("http://127.0.0.1:{port}");
1449 if let Err(error) = wait_for_health(client, &base_url).await {
1450 let _ = terminate_opencode_server(&mut child).await;
1451 return Err(error);
1452 }
1453 Ok((base_url, Some(child)))
1454 }
1455
1456 async fn open(
1457 &self,
1458 cwd: &Path,
1459 runtime_id: Option<String>,
1460 launch: Option<RuntimeLaunch>,
1461 ) -> Result<Box<dyn RuntimeConnection>> {
1462 let client = self.http_client()?;
1463 let (base_url, child) = self.service(&client, launch).await?;
1464 let cwd_string = cwd.to_string_lossy().to_string();
1465 let runtime_id = match runtime_id {
1466 Some(id) => {
1467 http_ok(
1468 client
1469 .get(format!("{base_url}/session/{id}"))
1470 .query(&[("directory", &cwd_string)])
1471 .send()
1472 .await,
1473 )
1474 .await?;
1475 id
1476 }
1477 None => {
1478 let response = http_ok(
1479 client
1480 .post(format!("{base_url}/session"))
1481 .query(&[("directory", &cwd_string)])
1482 .json(&json!({}))
1483 .send()
1484 .await,
1485 )
1486 .await?;
1487 response
1488 .json::<Value>()
1489 .await
1490 .map_err(http_error)?
1491 .get("id")
1492 .and_then(Value::as_str)
1493 .ok_or_else(|| Error::Other("OpenCode create session omitted id".into()))?
1494 .to_string()
1495 }
1496 };
1497 let receiver = spawn_sse(
1498 client.clone(),
1499 format!("{base_url}/event"),
1500 cwd_string.clone(),
1501 );
1502 Ok(Box::new(OpenCodeRuntimeConnection {
1503 handle: RuntimeHandle {
1504 harness: HarnessId::from(HarnessId::OPENCODE),
1505 runtime_id,
1506 endpoint: RuntimeEndpoint::Http {
1507 base_url: base_url.clone(),
1508 protocol: "opencode-http-sse".into(),
1509 },
1510 },
1511 base_url,
1512 cwd: cwd_string,
1513 client,
1514 receiver,
1515 child,
1516 }))
1517 }
1518}
1519
1520#[async_trait]
1521impl RuntimeBackend for OpenCodeRuntimeBackend {
1522 fn harness(&self) -> HarnessId {
1523 HarnessId::from(HarnessId::OPENCODE)
1524 }
1525
1526 fn capabilities(&self) -> RuntimeCapabilities {
1527 RuntimeCapabilities {
1528 start_session: true,
1529 resume_session: true,
1530 attach_existing_process: self.base_url.is_some(),
1531 send_input: true,
1532 stream_events: true,
1533 interrupt: true,
1534 steer: false,
1535 respond_to_requests: true,
1536 }
1537 }
1538
1539 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
1540 self.open(&request.cwd, None, request.launch).await
1541 }
1542
1543 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
1544 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1545 self.open(&cwd, Some(request.runtime_id), request.launch)
1546 .await
1547 }
1548
1549 async fn attach_existing(
1550 &self,
1551 request: RuntimeAttachRequest,
1552 ) -> Result<Box<dyn RuntimeConnection>> {
1553 if self.base_url.is_none() {
1554 return Err(Error::Other(
1555 "OpenCode live attach requires the existing server's `base_url`".into(),
1556 ));
1557 }
1558 let cwd = request.cwd.unwrap_or(std::env::current_dir()?);
1559 self.open(&cwd, Some(request.runtime_id), request.launch)
1560 .await
1561 }
1562}
1563
1564struct OpenCodeRuntimeConnection {
1565 handle: RuntimeHandle,
1566 base_url: String,
1567 cwd: String,
1568 client: reqwest::Client,
1569 receiver: mpsc::UnboundedReceiver<Value>,
1570 child: Option<super::GroupLeader>,
1571}
1572
1573#[async_trait]
1574impl RuntimeConnection for OpenCodeRuntimeConnection {
1575 fn handle(&self) -> &RuntimeHandle {
1576 &self.handle
1577 }
1578
1579 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
1580 let mut parts = Vec::new();
1581 if !input.text.is_empty() {
1582 parts.push(json!({"type": "text", "text": input.text}));
1583 }
1584 for url in input.image_urls {
1585 let mime = image_mime_type(&url).ok_or_else(|| {
1586 Error::Other("OpenCode image prompts require a recognizable image MIME type".into())
1587 })?;
1588 parts.push(json!({"type":"file", "mime":mime, "url":url}));
1589 }
1590 http_ok(
1591 self.client
1592 .post(format!(
1593 "{}/session/{}/prompt_async",
1594 self.base_url, self.handle.runtime_id
1595 ))
1596 .query(&[("directory", &self.cwd)])
1597 .json(&json!({"parts": parts}))
1598 .send()
1599 .await,
1600 )
1601 .await?;
1602 Ok(None)
1603 }
1604
1605 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
1606 loop {
1607 let Some(payload) = self.receiver.recv().await else {
1608 return Ok(None);
1609 };
1610 if opencode_event_session_id(&payload)
1611 .is_some_and(|session_id| session_id != self.handle.runtime_id)
1612 {
1613 continue;
1614 }
1615 let kind = payload
1616 .get("type")
1617 .and_then(Value::as_str)
1618 .unwrap_or("event")
1619 .to_string();
1620 return Ok(Some(HarnessEvent {
1621 sequence: None,
1622 kind,
1623 payload,
1624 }));
1625 }
1626 }
1627
1628 async fn interrupt(&mut self) -> Result<()> {
1629 http_ok(
1630 self.client
1631 .post(format!(
1632 "{}/session/{}/abort",
1633 self.base_url, self.handle.runtime_id
1634 ))
1635 .query(&[("directory", &self.cwd)])
1636 .send()
1637 .await,
1638 )
1639 .await?;
1640 Ok(())
1641 }
1642
1643 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
1644 let permission = request_id.as_str().ok_or_else(|| {
1645 Error::Other("OpenCode permission request id must be a string".into())
1646 })?;
1647 http_ok(
1648 self.client
1649 .post(format!(
1650 "{}/session/{}/permissions/{permission}",
1651 self.base_url, self.handle.runtime_id
1652 ))
1653 .query(&[("directory", &self.cwd)])
1654 .json(&response)
1655 .send()
1656 .await,
1657 )
1658 .await?;
1659 Ok(())
1660 }
1661
1662 async fn close(&mut self) -> Result<()> {
1663 if let Some(child) = &mut self.child {
1664 terminate_opencode_server(child).await?;
1665 }
1666 Ok(())
1667 }
1668}
1669
1670fn data_image_parts(url: &str) -> Option<(&str, &str)> {
1671 let rest = url.strip_prefix("data:")?;
1672 let (mime_type, data) = rest.split_once(";base64,")?;
1673 mime_type.starts_with("image/").then_some((mime_type, data))
1674}
1675
1676fn image_mime_type(url: &str) -> Option<&str> {
1677 if let Some((mime_type, _)) = data_image_parts(url) {
1678 return Some(mime_type);
1679 }
1680 let path = url.split(['?', '#']).next()?.to_ascii_lowercase();
1681 if path.ends_with(".png") {
1682 Some("image/png")
1683 } else if path.ends_with(".jpg") || path.ends_with(".jpeg") {
1684 Some("image/jpeg")
1685 } else if path.ends_with(".gif") {
1686 Some("image/gif")
1687 } else if path.ends_with(".webp") {
1688 Some("image/webp")
1689 } else {
1690 None
1691 }
1692}
1693
1694fn claude_image_part(url: &str) -> Result<Value> {
1695 if let Some((media_type, data)) = data_image_parts(url) {
1696 return Ok(json!({
1697 "type":"image",
1698 "source":{"type":"base64", "media_type":media_type, "data":data}
1699 }));
1700 }
1701 if url.starts_with("https://") || url.starts_with("http://") {
1702 return Ok(json!({"type":"image", "source":{"type":"url", "url":url}}));
1703 }
1704 Err(Error::Other(
1705 "Claude image prompts require image data URLs or HTTP(S) URLs".into(),
1706 ))
1707}
1708
1709fn opencode_event_session_id(payload: &Value) -> Option<&str> {
1710 let properties = payload.get("properties").unwrap_or(payload);
1711 properties
1712 .get("sessionID")
1713 .and_then(Value::as_str)
1714 .or_else(|| {
1715 properties
1716 .get("part")
1717 .and_then(|part| part.get("sessionID"))
1718 .and_then(Value::as_str)
1719 })
1720 .or_else(|| {
1721 properties
1722 .get("info")
1723 .and_then(|info| info.get("sessionID"))
1724 .and_then(Value::as_str)
1725 })
1726}
1727
1728async fn terminate_opencode_server(child: &mut Child) -> Result<()> {
1729 #[cfg(unix)]
1730 let process_group = child.id();
1731 let leader_exited = child.try_wait()?.is_some();
1732 if leader_exited {
1733 #[cfg(unix)]
1734 if let Some(pid) = process_group.filter(|pid| process_group_exists(*pid)) {
1735 crate::lsp::kill_process_group(pid);
1736 wait_for_process_group_exit(pid, Duration::from_secs(3)).await?;
1737 }
1738 return Ok(());
1739 }
1740 #[cfg(unix)]
1746 if let Some(pid) = process_group {
1747 unsafe {
1748 libc::kill(-(pid as libc::pid_t), libc::SIGTERM);
1749 }
1750 let mut leader_reaped = false;
1751 if let Ok(status) = tokio::time::timeout(Duration::from_millis(500), child.wait()).await {
1752 status?;
1753 leader_reaped = true;
1754 if !process_group_exists(pid) {
1755 return Ok(());
1756 }
1757 }
1758 crate::lsp::kill_process_group(pid);
1761 if leader_reaped {
1762 return wait_for_process_group_exit(pid, Duration::from_secs(3)).await;
1763 }
1764 }
1765 #[cfg(not(unix))]
1766 child.start_kill()?;
1767 tokio::time::timeout(Duration::from_secs(3), child.wait())
1768 .await
1769 .map_err(|_| Error::Other("timed out reaping the OpenCode server".into()))??;
1770 #[cfg(unix)]
1771 if let Some(pid) = process_group {
1772 wait_for_process_group_exit(pid, Duration::from_secs(3)).await?;
1773 }
1774 Ok(())
1775}
1776
1777#[cfg(unix)]
1778fn process_group_exists(pid: u32) -> bool {
1779 let result = unsafe { libc::kill(-(pid as libc::pid_t), 0) };
1780 result == 0 || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
1781}
1782
1783#[cfg(unix)]
1784async fn wait_for_process_group_exit(pid: u32, timeout: Duration) -> Result<()> {
1785 let deadline = tokio::time::Instant::now() + timeout;
1786 while process_group_exists(pid) {
1787 if tokio::time::Instant::now() >= deadline {
1788 return Err(Error::Other(format!(
1789 "timed out stopping OpenCode process group {pid}"
1790 )));
1791 }
1792 tokio::time::sleep(Duration::from_millis(10)).await;
1793 }
1794 Ok(())
1795}
1796
1797struct RawLineTransport {
1798 stdin: Mutex<ChildStdin>,
1799 child: Mutex<super::GroupLeader>,
1800 receiver: mpsc::UnboundedReceiver<Value>,
1801 endpoint: RuntimeEndpoint,
1802}
1803
1804impl RawLineTransport {
1805 async fn spawn(launch: &RuntimeLaunch, cwd: Option<&Path>, protocol: &str) -> Result<Self> {
1806 let mut command = Command::new(&launch.program);
1807 command
1808 .args(&launch.arguments)
1809 .envs(&launch.env)
1810 .stdin(Stdio::piped())
1811 .stdout(Stdio::piped())
1812 .stderr(Stdio::inherit())
1813 .kill_on_drop(true);
1814 #[cfg(unix)]
1818 command.process_group(0);
1819 if let Some(cwd) = cwd {
1820 command.current_dir(cwd);
1821 }
1822 let mut child = command.spawn().map_err(|error| {
1823 Error::Other(format!("could not launch {}: {error}", launch.program))
1824 })?;
1825 let pid = child.id();
1826 let stdin = child
1827 .stdin
1828 .take()
1829 .ok_or_else(|| Error::Other("runtime child has no stdin".into()))?;
1830 let stdout = child
1831 .stdout
1832 .take()
1833 .ok_or_else(|| Error::Other("runtime child has no stdout".into()))?;
1834 let (sender, receiver) = mpsc::unbounded_channel();
1835 tokio::spawn(async move {
1836 let mut lines = BufReader::new(stdout).lines();
1837 while let Ok(Some(line)) = lines.next_line().await {
1838 let value = serde_json::from_str(&line)
1839 .unwrap_or_else(|_| json!({"type": "malformed_output", "line": line}));
1840 let _ = sender.send(value);
1841 }
1842 });
1843 Ok(Self {
1844 stdin: Mutex::new(stdin),
1845 child: Mutex::new(super::GroupLeader(child)),
1846 receiver,
1847 endpoint: RuntimeEndpoint::LocalProcess {
1848 pid,
1849 command: std::iter::once(launch.program.clone())
1850 .chain(launch.arguments.iter().cloned())
1851 .collect(),
1852 protocol: protocol.into(),
1853 },
1854 })
1855 }
1856
1857 async fn write(&self, value: Value) -> Result<()> {
1858 let mut stdin = self.stdin.lock().await;
1859 stdin.write_all(value.to_string().as_bytes()).await?;
1860 stdin.write_all(b"\n").await?;
1861 stdin.flush().await?;
1862 Ok(())
1863 }
1864
1865 async fn close(&self) -> Result<()> {
1866 let mut child = self.child.lock().await;
1867 if child.try_wait()?.is_some() {
1868 return Ok(());
1869 }
1870 #[cfg(unix)]
1873 if let Some(pid) = child.id() {
1874 crate::lsp::kill_process_group(pid);
1875 tokio::time::timeout(Duration::from_secs(3), child.wait())
1876 .await
1877 .map_err(|_| Error::Other("timed out reaping runtime process group".into()))??;
1878 return Ok(());
1879 }
1880 #[cfg(not(unix))]
1881 child.kill().await?;
1882 Ok(())
1883 }
1884}
1885
1886async fn raw_next_event(
1887 receiver: &mut mpsc::UnboundedReceiver<Value>,
1888) -> Result<Option<HarnessEvent>> {
1889 let Some(payload) = receiver.recv().await else {
1890 return Ok(None);
1891 };
1892 Ok(Some(harness_event(payload)))
1893}
1894
1895fn harness_event(payload: Value) -> HarnessEvent {
1896 let kind = payload
1897 .get("type")
1898 .and_then(Value::as_str)
1899 .unwrap_or("event")
1900 .to_string();
1901 HarnessEvent {
1902 sequence: None,
1903 kind,
1904 payload,
1905 }
1906}
1907
1908fn broken_pipe(error: &Error) -> bool {
1915 matches!(error, Error::Io(io) if io.kind() == std::io::ErrorKind::BrokenPipe)
1916}
1917
1918fn claude_interrupt_timeout(bound: Duration) -> Error {
1919 Error::Other(format!(
1920 "Claude Code did not acknowledge the interrupt control request within {}s",
1921 bound.as_secs_f32()
1922 ))
1923}
1924
1925pub(crate) fn generated_session_id() -> String {
1926 let mut bytes = [0_u8; 16];
1927 if getrandom::getrandom(&mut bytes).is_err() {
1928 let nanos = SystemTime::now()
1929 .duration_since(UNIX_EPOCH)
1930 .unwrap_or_default()
1931 .as_nanos()
1932 .to_le_bytes();
1933 bytes.copy_from_slice(&nanos);
1934 }
1935 bytes[6] = (bytes[6] & 0x0f) | 0x40;
1936 bytes[8] = (bytes[8] & 0x3f) | 0x80;
1937 format!(
1938 "{:02x}{:02x}{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}{:02x}{:02x}{:02x}{:02x}",
1939 bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5], bytes[6], bytes[7],
1940 bytes[8], bytes[9], bytes[10], bytes[11], bytes[12], bytes[13], bytes[14], bytes[15]
1941 )
1942}
1943
1944async fn wait_for_health(client: &reqwest::Client, base_url: &str) -> Result<()> {
1945 wait_for_health_for(client, base_url, Duration::from_secs(10)).await
1946}
1947
1948async fn wait_for_health_for(
1949 client: &reqwest::Client,
1950 base_url: &str,
1951 total_timeout: Duration,
1952) -> Result<()> {
1953 let url = format!("{base_url}/global/health");
1954 let mut last = None;
1955 let deadline = tokio::time::Instant::now() + total_timeout;
1956 loop {
1961 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
1962 if remaining.is_zero() {
1963 break;
1964 }
1965 let request_timeout = remaining.min(Duration::from_millis(500));
1966 match tokio::time::timeout(request_timeout, client.get(&url).send()).await {
1967 Ok(Ok(response)) if response.status().is_success() => return Ok(()),
1968 Ok(Ok(response)) => last = Some(format!("HTTP {}", response.status())),
1969 Ok(Err(error)) => last = Some(error.to_string()),
1970 Err(_) => last = Some("health request timed out".into()),
1971 }
1972 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
1973 if !remaining.is_zero() {
1974 tokio::time::sleep(remaining.min(Duration::from_millis(100))).await;
1975 }
1976 }
1977 Err(Error::Other(format!(
1978 "OpenCode server at {base_url} did not become healthy: {}",
1979 last.unwrap_or_else(|| "no response".into())
1980 )))
1981}
1982
1983async fn http_ok(
1984 response: std::result::Result<reqwest::Response, reqwest::Error>,
1985) -> Result<reqwest::Response> {
1986 response
1987 .map_err(http_error)?
1988 .error_for_status()
1989 .map_err(http_error)
1990}
1991
1992fn http_error(error: reqwest::Error) -> Error {
1993 Error::Other(format!("runtime HTTP request failed: {error}"))
1994}
1995
1996fn spawn_sse(
1997 client: reqwest::Client,
1998 url: String,
1999 directory: String,
2000) -> mpsc::UnboundedReceiver<Value> {
2001 let (sender, receiver) = mpsc::unbounded_channel();
2002 tokio::spawn(async move {
2003 let response = client
2004 .get(url)
2005 .query(&[("directory", directory)])
2006 .send()
2007 .await;
2008 let Ok(response) = response.and_then(reqwest::Response::error_for_status) else {
2009 let _ = sender.send(
2010 json!({"type": "stream_error", "message": "could not open OpenCode SSE stream"}),
2011 );
2012 return;
2013 };
2014 let mut stream = response.bytes_stream();
2015 let mut buffer = String::new();
2016 while let Some(chunk) = stream.next().await {
2017 let Ok(chunk) = chunk else {
2018 break;
2019 };
2020 buffer.push_str(&String::from_utf8_lossy(&chunk));
2021 while let Some(newline) = buffer.find('\n') {
2022 let line = buffer[..newline].trim_end_matches('\r').to_string();
2023 buffer.drain(..=newline);
2024 if let Some(data) = line.strip_prefix("data:") {
2025 let data = data.trim();
2026 if let Ok(value) = serde_json::from_str(data) {
2027 let _ = sender.send(value);
2028 }
2029 }
2030 }
2031 }
2032 });
2033 receiver
2034}