1use 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
30pub 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 _ => 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
119pub 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 .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}