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