1use 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#[cfg(feature = "client")]
80const IDLE_DEFAULT_TIMEOUT: Duration = Duration::from_secs(29);
81
82#[derive(Clone, Debug, Error)]
84pub enum ImapIdleError {
85 #[error("IMAP IDLE failed: NO {0}")]
87 No(String),
88 #[error("IMAP IDLE failed: BAD {0}")]
90 Bad(String),
91 #[error("IMAP IDLE failed: BYE {0}")]
93 Bye(String),
94 #[error("IMAP IDLE failed: server returned a tagged response before the continuation request")]
96 UnexpectedTagged,
97 #[error("IMAP IDLE failed: server did not send the expected continuation request")]
99 ExpectedContinuationRequest,
100 #[error("IMAP IDLE failed: server did not return a tagged response to DONE")]
102 MissingTagged,
103 #[error("IMAP IDLE failed: reached unexpected EOF on stream")]
105 Eof,
106 #[error("IMAP IDLE failed: decode response error")]
108 DecodingFailure(Secret<Box<[u8]>>),
109 #[error("IMAP IDLE failed: parse response error: message is poisoned")]
111 MessageIsPoisoned(Secret<Box<[u8]>>),
112 #[error("IMAP IDLE failed: parse response error: message is too long")]
114 MessageTooLong(Secret<Box<[u8]>>),
115 #[error("IMAP IDLE failed: {0}")]
117 Send(#[from] ImapSendError),
118}
119
120#[derive(Debug)]
122pub struct ImapIdleEvent {
123 pub untagged: Vec<StatusBody<'static>>,
125 pub data: Vec<Data<'static>>,
128}
129
130#[derive(Debug)]
132pub enum ImapIdleYield {
133 WantsRead,
135 WantsWrite(Vec<u8>),
137 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#[derive(Clone, Debug, Default, Eq, PartialEq)]
152pub struct ImapIdleOptions {
153 pub timeout: Option<Duration>,
156}
157
158pub 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 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 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}