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::{Result, bail};
10use tokio::process::Command;
11use tokio::time::timeout;
12
13pub fn validate_thread_id(id: &str) -> Result<()> {
14    let valid = id.len() == 36
15        && id.bytes().enumerate().all(|(i, b)| {
16            if matches!(i, 8 | 13 | 18 | 23) {
17                b == b'-'
18            } else {
19                b.is_ascii_hexdigit()
20            }
21        });
22    if !valid {
23        bail!("expected the current Codex thread UUID, not a session name or prefix");
24    }
25    Ok(())
26}
27
28// Keep enough space for platform quoting, executable paths, and the fixed arguments.
29pub const MAX_QUEUE_BYTES: usize = 12 * 1024;
30
31#[derive(Debug)]
32pub struct DeliveryError {
33    pub retryable: bool,
34    message: String,
35}
36
37impl DeliveryError {
38    fn new(retryable: bool, message: impl Into<String>) -> Self {
39        Self {
40            retryable,
41            message: message.into(),
42        }
43    }
44}
45
46impl fmt::Display for DeliveryError {
47    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
48        f.write_str(&self.message)
49    }
50}
51impl std::error::Error for DeliveryError {}
52
53pub async fn deliver(executable: &Path, thread_id: &str, text: &str) -> Result<(), DeliveryError> {
54    validate_thread_id(thread_id).map_err(|e| DeliveryError::new(false, e.to_string()))?;
55    if text.len() > MAX_QUEUE_BYTES || text.contains('\0') {
56        return Err(DeliveryError::new(
57            false,
58            "message exceeds Interlink's 12 KiB Codex delivery limit or contains a NUL; recover it with failed_deliveries",
59        ));
60    }
61    let status = timeout(
62        Duration::from_secs(30),
63        Command::new(executable)
64            .arg("queue").arg("--thread").arg(thread_id).arg("--message").arg(text)
65            // CLI diagnostics may echo message contents; keep them out of logs and MCP stdout.
66            .stdin(Stdio::null()).stdout(Stdio::null()).stderr(Stdio::null())
67            .kill_on_drop(true).status(),
68    ).await.map_err(|_| DeliveryError::new(false, "Codex queue timed out; delivery outcome is unknown. A manual retry may duplicate the message"))?
69        .map_err(|e| DeliveryError::new(matches!(e.kind(), ErrorKind::Interrupted | ErrorKind::WouldBlock), format!("starting codex queue: {e}")))?;
70    if !status.success() {
71        return Err(DeliveryError::new(
72            true,
73            format!("codex queue failed with {status}"),
74        ));
75    }
76    Ok(())
77}
78
79#[cfg(test)]
80mod tests {
81    use super::*;
82
83    #[test]
84    fn accepts_only_full_thread_ids() {
85        assert!(validate_thread_id("01900000-1234-7000-8000-123456789abc").is_ok());
86        for id in [
87            "",
88            "main",
89            "01900000",
90            "../thread",
91            "01900000-1234-7000-8000-123456789abg",
92        ] {
93            assert!(validate_thread_id(id).is_err(), "accepted {id}");
94        }
95    }
96
97    #[tokio::test]
98    async fn missing_cli_is_a_delivery_failure() {
99        let dir = tempfile::tempdir().unwrap();
100        assert!(
101            deliver(
102                &dir.path().join("missing"),
103                "01900000-1234-7000-8000-123456789abc",
104                "hello"
105            )
106            .await
107            .is_err()
108        );
109    }
110}