Skip to main content

interlink/
inbox.rs

1//! A restart-safe JSONL inbox with a single reader and atomic cursor updates.
2
3use std::fs::{File, OpenOptions};
4use std::io::{BufRead, BufReader, ErrorKind, Read, Seek, SeekFrom, Write};
5use std::path::{Path, PathBuf};
6use std::time::Duration;
7
8use anyhow::Result;
9use fs2::FileExt;
10use tokio::time::{Instant, sleep};
11
12use crate::state::{atomic_write, lock};
13
14pub const RENEW_AFTER: Duration = Duration::from_secs(3000);
15pub const RENEW_NOTICE: &str = "[interlink listener renewal] No peer message arrived. End this turn without tools or a user-facing reply so the Stop hook can renew the inbox listener.";
16
17pub struct Inbox {
18    path: PathBuf,
19}
20
21pub struct Batch {
22    pub lines: Vec<String>,
23    end: u64,
24}
25
26pub enum Wake {
27    Messages(Batch),
28    Renew,
29}
30
31impl Inbox {
32    pub fn open(path: &Path) -> Result<Self> {
33        if let Some(parent) = path.parent() {
34            std::fs::create_dir_all(parent)?;
35        }
36        OpenOptions::new().create(true).append(true).open(path)?;
37        Ok(Self {
38            path: path.to_owned(),
39        })
40    }
41
42    pub fn listener_lock(&self) -> Result<Option<File>> {
43        let file = OpenOptions::new()
44            .create(true)
45            .truncate(false)
46            .write(true)
47            .open(self.path.with_extension("lock"))?;
48        match file.try_lock_exclusive() {
49            Ok(()) => Ok(Some(file)),
50            Err(e) if e.kind() == ErrorKind::WouldBlock => Ok(None),
51            Err(e) => Err(e.into()),
52        }
53    }
54
55    pub fn append(&self, record: &str) -> Result<()> {
56        let _lock = lock(&self.path.with_extension("io-lock"))?;
57        let mut file = OpenOptions::new()
58            .read(true)
59            .append(true)
60            .open(&self.path)?;
61        let len = file.metadata()?.len();
62        if len > 0 {
63            file.seek(SeekFrom::End(-1))?;
64            let mut last = [0];
65            file.read_exact(&mut last)?;
66            if last[0] != b'\n' {
67                // A partial append was never acknowledged to the bus. Remove its
68                // fragment before a redelivery appends the complete record.
69                let mut end = len;
70                let mut buffer = [0; 4096];
71                loop {
72                    let start = end.saturating_sub(buffer.len() as u64);
73                    let count = (end - start) as usize;
74                    file.seek(SeekFrom::Start(start))?;
75                    file.read_exact(&mut buffer[..count])?;
76                    if let Some(index) = buffer[..count].iter().rposition(|b| *b == b'\n') {
77                        file.set_len(start + index as u64 + 1)?;
78                        break;
79                    }
80                    if start == 0 {
81                        file.set_len(0)?;
82                        break;
83                    }
84                    end = start;
85                }
86            }
87        }
88        file.write_all(format!("{record}\n").as_bytes())?;
89        file.sync_data()?;
90        Ok(())
91    }
92
93    pub fn read_batch(&self) -> Result<Batch> {
94        let _lock = lock(&self.path.with_extension("io-lock"))?;
95        let cursor = match std::fs::read_to_string(self.path.with_extension("cursor")) {
96            Ok(raw) => raw.trim().parse::<u64>()?,
97            Err(e) if e.kind() == ErrorKind::NotFound => 0,
98            Err(e) => return Err(e.into()),
99        };
100        let file = File::open(&self.path)?;
101        // Older versions truncated inboxes on startup. Preserve migration behavior.
102        let start = if cursor > file.metadata()?.len() {
103            0
104        } else {
105            cursor
106        };
107        let mut reader = BufReader::new(file);
108        reader.seek(SeekFrom::Start(start))?;
109        let mut batch = Batch {
110            lines: Vec::new(),
111            end: start,
112        };
113        while batch.lines.len() < 64 && batch.end - start < 256 * 1024 {
114            let mut line = String::new();
115            let bytes = reader.read_line(&mut line)?;
116            // A crash during append can leave a partial record. Never consume it.
117            if bytes == 0 || !line.ends_with('\n') {
118                break;
119            }
120            batch.end += bytes as u64;
121            batch.lines.push(line);
122        }
123        Ok(batch)
124    }
125
126    pub fn commit(&self, batch: &Batch) -> Result<()> {
127        atomic_write(
128            &self.path.with_extension("cursor"),
129            batch.end.to_string().as_bytes(),
130        )
131    }
132
133    pub async fn wait(&self, renew_after: Duration) -> Result<Wake> {
134        let deadline = Instant::now() + renew_after;
135        loop {
136            let batch = self.read_batch()?;
137            if !batch.lines.is_empty() {
138                return Ok(Wake::Messages(batch));
139            }
140            if Instant::now() >= deadline {
141                return Ok(Wake::Renew);
142            }
143            sleep(Duration::from_millis(400)).await;
144        }
145    }
146}
147
148#[cfg(test)]
149mod tests {
150    use super::*;
151
152    #[test]
153    fn restart_preserves_unread_records_and_cursor_tracks_only_delivered_bytes() {
154        let dir = tempfile::tempdir().unwrap();
155        let path = dir.path().join("inbox.jsonl");
156        let first = Inbox::open(&path).unwrap();
157        first.append("first").unwrap();
158        let restarted = Inbox::open(&path).unwrap();
159        let batch = restarted.read_batch().unwrap();
160        assert_eq!(batch.lines, ["first\n"]);
161        first.append("second").unwrap();
162        restarted.commit(&batch).unwrap();
163        let batch = Inbox::open(&path).unwrap().read_batch().unwrap();
164        assert_eq!(batch.lines, ["second\n"]);
165    }
166
167    #[test]
168    fn partial_record_is_not_consumed_and_duplicate_listener_is_silent() {
169        let dir = tempfile::tempdir().unwrap();
170        let path = dir.path().join("inbox.jsonl");
171        let inbox = Inbox::open(&path).unwrap();
172        let guard = inbox.listener_lock().unwrap().unwrap();
173        assert!(inbox.listener_lock().unwrap().is_none());
174        std::fs::write(&path, "part").unwrap();
175        assert!(inbox.read_batch().unwrap().lines.is_empty());
176        OpenOptions::new()
177            .append(true)
178            .open(&path)
179            .unwrap()
180            .write_all(b"ial\n")
181            .unwrap();
182        assert_eq!(inbox.read_batch().unwrap().lines, ["partial\n"]);
183        drop(guard);
184        assert!(inbox.listener_lock().unwrap().is_some());
185    }
186
187    #[test]
188    fn append_repairs_incomplete_tail_without_losing_complete_records() {
189        let dir = tempfile::tempdir().unwrap();
190        let path = dir.path().join("inbox.jsonl");
191        let inbox = Inbox::open(&path).unwrap();
192        inbox.append("complete").unwrap();
193        OpenOptions::new()
194            .append(true)
195            .open(&path)
196            .unwrap()
197            .write_all(b"partial")
198            .unwrap();
199        inbox.append("redelivered").unwrap();
200        assert_eq!(
201            inbox.read_batch().unwrap().lines,
202            ["complete\n", "redelivered\n"]
203        );
204    }
205
206    #[tokio::test(start_paused = true)]
207    async fn long_idle_period_requests_renewal_instead_of_disarming() {
208        let dir = tempfile::tempdir().unwrap();
209        let inbox = Inbox::open(&dir.path().join("inbox.jsonl")).unwrap();
210        assert!(matches!(
211            inbox.wait(RENEW_AFTER).await.unwrap(),
212            Wake::Renew
213        ));
214        inbox.append("after renewal").unwrap();
215        assert!(matches!(
216            inbox.wait(RENEW_AFTER).await.unwrap(),
217            Wake::Messages(_)
218        ));
219    }
220}