Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_net/src/imap/server.rs

69.2 KiB, 175 runs

created by r1870400018:9844, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! IMAP4rev1 session loop.
2//!
3//! Drives one TLS-wrapped TCP connection through the IMAP command set
4//! defined in [`crate::imap`]. The session is parameterised over a
5//! `MailStore` and `UserStore` so the same loop serves any Hematite
6//! mailbox backend.
7//!
8//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
9//! Anthropic Claude
10
11use crate::mail::{
12 store::{
13 FolderName,
14 MailStore,
15 MailUser,
16 MessageFlags,
17 MessageMeta,
18 },
19 user::UserStore,
20};
21
22use oxedyne_fe2o3_core::prelude::*;
23
24use std::{
25 net::SocketAddr,
26 sync::Arc,
27 time::{
28 Duration,
29 SystemTime,
30 UNIX_EPOCH,
31 },
32};
33
34use tokio::{
35 io::{
36 AsyncRead,
37 AsyncReadExt,
38 AsyncWrite,
39 AsyncWriteExt,
40 },
41};
42
43
44//// IDLE.
45// A Maildir folder status is a directory listing, so the poll is cheap: one stat
46// per idle connection per interval, against a client that would otherwise ask
47// every ten minutes and be told nothing. RFC 2177 tells clients to re-issue
48// `IDLE` at least every 29 minutes, and ending a little earlier refreshes a
49// connection a NAT was about to drop silently. Only `DONE` is legal from an
50// idling client, so a partial line approaching the line cap is not a client to
51// indulge.
52pub const IDLE_POLL_INTERVAL: Duration = Duration::from_secs(5);
53pub const IDLE_MAX_DURATION: Duration = Duration::from_secs(28 * 60);
54pub const IDLE_MAX_LINE: usize = 1024; // bytes
55
56//// Session ceilings.
57// IMAP imposes no formal line limit; the cap is against a runaway client. The
58// literal cap mirrors the SMTP message size limit.
59pub const IMAP_MAX_LINE: usize = 8 * 1024; // bytes per command line
60pub const IMAP_MAX_LITERAL: usize = 20_480_000; // bytes per `APPEND`
61
62
63/// One IMAP listener configuration.
64///
65/// Cheaply cloneable -- inner state (handler, user store, hostname,
66/// folder list) is wrapped in `Arc`s so a single listener fan-outs
67/// across every accept loop without contention.
68///
69/// The server is generic over its transport: it is handed an established
70/// stream and speaks IMAP over it, whether that is a TLS stream from the
71/// listener, a plain socket on loopback, or an in-memory pipe in a test.
72#[derive(Clone)]
73pub struct ImapServer<M: MailStore, U: UserStore> {
74 pub store: M,
75 pub users: U,
76 pub hostname: Arc<String>, // advertised in the greeting and in BYE
77}
78
79/// Per-connection IMAP session state.
80struct ImapSession {
81 user: Option<MailUser>, // set on a successful LOGIN
82 selected: Option<FolderName>,
83 read_only: bool, // the selection came from EXAMINE
84 // In UID order, refreshed on SELECT, EXAMINE and NOOP, and after STORE and EXPUNGE.
85 messages: Vec<MessageMeta>,
86}
87
88impl ImapSession {
89 fn new() -> Self {
90 Self {
91 user: None,
92 selected: None,
93 read_only: false,
94 messages: Vec::new(),
95 }
96 }
97}
98
99
100impl<M: MailStore, U: UserStore> ImapServer<M, U> {
101
102 /// Runs until the client logs out or the connection drops.
103 pub async fn run<S: AsyncRead + AsyncWrite + Unpin + Send>(
104 &self,
105 mut stream: S,
106 peer: SocketAddr,
107 )
108 -> Outcome<()>
109 {
110 let mut session = ImapSession::new();
111
112 // Greeting.
113 let greet = fmt!(
114 "* OK [CAPABILITY {}] {} Hematite Steel IMAP ready\r\n",
115 capability_list(),
116 self.hostname,
117 );
118 if let Err(e) = stream.write_all(greet.as_bytes()).await {
119 return Err(err!(e, "Writing IMAP greeting."; IO, Network, Write));
120 }
121 if let Err(e) = stream.flush().await {
122 return Err(err!(e, "Flushing IMAP greeting."; IO, Network, Write));
123 }
124
125 loop {
126 let line = match res!(read_line(&mut stream).await) {
127 Some(l) => l,
128 None => break,
129 };
130 if line.trim().is_empty() {
131 continue;
132 }
133 let parsed = match parse_command(&line) {
134 Some(p) => p,
135 None => {
136 let bad = fmt!("* BAD Cannot parse command\r\n");
137 let _ = stream.write_all(bad.as_bytes()).await;
138 continue;
139 }
140 };
141 match self.dispatch(&mut stream, &mut session, parsed, peer).await {
142 Ok(true) => break,
143 Ok(false) => continue,
144 Err(e) => {
145 error!(err!(e,
146 "IMAP dispatch failure for {:?}.", peer;
147 IO));
148 break;
149 }
150 }
151 }
152
153 let _ = stream.shutdown().await;
154 Ok(())
155 }
156
157 /// `Ok(true)` ends the session: a LOGOUT, or a fatal protocol error.
158 async fn dispatch<S: AsyncRead + AsyncWrite + Unpin + Send>(
159 &self,
160 stream: &mut S,
161 session: &mut ImapSession,
162 parsed: ParsedCommand,
163 peer: SocketAddr,
164 )
165 -> Outcome<bool>
166 {
167 let tag = parsed.tag.clone();
168 let cmd = parsed.command.to_uppercase();
169 let args = parsed.args.clone();
170
171 // Trace every command so we can debug client behaviour.
172 // LOGIN argument is suppressed so the password never reaches
173 // the journal in cleartext -- the second arg of LOGIN is the
174 // password, the first is the username.
175 let log_args = if cmd == "LOGIN" {
176 "<redacted>".to_string()
177 } else {
178 args.clone()
179 };
180 info!("IMAP {} {} {} {}", peer, tag, cmd, log_args);
181
182 match cmd.as_str() {
183 "CAPABILITY" => {
184 let line = fmt!("* CAPABILITY {}\r\n", capability_list());
185 res!(write_all(stream, line.as_bytes()).await);
186 res!(write_ok(stream, &tag, "CAPABILITY completed").await);
187 Ok(false)
188 }
189 "NOOP" => {
190 if session.user.is_some() && session.selected.is_some() {
191 res!(self.refresh_selected(session, false).await);
192 res!(self.write_select_status(stream, session).await);
193 }
194 res!(write_ok(stream, &tag, "NOOP completed").await);
195 Ok(false)
196 }
197 "IDLE" => {
198 if session.user.is_none() || session.selected.is_none() {
199 res!(write_bad(stream, &tag,
200 "IDLE requires an authenticated session with a \
201 selected mailbox").await);
202 return Ok(false);
203 }
204 self.idle(stream, session, &tag).await
205 }
206 "LOGOUT" => {
207 res!(write_all(stream, b"* BYE Logging out\r\n").await);
208 res!(write_ok(stream, &tag, "LOGOUT completed").await);
209 Ok(true)
210 }
211 "LOGIN" => {
212 let mut it = ArgIter::new(&args);
213 let user = match it.next_string() {
214 Some(s) => s,
215 None => {
216 res!(write_bad(stream, &tag, "LOGIN requires user").await);
217 return Ok(false);
218 }
219 };
220 let pass = match it.next_string() {
221 Some(s) => s,
222 None => {
223 res!(write_bad(stream, &tag, "LOGIN requires password").await);
224 return Ok(false);
225 }
226 };
227 let result = self.users.authenticate(&user, &pass);
228 let mu = match result {
229 Ok(Some(u)) => u,
230 Ok(None) => {
231 res!(write_no(stream, &tag, "LOGIN failed").await);
232 return Ok(false);
233 }
234 Err(e) => {
235 warn!("LOGIN backend error: {}", e);
236 res!(write_no(stream, &tag, "LOGIN failed").await);
237 return Ok(false);
238 }
239 };
240 if let Err(e) = self.store.ensure_user(&mu) {
241 warn!("ensure_user error for {}: {}", mu.address(), e);
242 }
243 session.user = Some(mu);
244 res!(write_ok(stream, &tag, "LOGIN completed").await);
245 Ok(false)
246 }
247 "AUTHENTICATE" => {
248 // Not implemented in MVP -- LOGIN-only.
249 res!(write_no(stream, &tag, "AUTHENTICATE not supported, use LOGIN").await);
250 Ok(false)
251 }
252 "SELECT" | "EXAMINE" => {
253 if session.user.is_none() {
254 res!(write_no(stream, &tag, "Authenticate first").await);
255 return Ok(false);
256 }
257 let mut it = ArgIter::new(&args);
258 let folder = match it.next_string() {
259 Some(s) => FolderName::new(s),
260 None => {
261 res!(write_bad(stream, &tag, "SELECT requires folder").await);
262 return Ok(false);
263 }
264 };
265 session.selected = Some(folder);
266 session.read_only = cmd == "EXAMINE";
267 res!(self.refresh_selected(session, true).await);
268 res!(self.write_select_status(stream, session).await);
269 let suffix = if session.read_only { "[READ-ONLY]" } else { "[READ-WRITE]" };
270 res!(write_ok(stream, &tag, &fmt!("{} {} completed", suffix, cmd)).await);
271 Ok(false)
272 }
273 "CLOSE" => {
274 // RFC 3501 §6.4.2: if not read-only, expunge first.
275 if let (Some(user), Some(folder)) = (session.user.clone(), session.selected.clone()) {
276 if !session.read_only {
277 let _ = self.store.expunge(&user, &folder);
278 }
279 }
280 session.selected = None;
281 session.read_only = false;
282 session.messages.clear();
283 res!(write_ok(stream, &tag, "CLOSE completed").await);
284 Ok(false)
285 }
286 "LIST" | "LSUB" => {
287 let user = match &session.user {
288 Some(u) => u.clone(),
289 None => {
290 res!(write_no(stream, &tag, "Authenticate first").await);
291 return Ok(false);
292 }
293 };
294 let mut it = ArgIter::new(&args);
295 let _refname = it.next_string().unwrap_or_default();
296 let pattern = it.next_string().unwrap_or_default();
297 let folders = if cmd == "LSUB" {
298 res!(self.store.list_subscribed(&user))
299 } else {
300 res!(self.store.list_folders(&user))
301 };
302 for f in &folders {
303 if !match_imap_pattern(&pattern, f.as_str()) {
304 continue;
305 }
306 let attrs = special_use_attrs(f.as_str());
307 let line = fmt!("* {} ({}) \"/\" \"{}\"\r\n",
308 cmd, attrs, escape_quoted(f.as_str()));
309 res!(write_all(stream, line.as_bytes()).await);
310 }
311 res!(write_ok(stream, &tag, &fmt!("{} completed", cmd)).await);
312 Ok(false)
313 }
314 "SUBSCRIBE" => {
315 let user = match &session.user {
316 Some(u) => u.clone(),
317 None => {
318 res!(write_no(stream, &tag, "Authenticate first").await);
319 return Ok(false);
320 }
321 };
322 let mut it = ArgIter::new(&args);
323 let folder = match it.next_string() {
324 Some(s) => FolderName::new(s),
325 None => {
326 res!(write_bad(stream, &tag, "SUBSCRIBE requires folder").await);
327 return Ok(false);
328 }
329 };
330 let _ = self.store.subscribe(&user, &folder);
331 res!(write_ok(stream, &tag, "SUBSCRIBE completed").await);
332 Ok(false)
333 }
334 "UNSUBSCRIBE" => {
335 // No-op: the MVP store keeps the subscription set
336 // forever; Thunderbird does not rely on UNSUBSCRIBE.
337 res!(write_ok(stream, &tag, "UNSUBSCRIBE completed").await);
338 Ok(false)
339 }
340 "CREATE" => {
341 let user = match &session.user {
342 Some(u) => u.clone(),
343 None => {
344 res!(write_no(stream, &tag, "Authenticate first").await);
345 return Ok(false);
346 }
347 };
348 let mut it = ArgIter::new(&args);
349 let folder = match it.next_string() {
350 Some(s) => FolderName::new(s),
351 None => {
352 res!(write_bad(stream, &tag, "CREATE requires folder").await);
353 return Ok(false);
354 }
355 };
356 let _ = self.store.create_folder(&user, &folder);
357 let _ = self.store.subscribe(&user, &folder);
358 res!(write_ok(stream, &tag, "CREATE completed").await);
359 Ok(false)
360 }
361 "DELETE" => {
362 // Not supported on MVP -- folder deletion is a manual
363 // operation on disk.
364 res!(write_no(stream, &tag, "DELETE not supported").await);
365 Ok(false)
366 }
367 "STATUS" => {
368 let user = match &session.user {
369 Some(u) => u.clone(),
370 None => {
371 res!(write_no(stream, &tag, "Authenticate first").await);
372 return Ok(false);
373 }
374 };
375 let mut it = ArgIter::new(&args);
376 let folder = match it.next_string() {
377 Some(s) => FolderName::new(s),
378 None => {
379 res!(write_bad(stream, &tag, "STATUS requires folder").await);
380 return Ok(false);
381 }
382 };
383 let items = it.next_paren_list().unwrap_or_default();
384 let status = res!(self.store.folder_status(&user, &folder));
385 let mut parts: Vec<String> = Vec::new();
386 for item in items.split_whitespace() {
387 match item.to_uppercase().as_str() {
388 "MESSAGES" => parts.push(fmt!("MESSAGES {}", status.exists)),
389 "RECENT" => parts.push(fmt!("RECENT {}", status.recent)),
390 "UIDNEXT" => parts.push(fmt!("UIDNEXT {}", status.uid_next)),
391 "UIDVALIDITY" => parts.push(fmt!("UIDVALIDITY {}", status.uid_validity)),
392 "UNSEEN" => parts.push(fmt!("UNSEEN {}", status.unseen)),
393 _ => (),
394 }
395 }
396 let line = fmt!("* STATUS \"{}\" ({})\r\n",
397 escape_quoted(folder.as_str()),
398 parts.join(" "),
399 );
400 res!(write_all(stream, line.as_bytes()).await);
401 res!(write_ok(stream, &tag, "STATUS completed").await);
402 Ok(false)
403 }
404 "FETCH" | "UID" | "STORE" | "SEARCH" | "EXPUNGE" | "APPEND" => {
405 // Group the remaining commands together so the
406 // borrow checker keeps `session` mutably available.
407 self.dispatch_data_command(stream, session, &tag, &cmd, &args).await
408 }
409 _ => {
410 res!(write_bad(stream, &tag, "Unknown command").await);
411 Ok(false)
412 }
413 }
414 }
415
416 async fn dispatch_data_command<S: AsyncRead + AsyncWrite + Unpin + Send>(
417 &self,
418 stream: &mut S,
419 session: &mut ImapSession,
420 tag: &str,
421 cmd: &str,
422 args: &str,
423 )
424 -> Outcome<bool>
425 {
426 if session.user.is_none() {
427 res!(write_no(stream, tag, "Authenticate first").await);
428 return Ok(false);
429 }
430 match cmd {
431 "FETCH" => {
432 if session.selected.is_none() {
433 res!(write_no(stream, tag, "No mailbox selected").await);
434 return Ok(false);
435 }
436 let mut it = ArgIter::new(args);
437 let seq_set = match it.next_atom() {
438 Some(s) => s,
439 None => {
440 res!(write_bad(stream, tag, "FETCH requires seq set").await);
441 return Ok(false);
442 }
443 };
444 let items = it.rest().to_string();
445 res!(self.do_fetch(stream, session, &seq_set, &items, false).await);
446 res!(write_ok(stream, tag, "FETCH completed").await);
447 Ok(false)
448 }
449 "STORE" => {
450 if session.selected.is_none() {
451 res!(write_no(stream, tag, "No mailbox selected").await);
452 return Ok(false);
453 }
454 let mut it = ArgIter::new(args);
455 let seq_set = match it.next_atom() {
456 Some(s) => s,
457 None => {
458 res!(write_bad(stream, tag, "STORE requires seq set").await);
459 return Ok(false);
460 }
461 };
462 let op = match it.next_atom() {
463 Some(s) => s,
464 None => {
465 res!(write_bad(stream, tag, "STORE requires op").await);
466 return Ok(false);
467 }
468 };
469 let flags = it.next_paren_list().unwrap_or_default();
470 res!(self.do_store(stream, session, &seq_set, &op, &flags, false).await);
471 res!(write_ok(stream, tag, "STORE completed").await);
472 Ok(false)
473 }
474 "SEARCH" => {
475 if session.selected.is_none() {
476 res!(write_no(stream, tag, "No mailbox selected").await);
477 return Ok(false);
478 }
479 res!(self.do_search(stream, session, args, false).await);
480 res!(write_ok(stream, tag, "SEARCH completed").await);
481 Ok(false)
482 }
483 "EXPUNGE" => {
484 let user = match &session.user {
485 Some(u) => u.clone(),
486 None => {
487 res!(write_no(stream, tag, "Authenticate first").await);
488 return Ok(false);
489 }
490 };
491 let folder = match &session.selected {
492 Some(f) => f.clone(),
493 None => {
494 res!(write_no(stream, tag, "No mailbox selected").await);
495 return Ok(false);
496 }
497 };
498 let removed = res!(self.store.expunge(&user, &folder));
499 // For each expunged UID we send an untagged EXPUNGE
500 // with the *sequence number* (1-based index in the
501 // pre-expunge list). We use the cached message list.
502 for uid in &removed {
503 if let Some(idx) = session.messages.iter().position(|m| m.uid == *uid) {
504 let line = fmt!("* {} EXPUNGE\r\n", idx + 1);
505 res!(write_all(stream, line.as_bytes()).await);
506 session.messages.remove(idx);
507 }
508 }
509 res!(write_ok(stream, tag, "EXPUNGE completed").await);
510 Ok(false)
511 }
512 "APPEND" => {
513 let user = match &session.user {
514 Some(u) => u.clone(),
515 None => {
516 res!(write_no(stream, tag, "Authenticate first").await);
517 return Ok(false);
518 }
519 };
520 res!(self.do_append(stream, &user, tag, args).await);
521 Ok(false)
522 }
523 "UID" => {
524 // UID FETCH / UID STORE / UID SEARCH / UID COPY.
525 let mut it = ArgIter::new(args);
526 let sub = match it.next_atom() {
527 Some(s) => s.to_uppercase(),
528 None => {
529 res!(write_bad(stream, tag, "UID requires subcommand").await);
530 return Ok(false);
531 }
532 };
533 match sub.as_str() {
534 "FETCH" => {
535 let seq_set = match it.next_atom() {
536 Some(s) => s,
537 None => {
538 res!(write_bad(stream, tag, "UID FETCH requires seq set").await);
539 return Ok(false);
540 }
541 };
542 let items = it.rest().to_string();
543 res!(self.do_fetch(stream, session, &seq_set, &items, true).await);
544 res!(write_ok(stream, tag, "UID FETCH completed").await);
545 }
546 "STORE" => {
547 let seq_set = match it.next_atom() {
548 Some(s) => s,
549 None => {
550 res!(write_bad(stream, tag, "UID STORE requires seq set").await);
551 return Ok(false);
552 }
553 };
554 let op = match it.next_atom() {
555 Some(s) => s,
556 None => {
557 res!(write_bad(stream, tag, "UID STORE requires op").await);
558 return Ok(false);
559 }
560 };
561 let flags = it.next_paren_list().unwrap_or_default();
562 res!(self.do_store(stream, session, &seq_set, &op, &flags, true).await);
563 res!(write_ok(stream, tag, "UID STORE completed").await);
564 }
565 "SEARCH" => {
566 res!(self.do_search(stream, session, it.rest(), true).await);
567 res!(write_ok(stream, tag, "UID SEARCH completed").await);
568 }
569 _ => {
570 res!(write_bad(stream, tag, "UID subcommand not supported").await);
571 }
572 }
573 Ok(false)
574 }
575 _ => {
576 res!(write_bad(stream, tag, "Unknown command").await);
577 Ok(false)
578 }
579 }
580 }
581
582 /// RFC 2177 `IDLE`: hold the connection open and tell the client the
583 /// moment the mailbox changes, instead of making it ask.
584 ///
585 /// Without this a client polls -- Thunderbird every ten minutes by
586 /// default -- so mail that has already arrived sits unannounced for up to
587 /// ten minutes, and the client wakes the server up all day to be told
588 /// nothing has happened. IDLE inverts both halves of that bargain.
589 ///
590 /// The client sends `IDLE`, we answer `+ idling`, and from then until it
591 /// sends `DONE` the only thing it may send is `DONE`. We meanwhile watch
592 /// the mailbox and push an untagged `EXISTS` when the count moves.
593 ///
594 /// # Why a timeout around the read, and not `select!`
595 ///
596 /// The obvious shape -- `select!` between the socket and a ticker -- wants
597 /// a mutable borrow of the stream in one branch and another in the handler
598 /// of the other, which does not compile, and inviting people to work around
599 /// that by reading a byte at a time invites losing bytes: a future dropped
600 /// mid-line has already taken those bytes off the socket. `AsyncReadExt::read`
601 /// is cancel-safe -- if it is dropped, nothing was consumed -- so a timeout
602 /// around it is both correct and simple, and the partial line lives in a
603 /// buffer that outlives each attempt.
604 async fn idle<S: AsyncRead + AsyncWrite + Unpin + Send>(
605 &self,
606 stream: &mut S,
607 session: &mut ImapSession,
608 tag: &str,
609 )
610 -> Outcome<bool>
611 {
612 let user = match session.user.clone() {
613 Some(u) => u,
614 None => return Ok(false),
615 };
616 let folder = match session.selected.clone() {
617 Some(f) => f,
618 None => return Ok(false),
619 };
620
621 res!(write_all(stream, b"+ idling\r\n").await);
622
623 let mut last_exists = res!(self.store.folder_status(&user, &folder)).exists;
624 let mut pending: Vec<u8> = Vec::new();
625 let mut chunk = [0u8; 256];
626 let mut waited = Duration::from_secs(0);
627
628 loop {
629 let read = tokio::time::timeout(
630 IDLE_POLL_INTERVAL,
631 stream.read(&mut chunk),
632 ).await;
633
634 match read {
635 // The client hung up. Not an error: closing an idle
636 // connection is how clients end a session all the time.
637 Ok(Ok(0)) => return Ok(true),
638 Ok(Ok(n)) => {
639 pending.extend_from_slice(&chunk[..n]);
640 // Only `DONE` is legal here, so a line that is not `DONE`
641 // ends the IDLE rather than being silently swallowed --
642 // a client that thinks it issued a command and got no
643 // answer will hang for ever.
644 if let Some(i) = pending.iter().position(|b| *b == b'\n') {
645 let line: Vec<u8> = pending.drain(..=i).collect();
646 let text = String::from_utf8_lossy(&line)
647 .trim()
648 .to_uppercase();
649 if text == "DONE" {
650 res!(write_ok(stream, tag, "IDLE terminated").await);
651 } else {
652 res!(write_bad(stream, tag,
653 "Only DONE is valid while idling").await);
654 }
655 return Ok(false);
656 }
657 // A line this long is not a client we want to keep
658 // buffering for.
659 if pending.len() > IDLE_MAX_LINE {
660 res!(write_bad(stream, tag,
661 "Line too long while idling").await);
662 return Ok(false);
663 }
664 }
665 Ok(Err(e)) => return Err(err!(e,
666 "Reading from an idling IMAP client.";
667 IO, Network, Read)),
668 // Nothing from the client: look at the mailbox.
669 Err(_elapsed) => {
670 waited = waited.saturating_add(IDLE_POLL_INTERVAL);
671
672 let status = res!(self.store.folder_status(&user, &folder));
673 if status.exists != last_exists {
674 last_exists = status.exists;
675 res!(self.refresh_selected(session, false).await);
676 res!(write_all(stream,
677 fmt!("* {} EXISTS\r\n", status.exists).as_bytes()).await);
678 res!(write_all(stream,
679 fmt!("* {} RECENT\r\n", status.recent).as_bytes()).await);
680 }
681
682 // RFC 2177 tells clients to re-issue IDLE at least every
683 // 29 minutes, and warns servers to expect stale ones. End
684 // it ourselves a little before that: a well-behaved client
685 // simply idles again, and a NAT that would have dropped
686 // the connection silently never gets the chance.
687 if waited >= IDLE_MAX_DURATION {
688 res!(write_all(stream,
689 b"* OK Idle time expired, please re-issue IDLE\r\n").await);
690 res!(write_ok(stream, tag, "IDLE terminated").await);
691 return Ok(false);
692 }
693 }
694 }
695 }
696 }
697
698 async fn refresh_selected(
699 &self,
700 session: &mut ImapSession,
701 clear_rec: bool,
702 )
703 -> Outcome<()>
704 {
705 let user = match session.user.clone() {
706 Some(u) => u,
707 None => return Ok(()),
708 };
709 let folder = match session.selected.clone() {
710 Some(f) => f,
711 None => return Ok(()),
712 };
713 let read_only = !clear_rec || session.read_only;
714 let messages = res!(self.store.list_messages(&user, &folder, read_only));
715 session.messages = messages;
716 Ok(())
717 }
718
719 /// EXISTS, RECENT, UIDVALIDITY, UIDNEXT, FLAGS, PERMANENTFLAGS and [UNSEEN N],
720 /// all of which SELECT and EXAMINE are required to send.
721 async fn write_select_status<S: AsyncRead + AsyncWrite + Unpin + Send>(
722 &self,
723 stream: &mut S,
724 session: &mut ImapSession,
725 )
726 -> Outcome<()>
727 {
728 let user = match &session.user {
729 Some(u) => u.clone(),
730 None => return Ok(()),
731 };
732 let folder = match &session.selected {
733 Some(f) => f.clone(),
734 None => return Ok(()),
735 };
736 let status = res!(self.store.folder_status(&user, &folder));
737 let lines = [
738 fmt!("* {} EXISTS\r\n", status.exists),
739 fmt!("* {} RECENT\r\n", status.recent),
740 fmt!("* OK [UIDVALIDITY {}] UIDs valid\r\n", status.uid_validity),
741 fmt!("* OK [UIDNEXT {}] Next UID\r\n", status.uid_next),
742 fmt!("* FLAGS (\\Answered \\Flagged \\Deleted \\Seen \\Draft)\r\n"),
743 fmt!("* OK [PERMANENTFLAGS (\\Answered \\Flagged \\Deleted \\Seen \\Draft)] OK\r\n"),
744 ];
745 for line in &lines {
746 res!(write_all(stream, line.as_bytes()).await);
747 }
748 // UNSEEN: position of the first unseen message, if any.
749 if let Some(idx) = session.messages.iter().position(|m| !m.flags.seen) {
750 let line = fmt!("* OK [UNSEEN {}] First unseen\r\n", idx + 1);
751 res!(write_all(stream, line.as_bytes()).await);
752 }
753 Ok(())
754 }
755
756 async fn do_fetch<S: AsyncRead + AsyncWrite + Unpin + Send>(
757 &self,
758 stream: &mut S,
759 session: &mut ImapSession,
760 seq_set: &str,
761 items_raw: &str,
762 by_uid: bool,
763 )
764 -> Outcome<()>
765 {
766 let user = match &session.user {
767 Some(u) => u.clone(),
768 None => return Ok(()),
769 };
770 let folder = match &session.selected {
771 Some(f) => f.clone(),
772 None => return Ok(()),
773 };
774 res!(self.refresh_selected(session, false).await);
775
776 let items = parse_fetch_items(items_raw);
777 let indices = resolve_set(seq_set, &session.messages, by_uid);
778
779 for &idx in &indices {
780 let meta = match session.messages.get(idx) {
781 Some(m) => m.clone(),
782 None => continue,
783 };
784 let seq_no = idx + 1;
785
786 // Read raw bytes once if any item needs them.
787 let mut needs_bytes = false;
788 for it in &items {
789 match it {
790 FetchItem::Body { .. } |
791 FetchItem::Rfc822 |
792 FetchItem::Rfc822Header |
793 FetchItem::Rfc822Text |
794 FetchItem::Envelope => needs_bytes = true,
795 _ => (),
796 }
797 }
798 let raw: Option<Vec<u8>> = if needs_bytes {
799 Some(res!(self.store.fetch_bytes(&user, &folder, meta.uid)))
800 } else {
801 None
802 };
803
804 let mut implicit_seen = false;
805
806 // Encode each item as a sequence of segments. Text segments
807 // may be joined with a separating space; literal segments
808 // contain raw byte payloads that must be written
809 // unmodified. We collect segments first and decide whether
810 // to also append a synthetic FLAGS segment after evaluating
811 // implicit \Seen.
812 let mut segments: Vec<FetchSegment> = Vec::new();
813 for it in &items {
814 match it {
815 FetchItem::Uid => {
816 segments.push(FetchSegment::Text(fmt!("UID {}", meta.uid.0)));
817 }
818 FetchItem::Flags => {
819 segments.push(FetchSegment::Text(fmt!(
820 "FLAGS ({})", meta.flags.to_imap_list())));
821 }
822 FetchItem::Rfc822Size => {
823 segments.push(FetchSegment::Text(fmt!(
824 "RFC822.SIZE {}", meta.size)));
825 }
826 FetchItem::InternalDate => {
827 segments.push(FetchSegment::Text(fmt!(
828 "INTERNALDATE \"{}\"",
829 format_internal_date(meta.internal),
830 )));
831 }
832 FetchItem::Envelope => {
833 if let Some(ref bytes) = raw {
834 segments.push(FetchSegment::Text(fmt!(
835 "ENVELOPE {}", build_envelope(bytes))));
836 }
837 }
838 FetchItem::Body { peek, section } => {
839 if !peek { implicit_seen = true; }
840 let bytes = match raw {
841 Some(ref b) => extract_section(b, section),
842 None => Vec::new(),
843 };
844 let label = section_label(section);
845 segments.push(FetchSegment::Literal {
846 prefix: fmt!("BODY[{}] ", label),
847 payload: bytes,
848 });
849 }
850 FetchItem::Rfc822 => {
851 implicit_seen = true;
852 let bytes = match raw {
853 Some(ref b) => b.clone(),
854 None => Vec::new(),
855 };
856 segments.push(FetchSegment::Literal {
857 prefix: "RFC822 ".to_string(),
858 payload: bytes,
859 });
860 }
861 FetchItem::Rfc822Header => {
862 let bytes = match raw {
863 Some(ref b) => extract_section(b, &Section::Header),
864 None => Vec::new(),
865 };
866 segments.push(FetchSegment::Literal {
867 prefix: "RFC822.HEADER ".to_string(),
868 payload: bytes,
869 });
870 }
871 FetchItem::Rfc822Text => {
872 implicit_seen = true;
873 let bytes = match raw {
874 Some(ref b) => extract_section(b, &Section::Text),
875 None => Vec::new(),
876 };
877 segments.push(FetchSegment::Literal {
878 prefix: "RFC822.TEXT ".to_string(),
879 payload: bytes,
880 });
881 }
882 }
883 }
884
885 // RFC 3501 §6.4.8: every FETCH response triggered by a
886 // UID command must include the UID, even when the client
887 // did not request it. Inject one if it's missing.
888 if by_uid {
889 let already_uid = segments.iter().any(|s| {
890 matches!(s, FetchSegment::Text(t) if t.starts_with("UID "))
891 });
892 if !already_uid {
893 segments.insert(0, FetchSegment::Text(fmt!("UID {}", meta.uid.0)));
894 }
895 }
896
897 // Implicit \Seen.
898 if implicit_seen && !session.read_only && !meta.flags.seen {
899 let mut new_flags = meta.flags;
900 new_flags.seen = true;
901 let _ = self.store.set_flags(&user, &folder, meta.uid, new_flags);
902 if let Some(m) = session.messages.iter_mut().find(|m| m.uid == meta.uid) {
903 m.flags.seen = true;
904 }
905 let already = segments.iter().any(|s| {
906 matches!(s, FetchSegment::Text(t) if t.starts_with("FLAGS"))
907 });
908 if !already {
909 segments.push(FetchSegment::Text(fmt!(
910 "FLAGS ({})", new_flags.to_imap_list())));
911 }
912 }
913
914 // Emit `* SEQ FETCH ( ... )\r\n` with a SP between segments.
915 let opener = fmt!("* {} FETCH (", seq_no);
916 res!(write_all(stream, opener.as_bytes()).await);
917 for (i, seg) in segments.iter().enumerate() {
918 if i > 0 {
919 res!(write_all(stream, b" ").await);
920 }
921 match seg {
922 FetchSegment::Text(t) => {
923 res!(write_all(stream, t.as_bytes()).await);
924 }
925 FetchSegment::Literal { prefix, payload } => {
926 let head = fmt!("{}{{{}}}\r\n", prefix, payload.len());
927 res!(write_all(stream, head.as_bytes()).await);
928 res!(write_all(stream, payload).await);
929 }
930 }
931 }
932 res!(write_all(stream, b")\r\n").await);
933 }
934 Ok(())
935 }
936
937 async fn do_store<S: AsyncRead + AsyncWrite + Unpin + Send>(
938 &self,
939 stream: &mut S,
940 session: &mut ImapSession,
941 seq_set: &str,
942 op: &str,
943 flags_raw: &str,
944 by_uid: bool,
945 )
946 -> Outcome<()>
947 {
948 let user = match &session.user {
949 Some(u) => u.clone(),
950 None => return Ok(()),
951 };
952 let folder = match &session.selected {
953 Some(f) => f.clone(),
954 None => return Ok(()),
955 };
956 res!(self.refresh_selected(session, false).await);
957 let indices = resolve_set(seq_set, &session.messages, by_uid);
958 let silent = op.to_uppercase().ends_with(".SILENT");
959 let base = op.trim_end_matches(".SILENT").trim_end_matches(".silent");
960 let new_flag_names: Vec<&str> = flags_raw
961 .split_whitespace()
962 .collect();
963
964 for &idx in &indices {
965 let meta = match session.messages.get(idx) {
966 Some(m) => m.clone(),
967 None => continue,
968 };
969 let mut flags = meta.flags;
970 match base {
971 "+FLAGS" | "+flags" => {
972 for f in &new_flag_names { flags.set(f, true); }
973 }
974 "-FLAGS" | "-flags" => {
975 for f in &new_flag_names { flags.set(f, false); }
976 }
977 "FLAGS" | "flags" => {
978 flags = MessageFlags::default();
979 for f in &new_flag_names { flags.set(f, true); }
980 }
981 _ => {
982 return Err(err!(
983 "STORE op '{}' not understood.", op;
984 Invalid, Input));
985 }
986 }
987 let final_flags = res!(self.store.set_flags(&user, &folder, meta.uid, flags));
988 if let Some(m) = session.messages.iter_mut().find(|m| m.uid == meta.uid) {
989 m.flags = final_flags;
990 }
991 if !silent {
992 let mut parts = vec![fmt!("FLAGS ({})", final_flags.to_imap_list())];
993 if by_uid {
994 parts.push(fmt!("UID {}", meta.uid.0));
995 }
996 let line = fmt!("* {} FETCH ({})\r\n", idx + 1, parts.join(" "));
997 res!(write_all(stream, line.as_bytes()).await);
998 }
999 }
1000 Ok(())
1001 }
1002
1003 async fn do_search<S: AsyncRead + AsyncWrite + Unpin + Send>(
1004 &self,
1005 stream: &mut S,
1006 session: &mut ImapSession,
1007 args: &str,
1008 by_uid: bool,
1009 )
1010 -> Outcome<()>
1011 {
1012 // Minimal SEARCH: ALL, UID <set>, SEEN/UNSEEN, ANSWERED,
1013 // UNANSWERED, FLAGGED, DELETED. Anything else returns the
1014 // entire mailbox.
1015 res!(self.refresh_selected(session, false).await);
1016 let trimmed = args.trim();
1017 let upper = trimmed.to_uppercase();
1018 let mut matched: Vec<&MessageMeta> = Vec::new();
1019 if upper == "ALL" || upper.is_empty() {
1020 for m in &session.messages { matched.push(m); }
1021 } else if upper == "UNSEEN" {
1022 for m in &session.messages { if !m.flags.seen { matched.push(m); } }
1023 } else if upper == "SEEN" {
1024 for m in &session.messages { if m.flags.seen { matched.push(m); } }
1025 } else if upper == "FLAGGED" {
1026 for m in &session.messages { if m.flags.flagged { matched.push(m); } }
1027 } else if upper == "DELETED" {
1028 for m in &session.messages { if m.flags.deleted { matched.push(m); } }
1029 } else if upper.starts_with("UID ") {
1030 let set = &trimmed[4..];
1031 let indices = resolve_set(set, &session.messages, true);
1032 for &i in &indices {
1033 if let Some(m) = session.messages.get(i) { matched.push(m); }
1034 }
1035 } else {
1036 // Fall back to ALL.
1037 for m in &session.messages { matched.push(m); }
1038 }
1039 let mut parts: Vec<String> = Vec::new();
1040 for m in &matched {
1041 if by_uid {
1042 parts.push(fmt!("{}", m.uid.0));
1043 } else {
1044 if let Some(idx) = session.messages.iter().position(|x| x.uid == m.uid) {
1045 parts.push(fmt!("{}", idx + 1));
1046 }
1047 }
1048 }
1049 let line = fmt!("* SEARCH {}\r\n", parts.join(" "));
1050 res!(write_all(stream, line.as_bytes()).await);
1051 Ok(())
1052 }
1053
1054 async fn do_append<S: AsyncRead + AsyncWrite + Unpin + Send>(
1055 &self,
1056 stream: &mut S,
1057 user: &MailUser,
1058 tag: &str,
1059 args: &str,
1060 )
1061 -> Outcome<bool>
1062 {
1063 // Parse: <folder> [(flags)] [date-time] {literal_size}
1064 let mut it = ArgIter::new(args);
1065 let folder = match it.next_string() {
1066 Some(s) => FolderName::new(s),
1067 None => {
1068 res!(write_bad(stream, tag, "APPEND requires folder").await);
1069 return Ok(false);
1070 }
1071 };
1072 let mut flags = MessageFlags::default();
1073 let mut internal: Option<SystemTime> = None;
1074
1075 loop {
1076 it.skip_whitespace();
1077 if it.peek() == Some('(') {
1078 let raw_flags = it.next_paren_list().unwrap_or_default();
1079 for f in raw_flags.split_whitespace() {
1080 flags.set(f, true);
1081 }
1082 } else if it.peek() == Some('"') {
1083 let dt = it.next_string().unwrap_or_default();
1084 internal = parse_internal_date(&dt);
1085 } else if it.peek() == Some('{') {
1086 break;
1087 } else {
1088 break;
1089 }
1090 }
1091
1092 // Literal size.
1093 let lit = match it.next_literal_marker() {
1094 Some(n) => n,
1095 None => {
1096 res!(write_bad(stream, tag, "APPEND requires literal").await);
1097 return Ok(false);
1098 }
1099 };
1100 if lit > IMAP_MAX_LITERAL {
1101 res!(write_no(stream, tag, "Literal too large").await);
1102 return Ok(false);
1103 }
1104 // Synchronising literal: send continuation if not LITERAL+.
1105 if !it.literal_is_non_sync {
1106 res!(write_all(stream, b"+ Ready for literal data\r\n").await);
1107 }
1108 // Read exactly lit bytes.
1109 let mut buf = vec![0u8; lit];
1110 if let Err(e) = stream.read_exact(&mut buf).await {
1111 return Err(err!(e, "Reading APPEND literal."; IO, Network, Read));
1112 }
1113 // Read trailing CRLF after the literal.
1114 let mut tail = [0u8; 2];
1115 let _ = stream.read_exact(&mut tail).await;
1116
1117 let uid = res!(self.store.append(user, &folder, &buf, flags, internal));
1118 let _ = uid;
1119 res!(write_ok(stream, tag, "APPEND completed").await);
1120 Ok(false)
1121 }
1122}
1123
1124
1125// ┌───────────────────────────────────────────────────────────────────────────┐
1126// │ COMMAND PARSING │
1127// └───────────────────────────────────────────────────────────────────────────┘
1128
1129#[derive(Clone, Debug)]
1130struct ParsedCommand {
1131 tag: String,
1132 command: String,
1133 args: String,
1134}
1135
1136fn parse_command(line: &str) -> Option<ParsedCommand> {
1137 let trimmed = line.trim_end_matches(|c: char| c == '\r' || c == '\n');
1138 let mut parts = trimmed.splitn(3, ' ');
1139 let tag = ok!(parts.next()).to_string();
1140 let cmd = ok!(parts.next()).to_string();
1141 let args = parts.next().unwrap_or("").to_string();
1142 if tag.is_empty() || cmd.is_empty() {
1143 return None;
1144 }
1145 Some(ParsedCommand { tag, command: cmd, args })
1146}
1147
1148struct ArgIter<'a> {
1149 s: &'a str,
1150 pos: usize,
1151 // Set by next_literal_marker where the marker carried a `+` suffix, the LITERAL+ extension.
1152 literal_is_non_sync: bool,
1153}
1154
1155impl<'a> ArgIter<'a> {
1156 fn new(s: &'a str) -> Self { Self { s, pos: 0, literal_is_non_sync: false } }
1157
1158 fn peek(&self) -> Option<char> {
1159 self.s[self.pos..].chars().next()
1160 }
1161
1162 fn skip_whitespace(&mut self) {
1163 while let Some(c) = self.peek() {
1164 if c == ' ' || c == '\t' { self.pos += c.len_utf8(); }
1165 else { break; }
1166 }
1167 }
1168
1169 fn rest(&self) -> &str { &self.s[self.pos..] }
1170
1171 fn next_atom(&mut self) -> Option<String> {
1172 self.skip_whitespace();
1173 let start = self.pos;
1174 while let Some(c) = self.peek() {
1175 if c == ' ' || c == '\t' { break; }
1176 self.pos += c.len_utf8();
1177 }
1178 if start == self.pos { None } else { Some(self.s[start..self.pos].to_string()) }
1179 }
1180
1181 /// A quoted string (`"..."`) or an atom, which is what a folder name and other simple
1182 /// string arguments arrive as.
1183 fn next_string(&mut self) -> Option<String> {
1184 self.skip_whitespace();
1185 if self.peek() == Some('"') {
1186 self.pos += 1;
1187 let start = self.pos;
1188 let mut escaped = false;
1189 let mut out = String::new();
1190 while let Some(c) = self.peek() {
1191 self.pos += c.len_utf8();
1192 if escaped {
1193 out.push(c);
1194 escaped = false;
1195 continue;
1196 }
1197 if c == '\\' {
1198 escaped = true;
1199 continue;
1200 }
1201 if c == '"' {
1202 return Some(out);
1203 }
1204 out.push(c);
1205 }
1206 // Unterminated.
1207 Some(self.s[start..].to_string())
1208 } else {
1209 self.next_atom()
1210 }
1211 }
1212
1213 /// The inner text of a parenthesised flag list, without the parentheses.
1214 fn next_paren_list(&mut self) -> Option<String> {
1215 self.skip_whitespace();
1216 if self.peek() != Some('(') { return None; }
1217 self.pos += 1;
1218 let start = self.pos;
1219 let mut depth = 1usize;
1220 while let Some(c) = self.peek() {
1221 self.pos += c.len_utf8();
1222 if c == '(' { depth += 1; }
1223 if c == ')' {
1224 depth -= 1;
1225 if depth == 0 {
1226 return Some(self.s[start..self.pos - 1].to_string());
1227 }
1228 }
1229 }
1230 Some(self.s[start..self.pos].to_string())
1231 }
1232
1233 /// `{N}` or `{N+}`; the second form sets `literal_is_non_sync`.
1234 fn next_literal_marker(&mut self) -> Option<usize> {
1235 self.skip_whitespace();
1236 if self.peek() != Some('{') { return None; }
1237 self.pos += 1;
1238 let start = self.pos;
1239 while let Some(c) = self.peek() {
1240 if c == '}' { break; }
1241 self.pos += c.len_utf8();
1242 }
1243 let inner = &self.s[start..self.pos];
1244 if self.peek() == Some('}') { self.pos += 1; }
1245 let (digits, plus) = if let Some(d) = inner.strip_suffix('+') {
1246 (d, true)
1247 } else {
1248 (inner, false)
1249 };
1250 let n: usize = ok!(digits.parse().ok());
1251 self.literal_is_non_sync = plus;
1252 Some(n)
1253 }
1254}
1255
1256
1257// ┌───────────────────────────────────────────────────────────────────────────┐
1258// │ FETCH ITEM PARSING │
1259// └───────────────────────────────────────────────────────────────────────────┘
1260
1261#[derive(Clone, Debug)]
1262enum FetchItem {
1263 Uid,
1264 Flags,
1265 Rfc822Size,
1266 InternalDate,
1267 Envelope,
1268 Body { peek: bool, section: Section },
1269 Rfc822,
1270 Rfc822Header,
1271 Rfc822Text,
1272}
1273
1274#[derive(Clone, Debug)]
1275enum Section {
1276 Whole,
1277 Header,
1278 Text,
1279 HeaderFields(Vec<String>),
1280 HeaderFieldsNot(Vec<String>),
1281}
1282
1283fn parse_fetch_items(raw: &str) -> Vec<FetchItem> {
1284 // Strip enclosing parens if present. A single item may also appear
1285 // without parens.
1286 let inner = raw.trim();
1287 let inner = if inner.starts_with('(') && inner.ends_with(')') {
1288 &inner[1..inner.len() - 1]
1289 } else {
1290 inner
1291 };
1292 // Walk tokens. We need to handle BODY[...] and BODY.PEEK[...]
1293 // including their internal whitespace.
1294 let mut out = Vec::new();
1295 let bytes = inner.as_bytes();
1296 let mut i = 0;
1297 while i < bytes.len() {
1298 // Skip whitespace.
1299 while i < bytes.len() && (bytes[i] == b' ' || bytes[i] == b'\t') { i += 1; }
1300 if i >= bytes.len() { break; }
1301 // Read up to the next whitespace, or until a `[` (BODY[]).
1302 let start = i;
1303 while i < bytes.len() && bytes[i] != b' ' && bytes[i] != b'\t' && bytes[i] != b'[' {
1304 i += 1;
1305 }
1306 let mut name = inner[start..i].to_string();
1307 let mut section_text = String::new();
1308 if i < bytes.len() && bytes[i] == b'[' {
1309 // Read up to matching `]`.
1310 i += 1;
1311 let s = i;
1312 let mut depth = 1usize;
1313 while i < bytes.len() && depth > 0 {
1314 if bytes[i] == b'[' { depth += 1; }
1315 if bytes[i] == b']' { depth -= 1; if depth == 0 { break; } }
1316 i += 1;
1317 }
1318 section_text = inner[s..i].to_string();
1319 if i < bytes.len() && bytes[i] == b']' { i += 1; }
1320 name.push_str(&fmt!("[{}]", section_text));
1321 }
1322 let token_upper = name.to_uppercase();
1323 match token_upper.as_str() {
1324 "UID" => out.push(FetchItem::Uid),
1325 "FLAGS" => out.push(FetchItem::Flags),
1326 "RFC822.SIZE" => out.push(FetchItem::Rfc822Size),
1327 "INTERNALDATE" => out.push(FetchItem::InternalDate),
1328 "ENVELOPE" => out.push(FetchItem::Envelope),
1329 "RFC822" => out.push(FetchItem::Rfc822),
1330 "RFC822.HEADER" => out.push(FetchItem::Rfc822Header),
1331 "RFC822.TEXT" => out.push(FetchItem::Rfc822Text),
1332 _ => {
1333 if token_upper.starts_with("BODY[")
1334 || token_upper.starts_with("BODY.PEEK[")
1335 {
1336 let peek = token_upper.starts_with("BODY.PEEK[");
1337 let section = parse_section(&section_text);
1338 out.push(FetchItem::Body { peek, section });
1339 } else if token_upper == "BODY" {
1340 out.push(FetchItem::Body { peek: false, section: Section::Whole });
1341 }
1342 }
1343 }
1344 }
1345 out
1346}
1347
1348fn parse_section(text: &str) -> Section {
1349 let upper = text.to_uppercase();
1350 let trimmed = upper.trim();
1351 if trimmed.is_empty() {
1352 return Section::Whole;
1353 }
1354 if trimmed == "HEADER" {
1355 return Section::Header;
1356 }
1357 if trimmed == "TEXT" {
1358 return Section::Text;
1359 }
1360 if trimmed.starts_with("HEADER.FIELDS.NOT") {
1361 let rest = &text[text.find('(').map(|i| i + 1).unwrap_or(0)..];
1362 let rest = rest.trim_end_matches(')');
1363 let names = rest.split_whitespace().map(|s| s.to_string()).collect();
1364 return Section::HeaderFieldsNot(names);
1365 }
1366 if trimmed.starts_with("HEADER.FIELDS") {
1367 let rest = &text[text.find('(').map(|i| i + 1).unwrap_or(0)..];
1368 let rest = rest.trim_end_matches(')');
1369 let names = rest.split_whitespace().map(|s| s.to_string()).collect();
1370 return Section::HeaderFields(names);
1371 }
1372 Section::Whole
1373}
1374
1375fn section_label(section: &Section) -> String {
1376 match section {
1377 Section::Whole => String::new(),
1378 Section::Header => "HEADER".to_string(),
1379 Section::Text => "TEXT".to_string(),
1380 Section::HeaderFields(names) => fmt!(
1381 "HEADER.FIELDS ({})",
1382 names.iter().map(|s| s.as_str()).collect::<Vec<_>>().join(" "),
1383 ),
1384 Section::HeaderFieldsNot(names) => fmt!(
1385 "HEADER.FIELDS.NOT ({})",
1386 names.iter().map(|s| s.as_str()).collect::<Vec<_>>().join(" "),
1387 ),
1388 }
1389}
1390
1391fn extract_section(bytes: &[u8], section: &Section) -> Vec<u8> {
1392 let (head, body) = split_msg(bytes);
1393 match section {
1394 Section::Whole => bytes.to_vec(),
1395 Section::Header => {
1396 let mut h = head.to_vec();
1397 h.extend_from_slice(b"\r\n");
1398 h
1399 }
1400 Section::Text => body.to_vec(),
1401 Section::HeaderFields(names) => filter_headers(head, names, false),
1402 Section::HeaderFieldsNot(names) => filter_headers(head, names, true),
1403 }
1404}
1405
1406fn split_msg(bytes: &[u8]) -> (&[u8], &[u8]) {
1407 if let Some(i) = find_subseq(bytes, b"\r\n\r\n") {
1408 return (&bytes[..i], &bytes[i + 4..]);
1409 }
1410 if let Some(i) = find_subseq(bytes, b"\n\n") {
1411 return (&bytes[..i], &bytes[i + 2..]);
1412 }
1413 (bytes, &[])
1414}
1415
1416fn find_subseq(hay: &[u8], needle: &[u8]) -> Option<usize> {
1417 if needle.is_empty() || needle.len() > hay.len() { return None; }
1418 for i in 0..=hay.len() - needle.len() {
1419 if &hay[i..i + needle.len()] == needle { return Some(i); }
1420 }
1421 None
1422}
1423
1424fn filter_headers(head: &[u8], names: &[String], invert: bool) -> Vec<u8> {
1425 let text = String::from_utf8_lossy(head);
1426 let upper_names: Vec<String> = names.iter().map(|s| s.to_uppercase()).collect();
1427 let mut out = String::new();
1428 let mut keep_current = false;
1429 for line in text.split('\n') {
1430 let line = line.strip_suffix('\r').unwrap_or(line);
1431 if line.starts_with(' ') || line.starts_with('\t') {
1432 // Continuation; if we are keeping the previous header,
1433 // keep the continuation too.
1434 if keep_current {
1435 out.push_str(line);
1436 out.push_str("\r\n");
1437 }
1438 continue;
1439 }
1440 if let Some(i) = line.find(':') {
1441 let name = line[..i].trim().to_uppercase();
1442 let in_list = upper_names.iter().any(|n| n == &name);
1443 keep_current = if invert { !in_list } else { in_list };
1444 if keep_current {
1445 out.push_str(line);
1446 out.push_str("\r\n");
1447 }
1448 } else {
1449 keep_current = false;
1450 }
1451 }
1452 out.push_str("\r\n");
1453 out.into_bytes()
1454}
1455
1456fn build_envelope(bytes: &[u8]) -> String {
1457 let (head, _) = split_msg(bytes);
1458 let headers = parse_headers(head);
1459 let date = nstring(headers.get("date"));
1460 let subj = nstring(headers.get("subject"));
1461 let from = address_list(headers.get("from"));
1462 let sender = headers.get("sender")
1463 .map(|s| address_list(Some(s)))
1464 .unwrap_or_else(|| from.clone());
1465 let reply_to = headers.get("reply-to")
1466 .map(|s| address_list(Some(s)))
1467 .unwrap_or_else(|| from.clone());
1468 let to = address_list(headers.get("to"));
1469 let cc = address_list(headers.get("cc"));
1470 let bcc = address_list(headers.get("bcc"));
1471 let in_reply_to = nstring(headers.get("in-reply-to"));
1472 let message_id = nstring(headers.get("message-id"));
1473 fmt!(
1474 "({} {} {} {} {} {} {} {} {} {})",
1475 date, subj, from, sender, reply_to, to, cc, bcc, in_reply_to, message_id,
1476 )
1477}
1478
1479fn parse_headers(head: &[u8]) -> std::collections::HashMap<String, String> {
1480 let text = String::from_utf8_lossy(head);
1481 let mut out = std::collections::HashMap::new();
1482 let mut name: Option<String> = None;
1483 let mut value = String::new();
1484 for raw in text.split('\n') {
1485 let line = raw.strip_suffix('\r').unwrap_or(raw);
1486 if line.is_empty() { continue; }
1487 if line.starts_with(' ') || line.starts_with('\t') {
1488 value.push(' ');
1489 value.push_str(line.trim_start());
1490 continue;
1491 }
1492 if let Some(n) = name.take() {
1493 out.insert(n.to_lowercase(), value.trim().to_string());
1494 value.clear();
1495 }
1496 if let Some(i) = line.find(':') {
1497 name = Some(line[..i].trim().to_string());
1498 value = line[i + 1..].trim().to_string();
1499 }
1500 }
1501 if let Some(n) = name {
1502 out.insert(n.to_lowercase(), value.trim().to_string());
1503 }
1504 out
1505}
1506
1507fn nstring(s: Option<&String>) -> String {
1508 match s {
1509 Some(s) => fmt!("\"{}\"", escape_quoted(s)),
1510 None => "NIL".to_string(),
1511 }
1512}
1513
1514/// `((name nil mailbox host) ...)`, which is the ENVELOPE form. The MVP parser handles only the
1515/// simple `local@domain` and `Name <local@domain>` inputs.
1516fn address_list(s: Option<&String>) -> String {
1517 let s = match s { Some(s) => s.as_str(), None => return "NIL".to_string() };
1518 let mut entries: Vec<String> = Vec::new();
1519 for part in s.split(',') {
1520 let part = part.trim();
1521 if part.is_empty() { continue; }
1522 let (name, addr) = if let Some(start) = part.find('<') {
1523 let end = part.find('>').unwrap_or(part.len());
1524 let name = part[..start].trim().trim_matches('"').to_string();
1525 (name, part[start + 1..end].to_string())
1526 } else {
1527 (String::new(), part.to_string())
1528 };
1529 let (local, host) = match addr.rfind('@') {
1530 Some(i) => (addr[..i].to_string(), addr[i + 1..].to_string()),
1531 None => (addr, String::new()),
1532 };
1533 let name_field = if name.is_empty() {
1534 "NIL".to_string()
1535 } else {
1536 fmt!("\"{}\"", escape_quoted(&name))
1537 };
1538 entries.push(fmt!(
1539 "({} NIL \"{}\" \"{}\")",
1540 name_field,
1541 escape_quoted(&local),
1542 escape_quoted(&host),
1543 ));
1544 }
1545 if entries.is_empty() {
1546 "NIL".to_string()
1547 } else {
1548 fmt!("({})", entries.join(""))
1549 }
1550}
1551
1552fn escape_quoted(s: &str) -> String {
1553 s.replace('\\', "\\\\").replace('"', "\\\"")
1554}
1555
1556
1557// ┌───────────────────────────────────────────────────────────────────────────┐
1558// │ SET RESOLUTION │
1559// └───────────────────────────────────────────────────────────────────────────┘
1560
1561/// A message set is `1`, `1:5`, `*` or `1,3,5:7`, and resolves to indices into the cached
1562/// message list. `by_uid` reads it as a UID range rather than a sequence-number range.
1563fn resolve_set(set: &str, msgs: &[MessageMeta], by_uid: bool) -> Vec<usize> {
1564 let mut out: Vec<usize> = Vec::new();
1565 if msgs.is_empty() { return out; }
1566 let max_seq = msgs.len() as u32;
1567 let max_uid = msgs.last().map(|m| m.uid.0).unwrap_or(0);
1568 for piece in set.split(',') {
1569 let piece = piece.trim();
1570 if piece.is_empty() { continue; }
1571 let (lo_str, hi_str) = match piece.find(':') {
1572 Some(i) => (&piece[..i], &piece[i + 1..]),
1573 None => (piece, piece),
1574 };
1575 let lo = parse_set_number(lo_str, by_uid, max_seq, max_uid);
1576 let hi = parse_set_number(hi_str, by_uid, max_seq, max_uid);
1577 let (lo, hi) = if lo > hi { (hi, lo) } else { (lo, hi) };
1578 for i in 0..msgs.len() {
1579 let key = if by_uid { msgs[i].uid.0 } else { (i + 1) as u32 };
1580 if key >= lo && key <= hi {
1581 out.push(i);
1582 }
1583 }
1584 }
1585 out
1586}
1587
1588fn parse_set_number(s: &str, by_uid: bool, max_seq: u32, max_uid: u32) -> u32 {
1589 if s == "*" { if by_uid { max_uid } else { max_seq } }
1590 else { s.parse().unwrap_or(0) }
1591}
1592
1593
1594// ┌───────────────────────────────────────────────────────────────────────────┐
1595// │ RESPONSE WRITERS │
1596// └───────────────────────────────────────────────────────────────────────────┘
1597
1598async fn write_all<S: AsyncWrite + Unpin>(stream: &mut S, bytes: &[u8]) -> Outcome<()> {
1599 if let Err(e) = stream.write_all(bytes).await {
1600 return Err(err!(e, "Writing IMAP bytes."; IO, Network, Write));
1601 }
1602 if let Err(e) = stream.flush().await {
1603 return Err(err!(e, "Flushing IMAP bytes."; IO, Network, Write));
1604 }
1605 Ok(())
1606}
1607
1608async fn write_ok<S: AsyncWrite + Unpin>(stream: &mut S, tag: &str, text: &str) -> Outcome<()> {
1609 let line = fmt!("{} OK {}\r\n", tag, text);
1610 write_all(stream, line.as_bytes()).await
1611}
1612
1613async fn write_no<S: AsyncWrite + Unpin>(stream: &mut S, tag: &str, text: &str) -> Outcome<()> {
1614 let line = fmt!("{} NO {}\r\n", tag, text);
1615 write_all(stream, line.as_bytes()).await
1616}
1617
1618async fn write_bad<S: AsyncWrite + Unpin>(stream: &mut S, tag: &str, text: &str) -> Outcome<()> {
1619 let line = fmt!("{} BAD {}\r\n", tag, text);
1620 write_all(stream, line.as_bytes()).await
1621}
1622
1623async fn read_line<S: AsyncRead + Unpin>(stream: &mut S) -> Outcome<Option<String>> {
1624 let mut buf = Vec::with_capacity(128);
1625 let mut byte = [0u8; 1];
1626 loop {
1627 let n = match stream.read(&mut byte).await {
1628 Ok(n) => n,
1629 Err(e) => return Err(err!(e, "Reading IMAP line."; IO, Network, Read)),
1630 };
1631 if n == 0 {
1632 if buf.is_empty() { return Ok(None); }
1633 break;
1634 }
1635 buf.push(byte[0]);
1636 if byte[0] == b'\n' { break; }
1637 if buf.len() >= IMAP_MAX_LINE {
1638 return Err(err!(
1639 "IMAP line exceeded {} bytes.", IMAP_MAX_LINE;
1640 Invalid, Input, Excessive));
1641 }
1642 }
1643 while buf.last() == Some(&b'\n') || buf.last() == Some(&b'\r') {
1644 buf.pop();
1645 }
1646 Ok(Some(String::from_utf8_lossy(&buf).into_owned()))
1647}
1648
1649
1650// ┌───────────────────────────────────────────────────────────────────────────┐
1651// │ MISC HELPERS │
1652// └───────────────────────────────────────────────────────────────────────────┘
1653
1654fn capability_list() -> &'static str {
1655 "IMAP4rev1 LITERAL+ AUTH=PLAIN AUTH=LOGIN SPECIAL-USE IDLE"
1656}
1657
1658/// The RFC 6154 SPECIAL-USE attribute of a well-known folder name, with its leading backslash
1659/// and no trailing space, alongside the `\HasNoChildren` hint most clients rely on. An ordinary
1660/// folder gets `\HasNoChildren` alone.
1661fn special_use_attrs(name: &str) -> String {
1662 let su = match name {
1663 "Sent" => Some("\\Sent"),
1664 "Drafts" => Some("\\Drafts"),
1665 "Trash" => Some("\\Trash"),
1666 "Junk" => Some("\\Junk"),
1667 "Archive" => Some("\\Archive"),
1668 _ => None,
1669 };
1670 match su {
1671 Some(tag) => fmt!("{} \\HasNoChildren", tag),
1672 None => "\\HasNoChildren".to_string(),
1673 }
1674}
1675
1676/// In a `LIST` pattern, `*` matches any number of characters, the hierarchy separator included,
1677/// and `%` matches any number except `/`. An empty pattern matches everything.
1678fn match_imap_pattern(pattern: &str, name: &str) -> bool {
1679 if pattern.is_empty() { return true; }
1680 pattern_recurse(pattern.as_bytes(), name.as_bytes())
1681}
1682
1683fn pattern_recurse(p: &[u8], n: &[u8]) -> bool {
1684 if p.is_empty() { return n.is_empty(); }
1685 match p[0] {
1686 b'*' => {
1687 for i in 0..=n.len() {
1688 if pattern_recurse(&p[1..], &n[i..]) {
1689 return true;
1690 }
1691 }
1692 false
1693 }
1694 b'%' => {
1695 for i in 0..=n.len() {
1696 // % does not match the separator '/'.
1697 if i > 0 && n[i - 1] == b'/' { break; }
1698 if pattern_recurse(&p[1..], &n[i..]) {
1699 return true;
1700 }
1701 }
1702 false
1703 }
1704 c => {
1705 if !n.is_empty() && n[0].eq_ignore_ascii_case(&c) {
1706 pattern_recurse(&p[1..], &n[1..])
1707 } else {
1708 false
1709 }
1710 }
1711 }
1712}
1713
1714/// The INTERNALDATE form, e.g. `13-Apr-2026 10:00:00 +0000`, always at a UTC offset.
1715fn format_internal_date(t: SystemTime) -> String {
1716 let secs = t.duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0);
1717 let (y, mo, d, h, mi, s) = unix_to_civil(secs);
1718 let months = [
1719 "Jan","Feb","Mar","Apr","May","Jun",
1720 "Jul","Aug","Sep","Oct","Nov","Dec",
1721 ];
1722 let mon = months[(mo as usize - 1).min(11)];
1723 fmt!("{:02}-{}-{:04} {:02}:{:02}:{:02} +0000", d, mon, y, h, mi, s)
1724}
1725
1726/// Tolerant: what a client sends here varies more than the grammar says it should.
1727fn parse_internal_date(s: &str) -> Option<SystemTime> {
1728 // Format: "13-Apr-2026 10:00:00 +0000"
1729 let parts: Vec<&str> = s.split_whitespace().collect();
1730 if parts.len() < 2 { return None; }
1731 let date_parts: Vec<&str> = parts[0].split('-').collect();
1732 if date_parts.len() != 3 { return None; }
1733 let day: u32 = ok!(date_parts[0].parse().ok());
1734 let mon: u32 = match date_parts[1] {
1735 "Jan"=>1,"Feb"=>2,"Mar"=>3,"Apr"=>4,"May"=>5,"Jun"=>6,
1736 "Jul"=>7,"Aug"=>8,"Sep"=>9,"Oct"=>10,"Nov"=>11,"Dec"=>12,
1737 _ => return None,
1738 };
1739 let year: i32 = ok!(date_parts[2].parse().ok());
1740 let time_parts: Vec<&str> = parts[1].split(':').collect();
1741 if time_parts.len() != 3 { return None; }
1742 let h: u32 = ok!(time_parts[0].parse().ok());
1743 let mi: u32 = ok!(time_parts[1].parse().ok());
1744 let se: u32 = ok!(time_parts[2].parse().ok());
1745 let secs = civil_to_unix(year, mon, day, h, mi, se);
1746 Some(UNIX_EPOCH + std::time::Duration::from_secs(secs))
1747}
1748
1749/// Howard Hinnant's date algorithm, proleptic Gregorian and UTC.
1750fn unix_to_civil(secs: u64) -> (i32, u32, u32, u32, u32, u32) {
1751 let days = (secs / 86_400) as i64;
1752 let rem = (secs % 86_400) as u32;
1753 let h = rem / 3_600;
1754 let mi = (rem / 60) % 60;
1755 let s = rem % 60;
1756 let z = days + 719_468;
1757 let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
1758 let doe = (z - era * 146_097) as u32;
1759 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
1760 let y = yoe as i64 + era * 400;
1761 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
1762 let mp = (5 * doy + 2) / 153;
1763 let d = doy - (153 * mp + 2) / 5 + 1;
1764 let m = if mp < 10 { mp + 3 } else { mp - 9 };
1765 let y = if m <= 2 { y + 1 } else { y };
1766 (y as i32, m, d, h, mi, s)
1767}
1768
1769/// Inverse of `unix_to_civil`.
1770fn civil_to_unix(y: i32, m: u32, d: u32, h: u32, mi: u32, s: u32) -> u64 {
1771 let y = if m <= 2 { y - 1 } else { y };
1772 let era = if y >= 0 { y } else { y - 399 } / 400;
1773 let yoe = (y - era * 400) as u32;
1774 let doy = (153 * if m > 2 { m - 3 } else { m + 9 } + 2) / 5 + d - 1;
1775 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
1776 let days = era as i64 * 146_097 + doe as i64 - 719_468;
1777 let secs = (days as u64).wrapping_mul(86_400)
1778 + h as u64 * 3_600
1779 + mi as u64 * 60
1780 + s as u64;
1781 secs
1782}
1783
1784/// One emitted FETCH response component. A text segment is joined to its neighbours by a SP
1785/// inside the parenthesised response; a literal segment is raw bytes that must be written
1786/// byte-for-byte after a `prefix{N}\r\n` header.
1787enum FetchSegment {
1788 Text(String),
1789 Literal { prefix: String, payload: Vec<u8> },
1790}