1use std::collections::{BTreeMap, HashMap};
7use std::path::{Path, PathBuf};
8use std::process::Stdio;
9use std::sync::Arc;
10use std::time::Duration;
11
12use async_trait::async_trait;
13use serde::{Deserialize, Serialize};
14use serde_json::{json, Value};
15use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
16use tokio::process::{Child, ChildStdin, Command};
17use tokio::sync::{mpsc, oneshot, Mutex};
18
19use crate::{Error, HarnessId, Result};
20
21mod adapters;
22mod hosted;
23#[cfg(feature = "adapter-api")]
24mod supercode_http;
25pub(crate) use adapters::generated_session_id;
26pub use adapters::{
27 AcpRuntimeBackend, ClaudeCodeRuntimeBackend, OpenCodeRuntimeBackend, PiRuntimeBackend,
28};
29pub use hosted::{HostedHarnessConnection, HostedHarnessRuntime};
30#[cfg(feature = "adapter-api")]
31pub use supercode_http::SupercodeHttpRuntimeBackend;
32
33#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
35pub struct RuntimeCapabilities {
36 pub start_session: bool,
38 pub resume_session: bool,
40 pub attach_existing_process: bool,
42 pub send_input: bool,
44 pub stream_events: bool,
46 pub interrupt: bool,
48 #[serde(default)]
50 pub steer: bool,
51 pub respond_to_requests: bool,
53}
54
55#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
57pub struct RuntimeLaunch {
58 pub program: String,
60 pub arguments: Vec<String>,
62 pub env: BTreeMap<String, String>,
64}
65
66#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
73pub struct RuntimeConnectLaunch {
74 pub config_path: String,
77 pub address_pointer: String,
79 #[serde(default, skip_serializing_if = "Option::is_none")]
84 pub port_pointer: Option<String>,
85 #[serde(default, skip_serializing_if = "Option::is_none")]
90 pub default_address: Option<String>,
91 #[serde(default, skip_serializing_if = "Option::is_none")]
93 pub auth_pointer: Option<String>,
94 pub protocol: String,
96}
97
98#[derive(Clone, PartialEq, Eq)]
100pub struct BearerToken(String);
101
102impl BearerToken {
103 pub fn new(secret: impl Into<String>) -> Self {
105 Self(secret.into())
106 }
107
108 pub fn secret(&self) -> &str {
110 &self.0
111 }
112}
113
114impl std::fmt::Debug for BearerToken {
115 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
116 formatter.write_str("BearerToken(<redacted>)")
117 }
118}
119
120#[derive(Debug, Clone, PartialEq, Eq)]
122pub struct ResolvedRuntimeConnection {
123 pub address: String,
125 pub auth: Option<BearerToken>,
127}
128
129impl RuntimeConnectLaunch {
130 pub fn resolve(&self, home: &Path) -> Result<ResolvedRuntimeConnection> {
135 let path = match self.config_path.strip_prefix("~/") {
136 Some(rest) => home.join(rest),
137 None => PathBuf::from(&self.config_path),
138 };
139 let config: Value = match std::fs::read_to_string(&path) {
142 Ok(raw) => serde_json::from_str(&raw).map_err(|_| {
143 Error::Other(format!(
144 "connect-mode config {} is not valid JSON",
145 path.display()
146 ))
147 })?,
148 Err(error) => {
149 if self.default_address.is_some() {
150 Value::Object(Default::default())
151 } else {
152 return Err(Error::Other(format!(
153 "connect-mode config {} is unreadable: {error}",
154 path.display()
155 )));
156 }
157 }
158 };
159 let field = |pointer: &str, name: &str| -> Result<String> {
160 match config.pointer(pointer).and_then(Value::as_str) {
161 Some(value) if !value.trim().is_empty() => Ok(value.trim().to_string()),
162 _ => Err(Error::Other(format!(
163 "connect-mode {name} pointer `{pointer}` does not name a non-empty string in {}",
164 path.display()
165 ))),
166 }
167 };
168 let address = match config
171 .pointer(&self.address_pointer)
172 .and_then(Value::as_str)
173 {
174 Some(value) if !value.trim().is_empty() => value.trim().to_string(),
175 _ => {
176 let from_port = self
177 .port_pointer
178 .as_deref()
179 .and_then(|pointer| config.pointer(pointer))
180 .and_then(Value::as_u64)
181 .map(|port| {
182 let scheme = self
183 .default_address
184 .as_deref()
185 .and_then(|address| address.split_once("://"))
186 .map(|(scheme, _)| scheme)
187 .unwrap_or("ws");
188 format!("{scheme}://127.0.0.1:{port}")
189 });
190 match from_port.or_else(|| self.default_address.clone()) {
191 Some(address) => address,
192 None => {
193 return Err(Error::Other(format!(
194 "connect-mode address pointer `{}` does not name a non-empty string in {}",
195 self.address_pointer,
196 path.display()
197 )));
198 }
199 }
200 }
201 };
202 let mut address = address.trim_end_matches('/').to_string();
203 if !address.contains("://") {
207 let scheme = self
208 .default_address
209 .as_deref()
210 .and_then(|default| default.split_once("://"))
211 .map(|(scheme, _)| scheme)
212 .unwrap_or("ws");
213 address = format!("{scheme}://{address}");
214 }
215 let auth = match &self.auth_pointer {
219 Some(pointer) => match config.pointer(pointer).and_then(Value::as_str) {
220 Some(value) if !value.trim().is_empty() => {
221 Some(BearerToken::new(value.trim().to_string()))
222 }
223 _ if self.default_address.is_some() => None,
224 _ => Some(BearerToken::new(field(pointer, "auth")?)),
225 },
226 None => None,
227 };
228 Ok(ResolvedRuntimeConnection { address, auth })
229 }
230}
231
232#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
236pub struct McpServerLaunch {
237 pub name: String,
239 pub command: String,
241 #[serde(default)]
243 pub arguments: Vec<String>,
244 #[serde(default)]
246 pub env: BTreeMap<String, String>,
247}
248
249#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
251pub struct RuntimeStartRequest {
252 pub cwd: PathBuf,
254 pub launch: Option<RuntimeLaunch>,
256 #[serde(default)]
260 pub mcp_servers: Vec<McpServerLaunch>,
261 #[serde(default, skip_serializing_if = "Option::is_none")]
265 pub approval_policy: Option<String>,
266}
267
268#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
270pub struct RuntimeAttachRequest {
271 pub runtime_id: String,
273 pub cwd: Option<PathBuf>,
275 pub launch: Option<RuntimeLaunch>,
277 #[serde(default)]
281 pub mcp_servers: Vec<McpServerLaunch>,
282 #[serde(default, skip_serializing_if = "Option::is_none")]
286 pub approval_policy: Option<String>,
287}
288
289#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
291#[serde(tag = "kind", rename_all = "snake_case")]
292pub enum RuntimeEndpoint {
293 LocalProcess {
295 pid: Option<u32>,
297 command: Vec<String>,
299 protocol: String,
301 },
302 Http {
304 base_url: String,
306 protocol: String,
308 },
309}
310
311#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
313pub struct RuntimeHandle {
314 pub harness: HarnessId,
316 pub runtime_id: String,
318 pub endpoint: RuntimeEndpoint,
320}
321
322#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
324pub struct RuntimeInput {
325 pub text: String,
327 #[serde(default, skip_serializing_if = "Vec::is_empty")]
333 pub image_urls: Vec<String>,
334}
335
336#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
338pub struct HarnessEvent {
339 #[serde(default, skip_serializing_if = "Option::is_none")]
343 pub sequence: Option<u64>,
344 pub kind: String,
346 pub payload: Value,
348}
349
350#[async_trait]
352pub trait RuntimeConnection: Send {
353 fn handle(&self) -> &RuntimeHandle;
355 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>>;
358 async fn next_event(&mut self) -> Result<Option<HarnessEvent>>;
360 async fn interrupt(&mut self) -> Result<()>;
362 async fn steer(&mut self, _text: String) -> Result<()> {
364 Err(Error::Other(
365 "this runtime cannot steer an active turn".into(),
366 ))
367 }
368 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()>;
370 async fn acquire_control(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
373 Err(Error::Other(
374 "this runtime does not expose controller leases".into(),
375 ))
376 }
377 async fn heartbeat(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
379 Err(Error::Other(
380 "this runtime does not expose controller leases".into(),
381 ))
382 }
383 async fn detach(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
385 Err(Error::Other(
386 "this runtime does not expose detachable leases".into(),
387 ))
388 }
389 async fn close(&mut self) -> Result<()>;
391}
392
393#[async_trait]
396pub trait RuntimeBackend: Send + Sync {
397 fn harness(&self) -> HarnessId;
399 fn capabilities(&self) -> RuntimeCapabilities;
401 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>>;
403 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>>;
407 async fn attach_existing(
411 &self,
412 _request: RuntimeAttachRequest,
413 ) -> Result<Box<dyn RuntimeConnection>> {
414 Err(Error::Other(format!(
415 "{} cannot attach to an already-running process",
416 self.harness().as_str()
417 )))
418 }
419}
420
421#[derive(Debug, Clone)]
424pub struct CodexRuntimeBackend {
425 launch: RuntimeLaunch,
426}
427
428const CODEX_STARTUP_TIMEOUT: Duration = Duration::from_secs(10);
429
430#[derive(Debug)]
437struct CodexRuntimeHome {
438 root: PathBuf,
439 native_home: PathBuf,
440}
441
442impl CodexRuntimeHome {
443 fn prepare(launch: &mut RuntimeLaunch, runtime_id: Option<&str>) -> Result<Self> {
444 let native_home = codex_native_home(launch)?;
445 let root = supercode_runtime_root()
446 .join("codex")
447 .join(generated_session_id());
448 std::fs::create_dir_all(&root).map_err(|error| {
449 Error::Other(format!(
450 "could not create isolated Codex runtime home {}: {error}",
451 root.display()
452 ))
453 })?;
454 set_private_directory(&root)?;
455 let root = std::fs::canonicalize(&root)?;
456
457 for entry in [
458 "auth.json",
459 "config.toml",
460 "hooks.json",
461 "models_cache.json",
462 "installation_id",
463 ".personality_migration",
464 ".sandbox_migration",
465 "cache",
466 "generated_images",
467 "mcp-oauth-locks",
468 "memories",
469 "plugins",
470 "rules",
471 "shell_snapshots",
472 "skills",
473 "thread-writer-locks",
474 ] {
475 link_runtime_resource(&native_home.join(entry), &root.join(entry))?;
476 }
477
478 if let Some(runtime_id) = runtime_id {
479 let source = find_codex_rollout(&native_home.join("sessions"), runtime_id)?
480 .ok_or_else(|| {
481 Error::Other(format!(
482 "could not find Codex rollout `{runtime_id}` below {}",
483 native_home.join("sessions").display()
484 ))
485 })?;
486 let relative = source.strip_prefix(&native_home).map_err(|_| {
487 Error::Other(format!(
488 "Codex rollout {} is outside native home {}",
489 source.display(),
490 native_home.display()
491 ))
492 })?;
493 let projected = root.join(relative);
494 if let Some(parent) = projected.parent() {
495 std::fs::create_dir_all(parent)?;
496 }
497 std::fs::hard_link(&source, &projected).map_err(|error| {
498 Error::Other(format!(
499 "could not project Codex rollout {} into isolated runtime home: {error}",
500 source.display()
501 ))
502 })?;
503 }
504
505 launch
506 .env
507 .insert("CODEX_HOME".into(), root.to_string_lossy().into_owned());
508 Ok(Self { root, native_home })
509 }
510
511 fn started_rollout_path(&self, response: &Value) -> Result<PathBuf> {
512 let path = response
513 .pointer("/thread/path")
514 .and_then(Value::as_str)
515 .map(PathBuf::from)
516 .ok_or_else(|| {
517 Error::Other("Codex thread/start response omitted thread.path".into())
518 })?;
519 let relative = path.strip_prefix(&self.root).map_err(|_| {
520 Error::Other(format!(
521 "Codex created rollout {} outside isolated runtime home {}",
522 path.display(),
523 self.root.display()
524 ))
525 })?;
526 if !relative.starts_with("sessions") {
527 return Err(Error::Other(format!(
528 "Codex created non-session rollout {}",
529 path.display()
530 )));
531 }
532 Ok(path)
533 }
534
535 async fn publish_rollout(&self, path: &Path) -> Result<()> {
536 let relative = path.strip_prefix(&self.root).map_err(|_| {
537 Error::Other(format!(
538 "Codex created rollout {} outside isolated runtime home {}",
539 path.display(),
540 self.root.display()
541 ))
542 })?;
543 let publish_deadline = tokio::time::Instant::now() + Duration::from_secs(2);
544 while !path.is_file() {
545 if tokio::time::Instant::now() >= publish_deadline {
546 return Err(Error::Other(format!(
547 "Codex did not create promised rollout {} within 2s",
548 path.display()
549 )));
550 }
551 tokio::time::sleep(Duration::from_millis(10)).await;
552 }
553 let native = self.native_home.join(relative);
554 if let Some(parent) = native.parent() {
555 std::fs::create_dir_all(parent)?;
556 }
557 std::fs::hard_link(path, &native).map_err(|error| {
558 Error::Other(format!(
559 "could not publish Codex rollout {} to native home: {error}",
560 path.display()
561 ))
562 })
563 }
564
565 fn cleanup(&self) -> Result<()> {
566 match std::fs::remove_dir_all(&self.root) {
567 Ok(()) => Ok(()),
568 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
569 Err(error) => Err(Error::Other(format!(
570 "could not clean isolated Codex runtime home {}: {error}",
571 self.root.display()
572 ))),
573 }
574 }
575}
576
577impl Drop for CodexRuntimeHome {
578 fn drop(&mut self) {
579 let _ = self.cleanup();
580 }
581}
582
583const STDERR_TAIL_LINES: usize = 20;
585const STDERR_TAIL_CHARACTERS: usize = 2_000;
586
587fn closed_reason(recent_stderr: &std::collections::VecDeque<String>) -> String {
589 if recent_stderr.is_empty() {
590 return "runtime protocol closed".into();
591 }
592 let mut tail = recent_stderr
593 .iter()
594 .map(String::as_str)
595 .collect::<Vec<_>>()
596 .join(" | ");
597 if tail.chars().count() > STDERR_TAIL_CHARACTERS {
598 tail = tail
599 .chars()
600 .take(STDERR_TAIL_CHARACTERS)
601 .collect::<String>()
602 + "…";
603 }
604 format!("runtime protocol closed: {tail}")
605}
606
607fn is_stock_codex_launch(launch: &RuntimeLaunch) -> bool {
608 launch
609 .arguments
610 .iter()
611 .any(|argument| argument == "app-server")
612 && Path::new(&launch.program)
613 .file_name()
614 .and_then(|name| name.to_str())
615 .is_some_and(|name| name == "codex" || name == "codex.exe")
616}
617
618fn codex_native_home(launch: &RuntimeLaunch) -> Result<PathBuf> {
619 launch
620 .env
621 .get("CODEX_HOME")
622 .map(PathBuf::from)
623 .or_else(|| std::env::var_os("CODEX_HOME").map(PathBuf::from))
624 .or_else(|| {
625 supercode_interchange::user_home()
626 .map(std::path::PathBuf::into_os_string)
627 .map(PathBuf::from)
628 .map(|home| home.join(".codex"))
629 })
630 .ok_or_else(|| Error::Other("Codex runtime requires CODEX_HOME or HOME".into()))
631}
632
633fn supercode_runtime_root() -> PathBuf {
634 std::env::var_os("SUPERCODE_HOME")
635 .map(PathBuf::from)
636 .or_else(|| {
637 supercode_interchange::user_home()
638 .map(std::path::PathBuf::into_os_string)
639 .map(PathBuf::from)
640 .map(|home| home.join(".supercode"))
641 })
642 .unwrap_or_else(|| std::env::temp_dir().join("supercode"))
643 .join("runtime-homes")
644}
645
646fn find_codex_rollout(root: &Path, runtime_id: &str) -> Result<Option<PathBuf>> {
647 let entries = match std::fs::read_dir(root) {
648 Ok(entries) => entries,
649 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
650 Err(error) => return Err(error.into()),
651 };
652 let expected_suffix = format!("-{runtime_id}.jsonl");
653 for entry in entries {
654 let entry = entry?;
655 let kind = entry.file_type()?;
656 if kind.is_dir() {
657 if let Some(path) = find_codex_rollout(&entry.path(), runtime_id)? {
658 return Ok(Some(path));
659 }
660 } else if kind.is_file()
661 && entry
662 .file_name()
663 .to_str()
664 .is_some_and(|name| name.ends_with(&expected_suffix))
665 {
666 return Ok(Some(entry.path()));
667 }
668 }
669 Ok(None)
670}
671
672#[cfg(unix)]
673fn link_runtime_resource(source: &Path, target: &Path) -> Result<()> {
674 use std::os::unix::fs::symlink;
675
676 if source.exists() {
677 symlink(source, target)?;
678 }
679 Ok(())
680}
681
682#[cfg(not(unix))]
683fn link_runtime_resource(source: &Path, target: &Path) -> Result<()> {
684 if source.is_file() {
685 std::fs::copy(source, target)?;
686 }
687 Ok(())
688}
689
690#[cfg(unix)]
691fn set_private_directory(path: &Path) -> Result<()> {
692 use std::os::unix::fs::PermissionsExt;
693
694 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))?;
695 Ok(())
696}
697
698#[cfg(not(unix))]
699fn set_private_directory(_path: &Path) -> Result<()> {
700 Ok(())
701}
702
703impl Default for CodexRuntimeBackend {
704 fn default() -> Self {
705 Self::new()
706 }
707}
708
709impl CodexRuntimeBackend {
710 pub fn new() -> Self {
712 Self {
713 launch: RuntimeLaunch {
714 program: "codex".into(),
715 arguments: vec!["app-server".into()],
716 env: BTreeMap::new(),
717 },
718 }
719 }
720
721 pub fn with_launch(launch: RuntimeLaunch) -> Self {
723 Self { launch }
724 }
725
726 async fn connect(
727 &self,
728 launch: Option<RuntimeLaunch>,
729 runtime_id: Option<&str>,
730 ) -> Result<(
731 Arc<JsonLineClient>,
732 mpsc::UnboundedReceiver<Value>,
733 RuntimeEndpoint,
734 Option<CodexRuntimeHome>,
735 )> {
736 let mut launch = launch.unwrap_or_else(|| self.launch.clone());
737 let runtime_home = if is_stock_codex_launch(&launch) {
738 Some(CodexRuntimeHome::prepare(&mut launch, runtime_id)?)
739 } else {
740 None
741 };
742 let (client, receiver, endpoint) =
743 JsonLineClient::spawn(&launch, None, false, "codex-app-server-jsonl").await?;
744 tokio::time::timeout(
745 CODEX_STARTUP_TIMEOUT,
746 client.request(
747 "initialize",
748 json!({
749 "clientInfo": {
750 "name": "supercode",
751 "title": "Volter Harness",
752 "version": env!("CARGO_PKG_VERSION"),
753 }
754 }),
755 ),
756 )
757 .await
758 .map_err(|_| Error::Other("Codex app-server initialize timed out after 10s".into()))??;
759 client.notify("initialized", json!({})).await?;
760 Ok((client, receiver, endpoint, runtime_home))
761 }
762
763 async fn open_thread(
764 &self,
765 method: &str,
766 params: Value,
767 launch: Option<RuntimeLaunch>,
768 runtime_id: Option<&str>,
769 ) -> Result<Box<dyn RuntimeConnection>> {
770 let (client, receiver, endpoint, runtime_home) = self.connect(launch, runtime_id).await?;
771 let response = tokio::time::timeout(CODEX_STARTUP_TIMEOUT, client.request(method, params))
772 .await
773 .map_err(|_| Error::Other(format!("Codex {method} timed out after 10s")))??;
774 let thread_id = response
775 .pointer("/thread/id")
776 .and_then(Value::as_str)
777 .ok_or_else(|| Error::Other(format!("Codex {method} response omitted thread.id")))?
778 .to_string();
779 let unpublished_rollout = if method == "thread/start" {
780 runtime_home
781 .as_ref()
782 .map(|home| home.started_rollout_path(&response))
783 .transpose()?
784 } else {
785 None
786 };
787 Ok(Box::new(CodexRuntimeConnection {
788 handle: RuntimeHandle {
789 harness: HarnessId::from(HarnessId::CODEX),
790 runtime_id: thread_id,
791 endpoint,
792 },
793 client,
794 receiver,
795 active_turn: None,
796 runtime_home,
797 unpublished_rollout,
798 }))
799 }
800}
801
802#[async_trait]
803impl RuntimeBackend for CodexRuntimeBackend {
804 fn harness(&self) -> HarnessId {
805 HarnessId::from(HarnessId::CODEX)
806 }
807
808 fn capabilities(&self) -> RuntimeCapabilities {
809 RuntimeCapabilities {
810 start_session: true,
811 resume_session: true,
812 attach_existing_process: false,
816 send_input: true,
817 stream_events: true,
818 interrupt: true,
819 steer: true,
820 respond_to_requests: true,
821 }
822 }
823
824 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
825 let mut params = json!({"cwd": request.cwd});
826 if let Some(policy) = request.approval_policy {
827 params["approvalPolicy"] = json!(policy);
828 }
829 self.open_thread("thread/start", params, request.launch, None)
830 .await
831 }
832
833 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
834 let mut params = json!({"threadId": request.runtime_id});
835 if let Some(cwd) = request.cwd {
836 params["cwd"] = json!(cwd);
837 }
838 if let Some(policy) = request.approval_policy {
839 params["approvalPolicy"] = json!(policy);
840 }
841 let runtime_id = request.runtime_id.clone();
842 self.open_thread("thread/resume", params, request.launch, Some(&runtime_id))
843 .await
844 }
845}
846
847struct CodexRuntimeConnection {
848 handle: RuntimeHandle,
849 client: Arc<JsonLineClient>,
850 receiver: mpsc::UnboundedReceiver<Value>,
851 active_turn: Option<String>,
852 runtime_home: Option<CodexRuntimeHome>,
853 unpublished_rollout: Option<PathBuf>,
854}
855
856#[async_trait]
857impl RuntimeConnection for CodexRuntimeConnection {
858 fn handle(&self) -> &RuntimeHandle {
859 &self.handle
860 }
861
862 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
863 let mut parts = Vec::new();
864 if !input.text.is_empty() {
865 parts.push(json!({"type": "text", "text": input.text}));
866 }
867 parts.extend(
868 input
869 .image_urls
870 .into_iter()
871 .map(|url| json!({"type": "image", "url": url})),
872 );
873 let response = self
874 .client
875 .request(
876 "turn/start",
877 json!({
878 "threadId": self.handle.runtime_id,
879 "input": parts,
880 }),
881 )
882 .await?;
883 let turn_id = response
884 .pointer("/turn/id")
885 .and_then(Value::as_str)
886 .map(str::to_owned);
887 if let (Some(home), Some(path)) = (
888 self.runtime_home.as_ref(),
889 self.unpublished_rollout.as_ref(),
890 ) {
891 home.publish_rollout(path).await?;
892 self.unpublished_rollout = None;
893 }
894 self.active_turn = turn_id.clone();
895 Ok(turn_id)
896 }
897
898 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
899 let Some(payload) = self.receiver.recv().await else {
900 return Ok(None);
901 };
902 let kind = payload
903 .get("method")
904 .and_then(Value::as_str)
905 .map(str::to_owned)
906 .unwrap_or_else(|| "protocol".into());
907 if kind == "turn/completed" {
908 self.active_turn = None;
909 }
910 Ok(Some(HarnessEvent {
911 sequence: None,
912 kind,
913 payload,
914 }))
915 }
916
917 async fn interrupt(&mut self) -> Result<()> {
918 let Some(turn_id) = self.active_turn.as_ref() else {
919 return Err(Error::Other("Codex has no active turn to interrupt".into()));
920 };
921 self.client
922 .request(
923 "turn/interrupt",
924 json!({"threadId": self.handle.runtime_id, "turnId": turn_id}),
925 )
926 .await?;
927 Ok(())
928 }
929
930 async fn steer(&mut self, text: String) -> Result<()> {
931 let Some(turn_id) = self.active_turn.as_ref() else {
932 return Err(Error::Other("Codex has no active turn to steer".into()));
933 };
934 self.client
935 .request(
936 "turn/steer",
937 json!({
938 "threadId": self.handle.runtime_id,
939 "expectedTurnId": turn_id,
940 "input": [{"type":"text", "text":text}],
941 }),
942 )
943 .await?;
944 Ok(())
945 }
946
947 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
948 self.client.respond(request_id, response).await
949 }
950
951 async fn close(&mut self) -> Result<()> {
952 self.client.close().await?;
953 if let Some(home) = self.runtime_home.take() {
954 home.cleanup()?;
955 }
956 Ok(())
957 }
958}
959
960type PendingResponse = oneshot::Sender<std::result::Result<Value, String>>;
961type PendingResponses = Arc<Mutex<HashMap<u64, PendingResponse>>>;
962
963pub(super) struct GroupLeader(Child);
978
979impl std::ops::Deref for GroupLeader {
980 type Target = Child;
981
982 fn deref(&self) -> &Child {
983 &self.0
984 }
985}
986
987impl std::ops::DerefMut for GroupLeader {
988 fn deref_mut(&mut self) -> &mut Child {
989 &mut self.0
990 }
991}
992
993impl Drop for GroupLeader {
994 fn drop(&mut self) {
995 #[cfg(unix)]
996 if let Some(pid) = self.0.id() {
997 crate::lsp::kill_process_group(pid);
998 }
999 }
1000}
1001
1002pub(super) struct JsonLineClient {
1003 stdin: Mutex<ChildStdin>,
1004 child: Mutex<GroupLeader>,
1005 pending: PendingResponses,
1006 next_id: Mutex<u64>,
1007 include_jsonrpc: bool,
1008 events: mpsc::UnboundedSender<Value>,
1009 process_group: Option<u32>,
1010}
1011
1012impl JsonLineClient {
1013 pub(super) async fn spawn(
1014 launch: &RuntimeLaunch,
1015 cwd: Option<&std::path::Path>,
1016 include_jsonrpc: bool,
1017 protocol: &str,
1018 ) -> Result<(Arc<Self>, mpsc::UnboundedReceiver<Value>, RuntimeEndpoint)> {
1019 let mut command = Command::new(&launch.program);
1020 command
1021 .args(&launch.arguments)
1022 .envs(&launch.env)
1023 .stdin(Stdio::piped())
1024 .stdout(Stdio::piped())
1025 .stderr(Stdio::piped())
1026 .kill_on_drop(true);
1027 #[cfg(unix)]
1031 command.process_group(0);
1032 if let Some(cwd) = cwd {
1033 command.current_dir(cwd);
1034 }
1035 let mut child = command.spawn().map_err(|error| {
1036 Error::Other(format!("could not launch {}: {error}", launch.program))
1037 })?;
1038 let pid = child.id();
1039 let stdin = child
1040 .stdin
1041 .take()
1042 .ok_or_else(|| Error::Other("runtime child has no stdin".into()))?;
1043 let stdout = child
1044 .stdout
1045 .take()
1046 .ok_or_else(|| Error::Other("runtime child has no stdout".into()))?;
1047 let stderr = child
1048 .stderr
1049 .take()
1050 .ok_or_else(|| Error::Other("runtime child has no stderr".into()))?;
1051 let pending: PendingResponses = Arc::new(Mutex::new(HashMap::new()));
1052 let (events_tx, events_rx) = mpsc::unbounded_channel();
1053 let reader_events = events_tx.clone();
1054 let reader_pending = pending.clone();
1055 tokio::spawn(async move {
1056 let mut stdout_lines = BufReader::new(stdout).lines();
1057 let mut stderr_lines = BufReader::new(stderr).lines();
1058 let mut stdout_open = true;
1059 let mut stderr_open = true;
1060 let mut recent_stderr: std::collections::VecDeque<String> =
1065 std::collections::VecDeque::new();
1066 while stdout_open || stderr_open {
1067 tokio::select! {
1068 line = stdout_lines.next_line(), if stdout_open => match line {
1069 Ok(Some(line)) => {
1070 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1071 let _ = reader_events.send(json!({"type": "malformed_output", "line": line}));
1072 continue;
1073 };
1074 let response_id = value.get("id").and_then(Value::as_u64);
1075 let is_response = value.get("result").is_some() || value.get("error").is_some();
1076 if let Some(id) = response_id.filter(|_| is_response) {
1077 if let Some(sender) = reader_pending.lock().await.remove(&id) {
1078 let result = if let Some(error) = value.get("error") {
1079 Err(error.to_string())
1080 } else {
1081 Ok(value.get("result").cloned().unwrap_or(Value::Null))
1082 };
1083 let _ = sender.send(result);
1084 continue;
1085 }
1086 }
1087 let _ = reader_events.send(value);
1088 }
1089 Ok(None) => stdout_open = false,
1090 Err(error) => {
1091 let _ = reader_events.send(json!({"type": "transport_error", "message": error.to_string()}));
1092 stdout_open = false;
1093 }
1094 },
1095 line = stderr_lines.next_line(), if stderr_open => match line {
1096 Ok(Some(line)) => {
1097 if !line.trim().is_empty() {
1098 if recent_stderr.len() == STDERR_TAIL_LINES {
1099 recent_stderr.pop_front();
1100 }
1101 recent_stderr.push_back(line.clone());
1102 }
1103 let _ = reader_events.send(json!({"type": "transport_stderr", "line": line}));
1104 }
1105 Ok(None) => stderr_open = false,
1106 Err(error) => {
1107 let _ = reader_events.send(json!({"type": "transport_error", "message": error.to_string()}));
1108 stderr_open = false;
1109 }
1110 }
1111 }
1112 }
1113 let _ = reader_events.send(json!({"type": "transport_closed"}));
1114 let reason = closed_reason(&recent_stderr);
1115 let mut pending = reader_pending.lock().await;
1116 for (_, sender) in pending.drain() {
1117 let _ = sender.send(Err(reason.clone()));
1118 }
1119 });
1120 let endpoint = RuntimeEndpoint::LocalProcess {
1121 pid,
1122 command: std::iter::once(launch.program.clone())
1123 .chain(launch.arguments.iter().cloned())
1124 .collect(),
1125 protocol: protocol.into(),
1126 };
1127 Ok((
1128 Arc::new(Self {
1129 stdin: Mutex::new(stdin),
1130 child: Mutex::new(GroupLeader(child)),
1131 pending,
1132 next_id: Mutex::new(1),
1133 include_jsonrpc,
1134 events: events_tx,
1135 process_group: pid,
1136 }),
1137 events_rx,
1138 endpoint,
1139 ))
1140 }
1141
1142 pub(super) async fn request(&self, method: &str, params: Value) -> Result<Value> {
1143 let (_id, rx) = self.begin_request(method, params).await?;
1144 rx.await
1145 .map_err(|_| Error::Other("runtime response channel closed".into()))?
1146 .map_err(|message| {
1147 Error::Other(format!("runtime request `{method}` failed: {message}"))
1148 })
1149 }
1150
1151 pub(super) async fn begin_request(
1152 &self,
1153 method: &str,
1154 params: Value,
1155 ) -> Result<(u64, oneshot::Receiver<std::result::Result<Value, String>>)> {
1156 let id = {
1157 let mut next = self.next_id.lock().await;
1158 let id = *next;
1159 *next += 1;
1160 id
1161 };
1162 let (tx, rx) = oneshot::channel();
1163 self.pending.lock().await.insert(id, tx);
1164 let mut request = json!({"id": id, "method": method, "params": params});
1165 if self.include_jsonrpc {
1166 request["jsonrpc"] = json!("2.0");
1167 }
1168 if let Err(error) = self.write(&request).await {
1169 self.pending.lock().await.remove(&id);
1170 return Err(error);
1171 }
1172 Ok((id, rx))
1173 }
1174
1175 pub(super) async fn notify(&self, method: &str, params: Value) -> Result<()> {
1176 let mut notification = json!({"method": method, "params": params});
1177 if self.include_jsonrpc {
1178 notification["jsonrpc"] = json!("2.0");
1179 }
1180 self.write(¬ification).await
1181 }
1182
1183 pub(super) async fn respond(&self, id: Value, result: Value) -> Result<()> {
1184 let mut response = json!({"id": id, "result": result});
1185 if self.include_jsonrpc {
1186 response["jsonrpc"] = json!("2.0");
1187 }
1188 self.write(&response).await
1189 }
1190
1191 async fn write(&self, value: &Value) -> Result<()> {
1192 let mut stdin = self.stdin.lock().await;
1193 stdin.write_all(value.to_string().as_bytes()).await?;
1194 stdin.write_all(b"\n").await?;
1195 stdin.flush().await?;
1196 Ok(())
1197 }
1198
1199 pub(super) fn emit(&self, value: Value) {
1200 let _ = self.events.send(value);
1201 }
1202
1203 pub(super) async fn close(&self) -> Result<()> {
1204 let mut child = self.child.lock().await;
1205 if child.try_wait()?.is_some() {
1212 return Ok(());
1213 }
1214 #[cfg(unix)]
1215 {
1216 match self.process_group {
1217 Some(pid) => crate::lsp::kill_process_group(pid),
1218 None => child.kill().await?,
1219 }
1220 tokio::time::timeout(Duration::from_secs(3), child.wait())
1221 .await
1222 .map_err(|_| Error::Other("timed out reaping runtime process group".into()))??;
1223 }
1224 #[cfg(not(unix))]
1225 child.kill().await?;
1226 Ok(())
1227 }
1228}
1229
1230#[cfg(test)]
1231mod tests {
1232 use super::*;
1233
1234 #[test]
1235 fn closed_reason_reports_the_runtime_last_words() {
1236 let mut stderr = std::collections::VecDeque::new();
1237 stderr.push_back("grok: unsupported syscall SYS_execve".to_string());
1238 assert_eq!(
1239 closed_reason(&stderr),
1240 "runtime protocol closed: grok: unsupported syscall SYS_execve",
1241 );
1242 }
1243
1244 #[test]
1245 fn closed_reason_stays_bare_without_stderr() {
1246 assert_eq!(
1247 closed_reason(&std::collections::VecDeque::new()),
1248 "runtime protocol closed",
1249 );
1250 }
1251
1252 #[test]
1253 fn closed_reason_truncates_a_long_tail() {
1254 let mut stderr = std::collections::VecDeque::new();
1255 stderr.push_back("x".repeat(STDERR_TAIL_CHARACTERS + 500));
1256 let reason = closed_reason(&stderr);
1257 assert!(reason.ends_with('…'), "{reason}");
1258 assert_eq!(
1259 reason.chars().count(),
1260 "runtime protocol closed: ".chars().count() + STDERR_TAIL_CHARACTERS + 1,
1261 );
1262 }
1263
1264 fn scratch_home(tag: &str) -> PathBuf {
1265 let dir = std::env::temp_dir().join(format!(
1266 "supercode-connect-launch-{tag}-{}-{}",
1267 std::process::id(),
1268 std::time::SystemTime::now()
1269 .duration_since(std::time::UNIX_EPOCH)
1270 .unwrap()
1271 .as_nanos()
1272 ));
1273 std::fs::create_dir_all(&dir).unwrap();
1274 dir
1275 }
1276
1277 #[test]
1278 fn connect_launch_resolves_address_and_auth_from_the_harness_config() {
1279 let home = scratch_home("resolve");
1280 std::fs::create_dir_all(home.join(".gateway")).unwrap();
1281 std::fs::write(
1282 home.join(".gateway/config.json"),
1283 r#"{"gateway": {"url": "ws://127.0.0.1:18789/", "auth": {"token": "secret-credential"}}}"#,
1284 )
1285 .unwrap();
1286 let launch = RuntimeConnectLaunch {
1287 config_path: "~/.gateway/config.json".into(),
1288 address_pointer: "/gateway/url".into(),
1289 port_pointer: None,
1290 default_address: None,
1291 auth_pointer: Some("/gateway/auth/token".into()),
1292 protocol: "acp-v1-jsonrpc".into(),
1293 };
1294 let resolved = launch.resolve(&home).unwrap();
1295 assert_eq!(resolved.address, "ws://127.0.0.1:18789");
1296 assert_eq!(
1297 resolved.auth.as_ref().unwrap().secret(),
1298 "secret-credential"
1299 );
1300 let debugged = format!("{resolved:?}");
1301 assert!(!debugged.contains("secret-credential"));
1302 assert!(debugged.contains("<redacted>"));
1303 }
1304
1305 #[test]
1306 fn connect_launch_resolution_fails_closed_without_echoing_config_contents() {
1307 let home = scratch_home("fail-closed");
1308 let launch = RuntimeConnectLaunch {
1309 config_path: "~/missing.json".into(),
1310 address_pointer: "/url".into(),
1311 port_pointer: None,
1312 default_address: None,
1313 auth_pointer: None,
1314 protocol: "acp-v1-jsonrpc".into(),
1315 };
1316 assert!(launch.resolve(&home).is_err());
1317
1318 std::fs::write(
1319 home.join("present.json"),
1320 r#"{"url": "", "auth": {"token": "secret-credential"}}"#,
1321 )
1322 .unwrap();
1323 let empty_address = RuntimeConnectLaunch {
1324 config_path: "~/present.json".into(),
1325 address_pointer: "/url".into(),
1326 port_pointer: None,
1327 default_address: None,
1328 auth_pointer: None,
1329 protocol: "acp-v1-jsonrpc".into(),
1330 };
1331 let error = empty_address.resolve(&home).unwrap_err();
1332 assert!(error.to_string().contains("/url"));
1333 assert!(!error.to_string().contains("secret-credential"));
1334
1335 let missing_auth = RuntimeConnectLaunch {
1336 config_path: "~/present.json".into(),
1337 address_pointer: "/auth/token".into(),
1338 port_pointer: None,
1339 default_address: None,
1340 auth_pointer: Some("/absent".into()),
1341 protocol: "acp-v1-jsonrpc".into(),
1342 };
1343 let error = missing_auth.resolve(&home).unwrap_err();
1344 assert!(error.to_string().contains("/absent"));
1345 assert!(!error.to_string().contains("secret-credential"));
1346 }
1347
1348 #[test]
1349 fn connect_launch_round_trips_through_json() {
1350 let launch = RuntimeConnectLaunch {
1351 config_path: "~/.openclaw/openclaw.json".into(),
1352 address_pointer: "/gateway/url".into(),
1353 port_pointer: None,
1354 default_address: None,
1355 auth_pointer: Some("/gateway/token".into()),
1356 protocol: "acp-v1-jsonrpc".into(),
1357 };
1358 let encoded = serde_json::to_value(&launch).unwrap();
1359 let decoded: RuntimeConnectLaunch = serde_json::from_value(encoded).unwrap();
1360 assert_eq!(decoded, launch);
1361 let minimal: RuntimeConnectLaunch = serde_json::from_value(json!({
1362 "config_path": "~/.gateway.json",
1363 "address_pointer": "/url",
1364 "protocol": "http",
1365 }))
1366 .unwrap();
1367 assert_eq!(minimal.auth_pointer, None);
1368 }
1369
1370 #[test]
1371 fn codex_capabilities_do_not_claim_arbitrary_process_attach() {
1372 let capabilities = CodexRuntimeBackend::new().capabilities();
1373 assert!(capabilities.start_session);
1374 assert!(capabilities.resume_session);
1375 assert!(!capabilities.attach_existing_process);
1376 assert!(capabilities.send_input);
1377 assert!(capabilities.stream_events);
1378 assert!(capabilities.interrupt);
1379 assert!(capabilities.steer);
1380 }
1381
1382 #[test]
1383 fn runtime_handle_is_language_neutral_json() {
1384 let handle = RuntimeHandle {
1385 harness: HarnessId::from(HarnessId::CODEX),
1386 runtime_id: "thread-1".into(),
1387 endpoint: RuntimeEndpoint::LocalProcess {
1388 pid: Some(42),
1389 command: vec!["codex".into(), "app-server".into()],
1390 protocol: "codex-app-server-jsonl".into(),
1391 },
1392 };
1393 let encoded = serde_json::to_string(&handle).unwrap();
1394 assert_eq!(
1395 serde_json::from_str::<RuntimeHandle>(&encoded).unwrap(),
1396 handle
1397 );
1398 }
1399
1400 #[cfg(unix)]
1401 #[tokio::test]
1402 async fn codex_adapter_performs_handshake_start_and_turn() {
1403 let script = r#"
1404 i=0
1405 while IFS= read -r line; do
1406 i=$((i + 1))
1407 case "$i" in
1408 1) printf '%s\n' '{"id":1,"result":{"userAgent":"mock"}}' ;;
1409 2) ;;
1410 3) printf '%s\n' '{"id":2,"result":{"thread":{"id":"thr_mock"}}}' ;;
1411 4)
1412 printf '%s\n' '{"id":3,"result":{"turn":{"id":"turn_mock"}}}'
1413 printf '%s\n' '{"method":"turn/started","params":{"turn":{"id":"turn_mock"}}}'
1414 ;;
1415 5) printf '%s\n' '{"id":4,"result":{"turnId":"turn_mock"}}' ;;
1416 esac
1417 done
1418 "#;
1419 let backend = CodexRuntimeBackend::with_launch(RuntimeLaunch {
1420 program: "/bin/sh".into(),
1421 arguments: vec!["-c".into(), script.into()],
1422 env: BTreeMap::new(),
1423 });
1424 let mut connection = backend
1425 .start(RuntimeStartRequest {
1426 cwd: std::env::current_dir().unwrap(),
1427 launch: None,
1428 mcp_servers: Vec::new(),
1429 approval_policy: None,
1430 })
1431 .await
1432 .unwrap();
1433 assert_eq!(connection.handle().runtime_id, "thr_mock");
1434 assert_eq!(
1435 connection
1436 .send_input(RuntimeInput {
1437 text: "hi".into(),
1438 image_urls: Vec::new(),
1439 })
1440 .await
1441 .unwrap()
1442 .as_deref(),
1443 Some("turn_mock")
1444 );
1445 connection.steer("focus on tests".into()).await.unwrap();
1446 assert_eq!(
1447 connection.next_event().await.unwrap().unwrap().kind,
1448 "turn/started"
1449 );
1450 connection.close().await.unwrap();
1451 }
1452}