Skip to main content

interlink/
mailbox.rs

1//! Shared consumption and notification state for every host adapter.
2
3use 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        // Keep replay tombstones throughout the gate's acceptance window. Unread
121        // requests are never evicted just to make space for a newer message.
122        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                // A delayed progress message must not resurrect a completed task.
143                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        // Later state changes and notification IDs can grow the metadata. Never
161        // apply the admission cap to consumption or recovery writes.
162        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            // A backward clock jump must not leave a reservation stuck indefinitely.
193            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        // Free the reservation only when none of its covered messages still need
250        // attention. New arrivals are left unread and can immediately wake the host.
251        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    /// Notification IDs are advisory: stale notices always fetch current unread state.
270    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        // Stable ordering keeps questions and failures ahead of routine progress.
274        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(&notice.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(&notice.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(&notice.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(&notice.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(&notice.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}