1use std::collections::HashSet;
4use std::fs::{File, OpenOptions};
5use std::io::ErrorKind;
6use std::path::{Path, PathBuf};
7
8use anyhow::{Result, bail};
9use fs2::FileExt;
10use serde::{Deserialize, Serialize};
11
12use crate::agent::MAX_PAST_MS;
13use crate::identity::{SignedMessage, TaskStatus, mint_session_id};
14use crate::now_ms;
15use crate::state::{atomic_write, lock};
16
17const NOTICE_RETRY_MS: u64 = 30_000;
18const MAX_NOTICE_RETRY_MS: u64 = 300_000;
19const MAX_RECORDS: usize = 4096;
20const MAX_BYTES: usize = 32 * 1024 * 1024;
21
22#[derive(Clone, Debug, Serialize, Deserialize)]
23pub struct Received {
24 pub message: SignedMessage,
25 pub peer: String,
26 pub state: String,
27 pub superseded_by: Option<String>,
28}
29
30impl Received {
31 pub fn unread(&self) -> bool {
32 self.state != "receiver_acknowledged" && self.superseded_by.is_none()
33 }
34}
35
36#[derive(Clone, Debug, Serialize, Deserialize)]
37pub struct Notice {
38 pub id: String,
39 pub state: String,
40 pub messages: Vec<(String, String)>,
41 #[serde(default)]
42 pub retry_at: u64,
43 #[serde(default)]
44 pub attempt: u32,
45}
46
47impl Notice {
48 pub fn text(&self) -> String {
49 format!(
50 "[Interlink inbox notice] Check the inbox before ending this turn.\n\
51 1. Call Interlink's receive_messages(notification_id=\"{}\"). If necessary, discover the Interlink MCP tools first.\n\
52 2. Read the returned messages, then call acknowledge_messages(messages=[the exact returned receipt objects, each containing sender and msg_id]). Acknowledge before acting on a request or fetching another batch.\n\
53 3. Only after fetching: if the inbox is empty, end silently. After acknowledgement, handle messages within the operator's authorized scope; routine progress needs no user-facing reply.\n\
54 If either tool is unavailable or fails, report that blocker. Do not silently skip the inbox check.",
55 self.id
56 )
57 }
58}
59
60#[derive(Default, Serialize, Deserialize)]
61struct Data {
62 records: Vec<Received>,
63 notice: Option<Notice>,
64}
65
66pub struct Mailbox {
67 path: PathBuf,
68}
69
70impl Mailbox {
71 pub fn new(path: &Path) -> Self {
72 Self {
73 path: path.to_owned(),
74 }
75 }
76
77 pub fn notifier_lock(&self) -> Result<Option<File>> {
78 if let Some(parent) = self.path.parent() {
79 std::fs::create_dir_all(parent)?;
80 }
81 let file = OpenOptions::new()
82 .create(true)
83 .truncate(false)
84 .write(true)
85 .open(self.path.with_extension("notifier-lock"))?;
86 match file.try_lock_exclusive() {
87 Ok(()) => Ok(Some(file)),
88 Err(e) if e.kind() == ErrorKind::WouldBlock => Ok(None),
89 Err(e) => Err(e.into()),
90 }
91 }
92
93 fn load(&self) -> Result<Data> {
94 match std::fs::read(&self.path) {
95 Ok(bytes) => Ok(serde_json::from_slice(&bytes)?),
96 Err(e) if e.kind() == ErrorKind::NotFound => Ok(Data::default()),
97 Err(e) => Err(e.into()),
98 }
99 }
100
101 fn save(&self, data: &Data) -> Result<()> {
102 let bytes = serde_json::to_vec(data)?;
103 atomic_write(&self.path, &bytes)
104 }
105
106 pub fn records(&self) -> Result<Vec<Received>> {
107 Ok(self.load()?.records)
108 }
109
110 pub fn retain(&self, message: &SignedMessage, peer: &str) -> Result<()> {
111 let _lock = lock(&self.path.with_extension("lock"))?;
112 let mut data = self.load()?;
113 if data
114 .records
115 .iter()
116 .any(|r| r.message.from == message.from && r.message.msg_id == message.msg_id)
117 {
118 return Ok(());
119 }
120 data.records
123 .retain(|r| r.unread() || now_ms().saturating_sub(r.message.ts) <= MAX_PAST_MS);
124 let mut incoming = Received {
125 message: message.clone(),
126 peer: peer.into(),
127 state: "receiver_stored".into(),
128 superseded_by: None,
129 };
130 if message.task_id.is_some()
131 && message
132 .status
133 .is_some_and(|s| s == TaskStatus::Update || s.is_terminal())
134 {
135 for prior in &mut data.records {
136 let same_task = prior.message.from == message.from
137 && prior.message.reply_to == message.reply_to
138 && prior.message.task_id == message.task_id;
139 if !same_task {
140 continue;
141 }
142 if message.status == Some(TaskStatus::Update)
144 && (prior.message.status.is_some_and(TaskStatus::is_terminal)
145 || (prior.message.status == Some(TaskStatus::Update)
146 && prior.message.ts > message.ts))
147 {
148 incoming.superseded_by = Some(prior.message.msg_id.clone());
149 }
150 if prior.message.status == Some(TaskStatus::Update)
151 && prior.superseded_by.is_none()
152 && (message.status.is_some_and(TaskStatus::is_terminal)
153 || prior.message.ts <= message.ts)
154 {
155 prior.superseded_by = Some(message.msg_id.clone());
156 }
157 }
158 }
159 data.records.push(incoming);
160 if data.records.len() > MAX_RECORDS || serde_json::to_vec(&data.records)?.len() > MAX_BYTES
163 {
164 bail!(
165 "Interlink mailbox is full; messages remain on the broker until space is available"
166 );
167 }
168 self.save(&data)
169 }
170
171 pub fn reserve_notice(&self) -> Result<Option<Notice>> {
172 self.reserve_notice_at(now_ms())
173 }
174
175 fn reserve_notice_at(&self, now: u64) -> Result<Option<Notice>> {
176 let _lock = lock(&self.path.with_extension("lock"))?;
177 let mut data = self.load()?;
178 let messages: Vec<(String, String)> = data
179 .records
180 .iter()
181 .filter(|r| r.unread() && r.message.status != Some(TaskStatus::Update))
182 .map(|r| (r.message.from.clone(), r.message.msg_id.clone()))
183 .collect();
184 if messages.is_empty() {
185 if data.notice.take().is_some() {
186 self.save(&data)?;
187 }
188 return Ok(None);
189 }
190 let mut attempt = 0;
191 if let Some(notice) = &data.notice {
192 let remaining = notice.retry_at.saturating_sub(now);
194 if remaining > 0 && remaining <= MAX_NOTICE_RETRY_MS {
195 return Ok(None);
196 }
197 attempt = notice.attempt.saturating_add(1).min(4);
198 }
199 let delay = (NOTICE_RETRY_MS << attempt).min(MAX_NOTICE_RETRY_MS);
200 let notice = Notice {
201 id: mint_session_id()?,
202 state: "preparing".into(),
203 messages,
204 retry_at: now.saturating_add(delay),
205 attempt,
206 };
207 data.notice = Some(notice.clone());
208 self.save(&data)?;
209 Ok(Some(notice))
210 }
211
212 pub fn finish_notice(&self, id: &str, state: &str) -> Result<()> {
213 let _lock = lock(&self.path.with_extension("lock"))?;
214 let mut data = self.load()?;
215 if let Some(notice) = &mut data.notice
216 && notice.id == id
217 {
218 notice.state = state.into();
219 for record in &mut data.records {
220 if record.unread()
221 && notice.messages.iter().any(|(sender, id)| {
222 sender == &record.message.from && id == &record.message.msg_id
223 })
224 {
225 record.state = state.into();
226 }
227 }
228 self.save(&data)?;
229 }
230 Ok(())
231 }
232
233 pub fn acknowledge(&self, ids: &[(String, String)]) -> Result<usize> {
234 let _lock = lock(&self.path.with_extension("lock"))?;
235 let mut data = self.load()?;
236 let mut acknowledged = 0;
237 let ids: HashSet<_> = ids
238 .iter()
239 .map(|(sender, id)| (sender.as_str(), id.as_str()))
240 .collect();
241 for record in &mut data.records {
242 if record.state != "receiver_acknowledged"
243 && ids.contains(&(record.message.from.as_str(), record.message.msg_id.as_str()))
244 {
245 record.state = "receiver_acknowledged".into();
246 acknowledged += 1;
247 }
248 }
249 if data.notice.as_ref().is_some_and(|notice| {
252 let covered: HashSet<_> = notice
253 .messages
254 .iter()
255 .map(|(sender, id)| (sender.as_str(), id.as_str()))
256 .collect();
257 !data.records.iter().any(|record| {
258 record.unread()
259 && covered
260 .contains(&(record.message.from.as_str(), record.message.msg_id.as_str()))
261 })
262 }) {
263 data.notice = None;
264 }
265 self.save(&data)?;
266 Ok(acknowledged)
267 }
268
269 pub fn receive(&self, _notification_id: Option<&str>, limit: usize) -> Result<Vec<Received>> {
271 let data = self.load()?;
272 let mut received: Vec<_> = data.records.into_iter().filter(Received::unread).collect();
273 received.sort_by_key(|r| r.message.status == Some(TaskStatus::Update));
275 received.truncate(limit);
276 Ok(received)
277 }
278}
279
280#[cfg(test)]
281mod tests {
282 use super::*;
283 use crate::identity::{AgentKey, MessageKind};
284
285 fn message(id: &str, task: &str, status: Option<TaskStatus>, ts: u64) -> SignedMessage {
286 let key = AgentKey::from_b64("AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA").unwrap();
287 key.sign_full(
288 key.id(),
289 id,
290 ts,
291 id,
292 MessageKind::Message,
293 Some(task),
294 status,
295 None,
296 )
297 }
298
299 fn ids(records: &[Received]) -> Vec<(String, String)> {
300 records
301 .iter()
302 .map(|r| (r.message.from.clone(), r.message.msg_id.clone()))
303 .collect()
304 }
305
306 #[test]
307 fn lost_notices_recover_in_every_state_and_backoff_is_bounded() {
308 for state in [
309 "preparing",
310 "host_queued",
311 "inbox_queued",
312 "notification_sent",
313 "delivery_failed",
314 ] {
315 let dir = tempfile::tempdir().unwrap();
316 let path = dir.path().join("mail.json");
317 let mailbox = Mailbox::new(&path);
318 mailbox
319 .retain(&message("a", "task", None, now_ms()), "peer")
320 .unwrap();
321 let mut notice = mailbox.reserve_notice_at(1_000).unwrap().unwrap();
322 mailbox.finish_notice(¬ice.id, state).unwrap();
323 let restarted = Mailbox::new(&path);
324 for _ in 0..10 {
325 assert!(
326 restarted
327 .reserve_notice_at(notice.retry_at - 1)
328 .unwrap()
329 .is_none()
330 );
331 let next = restarted
332 .reserve_notice_at(notice.retry_at)
333 .unwrap()
334 .unwrap();
335 assert_ne!(next.id, notice.id);
336 assert!(next.retry_at - notice.retry_at <= MAX_NOTICE_RETRY_MS);
337 notice = next;
338 }
339 }
340 }
341
342 #[test]
343 fn dropped_fetch_and_partial_ack_leave_remaining_messages_recoverable() {
344 let dir = tempfile::tempdir().unwrap();
345 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
346 for id in ["a", "b"] {
347 mailbox
348 .retain(&message(id, "task", None, now_ms()), "peer")
349 .unwrap();
350 }
351 let first = mailbox.reserve_notice_at(1_000).unwrap().unwrap();
352 let fetched = mailbox.receive(None, 1).unwrap();
353 assert_eq!(
354 mailbox.receive(None, 1).unwrap()[0].message.msg_id,
355 fetched[0].message.msg_id
356 );
357 mailbox
358 .retain(&message("c", "task", None, now_ms()), "peer")
359 .unwrap();
360 assert_eq!(mailbox.acknowledge(&ids(&fetched)).unwrap(), 1);
361 assert_eq!(mailbox.acknowledge(&ids(&fetched)).unwrap(), 0);
362 let retry = mailbox.reserve_notice_at(first.retry_at).unwrap().unwrap();
363 mailbox.finish_notice(&first.id, "host_queued").unwrap();
364 assert_eq!(mailbox.load().unwrap().notice.unwrap().id, retry.id);
365 assert_eq!(mailbox.receive(None, 20).unwrap().len(), 2);
366 mailbox
367 .acknowledge(&ids(&mailbox.records().unwrap()))
368 .unwrap();
369 assert!(mailbox.reserve_notice_at(retry.retry_at).unwrap().is_none());
370 }
371
372 #[test]
373 fn history_consumption_without_notice_id_unblocks_future_messages() {
374 let dir = tempfile::tempdir().unwrap();
375 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
376 mailbox
377 .retain(&message("a", "task", None, now_ms()), "peer")
378 .unwrap();
379 let old = mailbox.reserve_notice().unwrap().unwrap();
380 mailbox.finish_notice(&old.id, "notification_sent").unwrap();
381 mailbox
382 .retain(&message("b", "task", None, now_ms()), "peer")
383 .unwrap();
384 let fetched = mailbox.receive(None, 1).unwrap();
385 mailbox.acknowledge(&ids(&fetched)).unwrap();
386 let next = mailbox.reserve_notice().unwrap().unwrap();
387 assert_ne!(next.id, old.id);
388 assert_eq!(
389 mailbox.receive(Some(&old.id), 20).unwrap()[0]
390 .message
391 .msg_id,
392 "b"
393 );
394 }
395
396 #[test]
397 fn old_mailbox_notices_and_backward_clock_jumps_do_not_block_recovery() {
398 let dir = tempfile::tempdir().unwrap();
399 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
400 mailbox
401 .retain(&message("a", "task", None, now_ms()), "peer")
402 .unwrap();
403 let first = mailbox.reserve_notice().unwrap().unwrap();
404 let mut legacy = serde_json::to_value(mailbox.load().unwrap()).unwrap();
405 legacy["notice"].as_object_mut().unwrap().remove("retry_at");
406 legacy["notice"].as_object_mut().unwrap().remove("attempt");
407 atomic_write(&mailbox.path, &serde_json::to_vec(&legacy).unwrap()).unwrap();
408 let restored = mailbox.reserve_notice().unwrap().unwrap();
409 assert_ne!(first.id, restored.id);
410 assert!(mailbox.reserve_notice_at(0).unwrap().is_some());
411 }
412
413 #[test]
414 fn history_consumption_survives_restart_and_an_already_queued_notice() {
415 let dir = tempfile::tempdir().unwrap();
416 let path = dir.path().join("mail.json");
417 let mailbox = Mailbox::new(&path);
418 let msg = message("one", "task", None, now_ms());
419 mailbox.retain(&msg, "peer").unwrap();
420 let notice = mailbox.reserve_notice().unwrap().unwrap();
421 mailbox.finish_notice(¬ice.id, "host_queued").unwrap();
422 mailbox
423 .acknowledge(&[(msg.from.clone(), "one".into())])
424 .unwrap();
425 let restarted = Mailbox::new(&path);
426 restarted.retain(&msg, "peer").unwrap();
427 assert_eq!(restarted.records().unwrap().len(), 1);
428 assert!(restarted.reserve_notice().unwrap().is_none());
429 assert!(restarted.receive(Some(¬ice.id), 20).unwrap().is_empty());
430 assert!(restarted.reserve_notice().unwrap().is_none());
431 assert_eq!(
432 restarted.records().unwrap()[0].state,
433 "receiver_acknowledged"
434 );
435 }
436
437 #[test]
438 fn progress_is_quiet_and_superseded_but_questions_and_failures_survive() {
439 let dir = tempfile::tempdir().unwrap();
440 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
441 let ts = now_ms();
442 for (id, status, offset) in [
443 ("old", TaskStatus::Update, 0),
444 ("new", TaskStatus::Update, 2),
445 ("delayed", TaskStatus::Update, 1),
446 ] {
447 mailbox
448 .retain(&message(id, "task", Some(status), ts + offset), "peer")
449 .unwrap();
450 }
451 assert!(mailbox.reserve_notice().unwrap().is_none());
452 mailbox
453 .retain(
454 &message("question", "task", Some(TaskStatus::NeedsInput), ts + 3),
455 "peer",
456 )
457 .unwrap();
458 mailbox
459 .retain(
460 &message("failure", "task", Some(TaskStatus::Failed), ts + 4),
461 "peer",
462 )
463 .unwrap();
464 mailbox
465 .retain(
466 &message("after-terminal", "task", Some(TaskStatus::Update), ts + 5),
467 "peer",
468 )
469 .unwrap();
470 let notice = mailbox.reserve_notice().unwrap().unwrap();
471 let records = mailbox.receive(Some(¬ice.id), 20).unwrap();
472 assert_eq!(
473 records
474 .iter()
475 .map(|r| r.message.msg_id.as_str())
476 .collect::<Vec<_>>(),
477 ["question", "failure"]
478 );
479 assert_eq!(mailbox.records().unwrap().len(), 6);
480 }
481
482 #[test]
483 fn coalescing_is_scoped_to_sender_session_and_task() {
484 let dir = tempfile::tempdir().unwrap();
485 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
486 let ts = now_ms();
487 let mut a = message("a", "task", Some(TaskStatus::Update), ts);
488 a.reply_to = Some("session-a".into());
489 let mut b = message("b", "task", Some(TaskStatus::Result), ts + 1);
490 b.reply_to = Some("session-b".into());
491 mailbox.retain(&a, "peer").unwrap();
492 mailbox.retain(&b, "peer").unwrap();
493 mailbox
494 .retain(
495 &message("c", "other", Some(TaskStatus::Result), ts + 2),
496 "peer",
497 )
498 .unwrap();
499 assert_eq!(mailbox.receive(None, 20).unwrap().len(), 3);
500 }
501
502 #[test]
503 fn notification_races_do_not_regress_acknowledgement_or_clear_new_wakeups() {
504 let dir = tempfile::tempdir().unwrap();
505 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
506 mailbox
507 .retain(&message("a", "task", None, now_ms()), "peer")
508 .unwrap();
509 let first = mailbox.reserve_notice().unwrap().unwrap();
510 let fetched = mailbox.receive(Some(&first.id), 20).unwrap();
511 mailbox.acknowledge(&ids(&fetched)).unwrap();
512 mailbox
513 .retain(&message("b", "task", None, now_ms()), "peer")
514 .unwrap();
515 let second = mailbox.reserve_notice().unwrap().unwrap();
516 mailbox.finish_notice(&first.id, "host_queued").unwrap();
517 let fetched = mailbox.receive(Some(&first.id), 20).unwrap();
518 assert_eq!(mailbox.load().unwrap().notice.unwrap().id, second.id);
519 mailbox.acknowledge(&ids(&fetched)).unwrap();
520 assert!(
521 mailbox
522 .records()
523 .unwrap()
524 .iter()
525 .all(|r| r.state == "receiver_acknowledged")
526 );
527 }
528
529 #[test]
530 fn only_one_notifier_and_acknowledgement_is_idempotent_between_handles() {
531 let dir = tempfile::tempdir().unwrap();
532 let path = dir.path().join("mail.json");
533 let a = Mailbox::new(&path);
534 let b = Mailbox::new(&path);
535 let guard = a.notifier_lock().unwrap().unwrap();
536 assert!(b.notifier_lock().unwrap().is_none());
537 a.retain(&message("a", "task", None, now_ms()), "peer")
538 .unwrap();
539 let workers: Vec<_> = (0..2)
540 .map(|_| {
541 let path = path.clone();
542 std::thread::spawn(move || {
543 let mailbox = Mailbox::new(&path);
544 let records = mailbox.records().unwrap();
545 mailbox.acknowledge(&ids(&records)).unwrap()
546 })
547 })
548 .collect();
549 assert_eq!(
550 workers
551 .into_iter()
552 .map(|w| w.join().unwrap())
553 .sum::<usize>(),
554 1
555 );
556 drop(guard);
557 assert!(b.notifier_lock().unwrap().is_some());
558 }
559 #[test]
560 fn full_mailbox_can_be_consumed_and_expired_tombstones_make_room() {
561 let dir = tempfile::tempdir().unwrap();
562 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
563 let base = message("base", "task", None, now_ms());
564 let mut data = Data::default();
565 for i in 0..MAX_RECORDS {
566 let mut msg = base.clone();
567 msg.msg_id = i.to_string();
568 data.records.push(Received {
569 message: msg,
570 peer: "peer".into(),
571 state: "receiver_stored".into(),
572 superseded_by: None,
573 });
574 }
575 mailbox.save(&data).unwrap();
576 let next = message("next", "task", None, now_ms());
577 assert!(mailbox.retain(&next, "peer").is_err());
578 assert_eq!(mailbox.records().unwrap().len(), MAX_RECORDS);
579 let notice = mailbox.reserve_notice().unwrap().unwrap();
580 assert_eq!(
581 mailbox
582 .receive(Some(¬ice.id), MAX_RECORDS)
583 .unwrap()
584 .len(),
585 MAX_RECORDS
586 );
587 mailbox
588 .acknowledge(&ids(&mailbox.records().unwrap()))
589 .unwrap();
590 let mut data = mailbox.load().unwrap();
591 for record in &mut data.records {
592 record.message.ts = now_ms() - MAX_PAST_MS - 1;
593 }
594 mailbox.save(&data).unwrap();
595 mailbox.retain(&next, "peer").unwrap();
596 assert_eq!(mailbox.records().unwrap().len(), 1);
597 }
598
599 #[test]
600 fn history_acknowledgement_is_scoped_by_identity_and_questions_take_priority() {
601 let dir = tempfile::tempdir().unwrap();
602 let mailbox = Mailbox::new(&dir.path().join("mail.json"));
603 let progress = message("same-id", "task", Some(TaskStatus::Update), now_ms());
604 let other = AgentKey::generate().unwrap();
605 let question = other.sign_full(
606 other.id(),
607 "question",
608 now_ms(),
609 "same-id",
610 MessageKind::Message,
611 Some("task"),
612 Some(TaskStatus::NeedsInput),
613 None,
614 );
615 mailbox.retain(&progress, "peer").unwrap();
616 mailbox.retain(&question, "other").unwrap();
617 mailbox
618 .acknowledge(&[(progress.from, progress.msg_id)])
619 .unwrap();
620 let received = mailbox.receive(None, 1).unwrap();
621 assert_eq!(received.len(), 1);
622 assert_eq!(received[0].peer, "other");
623 mailbox.acknowledge(&ids(&received)).unwrap();
624 mailbox
625 .retain(
626 &message("progress", "other-task", Some(TaskStatus::Update), now_ms()),
627 "peer",
628 )
629 .unwrap();
630 mailbox
631 .retain(&message("request", "task", None, now_ms()), "peer")
632 .unwrap();
633 assert_eq!(
634 mailbox.receive(None, 1).unwrap()[0].message.msg_id,
635 "request"
636 );
637 }
638}