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