1use 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 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 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 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}