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, NativeProjection};
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 #[serde(default, skip_serializing_if = "Option::is_none")]
270 pub model: Option<String>,
271}
272
273#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
275pub struct RuntimeAttachRequest {
276 pub runtime_id: String,
278 pub cwd: Option<PathBuf>,
280 pub launch: Option<RuntimeLaunch>,
282 #[serde(default)]
286 pub mcp_servers: Vec<McpServerLaunch>,
287 #[serde(default, skip_serializing_if = "Option::is_none")]
291 pub approval_policy: Option<String>,
292}
293
294#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
296#[serde(tag = "kind", rename_all = "snake_case")]
297pub enum RuntimeEndpoint {
298 LocalProcess {
300 pid: Option<u32>,
302 command: Vec<String>,
304 protocol: String,
306 },
307 Http {
309 base_url: String,
311 protocol: String,
313 },
314}
315
316#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
318pub struct RuntimeHandle {
319 pub harness: HarnessId,
321 pub runtime_id: String,
323 pub endpoint: RuntimeEndpoint,
325}
326
327#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
329pub struct RuntimeInput {
330 pub text: String,
332 #[serde(default, skip_serializing_if = "Vec::is_empty")]
338 pub image_urls: Vec<String>,
339}
340
341#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
343pub struct HarnessEvent {
344 #[serde(default, skip_serializing_if = "Option::is_none")]
348 pub sequence: Option<u64>,
349 pub kind: String,
351 pub payload: Value,
353}
354
355#[async_trait]
357pub trait RuntimeConnection: Send {
358 fn handle(&self) -> &RuntimeHandle;
360 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>>;
363 async fn next_event(&mut self) -> Result<Option<HarnessEvent>>;
365 async fn interrupt(&mut self) -> Result<()>;
367 async fn steer(&mut self, _text: String) -> Result<()> {
369 Err(Error::Other(
370 "this runtime cannot steer an active turn".into(),
371 ))
372 }
373 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()>;
375 async fn acquire_control(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
378 Err(Error::Other(
379 "this runtime does not expose controller leases".into(),
380 ))
381 }
382 async fn heartbeat(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
384 Err(Error::Other(
385 "this runtime does not expose controller leases".into(),
386 ))
387 }
388 async fn detach(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
390 Err(Error::Other(
391 "this runtime does not expose detachable leases".into(),
392 ))
393 }
394 async fn close(&mut self) -> Result<()>;
396}
397
398#[async_trait]
401pub trait RuntimeBackend: Send + Sync {
402 fn harness(&self) -> HarnessId;
404 fn capabilities(&self) -> RuntimeCapabilities;
406 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>>;
408 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>>;
412 async fn attach_existing(
416 &self,
417 _request: RuntimeAttachRequest,
418 ) -> Result<Box<dyn RuntimeConnection>> {
419 Err(Error::Other(format!(
420 "{} cannot attach to an already-running process",
421 self.harness().as_str()
422 )))
423 }
424}
425
426#[derive(Debug, Clone)]
429pub struct CodexRuntimeBackend {
430 launch: RuntimeLaunch,
431}
432
433const CODEX_STARTUP_TIMEOUT: Duration = Duration::from_secs(10);
434
435#[derive(Debug)]
442struct CodexRuntimeHome {
443 root: PathBuf,
444 native_home: PathBuf,
445 launch_env: BTreeMap<String, String>,
447}
448
449impl CodexRuntimeHome {
450 fn prepare(launch: &mut RuntimeLaunch, runtime_id: Option<&str>) -> Result<Self> {
451 let native_home = codex_native_home(launch)?;
452 let root = supercode_runtime_root()
453 .join("codex")
454 .join(generated_session_id());
455 std::fs::create_dir_all(&root).map_err(|error| {
456 Error::Other(format!(
457 "could not create isolated Codex runtime home {}: {error}",
458 root.display()
459 ))
460 })?;
461 set_private_directory(&root)?;
462 let root = std::fs::canonicalize(&root)?;
463
464 for entry in [
465 "auth.json",
466 "config.toml",
467 "hooks.json",
468 "models_cache.json",
469 "installation_id",
470 ".personality_migration",
471 ".sandbox_migration",
472 "cache",
473 "generated_images",
474 "mcp-oauth-locks",
475 "memories",
476 "plugins",
477 "rules",
478 "shell_snapshots",
479 "skills",
480 "thread-writer-locks",
481 ] {
482 link_runtime_resource(&native_home.join(entry), &root.join(entry))?;
483 }
484
485 if let Some(runtime_id) = runtime_id {
486 let source = find_codex_rollout(&native_home.join("sessions"), runtime_id)?
487 .ok_or_else(|| {
488 Error::Other(format!(
489 "could not find Codex rollout `{runtime_id}` below {}",
490 native_home.join("sessions").display()
491 ))
492 })?;
493 let relative = source.strip_prefix(&native_home).map_err(|_| {
494 Error::Other(format!(
495 "Codex rollout {} is outside native home {}",
496 source.display(),
497 native_home.display()
498 ))
499 })?;
500 let projected = root.join(relative);
501 if let Some(parent) = projected.parent() {
502 std::fs::create_dir_all(parent)?;
503 }
504 std::fs::hard_link(&source, &projected).map_err(|error| {
505 Error::Other(format!(
506 "could not project Codex rollout {} into isolated runtime home: {error}",
507 source.display()
508 ))
509 })?;
510 }
511
512 launch
513 .env
514 .insert("CODEX_HOME".into(), root.to_string_lossy().into_owned());
515 Ok(Self {
516 root,
517 native_home,
518 launch_env: launch.env.clone(),
519 })
520 }
521
522 fn started_rollout_path(&self, response: &Value) -> Result<PathBuf> {
523 let path = response
524 .pointer("/thread/path")
525 .and_then(Value::as_str)
526 .map(PathBuf::from)
527 .ok_or_else(|| {
528 Error::Other("Codex thread/start response omitted thread.path".into())
529 })?;
530 let relative = path.strip_prefix(&self.root).map_err(|_| {
531 Error::Other(format!(
532 "Codex created rollout {} outside isolated runtime home {}",
533 path.display(),
534 self.root.display()
535 ))
536 })?;
537 if !relative.starts_with("sessions") {
538 return Err(Error::Other(format!(
539 "Codex created non-session rollout {}",
540 path.display()
541 )));
542 }
543 Ok(path)
544 }
545
546 async fn publish_rollout(&self, path: &Path, session_id: &str) -> Result<()> {
547 let relative = path.strip_prefix(&self.root).map_err(|_| {
548 Error::Other(format!(
549 "Codex created rollout {} outside isolated runtime home {}",
550 path.display(),
551 self.root.display()
552 ))
553 })?;
554 let publish_deadline = tokio::time::Instant::now() + Duration::from_secs(2);
555 while !path.is_file() {
556 if tokio::time::Instant::now() >= publish_deadline {
557 return Err(Error::Other(format!(
558 "Codex did not create promised rollout {} within 2s",
559 path.display()
560 )));
561 }
562 tokio::time::sleep(Duration::from_millis(10)).await;
563 }
564 let native = self.native_home.join(relative);
565 if let Some(parent) = native.parent() {
566 std::fs::create_dir_all(parent)?;
567 }
568 if let Err(error) =
571 crate::launch_agent::record(&self.native_home, session_id, &self.launch_env)
572 {
573 tracing::warn!(
574 "could not record the launch agent of Codex session {session_id}: {error}"
575 );
576 }
577 std::fs::hard_link(path, &native).map_err(|error| {
578 Error::Other(format!(
579 "could not publish Codex rollout {} to native home: {error}",
580 path.display()
581 ))
582 })
583 }
584
585 fn cleanup(&self) -> Result<()> {
586 match std::fs::remove_dir_all(&self.root) {
587 Ok(()) => Ok(()),
588 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
589 Err(error) => Err(Error::Other(format!(
590 "could not clean isolated Codex runtime home {}: {error}",
591 self.root.display()
592 ))),
593 }
594 }
595}
596
597impl Drop for CodexRuntimeHome {
598 fn drop(&mut self) {
599 let _ = self.cleanup();
600 }
601}
602
603const STDERR_TAIL_LINES: usize = 20;
605const STDERR_TAIL_CHARACTERS: usize = 2_000;
606
607fn closed_reason(recent_stderr: &std::collections::VecDeque<String>) -> String {
609 if recent_stderr.is_empty() {
610 return "runtime protocol closed".into();
611 }
612 let mut tail = recent_stderr
613 .iter()
614 .map(String::as_str)
615 .collect::<Vec<_>>()
616 .join(" | ");
617 if tail.chars().count() > STDERR_TAIL_CHARACTERS {
618 tail = tail
619 .chars()
620 .take(STDERR_TAIL_CHARACTERS)
621 .collect::<String>()
622 + "…";
623 }
624 format!("runtime protocol closed: {tail}")
625}
626
627fn is_stock_codex_launch(launch: &RuntimeLaunch) -> bool {
628 launch
629 .arguments
630 .iter()
631 .any(|argument| argument == "app-server")
632 && Path::new(&launch.program)
633 .file_name()
634 .and_then(|name| name.to_str())
635 .is_some_and(|name| matches!(name, "codex" | "codex.exe" | "codex.cmd"))
636}
637
638fn codex_native_home(launch: &RuntimeLaunch) -> Result<PathBuf> {
639 launch
640 .env
641 .get("CODEX_HOME")
642 .map(PathBuf::from)
643 .or_else(|| std::env::var_os("CODEX_HOME").map(PathBuf::from))
644 .or_else(|| {
645 supercode_interchange::user_home()
646 .map(std::path::PathBuf::into_os_string)
647 .map(PathBuf::from)
648 .map(|home| home.join(".codex"))
649 })
650 .ok_or_else(|| Error::Other("Codex runtime requires CODEX_HOME or HOME".into()))
651}
652
653fn supercode_runtime_root() -> PathBuf {
654 std::env::var_os("SUPERCODE_HOME")
655 .map(PathBuf::from)
656 .or_else(|| {
657 supercode_interchange::user_home()
658 .map(std::path::PathBuf::into_os_string)
659 .map(PathBuf::from)
660 .map(|home| home.join(".supercode"))
661 })
662 .unwrap_or_else(|| std::env::temp_dir().join("supercode"))
663 .join("runtime-homes")
664}
665
666fn find_codex_rollout(root: &Path, runtime_id: &str) -> Result<Option<PathBuf>> {
667 let entries = match std::fs::read_dir(root) {
668 Ok(entries) => entries,
669 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
670 Err(error) => return Err(error.into()),
671 };
672 let expected_suffix = format!("-{runtime_id}.jsonl");
673 for entry in entries {
674 let entry = entry?;
675 let kind = entry.file_type()?;
676 if kind.is_dir() {
677 if let Some(path) = find_codex_rollout(&entry.path(), runtime_id)? {
678 return Ok(Some(path));
679 }
680 } else if kind.is_file()
681 && entry
682 .file_name()
683 .to_str()
684 .is_some_and(|name| name.ends_with(&expected_suffix))
685 {
686 return Ok(Some(entry.path()));
687 }
688 }
689 Ok(None)
690}
691
692#[cfg(unix)]
693fn link_runtime_resource(source: &Path, target: &Path) -> Result<()> {
694 use std::os::unix::fs::symlink;
695
696 if source.exists() {
697 symlink(source, target)?;
698 }
699 Ok(())
700}
701
702#[cfg(not(unix))]
703fn link_runtime_resource(source: &Path, target: &Path) -> Result<()> {
704 if source.is_file() {
705 std::fs::copy(source, target)?;
706 }
707 Ok(())
708}
709
710#[cfg(unix)]
711fn set_private_directory(path: &Path) -> Result<()> {
712 use std::os::unix::fs::PermissionsExt;
713
714 std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o700))?;
715 Ok(())
716}
717
718#[cfg(not(unix))]
719fn set_private_directory(_path: &Path) -> Result<()> {
720 Ok(())
721}
722
723impl Default for CodexRuntimeBackend {
724 fn default() -> Self {
725 Self::new()
726 }
727}
728
729impl CodexRuntimeBackend {
730 pub fn new() -> Self {
732 Self {
733 launch: RuntimeLaunch {
734 program: "codex".into(),
735 arguments: vec!["app-server".into()],
736 env: BTreeMap::new(),
737 },
738 }
739 }
740
741 pub fn with_launch(launch: RuntimeLaunch) -> Self {
743 Self { launch }
744 }
745
746 async fn connect(
747 &self,
748 launch: Option<RuntimeLaunch>,
749 runtime_id: Option<&str>,
750 ) -> Result<(
751 Arc<JsonLineClient>,
752 mpsc::UnboundedReceiver<Value>,
753 RuntimeEndpoint,
754 Option<CodexRuntimeHome>,
755 )> {
756 let mut launch = launch.unwrap_or_else(|| self.launch.clone());
757 let stock_codex = is_stock_codex_launch(&launch);
758 if stock_codex {
759 launch.program = crate::startup_prompts::codex_program().map_err(Error::Other)?;
760 launch.arguments = crate::startup_prompts::codex_arguments(&launch.arguments);
761 }
762 let runtime_home = if stock_codex {
763 Some(CodexRuntimeHome::prepare(&mut launch, runtime_id)?)
764 } else {
765 None
766 };
767 let (client, receiver, endpoint) =
768 JsonLineClient::spawn(&launch, None, false, "codex-app-server-jsonl").await?;
769 tokio::time::timeout(
770 CODEX_STARTUP_TIMEOUT,
771 client.request(
772 "initialize",
773 json!({
774 "clientInfo": {
775 "name": "supercode",
776 "title": "Volter Harness",
777 "version": env!("CARGO_PKG_VERSION"),
778 }
779 }),
780 ),
781 )
782 .await
783 .map_err(|_| Error::Other("Codex app-server initialize timed out after 10s".into()))??;
784 client.notify("initialized", json!({})).await?;
785 Ok((client, receiver, endpoint, runtime_home))
786 }
787
788 async fn open_thread(
789 &self,
790 method: &str,
791 params: Value,
792 launch: Option<RuntimeLaunch>,
793 runtime_id: Option<&str>,
794 ) -> Result<Box<dyn RuntimeConnection>> {
795 let (client, receiver, endpoint, runtime_home) = self.connect(launch, runtime_id).await?;
796 let response = tokio::time::timeout(CODEX_STARTUP_TIMEOUT, client.request(method, params))
797 .await
798 .map_err(|_| Error::Other(format!("Codex {method} timed out after 10s")))??;
799 let thread_id = response
800 .pointer("/thread/id")
801 .and_then(Value::as_str)
802 .ok_or_else(|| Error::Other(format!("Codex {method} response omitted thread.id")))?
803 .to_string();
804 let unpublished_rollout = if method == "thread/start" {
805 runtime_home
806 .as_ref()
807 .map(|home| home.started_rollout_path(&response))
808 .transpose()?
809 } else {
810 None
811 };
812 Ok(Box::new(CodexRuntimeConnection {
813 handle: RuntimeHandle {
814 harness: HarnessId::from(HarnessId::CODEX),
815 runtime_id: thread_id,
816 endpoint,
817 },
818 client,
819 receiver,
820 active_turn: None,
821 runtime_home,
822 unpublished_rollout,
823 }))
824 }
825}
826
827#[async_trait]
828impl RuntimeBackend for CodexRuntimeBackend {
829 fn harness(&self) -> HarnessId {
830 HarnessId::from(HarnessId::CODEX)
831 }
832
833 fn capabilities(&self) -> RuntimeCapabilities {
834 RuntimeCapabilities {
835 start_session: true,
836 resume_session: true,
837 attach_existing_process: false,
841 send_input: true,
842 stream_events: true,
843 interrupt: true,
844 steer: true,
845 respond_to_requests: true,
846 }
847 }
848
849 async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>> {
850 let mut params = json!({"cwd": request.cwd});
851 if let Some(policy) = request.approval_policy {
852 params["approvalPolicy"] = json!(policy);
853 }
854 if let Some(model) = request.model {
855 params["model"] = json!(model);
856 }
857 self.open_thread("thread/start", params, request.launch, None)
858 .await
859 }
860
861 async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>> {
862 let mut params = json!({"threadId": request.runtime_id});
863 if let Some(cwd) = request.cwd {
864 params["cwd"] = json!(cwd);
865 }
866 if let Some(policy) = request.approval_policy {
867 params["approvalPolicy"] = json!(policy);
868 }
869 let runtime_id = request.runtime_id.clone();
870 self.open_thread("thread/resume", params, request.launch, Some(&runtime_id))
871 .await
872 }
873}
874
875struct CodexRuntimeConnection {
876 handle: RuntimeHandle,
877 client: Arc<JsonLineClient>,
878 receiver: mpsc::UnboundedReceiver<Value>,
879 active_turn: Option<String>,
880 runtime_home: Option<CodexRuntimeHome>,
881 unpublished_rollout: Option<PathBuf>,
882}
883
884#[async_trait]
885impl RuntimeConnection for CodexRuntimeConnection {
886 fn handle(&self) -> &RuntimeHandle {
887 &self.handle
888 }
889
890 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
891 let mut parts = Vec::new();
892 if !input.text.is_empty() {
893 parts.push(json!({"type": "text", "text": input.text}));
894 }
895 parts.extend(
896 input
897 .image_urls
898 .into_iter()
899 .map(|url| json!({"type": "image", "url": url})),
900 );
901 let response = self
902 .client
903 .request(
904 "turn/start",
905 json!({
906 "threadId": self.handle.runtime_id,
907 "input": parts,
908 }),
909 )
910 .await?;
911 let turn_id = response
912 .pointer("/turn/id")
913 .and_then(Value::as_str)
914 .map(str::to_owned);
915 if let (Some(home), Some(path)) = (
916 self.runtime_home.as_ref(),
917 self.unpublished_rollout.as_ref(),
918 ) {
919 home.publish_rollout(path, &self.handle.runtime_id).await?;
920 self.unpublished_rollout = None;
921 }
922 self.active_turn = turn_id.clone();
923 Ok(turn_id)
924 }
925
926 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
927 let Some(payload) = self.receiver.recv().await else {
928 return Ok(None);
929 };
930 let kind = payload
931 .get("method")
932 .and_then(Value::as_str)
933 .map(str::to_owned)
934 .unwrap_or_else(|| "protocol".into());
935 if kind == "turn/completed" {
936 self.active_turn = None;
937 }
938 Ok(Some(HarnessEvent {
939 sequence: None,
940 kind,
941 payload,
942 }))
943 }
944
945 async fn interrupt(&mut self) -> Result<()> {
946 let Some(turn_id) = self.active_turn.as_ref() else {
947 return Err(Error::Other("Codex has no active turn to interrupt".into()));
948 };
949 self.client
950 .request(
951 "turn/interrupt",
952 json!({"threadId": self.handle.runtime_id, "turnId": turn_id}),
953 )
954 .await?;
955 Ok(())
956 }
957
958 async fn steer(&mut self, text: String) -> Result<()> {
959 let Some(turn_id) = self.active_turn.as_ref() else {
960 return Err(Error::Other("Codex has no active turn to steer".into()));
961 };
962 self.client
963 .request(
964 "turn/steer",
965 json!({
966 "threadId": self.handle.runtime_id,
967 "expectedTurnId": turn_id,
968 "input": [{"type":"text", "text":text}],
969 }),
970 )
971 .await?;
972 Ok(())
973 }
974
975 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
976 self.client.respond(request_id, response).await
977 }
978
979 async fn close(&mut self) -> Result<()> {
980 self.client.close().await?;
981 if let Some(home) = self.runtime_home.take() {
982 home.cleanup()?;
983 }
984 Ok(())
985 }
986}
987
988type PendingResponse = oneshot::Sender<std::result::Result<Value, String>>;
989type PendingResponses = Arc<Mutex<HashMap<u64, PendingResponse>>>;
990
991pub(super) struct GroupLeader(Child);
1006
1007impl std::ops::Deref for GroupLeader {
1008 type Target = Child;
1009
1010 fn deref(&self) -> &Child {
1011 &self.0
1012 }
1013}
1014
1015impl std::ops::DerefMut for GroupLeader {
1016 fn deref_mut(&mut self) -> &mut Child {
1017 &mut self.0
1018 }
1019}
1020
1021impl Drop for GroupLeader {
1022 fn drop(&mut self) {
1023 #[cfg(unix)]
1024 if let Some(pid) = self.0.id() {
1025 crate::lsp::kill_process_group(pid);
1026 }
1027 }
1028}
1029
1030pub(super) struct JsonLineClient {
1031 stdin: Mutex<ChildStdin>,
1032 child: Mutex<GroupLeader>,
1033 pending: PendingResponses,
1034 next_id: Mutex<u64>,
1035 include_jsonrpc: bool,
1036 events: mpsc::UnboundedSender<Value>,
1037 process_group: Option<u32>,
1038}
1039
1040impl JsonLineClient {
1041 pub(super) async fn spawn(
1042 launch: &RuntimeLaunch,
1043 cwd: Option<&std::path::Path>,
1044 include_jsonrpc: bool,
1045 protocol: &str,
1046 ) -> Result<(Arc<Self>, mpsc::UnboundedReceiver<Value>, RuntimeEndpoint)> {
1047 let mut command = Command::new(&launch.program);
1048 command
1049 .args(&launch.arguments)
1050 .envs(&launch.env)
1051 .stdin(Stdio::piped())
1052 .stdout(Stdio::piped())
1053 .stderr(Stdio::piped())
1054 .kill_on_drop(true);
1055 #[cfg(unix)]
1059 command.process_group(0);
1060 if let Some(cwd) = cwd {
1061 command.current_dir(cwd);
1062 }
1063 let mut child = command.spawn().map_err(|error| {
1064 Error::Other(format!("could not launch {}: {error}", launch.program))
1065 })?;
1066 let pid = child.id();
1067 let stdin = child
1068 .stdin
1069 .take()
1070 .ok_or_else(|| Error::Other("runtime child has no stdin".into()))?;
1071 let stdout = child
1072 .stdout
1073 .take()
1074 .ok_or_else(|| Error::Other("runtime child has no stdout".into()))?;
1075 let stderr = child
1076 .stderr
1077 .take()
1078 .ok_or_else(|| Error::Other("runtime child has no stderr".into()))?;
1079 let pending: PendingResponses = Arc::new(Mutex::new(HashMap::new()));
1080 let (events_tx, events_rx) = mpsc::unbounded_channel();
1081 let reader_events = events_tx.clone();
1082 let reader_pending = pending.clone();
1083 tokio::spawn(async move {
1084 let mut stdout_lines = BufReader::new(stdout).lines();
1085 let mut stderr_lines = BufReader::new(stderr).lines();
1086 let mut stdout_open = true;
1087 let mut stderr_open = true;
1088 let mut recent_stderr: std::collections::VecDeque<String> =
1093 std::collections::VecDeque::new();
1094 while stdout_open || stderr_open {
1095 tokio::select! {
1096 line = stdout_lines.next_line(), if stdout_open => match line {
1097 Ok(Some(line)) => {
1098 let Ok(value) = serde_json::from_str::<Value>(&line) else {
1099 let _ = reader_events.send(json!({"type": "malformed_output", "line": line}));
1100 continue;
1101 };
1102 let response_id = value.get("id").and_then(Value::as_u64);
1103 let is_response = value.get("result").is_some() || value.get("error").is_some();
1104 if let Some(id) = response_id.filter(|_| is_response) {
1105 if let Some(sender) = reader_pending.lock().await.remove(&id) {
1106 let result = if let Some(error) = value.get("error") {
1107 Err(error.to_string())
1108 } else {
1109 Ok(value.get("result").cloned().unwrap_or(Value::Null))
1110 };
1111 let _ = sender.send(result);
1112 continue;
1113 }
1114 }
1115 let _ = reader_events.send(value);
1116 }
1117 Ok(None) => stdout_open = false,
1118 Err(error) => {
1119 let _ = reader_events.send(json!({"type": "transport_error", "message": error.to_string()}));
1120 stdout_open = false;
1121 }
1122 },
1123 line = stderr_lines.next_line(), if stderr_open => match line {
1124 Ok(Some(line)) => {
1125 if !line.trim().is_empty() {
1126 if recent_stderr.len() == STDERR_TAIL_LINES {
1127 recent_stderr.pop_front();
1128 }
1129 recent_stderr.push_back(line.clone());
1130 }
1131 let _ = reader_events.send(json!({"type": "transport_stderr", "line": line}));
1132 }
1133 Ok(None) => stderr_open = false,
1134 Err(error) => {
1135 let _ = reader_events.send(json!({"type": "transport_error", "message": error.to_string()}));
1136 stderr_open = false;
1137 }
1138 }
1139 }
1140 }
1141 let _ = reader_events.send(json!({"type": "transport_closed"}));
1142 let reason = closed_reason(&recent_stderr);
1143 let mut pending = reader_pending.lock().await;
1144 for (_, sender) in pending.drain() {
1145 let _ = sender.send(Err(reason.clone()));
1146 }
1147 });
1148 let endpoint = RuntimeEndpoint::LocalProcess {
1149 pid,
1150 command: std::iter::once(launch.program.clone())
1151 .chain(launch.arguments.iter().cloned())
1152 .collect(),
1153 protocol: protocol.into(),
1154 };
1155 Ok((
1156 Arc::new(Self {
1157 stdin: Mutex::new(stdin),
1158 child: Mutex::new(GroupLeader(child)),
1159 pending,
1160 next_id: Mutex::new(1),
1161 include_jsonrpc,
1162 events: events_tx,
1163 process_group: pid,
1164 }),
1165 events_rx,
1166 endpoint,
1167 ))
1168 }
1169
1170 pub(super) async fn request(&self, method: &str, params: Value) -> Result<Value> {
1171 let (_id, rx) = self.begin_request(method, params).await?;
1172 rx.await
1173 .map_err(|_| Error::Other("runtime response channel closed".into()))?
1174 .map_err(|message| {
1175 Error::Other(format!("runtime request `{method}` failed: {message}"))
1176 })
1177 }
1178
1179 pub(super) async fn begin_request(
1180 &self,
1181 method: &str,
1182 params: Value,
1183 ) -> Result<(u64, oneshot::Receiver<std::result::Result<Value, String>>)> {
1184 let id = {
1185 let mut next = self.next_id.lock().await;
1186 let id = *next;
1187 *next += 1;
1188 id
1189 };
1190 let (tx, rx) = oneshot::channel();
1191 self.pending.lock().await.insert(id, tx);
1192 let mut request = json!({"id": id, "method": method, "params": params});
1193 if self.include_jsonrpc {
1194 request["jsonrpc"] = json!("2.0");
1195 }
1196 if let Err(error) = self.write(&request).await {
1197 self.pending.lock().await.remove(&id);
1198 return Err(error);
1199 }
1200 Ok((id, rx))
1201 }
1202
1203 pub(super) async fn notify(&self, method: &str, params: Value) -> Result<()> {
1204 let mut notification = json!({"method": method, "params": params});
1205 if self.include_jsonrpc {
1206 notification["jsonrpc"] = json!("2.0");
1207 }
1208 self.write(¬ification).await
1209 }
1210
1211 pub(super) async fn respond(&self, id: Value, result: Value) -> Result<()> {
1212 let mut response = json!({"id": id, "result": result});
1213 if self.include_jsonrpc {
1214 response["jsonrpc"] = json!("2.0");
1215 }
1216 self.write(&response).await
1217 }
1218
1219 async fn write(&self, value: &Value) -> Result<()> {
1220 let mut stdin = self.stdin.lock().await;
1221 stdin.write_all(value.to_string().as_bytes()).await?;
1222 stdin.write_all(b"\n").await?;
1223 stdin.flush().await?;
1224 Ok(())
1225 }
1226
1227 pub(super) fn emit(&self, value: Value) {
1228 let _ = self.events.send(value);
1229 }
1230
1231 pub(super) async fn close(&self) -> Result<()> {
1232 let mut child = self.child.lock().await;
1233 if child.try_wait()?.is_some() {
1240 return Ok(());
1241 }
1242 #[cfg(unix)]
1243 {
1244 match self.process_group {
1245 Some(pid) => crate::lsp::kill_process_group(pid),
1246 None => child.kill().await?,
1247 }
1248 tokio::time::timeout(Duration::from_secs(3), child.wait())
1249 .await
1250 .map_err(|_| Error::Other("timed out reaping runtime process group".into()))??;
1251 }
1252 #[cfg(not(unix))]
1253 child.kill().await?;
1254 Ok(())
1255 }
1256}