Skip to main content

interlink/
codex.rs

1//! Delivery to an existing local Codex CLI thread through its shared daemon.
2
3use std::fmt;
4use std::io::ErrorKind;
5use std::path::Path;
6use std::process::Stdio;
7use std::time::Duration;
8
9use anyhow::{Context, Result, bail};
10use serde_json::{Value, json};
11use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader};
12use tokio::process::{Child, ChildStdin, ChildStdout, Command};
13use tokio::time::timeout;
14
15pub fn validate_thread_id(id: &str) -> Result<()> {
16    let valid = id.len() == 36
17        && id.bytes().enumerate().all(|(i, b)| {
18            if matches!(i, 8 | 13 | 18 | 23) {
19                b == b'-'
20            } else {
21                b.is_ascii_hexdigit()
22            }
23        });
24    if !valid {
25        bail!("expected the current Codex thread UUID, not a session name or prefix");
26    }
27    Ok(())
28}
29
30/// A separate metadata connection never resumes or subscribes to the owning
31/// thread, so keeping it open cannot keep that conversation artificially alive.
32pub struct TitleReader {
33    _child: Child,
34    input: ChildStdin,
35    output: BufReader<ChildStdout>,
36    next_id: u64,
37}
38
39impl TitleReader {
40    pub async fn connect(executable: &Path) -> Result<Self> {
41        let mut child = Command::new(executable)
42            .args(["app-server", "--stdio"])
43            .stdin(Stdio::piped())
44            .stdout(Stdio::piped())
45            .stderr(Stdio::null())
46            .kill_on_drop(true)
47            .spawn()?;
48        let input = child.stdin.take().context("missing metadata stdin")?;
49        let output = BufReader::new(child.stdout.take().context("missing metadata stdout")?);
50        let mut reader = Self {
51            _child: child,
52            input,
53            output,
54            next_id: 0,
55        };
56        reader.request("initialize", json!({"clientInfo":{"name":"interlink-title-reader","version":env!("CARGO_PKG_VERSION")}})).await?;
57        reader
58            .send(json!({"jsonrpc":"2.0","method":"initialized"}))
59            .await?;
60        Ok(reader)
61    }
62
63    pub async fn read_title(&mut self, thread_id: &str) -> Result<Option<String>> {
64        validate_thread_id(thread_id)?;
65        let response = self
66            .request(
67                "thread/read",
68                json!({"threadId":thread_id,"includeTurns":false}),
69            )
70            .await?;
71        let thread = &response["thread"];
72        if thread["id"].as_str() != Some(thread_id) {
73            bail!("metadata returned a different thread");
74        }
75        match thread.get("name") {
76            Some(Value::String(name)) => Ok(Some(name.clone())),
77            Some(Value::Null) => Ok(None),
78            // Missing fields on older hosts must not erase a previously known title.
79            _ => bail!("host did not provide title metadata"),
80        }
81    }
82
83    async fn send(&mut self, message: Value) -> Result<()> {
84        self.input
85            .write_all(format!("{message}\n").as_bytes())
86            .await?;
87        self.input.flush().await?;
88        Ok(())
89    }
90
91    async fn request(&mut self, method: &str, params: Value) -> Result<Value> {
92        self.next_id += 1;
93        let id = self.next_id;
94        self.send(json!({"jsonrpc":"2.0","id":id,"method":method,"params":params}))
95            .await?;
96        loop {
97            let mut line = Vec::new();
98            (&mut self.output)
99                .take(1024 * 1024)
100                .read_until(b'\n', &mut line)
101                .await?;
102            if line.last() != Some(&b'\n') {
103                bail!("metadata response missing or too large");
104            }
105            let message: Value = serde_json::from_slice(&line)?;
106            if message["id"] == id {
107                if message.get("error").is_some() {
108                    bail!("Codex metadata request failed");
109                }
110                return message
111                    .get("result")
112                    .cloned()
113                    .context("missing metadata result");
114            }
115        }
116    }
117}
118
119// Keep enough space for platform quoting, executable paths, and the fixed arguments.
120pub const MAX_QUEUE_BYTES: usize = 12 * 1024;
121
122#[derive(Debug)]
123pub struct DeliveryError {
124    pub retryable: bool,
125    message: String,
126}
127
128impl DeliveryError {
129    fn new(retryable: bool, message: impl Into<String>) -> Self {
130        Self {
131            retryable,
132            message: message.into(),
133        }
134    }
135}
136
137impl fmt::Display for DeliveryError {
138    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
139        f.write_str(&self.message)
140    }
141}
142impl std::error::Error for DeliveryError {}
143
144pub async fn deliver(executable: &Path, thread_id: &str, text: &str) -> Result<(), DeliveryError> {
145    validate_thread_id(thread_id).map_err(|e| DeliveryError::new(false, e.to_string()))?;
146    if text.len() > MAX_QUEUE_BYTES || text.contains('\0') {
147        return Err(DeliveryError::new(
148            false,
149            "message exceeds Interlink's 12 KiB Codex delivery limit or contains a NUL; recover it with failed_deliveries",
150        ));
151    }
152    let status = timeout(
153        Duration::from_secs(30),
154        Command::new(executable)
155            .arg("queue").arg("--thread").arg(thread_id).arg("--message").arg(text)
156            // CLI diagnostics may echo message contents; keep them out of logs and MCP stdout.
157            .stdin(Stdio::null()).stdout(Stdio::null()).stderr(Stdio::null())
158            .kill_on_drop(true).status(),
159    ).await.map_err(|_| DeliveryError::new(false, "Codex queue timed out; delivery outcome is unknown. A manual retry may duplicate the message"))?
160        .map_err(|e| DeliveryError::new(matches!(e.kind(), ErrorKind::Interrupted | ErrorKind::WouldBlock), format!("starting codex queue: {e}")))?;
161    if !status.success() {
162        return Err(DeliveryError::new(
163            true,
164            format!("codex queue failed with {status}"),
165        ));
166    }
167    Ok(())
168}
169
170#[cfg(test)]
171mod tests {
172    use super::*;
173
174    #[test]
175    fn accepts_only_full_thread_ids() {
176        assert!(validate_thread_id("01900000-1234-7000-8000-123456789abc").is_ok());
177        for id in [
178            "",
179            "main",
180            "01900000",
181            "../thread",
182            "01900000-1234-7000-8000-123456789abg",
183        ] {
184            assert!(validate_thread_id(id).is_err(), "accepted {id}");
185        }
186    }
187
188    #[tokio::test]
189    async fn missing_cli_is_a_delivery_failure() {
190        let dir = tempfile::tempdir().unwrap();
191        assert!(
192            deliver(
193                &dir.path().join("missing"),
194                "01900000-1234-7000-8000-123456789abc",
195                "hello"
196            )
197            .await
198            .is_err()
199        );
200    }
201}