1use 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
28pub 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 .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}