Skip to main content

io_imap/
watch.rs

1//! IMAP single-mailbox watcher: IDLE (RFC 2177) for the wake signal,
2//! SELECT (QRESYNC) (RFC 7162) for UID-keyed deltas.
3//!
4//! QRESYNC: <https://www.rfc-editor.org/rfc/rfc7162>
5//!
6//! ```text
7//! SELECT (CONDSTORE) → FETCH 1:* (UID FLAGS) [seed shadow]
8//!     → IDLE → SELECT (QRESYNC) → emit deltas → IDLE → ...
9//! ```
10//!
11//! Connection is dedicated. Flip the shared [`AtomicBool`] to wind
12//! down cleanly.
13//!
14//! # Example
15//!
16//! ```rust,no_run
17//! use core::sync::atomic::AtomicBool;
18//! use std::{
19//!     io::{Read, Write},
20//!     net::TcpStream,
21//!     sync::Arc,
22//! };
23//!
24//! use io_imap::{
25//!     codec::fragmentizer::Fragmentizer,
26//!     coroutine::{ImapCoroutine, ImapCoroutineState},
27//!     types::response::Capability,
28//!     watch::{ImapMailboxWatch, ImapMailboxWatchYield},
29//! };
30//!
31//! // Ready stream needed (TCP-connected, TLS-negotiated, IMAP-authenticated)
32//! let mut stream = TcpStream::connect("localhost:143").unwrap();
33//!
34//! let mut fragmentizer = Fragmentizer::new(50 * 1024 * 1024);
35//! let mut buf = [0u8; 4096];
36//!
37//! let capability = [Capability::QResync];
38//! let mailbox = "INBOX".try_into().unwrap();
39//! let shutdown = Arc::new(AtomicBool::new(false));
40//! let mut coroutine =
41//!     ImapMailboxWatch::new(&capability, mailbox, shutdown.clone()).unwrap();
42//! let mut arg = None;
43//!
44//! loop {
45//!     match coroutine.resume(&mut fragmentizer, arg.take()) {
46//!         ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(bytes)) => {
47//!             stream.write_all(&bytes).unwrap();
48//!         }
49//!         ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead) => {
50//!             let n = stream.read(&mut buf).unwrap();
51//!             arg = Some(&buf[..n]);
52//!         }
53//!         ImapCoroutineState::Yielded(ImapMailboxWatchYield::Event(event)) => {
54//!             println!("{event:?}");
55//!         }
56//!         ImapCoroutineState::Complete(Ok(())) => break,
57//!         ImapCoroutineState::Complete(Err(err)) => panic!("{err}"),
58//!     }
59//! }
60//! ```
61
62use core::{
63    mem,
64    num::{NonZeroU32, NonZeroU64},
65    sync::atomic::{AtomicBool, Ordering},
66};
67
68use alloc::{
69    collections::{BTreeMap, VecDeque},
70    string::String,
71    sync::Arc,
72    vec,
73    vec::Vec,
74};
75
76use imap_codec::{
77    fragmentizer::Fragmentizer,
78    imap_types::{
79        command::SelectParameter,
80        core::{Atom, Vec1},
81        extensions::enable::CapabilityEnable,
82        fetch::{MacroOrMessageDataItemNames, MessageDataItem, MessageDataItemName},
83        flag::{Flag, FlagFetch},
84        mailbox::Mailbox,
85        response::Capability,
86        sequence::SequenceSet,
87    },
88};
89use log::{debug, trace};
90use thiserror::Error;
91
92use crate::{
93    coroutine::*,
94    rfc2177::idle::{ImapIdle, ImapIdleError, ImapIdleOptions, ImapIdleYield},
95    rfc3501::{
96        fetch::{ImapMessageFetch, ImapMessageFetchError, ImapMessageFetchOptions},
97        select::{
98            ImapMailboxSelect, ImapMailboxSelectData, ImapMailboxSelectError,
99            ImapMailboxSelectOptions,
100        },
101    },
102    rfc5161::enable::{ImapExtensionEnable, ImapExtensionEnableError},
103};
104
105/// UID-keyed mailbox change emitted by the watcher.
106///
107/// `FlagsAdded`/`FlagsRemoved` are pre-diffed against the internal
108/// shadow; each `flags` vector lists only the changed flags.
109#[derive(Clone, Debug)]
110pub enum ImapMailboxWatchEvent {
111    /// A message appeared in the mailbox.
112    EnvelopeAdded {
113        /// The UID of the new message.
114        uid: NonZeroU32,
115        /// The FETCH items announcing the message.
116        items: Vec<MessageDataItem<'static>>,
117    },
118    /// Flags were set on an existing message.
119    FlagsAdded {
120        /// The UID of the changed message.
121        uid: NonZeroU32,
122        /// The flags that were added.
123        flags: Vec<Flag<'static>>,
124    },
125    /// Flags were cleared on an existing message.
126    FlagsRemoved {
127        /// The UID of the changed message.
128        uid: NonZeroU32,
129        /// The flags that were removed.
130        flags: Vec<Flag<'static>>,
131    },
132    /// A message left the mailbox (expunged or moved away).
133    EnvelopeRemoved {
134        /// The UID of the removed message.
135        uid: NonZeroU32,
136    },
137}
138
139/// Failure causes during the mailbox watch flow.
140#[derive(Debug, Error)]
141pub enum ImapMailboxWatchError {
142    /// The capability list given to `new` lacks QRESYNC.
143    #[error("IMAP server does not advertise QRESYNC")]
144    QresyncUnsupported,
145    /// The SELECT response carried no UIDVALIDITY, so deltas cannot be
146    /// keyed safely.
147    #[error("IMAP server did not return UIDVALIDITY in SELECT response")]
148    MissingUidValidity,
149    /// The SELECT response carried no HIGHESTMODSEQ, so there is no
150    /// resync point.
151    #[error("IMAP server did not return HIGHESTMODSEQ in SELECT response")]
152    MissingHighestModSeq,
153    /// The baseline `1:*` sequence set failed to parse.
154    #[error("Invalid `1:*` sequence set: {0}")]
155    InvalidSequenceSet(String),
156    /// The initial or QRESYNC SELECT failed.
157    #[error("IMAP SELECT error")]
158    Select(#[from] ImapMailboxSelectError),
159    /// The baseline FETCH failed.
160    #[error("IMAP FETCH error")]
161    Fetch(#[from] ImapMessageFetchError),
162    /// The IDLE wake-loop failed.
163    #[error("IMAP IDLE error")]
164    Idle(#[from] ImapIdleError),
165    /// The ENABLE QRESYNC round failed.
166    #[error("IMAP ENABLE error")]
167    Enable(#[from] ImapExtensionEnableError),
168}
169
170/// Yield variants from the mailbox watcher.
171#[derive(Debug)]
172pub enum ImapMailboxWatchYield {
173    /// The caller reads from its stream and resumes with the bytes.
174    WantsRead,
175    /// The caller writes the given bytes to its stream and resumes.
176    WantsWrite(Vec<u8>),
177    /// A mailbox change to consume; the watcher keeps running.
178    Event(ImapMailboxWatchEvent),
179}
180
181enum State {
182    EnableQresync(ImapExtensionEnable),
183    SelectInitial(ImapMailboxSelect),
184    FetchBaseline(ImapMessageFetch),
185    BeginIdle,
186    Idle(ImapIdle),
187    SelectQresync(ImapMailboxSelect),
188    EmitDeltas,
189    Terminal,
190}
191
192/// I/O-free IDLE+QRESYNC mailbox watcher.
193pub struct ImapMailboxWatch {
194    state: State,
195    shutdown: Arc<AtomicBool>,
196    idle_done: Arc<AtomicBool>,
197    idle_saw_data: bool,
198    mailbox: Mailbox<'static>,
199    uid_validity: Option<NonZeroU32>,
200    highest_mod_seq: u64,
201    shadow: BTreeMap<NonZeroU32, Vec<Flag<'static>>>,
202    pending: VecDeque<ImapMailboxWatchEvent>,
203}
204
205impl ImapMailboxWatch {
206    /// Errors with `QresyncUnsupported` when `capability` lacks QRESYNC.
207    pub fn new(
208        capability: &[Capability<'static>],
209        mailbox: Mailbox<'static>,
210        shutdown: Arc<AtomicBool>,
211    ) -> Result<Self, ImapMailboxWatchError> {
212        if !capability.contains(&Capability::QResync) {
213            return Err(ImapMailboxWatchError::QresyncUnsupported);
214        }
215
216        // NOTE: RFC 7162 §3.1: QRESYNC implies CONDSTORE, but pass
217        // both since some servers only echo CONDSTORE in ENABLED.
218        let condstore = CapabilityEnable::CondStore;
219        // NOTE: QRESYNC is not in the typed enum, route via Atom.
220        let qresync = CapabilityEnable::from(
221            Atom::try_from("QRESYNC").expect("`QRESYNC` is a syntactically valid IMAP atom"),
222        );
223        let capabilities =
224            Vec1::try_from(vec![condstore, qresync]).expect("two capabilities is non-empty");
225        let enable = ImapExtensionEnable::new(capabilities);
226
227        Ok(Self {
228            state: State::EnableQresync(enable),
229            shutdown,
230            idle_done: Arc::new(AtomicBool::new(false)),
231            idle_saw_data: false,
232            mailbox,
233            uid_validity: None,
234            highest_mod_seq: 0,
235            shadow: BTreeMap::new(),
236            pending: VecDeque::new(),
237        })
238    }
239
240    fn compute_deltas(&mut self, data: &ImapMailboxSelectData) {
241        for uid in &data.vanished_earlier {
242            if self.shadow.remove(uid).is_some() {
243                self.pending
244                    .push_back(ImapMailboxWatchEvent::EnvelopeRemoved { uid: *uid });
245            }
246        }
247
248        for fetch in &data.changed {
249            let items_vec: Vec<MessageDataItem<'static>> =
250                fetch.items.clone().into_inner().into_iter().collect();
251            let (uid_opt, new_flags) = extract_uid_flags(&items_vec);
252            let Some(uid) = uid_opt else {
253                continue;
254            };
255
256            match self.shadow.get(&uid).cloned() {
257                None => {
258                    self.shadow.insert(uid, new_flags);
259                    self.pending
260                        .push_back(ImapMailboxWatchEvent::EnvelopeAdded {
261                            uid,
262                            items: items_vec,
263                        });
264                }
265                Some(old_flags) => {
266                    let added: Vec<Flag<'static>> = new_flags
267                        .iter()
268                        .filter(|f| !old_flags.contains(f))
269                        .cloned()
270                        .collect();
271                    let removed: Vec<Flag<'static>> = old_flags
272                        .iter()
273                        .filter(|f| !new_flags.contains(f))
274                        .cloned()
275                        .collect();
276                    self.shadow.insert(uid, new_flags);
277                    if !added.is_empty() {
278                        self.pending
279                            .push_back(ImapMailboxWatchEvent::FlagsAdded { uid, flags: added });
280                    }
281                    if !removed.is_empty() {
282                        self.pending.push_back(ImapMailboxWatchEvent::FlagsRemoved {
283                            uid,
284                            flags: removed,
285                        });
286                    }
287                }
288            }
289        }
290    }
291}
292
293impl ImapCoroutine for ImapMailboxWatch {
294    type Yield = ImapMailboxWatchYield;
295    type Return = Result<(), ImapMailboxWatchError>;
296
297    fn resume(
298        &mut self,
299        fragmentizer: &mut Fragmentizer,
300        mut arg: Option<&[u8]>,
301    ) -> ImapCoroutineState<Self::Yield, Self::Return> {
302        if self.shutdown.load(Ordering::SeqCst) {
303            self.idle_done.store(true, Ordering::SeqCst);
304        }
305
306        loop {
307            let state = mem::replace(&mut self.state, State::Terminal);
308
309            match state {
310                State::EnableQresync(mut enable) => match enable.resume(fragmentizer, arg.take()) {
311                    ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
312                        self.state = State::EnableQresync(enable);
313                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
314                    }
315                    ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
316                        self.state = State::EnableQresync(enable);
317                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
318                            bytes,
319                        ));
320                    }
321                    ImapCoroutineState::Complete(Ok(enabled)) => {
322                        debug!("enabled qresync");
323                        trace!("{enabled:?}");
324                        let parameters = vec![SelectParameter::CondStore];
325                        let select = ImapMailboxSelect::new(
326                            self.mailbox.clone(),
327                            ImapMailboxSelectOptions { parameters },
328                        );
329                        self.state = State::SelectInitial(select);
330                    }
331                    ImapCoroutineState::Complete(Err(err)) => {
332                        return ImapCoroutineState::Complete(Err(err.into()));
333                    }
334                },
335
336                State::SelectInitial(mut select) => match select.resume(fragmentizer, arg.take()) {
337                    ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
338                        self.state = State::SelectInitial(select);
339                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
340                    }
341                    ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
342                        self.state = State::SelectInitial(select);
343                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
344                            bytes,
345                        ));
346                    }
347                    ImapCoroutineState::Complete(Ok(data)) => {
348                        let Some(uid_validity) = data.uid_validity else {
349                            return ImapCoroutineState::Complete(Err(
350                                ImapMailboxWatchError::MissingUidValidity,
351                            ));
352                        };
353                        let Some(highest_mod_seq) = data.highest_mod_seq else {
354                            return ImapCoroutineState::Complete(Err(
355                                ImapMailboxWatchError::MissingHighestModSeq,
356                            ));
357                        };
358
359                        self.uid_validity = Some(uid_validity);
360                        self.highest_mod_seq = highest_mod_seq;
361                        debug!("selected mailbox with condstore");
362                        trace!("uid_validity: {uid_validity}");
363                        trace!("highest_mod_seq: {highest_mod_seq}");
364
365                        let sequence_set: SequenceSet = match "1:*".try_into() {
366                            Ok(s) => s,
367                            Err(_) => {
368                                return ImapCoroutineState::Complete(Err(
369                                    ImapMailboxWatchError::InvalidSequenceSet("1:*".into()),
370                                ));
371                            }
372                        };
373                        let item_names = MacroOrMessageDataItemNames::MessageDataItemNames(vec![
374                            MessageDataItemName::Uid,
375                            MessageDataItemName::Flags,
376                        ]);
377                        let fetch = ImapMessageFetch::new(
378                            sequence_set,
379                            item_names,
380                            ImapMessageFetchOptions::default(),
381                        );
382                        self.state = State::FetchBaseline(fetch);
383                    }
384                    ImapCoroutineState::Complete(Err(err)) => {
385                        return ImapCoroutineState::Complete(Err(err.into()));
386                    }
387                },
388
389                State::FetchBaseline(mut fetch) => match fetch.resume(fragmentizer, arg.take()) {
390                    ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
391                        self.state = State::FetchBaseline(fetch);
392                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
393                    }
394                    ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
395                        self.state = State::FetchBaseline(fetch);
396                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
397                            bytes,
398                        ));
399                    }
400                    ImapCoroutineState::Complete(Ok(data)) => {
401                        for (_seq, items) in data {
402                            let items_vec = items.into_inner();
403                            if let (Some(uid), flags) = extract_uid_flags(&items_vec) {
404                                self.shadow.insert(uid, flags);
405                            }
406                        }
407                        debug!("seeded baseline shadow");
408                        trace!("uids: {}", self.shadow.len());
409                        self.state = State::BeginIdle;
410                    }
411                    ImapCoroutineState::Complete(Err(err)) => {
412                        return ImapCoroutineState::Complete(Err(err.into()));
413                    }
414                },
415
416                State::BeginIdle => {
417                    if self.shutdown.load(Ordering::SeqCst) {
418                        return ImapCoroutineState::Complete(Ok(()));
419                    }
420
421                    self.idle_done.store(false, Ordering::SeqCst);
422                    self.idle_saw_data = false;
423                    let idle = ImapIdle::new(self.idle_done.clone(), ImapIdleOptions::default());
424                    self.state = State::Idle(idle);
425                }
426
427                State::Idle(mut idle) => match idle.resume(fragmentizer, arg.take()) {
428                    ImapCoroutineState::Yielded(ImapIdleYield::Event(_)) => {
429                        debug!("idle saw untagged data");
430                        self.idle_saw_data = true;
431                        self.idle_done.store(true, Ordering::SeqCst);
432                        self.state = State::Idle(idle);
433                    }
434                    ImapCoroutineState::Yielded(ImapIdleYield::WantsRead) => {
435                        self.state = State::Idle(idle);
436                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
437                    }
438                    ImapCoroutineState::Yielded(ImapIdleYield::WantsWrite(bytes)) => {
439                        self.state = State::Idle(idle);
440                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
441                            bytes,
442                        ));
443                    }
444                    ImapCoroutineState::Complete(Ok(())) => {
445                        if self.shutdown.load(Ordering::SeqCst) {
446                            return ImapCoroutineState::Complete(Ok(()));
447                        }
448
449                        if self.idle_saw_data {
450                            // NOTE: uid_validity is set by SelectInitial.
451                            let uid_validity = self.uid_validity.unwrap();
452                            let modseq = NonZeroU64::new(self.highest_mod_seq)
453                                .unwrap_or_else(|| NonZeroU64::new(1).expect("1 is non-zero"));
454                            let parameters = vec![SelectParameter::QResync {
455                                uid_validity,
456                                mod_sequence_value: modseq,
457                                known_uids: None,
458                                seq_match_data: None,
459                            }];
460                            let select = ImapMailboxSelect::new(
461                                self.mailbox.clone(),
462                                ImapMailboxSelectOptions { parameters },
463                            );
464                            self.state = State::SelectQresync(select);
465                        } else {
466                            debug!("idle timed out with no data, restarting");
467                            self.state = State::BeginIdle;
468                        }
469                    }
470                    ImapCoroutineState::Complete(Err(err)) => {
471                        return ImapCoroutineState::Complete(Err(err.into()));
472                    }
473                },
474
475                State::SelectQresync(mut select) => match select.resume(fragmentizer, arg.take()) {
476                    ImapCoroutineState::Yielded(ImapYield::WantsRead) => {
477                        self.state = State::SelectQresync(select);
478                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsRead);
479                    }
480                    ImapCoroutineState::Yielded(ImapYield::WantsWrite(bytes)) => {
481                        self.state = State::SelectQresync(select);
482                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::WantsWrite(
483                            bytes,
484                        ));
485                    }
486                    ImapCoroutineState::Complete(Ok(data)) => {
487                        self.compute_deltas(&data);
488                        if let Some(new_modseq) = data.highest_mod_seq {
489                            self.highest_mod_seq = new_modseq;
490                        }
491                        self.state = State::EmitDeltas;
492                    }
493                    ImapCoroutineState::Complete(Err(err)) => {
494                        return ImapCoroutineState::Complete(Err(err.into()));
495                    }
496                },
497
498                State::EmitDeltas => {
499                    if let Some(event) = self.pending.pop_front() {
500                        self.state = State::EmitDeltas;
501                        return ImapCoroutineState::Yielded(ImapMailboxWatchYield::Event(event));
502                    }
503                    self.state = State::BeginIdle;
504                }
505
506                State::Terminal => {
507                    self.state = State::Terminal;
508                    return ImapCoroutineState::Complete(Ok(()));
509                }
510            }
511        }
512    }
513}
514
515/// Extract the UID and flag list from a single FETCH; preserves wire
516/// order, drops non-`Flag` variants of [`FlagFetch`].
517fn extract_uid_flags(
518    items: &[MessageDataItem<'static>],
519) -> (Option<NonZeroU32>, Vec<Flag<'static>>) {
520    let mut uid = None;
521    let mut flags = Vec::new();
522    for item in items {
523        match item {
524            MessageDataItem::Uid(u) => uid = Some(*u),
525            MessageDataItem::Flags(fs) => {
526                flags = fs
527                    .iter()
528                    .filter_map(|f| match f {
529                        FlagFetch::Flag(flag) => Some(flag.clone()),
530                        _ => None,
531                    })
532                    .collect();
533            }
534            _ => {}
535        }
536    }
537    (uid, flags)
538}