Skip to main content

io_imap/rfc2177/
idle.rs

1//! IMAP IDLE coroutine yielding mailbox change events.
2//!
3//! # Example
4//!
5//! ```rust,no_run
6//! use core::sync::atomic::AtomicBool;
7//! use std::{
8//!     io::{Read, Write},
9//!     net::TcpStream,
10//!     sync::Arc,
11//! };
12//!
13//! use io_imap::{
14//!     codec::fragmentizer::Fragmentizer,
15//!     coroutine::{ImapCoroutine, ImapCoroutineState},
16//!     rfc2177::idle::{ImapIdle, ImapIdleOptions, ImapIdleYield},
17//! };
18//!
19//! // Ready stream needed (TCP-connected, TLS-negotiated, IMAP-authenticated)
20//! let mut stream = TcpStream::connect("localhost:143").unwrap();
21//!
22//! let mut fragmentizer = Fragmentizer::new(50 * 1024 * 1024);
23//! let mut buf = [0u8; 4096];
24//!
25//! let shutdown = Arc::new(AtomicBool::new(false));
26//! let mut coroutine =
27//!     ImapIdle::new(shutdown.clone(), ImapIdleOptions::default());
28//! let mut arg = None;
29//!
30//! loop {
31//!     match coroutine.resume(&mut fragmentizer, arg.take()) {
32//!         ImapCoroutineState::Yielded(ImapIdleYield::WantsWrite(bytes)) => {
33//!             stream.write_all(&bytes).unwrap();
34//!         }
35//!         ImapCoroutineState::Yielded(ImapIdleYield::WantsRead) => {
36//!             let n = stream.read(&mut buf).unwrap();
37//!             arg = Some(&buf[..n]);
38//!         }
39//!         ImapCoroutineState::Yielded(ImapIdleYield::Event(event)) => {
40//!             println!("{event:?}");
41//!         }
42//!         ImapCoroutineState::Complete(Ok(())) => break,
43//!         ImapCoroutineState::Complete(Err(err)) => panic!("{err}"),
44//!     }
45//! }
46//! ```
47
48use core::{
49    fmt, mem,
50    sync::atomic::{AtomicBool, Ordering},
51    time::Duration,
52};
53
54use alloc::{boxed::Box, string::String, string::ToString, sync::Arc, vec::Vec};
55
56#[cfg(feature = "client")]
57use std::time::Instant;
58
59use imap_codec::{
60    CommandCodec, IdleDoneCodec, ResponseCodec,
61    fragmentizer::{DecodeMessageError, FragmentInfo, Fragmentizer},
62    imap_types::{
63        IntoStatic,
64        command::{Command, CommandBody},
65        core::TagGenerator,
66        extensions::idle::IdleDone,
67        response::{Bye, Data, Response, Status, StatusBody, StatusKind, Tagged},
68        secret::Secret,
69        utils::escape_byte_string,
70    },
71};
72use log::{debug, trace};
73use thiserror::Error;
74
75use crate::{coroutine::*, imap_try, send::*};
76
77/// Default refresh interval: 29 s survives NAT middle-boxes and stays
78/// well under the 29-minute RFC 2177 ยง3 cap.
79#[cfg(feature = "client")]
80const IDLE_DEFAULT_TIMEOUT: Duration = Duration::from_secs(29);
81
82/// Failure causes during the IMAP IDLE flow.
83#[derive(Clone, Debug, Error)]
84pub enum ImapIdleError {
85    /// The server rejected the IDLE command with a NO response.
86    #[error("IMAP IDLE failed: NO {0}")]
87    No(String),
88    /// The server rejected the IDLE command with a BAD response.
89    #[error("IMAP IDLE failed: BAD {0}")]
90    Bad(String),
91    /// The server closed the connection with a BYE response.
92    #[error("IMAP IDLE failed: BYE {0}")]
93    Bye(String),
94    /// The server sent a tagged OK before the continuation request.
95    #[error("IMAP IDLE failed: server returned a tagged response before the continuation request")]
96    UnexpectedTagged,
97    /// The server never sent the continuation request after IDLE.
98    #[error("IMAP IDLE failed: server did not send the expected continuation request")]
99    ExpectedContinuationRequest,
100    /// The server never answered DONE with a tagged response.
101    #[error("IMAP IDLE failed: server did not return a tagged response to DONE")]
102    MissingTagged,
103    /// The stream reached EOF while waiting for server responses.
104    #[error("IMAP IDLE failed: reached unexpected EOF on stream")]
105    Eof,
106    /// A server response could not be decoded.
107    #[error("IMAP IDLE failed: decode response error")]
108    DecodingFailure(Secret<Box<[u8]>>),
109    /// A server response was flagged as poisoned by the fragmentizer.
110    #[error("IMAP IDLE failed: parse response error: message is poisoned")]
111    MessageIsPoisoned(Secret<Box<[u8]>>),
112    /// A server response exceeded the fragmentizer's maximum size.
113    #[error("IMAP IDLE failed: parse response error: message is too long")]
114    MessageTooLong(Secret<Box<[u8]>>),
115    /// The underlying send sub-coroutine failed.
116    #[error("IMAP IDLE failed: {0}")]
117    Send(#[from] ImapSendError),
118}
119
120/// Batch of unilateral untagged responses received during an IDLE.
121#[derive(Debug)]
122pub struct ImapIdleEvent {
123    /// Untagged status responses received while idling.
124    pub untagged: Vec<StatusBody<'static>>,
125    /// Mailbox data updates received while idling, such as EXISTS,
126    /// EXPUNGE or FETCH.
127    pub data: Vec<Data<'static>>,
128}
129
130/// Yield variants from the IDLE coroutine.
131#[derive(Debug)]
132pub enum ImapIdleYield {
133    /// The caller reads bytes from the stream and resumes with them.
134    WantsRead,
135    /// The caller writes the given bytes to the stream and resumes.
136    WantsWrite(Vec<u8>),
137    /// A mailbox change event to consume; the coroutine keeps running.
138    Event(ImapIdleEvent),
139}
140
141impl From<ImapYield> for ImapIdleYield {
142    fn from(y: ImapYield) -> Self {
143        match y {
144            ImapYield::WantsRead => ImapIdleYield::WantsRead,
145            ImapYield::WantsWrite(bytes) => ImapIdleYield::WantsWrite(bytes),
146        }
147    }
148}
149
150/// Options for [`ImapIdle::new`].
151#[derive(Clone, Debug, Default, Eq, PartialEq)]
152pub struct ImapIdleOptions {
153    /// Refresh interval; defaults to 29 s so the connection survives
154    /// middle-boxes. Unused without the `client` feature.
155    pub timeout: Option<Duration>,
156}
157
158/// I/O-free IMAP IDLE coroutine yielding mailbox change events.
159pub struct ImapIdle {
160    tag: TagGenerator,
161    state: State,
162    wants_read: bool,
163    codec: ResponseCodec,
164    data: Vec<Data<'static>>,
165    untagged: Vec<StatusBody<'static>>,
166    bye: Option<Bye<'static>>,
167    done: Arc<AtomicBool>,
168    #[cfg_attr(not(feature = "client"), allow(dead_code))]
169    opts: ImapIdleOptions,
170    #[cfg(feature = "client")]
171    timer: Option<Instant>,
172}
173
174impl ImapIdle {
175    /// Creates a coroutine that keeps an IDLE session open on the
176    /// selected mailbox and yields mailbox change events.
177    ///
178    /// Flip `done` to `true` to wind down with a clean DONE; `opts`
179    /// tunes the refresh interval.
180    pub fn new(done: Arc<AtomicBool>, opts: ImapIdleOptions) -> Self {
181        let mut tag = TagGenerator::new();
182
183        let command = Command {
184            tag: tag.generate(),
185            body: CommandBody::Idle,
186        };
187
188        trace!("send IMAP command {command:?}");
189
190        let state = State::Idle(ImapSend::new(CommandCodec::new(), command));
191
192        Self {
193            tag,
194            state,
195            wants_read: false,
196            codec: ResponseCodec::new(),
197            data: Vec::new(),
198            untagged: Vec::new(),
199            bye: None,
200            done,
201            opts,
202            #[cfg(feature = "client")]
203            timer: None,
204        }
205    }
206
207    #[cfg(feature = "client")]
208    fn timeout(&self) -> Duration {
209        self.opts.timeout.unwrap_or(IDLE_DEFAULT_TIMEOUT)
210    }
211
212    #[cfg(feature = "client")]
213    fn timed_out(&self) -> bool {
214        self.timer
215            .as_ref()
216            .map(|t| t.elapsed() >= self.timeout())
217            .unwrap_or(false)
218    }
219}
220
221impl ImapCoroutine for ImapIdle {
222    type Yield = ImapIdleYield;
223    type Return = Result<(), ImapIdleError>;
224
225    fn resume(
226        &mut self,
227        fragmentizer: &mut Fragmentizer,
228        mut arg: Option<&[u8]>,
229    ) -> ImapCoroutineState<Self::Yield, Self::Return> {
230        #[cfg(feature = "client")]
231        if self.timer.is_none() {
232            self.timer = Some(Instant::now());
233        }
234
235        loop {
236            if mem::take(&mut self.wants_read) {
237                return ImapCoroutineState::Yielded(ImapIdleYield::WantsRead);
238            }
239
240            match &mut self.state {
241                State::Idle(send) => {
242                    // NOTE: servers may pack untagged responses into the same
243                    // frame as `+ idling`; surface them immediately.
244                    let out = imap_try!(send, fragmentizer, arg.take());
245
246                    if let Some(bye) = out.bye {
247                        let err = ImapIdleError::Bye(bye.text.to_string());
248                        return ImapCoroutineState::Complete(Err(err));
249                    }
250
251                    if let Some(Tagged { body, .. }) = out.tagged {
252                        let err = match body.kind {
253                            StatusKind::Ok => ImapIdleError::UnexpectedTagged,
254                            StatusKind::No => ImapIdleError::No(body.text.to_string()),
255                            StatusKind::Bad => ImapIdleError::Bad(body.text.to_string()),
256                        };
257
258                        return ImapCoroutineState::Complete(Err(err));
259                    }
260
261                    if out.continuation_request.is_none() {
262                        let err = ImapIdleError::ExpectedContinuationRequest;
263                        return ImapCoroutineState::Complete(Err(err));
264                    }
265
266                    self.state = State::Read;
267                    debug!("{}", self.state);
268
269                    if !out.data.is_empty() || !out.untagged.is_empty() {
270                        let event = ImapIdleEvent {
271                            data: out.data,
272                            untagged: out.untagged,
273                        };
274
275                        return ImapCoroutineState::Yielded(ImapIdleYield::Event(event));
276                    }
277                }
278                State::Read => {
279                    let done = self.done.load(Ordering::SeqCst);
280                    #[cfg(feature = "client")]
281                    let timed_out = self.timed_out();
282                    #[cfg(not(feature = "client"))]
283                    let timed_out = false;
284
285                    if done || timed_out {
286                        trace!("idle done: {done}");
287                        trace!("idle timed out: {timed_out}");
288                        let send = ImapSend::new(IdleDoneCodec::new(), IdleDone);
289                        self.state = State::Done(send);
290                        debug!("{}", self.state);
291                        continue;
292                    }
293
294                    match arg.take() {
295                        Some(&[]) => {
296                            return ImapCoroutineState::Complete(Err(ImapIdleError::Eof));
297                        }
298                        Some(bytes) => {
299                            trace!("read bytes: {}", escape_byte_string(bytes));
300                            fragmentizer.enqueue_bytes(bytes);
301                        }
302                        None => {
303                            self.wants_read = true;
304                            continue;
305                        }
306                    }
307
308                    loop {
309                        match fragmentizer.progress() {
310                            Some(info @ FragmentInfo::Line { .. }) => {
311                                let bytes = fragmentizer.fragment_bytes(info);
312                                trace!("read line fragment: {}", escape_byte_string(bytes));
313
314                                if !fragmentizer.is_message_complete() {
315                                    continue;
316                                }
317
318                                match fragmentizer.decode_message(&self.codec) {
319                                    Ok(Response::Data(data)) => {
320                                        self.data.push(data.into_static());
321                                    }
322                                    Ok(Response::Status(Status::Untagged(status))) => {
323                                        self.untagged.push(status.into_static());
324                                    }
325                                    Ok(Response::Status(Status::Tagged(_))) => {}
326                                    Ok(Response::Status(Status::Bye(bye))) => {
327                                        self.bye.replace(bye.into_static());
328                                    }
329                                    Ok(Response::CommandContinuationRequest(_)) => {}
330                                    Err(decode_err) => {
331                                        let bytes = fragmentizer.message_bytes();
332                                        let bytes = Secret::new(bytes.into());
333                                        let err = match decode_err {
334                                            DecodeMessageError::DecodingFailure(_)
335                                            | DecodeMessageError::DecodingRemainder { .. } => {
336                                                ImapIdleError::DecodingFailure(bytes)
337                                            }
338                                            DecodeMessageError::MessageTooLong { .. } => {
339                                                ImapIdleError::MessageTooLong(bytes)
340                                            }
341                                            DecodeMessageError::MessagePoisoned { .. } => {
342                                                ImapIdleError::MessageIsPoisoned(bytes)
343                                            }
344                                        };
345                                        return ImapCoroutineState::Complete(Err(err));
346                                    }
347                                }
348                            }
349                            Some(info @ FragmentInfo::Literal { .. }) => {
350                                let bytes = fragmentizer.fragment_bytes(info);
351                                trace!("read literal fragment ({} bytes)", bytes.len());
352                            }
353                            None => {
354                                let event = ImapIdleEvent {
355                                    data: mem::take(&mut self.data),
356                                    untagged: mem::take(&mut self.untagged),
357                                };
358
359                                return ImapCoroutineState::Yielded(ImapIdleYield::Event(event));
360                            }
361                        }
362                    }
363                }
364                State::Done(send) => {
365                    let out = imap_try!(send, fragmentizer, arg.take());
366
367                    if let Some(bye) = out.bye {
368                        let err = ImapIdleError::Bye(bye.text.to_string());
369                        return ImapCoroutineState::Complete(Err(err));
370                    }
371
372                    let Some(Tagged { body, .. }) = out.tagged else {
373                        return ImapCoroutineState::Complete(Err(ImapIdleError::MissingTagged));
374                    };
375
376                    #[cfg(feature = "client")]
377                    let timed_out = self
378                        .timer
379                        .take()
380                        .map(|t| t.elapsed() >= self.timeout())
381                        .unwrap_or(false);
382                    #[cfg(not(feature = "client"))]
383                    let timed_out = false;
384
385                    return match body.kind {
386                        StatusKind::Ok if timed_out => {
387                            trace!("reached timeout, starting a new IDLE command");
388                            let command = Command {
389                                tag: self.tag.generate(),
390                                body: CommandBody::Idle,
391                            };
392                            let send = ImapSend::new(CommandCodec::new(), command);
393                            self.state = State::Idle(send);
394                            debug!("{}", self.state);
395                            continue;
396                        }
397                        StatusKind::Ok => ImapCoroutineState::Complete(Ok(())),
398                        StatusKind::No => ImapCoroutineState::Complete(Err(ImapIdleError::No(
399                            body.text.to_string(),
400                        ))),
401                        StatusKind::Bad => ImapCoroutineState::Complete(Err(ImapIdleError::Bad(
402                            body.text.to_string(),
403                        ))),
404                    };
405                }
406            }
407        }
408    }
409}
410
411enum State {
412    Idle(ImapSend<CommandCodec>),
413    Read,
414    Done(ImapSend<IdleDoneCodec>),
415}
416
417impl fmt::Display for State {
418    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
419        match self {
420            Self::Idle(_) => f.write_str("send idle"),
421            Self::Read => f.write_str("read events"),
422            Self::Done(_) => f.write_str("send done"),
423        }
424    }
425}
426
427#[cfg(test)]
428mod tests {
429    use core::str;
430
431    use alloc::{borrow::ToOwned, format};
432
433    use crate::rfc2177::idle::*;
434
435    #[test]
436    fn shutdown_returns_ok() {
437        let done = Arc::new(AtomicBool::new(false));
438        let mut idle = ImapIdle::new(done.clone(), ImapIdleOptions::default());
439        let mut frag = Fragmentizer::new(50 * 1024 * 1024);
440
441        let bytes = expect_wants_write(&mut idle, &mut frag, None);
442        let line = str::from_utf8(&bytes).expect("utf8 command");
443        let tag = first_word(line).to_owned();
444        assert!(line.trim_end().ends_with("IDLE"));
445
446        expect_wants_read(&mut idle, &mut frag);
447        expect_wants_read_after(&mut idle, &mut frag, b"+ idling\r\n");
448
449        done.store(true, Ordering::SeqCst);
450        let bytes = expect_wants_write(&mut idle, &mut frag, None);
451        assert_eq!(b"DONE\r\n", &*bytes);
452
453        expect_wants_read(&mut idle, &mut frag);
454
455        let reply = format!("{tag} OK IDLE terminated\r\n");
456        expect_complete_ok(&mut idle, &mut frag, reply.as_bytes());
457    }
458
459    #[test]
460    fn unsolicited_during_read_yields_event() {
461        let done = Arc::new(AtomicBool::new(false));
462        let mut idle = ImapIdle::new(done, ImapIdleOptions::default());
463        let mut frag = Fragmentizer::new(50 * 1024 * 1024);
464
465        let _ = expect_wants_write(&mut idle, &mut frag, None);
466        expect_wants_read(&mut idle, &mut frag);
467        expect_wants_read_after(&mut idle, &mut frag, b"+ idling\r\n");
468
469        let event = expect_event(&mut idle, &mut frag, b"* 5 EXISTS\r\n");
470        assert_eq!(1, event.data.len());
471        assert!(event.untagged.is_empty());
472    }
473
474    #[test]
475    fn unsolicited_piggyback_on_continuation_yields_event() {
476        let done = Arc::new(AtomicBool::new(false));
477        let mut idle = ImapIdle::new(done, ImapIdleOptions::default());
478        let mut frag = Fragmentizer::new(50 * 1024 * 1024);
479
480        let _ = expect_wants_write(&mut idle, &mut frag, None);
481        expect_wants_read(&mut idle, &mut frag);
482
483        let event = expect_event(&mut idle, &mut frag, b"+ idling\r\n* 10 EXISTS\r\n");
484        assert_eq!(1, event.data.len());
485    }
486
487    #[test]
488    fn idle_tagged_bad_returns_bad_error() {
489        let done = Arc::new(AtomicBool::new(false));
490        let mut idle = ImapIdle::new(done, ImapIdleOptions::default());
491        let mut frag = Fragmentizer::new(50 * 1024 * 1024);
492
493        let bytes = expect_wants_write(&mut idle, &mut frag, None);
494        let tag = first_word(str::from_utf8(&bytes).expect("utf8 command")).to_owned();
495
496        expect_wants_read(&mut idle, &mut frag);
497
498        let reply = format!("{tag} BAD IDLE not supported\r\n");
499        let err = expect_complete_err(&mut idle, &mut frag, reply.as_bytes());
500        let ImapIdleError::Bad(text) = err else {
501            panic!("expected ImapIdleError::Bad, got {err:?}");
502        };
503        assert_eq!(text, "IDLE not supported");
504    }
505
506    #[test]
507    fn done_tagged_no_returns_no_error() {
508        let done = Arc::new(AtomicBool::new(false));
509        let mut idle = ImapIdle::new(done.clone(), ImapIdleOptions::default());
510        let mut frag = Fragmentizer::new(50 * 1024 * 1024);
511
512        let bytes = expect_wants_write(&mut idle, &mut frag, None);
513        let tag = first_word(str::from_utf8(&bytes).expect("utf8 command")).to_owned();
514
515        expect_wants_read(&mut idle, &mut frag);
516        expect_wants_read_after(&mut idle, &mut frag, b"+ idling\r\n");
517
518        done.store(true, Ordering::SeqCst);
519        let _ = expect_wants_write(&mut idle, &mut frag, None);
520        expect_wants_read(&mut idle, &mut frag);
521
522        let reply = format!("{tag} NO IDLE aborted\r\n");
523        let err = expect_complete_err(&mut idle, &mut frag, reply.as_bytes());
524        let ImapIdleError::No(text) = err else {
525            panic!("expected ImapIdleError::No, got {err:?}");
526        };
527        assert_eq!(text, "IDLE aborted");
528    }
529
530    fn expect_wants_write(
531        cor: &mut ImapIdle,
532        frag: &mut Fragmentizer,
533        arg: Option<&[u8]>,
534    ) -> Vec<u8> {
535        match cor.resume(frag, arg) {
536            ImapCoroutineState::Yielded(ImapIdleYield::WantsWrite(bytes)) => bytes,
537            state => panic!("expected WantsWrite, got {state:?}"),
538        }
539    }
540
541    fn expect_wants_read(cor: &mut ImapIdle, frag: &mut Fragmentizer) {
542        match cor.resume(frag, None) {
543            ImapCoroutineState::Yielded(ImapIdleYield::WantsRead) => {}
544            state => panic!("expected WantsRead, got {state:?}"),
545        }
546    }
547
548    fn expect_wants_read_after(cor: &mut ImapIdle, frag: &mut Fragmentizer, arg: &[u8]) {
549        match cor.resume(frag, Some(arg)) {
550            ImapCoroutineState::Yielded(ImapIdleYield::WantsRead) => {}
551            state => panic!("expected WantsRead, got {state:?}"),
552        }
553    }
554
555    fn expect_event(cor: &mut ImapIdle, frag: &mut Fragmentizer, arg: &[u8]) -> ImapIdleEvent {
556        match cor.resume(frag, Some(arg)) {
557            ImapCoroutineState::Yielded(ImapIdleYield::Event(event)) => event,
558            state => panic!("expected Event, got {state:?}"),
559        }
560    }
561
562    fn expect_complete_ok(cor: &mut ImapIdle, frag: &mut Fragmentizer, reply: &[u8]) {
563        match cor.resume(frag, Some(reply)) {
564            ImapCoroutineState::Complete(Ok(())) => {}
565            state => panic!("expected Complete(Ok), got {state:?}"),
566        }
567    }
568
569    fn expect_complete_err(
570        cor: &mut ImapIdle,
571        frag: &mut Fragmentizer,
572        reply: &[u8],
573    ) -> ImapIdleError {
574        match cor.resume(frag, Some(reply)) {
575            ImapCoroutineState::Complete(Err(err)) => err,
576            state => panic!("expected Complete(Err), got {state:?}"),
577        }
578    }
579
580    fn first_word(line: &str) -> &str {
581        line.split_whitespace()
582            .next()
583            .expect("first whitespace-separated token")
584    }
585}