Skip to main content

supercode_harness/
runtime.rs

1//! Primitive live-runtime contracts and the Codex app-server reference adapter.
2//!
3//! These APIs control harness-native sessions; they do not emulate terminal
4//! keystrokes and do not claim to attach to an arbitrary already-running TUI.
5
6use 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/// Mechanical facts an adapter can guarantee.
34#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
35pub struct RuntimeCapabilities {
36    /// Can create a fresh harness-native session.
37    pub start_session: bool,
38    /// Can resume a harness-native persisted session by id.
39    pub resume_session: bool,
40    /// Can join an arbitrary already-running harness process.
41    pub attach_existing_process: bool,
42    /// Can send user input through a structured protocol.
43    pub send_input: bool,
44    /// Can receive structured live events.
45    pub stream_events: bool,
46    /// Can interrupt an in-flight turn.
47    pub interrupt: bool,
48    /// Can redirect an in-flight turn without interrupting it.
49    #[serde(default)]
50    pub steer: bool,
51    /// Can answer protocol requests such as approvals or elicitation.
52    pub respond_to_requests: bool,
53}
54
55/// Executable configuration used to launch one adapter endpoint.
56#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
57pub struct RuntimeLaunch {
58    /// Executable name or path.
59    pub program: String,
60    /// Arguments passed before adapter-generated protocol arguments.
61    pub arguments: Vec<String>,
62    /// Extra environment variables.
63    pub env: BTreeMap<String, String>,
64}
65
66/// Connect to an already-running harness endpoint instead of spawning one.
67///
68/// The registry stores where the endpoint and its credential live — the
69/// harness's own config file — never the values themselves. The service
70/// resolves them when it opens the connection, so a rotated token or a moved
71/// gateway is picked up on the next open without a registry change.
72#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
73pub struct RuntimeConnectLaunch {
74    /// Harness config file holding the endpoint; a leading `~/` expands to the
75    /// caller's home directory at resolve time.
76    pub config_path: String,
77    /// JSON pointer to the endpoint address inside the config file.
78    pub address_pointer: String,
79    /// Optional JSON pointer to a PORT number in the config file, consulted
80    /// when `address_pointer` names nothing: the address becomes that port on
81    /// loopback under `default_address`'s scheme. Harnesses like openclaw
82    /// configure a bare `gateway.port`, never a full URL.
83    #[serde(default, skip_serializing_if = "Option::is_none")]
84    pub port_pointer: Option<String>,
85    /// Optional fallback endpoint when neither pointer resolves — the
86    /// harness's documented out-of-the-box endpoint. With this set, a missing
87    /// or pointer-less config is the harness "running on defaults", not an
88    /// error.
89    #[serde(default, skip_serializing_if = "Option::is_none")]
90    pub default_address: Option<String>,
91    /// Optional JSON pointer to the bearer credential inside the config file.
92    #[serde(default, skip_serializing_if = "Option::is_none")]
93    pub auth_pointer: Option<String>,
94    /// Protocol spoken at the endpoint.
95    pub protocol: String,
96}
97
98/// Bearer credential whose `Debug` output never contains the secret.
99#[derive(Clone, PartialEq, Eq)]
100pub struct BearerToken(String);
101
102impl BearerToken {
103    /// Wrap a resolved credential.
104    pub fn new(secret: impl Into<String>) -> Self {
105        Self(secret.into())
106    }
107
108    /// The secret itself, for constructing an Authorization header.
109    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/// Endpoint and credential resolved from a [`RuntimeConnectLaunch`].
121#[derive(Debug, Clone, PartialEq, Eq)]
122pub struct ResolvedRuntimeConnection {
123    /// Concrete endpoint address.
124    pub address: String,
125    /// Bearer credential when the launch declares one.
126    pub auth: Option<BearerToken>,
127}
128
129impl RuntimeConnectLaunch {
130    /// Resolve the endpoint address and credential from the harness's config
131    /// file. Fails closed: a declared pointer that does not resolve to a
132    /// non-empty string is an error, and diagnostics name the path and the
133    /// pointer without echoing config contents.
134    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        // A missing config file is the harness on documented defaults when
140        // the descriptor declares them; otherwise it stays an error.
141        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        // Address chain: explicit URL pointer → configured port on loopback →
169        // the descriptor's documented default endpoint.
170        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        // Normalize a bare host:port to the endpoint's scheme — configs
204        // routinely omit it and a scheme-less URL makes gateway clients fall
205        // back to their compiled-in default endpoint instead.
206        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        // Auth is optional exactly when the endpoint can run without it: a
216        // declared pointer that resolves to nothing is only an error when no
217        // default endpoint is declared (the original fail-closed contract).
218        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/// One stdio MCP server the caller wants mounted into the session it is
233/// starting. Uniform shape; each backend translates it into whatever its own
234/// harness accepts (the ACP backend into `session/new`'s `mcpServers`).
235#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
236pub struct McpServerLaunch {
237    /// Server name the harness registers the tools under.
238    pub name: String,
239    /// Executable to spawn.
240    pub command: String,
241    /// Arguments passed to it.
242    #[serde(default)]
243    pub arguments: Vec<String>,
244    /// Extra environment for the spawned server.
245    #[serde(default)]
246    pub env: BTreeMap<String, String>,
247}
248
249/// Request to create a fresh runtime session.
250#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
251pub struct RuntimeStartRequest {
252    /// Project working directory.
253    pub cwd: PathBuf,
254    /// Optional executable override, primarily for alternate installs/tests.
255    pub launch: Option<RuntimeLaunch>,
256    /// MCP servers to mount into the new session, where the harness's own
257    /// start door carries them. Backends that have no such door ignore it —
258    /// their caller mounts through a config file instead.
259    #[serde(default)]
260    pub mcp_servers: Vec<McpServerLaunch>,
261    /// The harness's approval policy for this session, where its start door takes one (Codex's
262    /// `approvalPolicy` on thread/start and thread/resume, e.g. `untrusted`: every command asks first).
263    /// Absent, the harness's own default.
264    #[serde(default, skip_serializing_if = "Option::is_none")]
265    pub approval_policy: Option<String>,
266    /// The model the session runs, in the harness's own naming, where its start door takes one
267    /// (Claude Code's `--model`, e.g. `opus` or a full model id; Codex's `model` on thread/start).
268    /// Absent, the harness's own default.
269    #[serde(default, skip_serializing_if = "Option::is_none")]
270    pub model: Option<String>,
271}
272
273/// Request to resume or attach through a new adapter connection.
274#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
275pub struct RuntimeAttachRequest {
276    /// Harness-native session/thread id.
277    pub runtime_id: String,
278    /// Optional cwd override accepted by the harness protocol.
279    pub cwd: Option<PathBuf>,
280    /// Optional executable override.
281    pub launch: Option<RuntimeLaunch>,
282    /// MCP servers to mount into the resumed session, exactly as a start
283    /// request mounts them: a session's tools do not survive its process, so
284    /// the caller that resumes it names them again.
285    #[serde(default)]
286    pub mcp_servers: Vec<McpServerLaunch>,
287    /// The harness's approval policy for this session, where its start door takes one (Codex's
288    /// `approvalPolicy` on thread/start and thread/resume, e.g. `untrusted`: every command asks first).
289    /// Absent, the harness's own default.
290    #[serde(default, skip_serializing_if = "Option::is_none")]
291    pub approval_policy: Option<String>,
292}
293
294/// Observable endpoint backing a runtime connection.
295#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
296#[serde(tag = "kind", rename_all = "snake_case")]
297pub enum RuntimeEndpoint {
298    /// Child process owned by this connection.
299    LocalProcess {
300        /// Process id when available.
301        pid: Option<u32>,
302        /// Executable plus arguments.
303        command: Vec<String>,
304        /// Native protocol spoken over stdio.
305        protocol: String,
306    },
307    /// Existing HTTP service.
308    Http {
309        /// Service base URL.
310        base_url: String,
311        /// Native protocol name.
312        protocol: String,
313    },
314}
315
316/// Identity returned after a live session is started or resumed.
317#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
318pub struct RuntimeHandle {
319    /// Runtime adapter/harness.
320    pub harness: HarnessId,
321    /// Harness-native live session identity.
322    pub runtime_id: String,
323    /// Concrete endpoint used by this connection.
324    pub endpoint: RuntimeEndpoint,
325}
326
327/// User input accepted by a live runtime.
328#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
329pub struct RuntimeInput {
330    /// Plain text prompt or steering instruction.
331    pub text: String,
332    /// Runtime-resolved image URLs or `data:image/...;base64,...` payloads.
333    ///
334    /// Adapters must either preserve these as native multimodal input or
335    /// reject the turn explicitly; they must never flatten image bytes into
336    /// the text prompt.
337    #[serde(default, skip_serializing_if = "Vec::is_empty")]
338    pub image_urls: Vec<String>,
339}
340
341/// Protocol-neutral envelope around a native live event.
342#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
343pub struct HarnessEvent {
344    /// Canonical SDK sequence when the event originated from an SDK runtime.
345    /// Native harness adapters leave this absent and the service sequences
346    /// their transport stream locally.
347    #[serde(default, skip_serializing_if = "Option::is_none")]
348    pub sequence: Option<u64>,
349    /// Native method/type name, or `request` for a server-initiated request.
350    pub kind: String,
351    /// Lossless native event/request value.
352    pub payload: Value,
353}
354
355/// One connected harness-native runtime session.
356#[async_trait]
357pub trait RuntimeConnection: Send {
358    /// Identity and endpoint of this connection.
359    fn handle(&self) -> &RuntimeHandle;
360    /// Submit structured user input and return the harness-native turn id when
361    /// one is allocated.
362    async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>>;
363    /// Wait for the next native live event.
364    async fn next_event(&mut self) -> Result<Option<HarnessEvent>>;
365    /// Interrupt the current turn, when supported.
366    async fn interrupt(&mut self) -> Result<()>;
367    /// Redirect the current turn, when supported.
368    async fn steer(&mut self, _text: String) -> Result<()> {
369        Err(Error::Other(
370            "this runtime cannot steer an active turn".into(),
371        ))
372    }
373    /// Answer a server-initiated protocol request by its native JSON id.
374    async fn respond(&mut self, request_id: Value, response: Value) -> Result<()>;
375    /// Acquire this connection's native controller lease without displacing
376    /// an existing controller.
377    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    /// Refresh this connection's observer/controller lease.
383    async fn heartbeat(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
384        Err(Error::Other(
385            "this runtime does not expose controller leases".into(),
386        ))
387    }
388    /// Detach this exact connection without stopping the native runtime.
389    async fn detach(&mut self) -> Result<crate::RuntimeLeaseSnapshot> {
390        Err(Error::Other(
391            "this runtime does not expose detachable leases".into(),
392        ))
393    }
394    /// Close the adapter-owned transport/process.
395    async fn close(&mut self) -> Result<()>;
396}
397
398/// Factory for starting, resuming, and (where the native protocol permits it)
399/// joining one harness's already-running runtime endpoint.
400#[async_trait]
401pub trait RuntimeBackend: Send + Sync {
402    /// Harness implemented by this backend.
403    fn harness(&self) -> HarnessId;
404    /// Honest mechanical capability report.
405    fn capabilities(&self) -> RuntimeCapabilities;
406    /// Create a fresh harness-native session.
407    async fn start(&self, request: RuntimeStartRequest) -> Result<Box<dyn RuntimeConnection>>;
408    /// Resume a persisted harness-native session through a new protocol
409    /// connection. This does not imply joining the process that originally
410    /// wrote the session.
411    async fn attach(&self, request: RuntimeAttachRequest) -> Result<Box<dyn RuntimeConnection>>;
412    /// Join an already-running harness process or server. Most stock harnesses
413    /// cannot do this; adapters must opt in rather than silently treating a
414    /// persisted resume as a live attach.
415    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/// Codex live-runtime backend using the official `codex app-server` JSONL
427/// protocol (`initialize`, `thread/start|resume`, `turn/start|interrupt`).
428#[derive(Debug, Clone)]
429pub struct CodexRuntimeBackend {
430    launch: RuntimeLaunch,
431}
432
433const CODEX_STARTUP_TIMEOUT: Duration = Duration::from_secs(10);
434
435/// A stock Codex app-server eagerly indexes everything below `CODEX_HOME`
436/// before answering `initialize`. That turns a runtime open into an unbounded
437/// corpus scan for long-time Codex users. Give each connection a private state
438/// database and project only the one native rollout it needs into that home.
439/// The rollout itself is hard-linked, so Codex continues the original inode
440/// rather than a copy that would need lossy reconciliation later.
441#[derive(Debug)]
442struct CodexRuntimeHome {
443    root: PathBuf,
444    native_home: PathBuf,
445    /// The launch's environment, whose agent variables are recorded with the published rollout.
446    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        // The agent the launch named is recorded in the native home before the rollout is visible there, so whoever
569        // discovers the session finds its agent with it (launch_agent.rs).
570        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
603/// Keeps a runtime's closing diagnostics short enough to read in an error.
604const STDERR_TAIL_LINES: usize = 20;
605const STDERR_TAIL_CHARACTERS: usize = 2_000;
606
607/// Reports a closed protocol together with whatever the runtime last said.
608fn 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    /// Use the executing user's consolidated Codex app server.
731    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    /// Use an explicit command prefix.
742    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            // A new app-server can resume the same stored thread, but stock
838            // Codex does not let it join an arbitrary already-running TUI's
839            // transport/event fanout.
840            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
991/// A spawned child that LEADS its own process group (`process_group(0)`), so
992/// dropping it signals the whole group rather than just the leader.
993///
994/// `kill_on_drop(true)` reaches the direct child only. Harness launchers are
995/// commonly package-manager shims that spawn the real worker — the worker
996/// holding the protocol pipes — so a dropped launcher leaves that worker
997/// running with nothing attached to it. Dropping is not a rare path: it is
998/// what a blown deadline does to a launch or a control call still in flight.
999///
1000/// The group is signalled only while `id()` still answers, i.e. while this
1001/// process has not been reaped here. A reaped leader's pid can be reused by
1002/// an unrelated group, and killing that group would be someone else's
1003/// outage; a graceful `close()` that reaped the group therefore makes this
1004/// drop a no-op.
1005pub(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        // Package-manager shims commonly spawn a native worker. Isolate the
1056        // complete adapter tree so close can reap it instead of orphaning the
1057        // worker with inherited protocol handles.
1058        #[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            // A runtime that dies mid-handshake explains itself on stderr and
1089            // nowhere else. Events reach only an already-started runtime, so
1090            // without this the caller is told the protocol closed and never
1091            // told why.
1092            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(&notification).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        // `process_group` is a pid COPY taken at spawn, and a reaped pid
1234        // belongs to whoever the OS hands it to next. Close is called more
1235        // than once — a hosted runtime closes on shutdown, on transport end,
1236        // and again when its host task exits — so the second call must find
1237        // this child still unreaped here before signalling anything, exactly
1238        // as the raw-line transport does.
1239        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}