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 | |
| 11 | use crate::mail::{ |
| 12 | store::{ |
| 13 | FolderName, |
| 14 | MailStore, |
| 15 | MailUser, |
| 16 | MessageFlags, |
| 17 | MessageMeta, |
| 18 | }, |
| 19 | user::UserStore, |
| 20 | }; |
| 21 | |
| 22 | use oxedyne_fe2o3_core::prelude::*; |
| 23 | |
| 24 | use std::{ |
| 25 | net::SocketAddr, |
| 26 | sync::Arc, |
| 27 | time::{ |
| 28 | Duration, |
| 29 | SystemTime, |
| 30 | UNIX_EPOCH, |
| 31 | }, |
| 32 | }; |
| 33 | |
| 34 | use 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. |
| 52 | pub const IDLE_POLL_INTERVAL: Duration = Duration::from_secs(5); |
| 53 | pub const IDLE_MAX_DURATION: Duration = Duration::from_secs(28 * 60); |
| 54 | pub 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. |
| 59 | pub const IMAP_MAX_LINE: usize = 8 * 1024; // bytes per command line |
| 60 | pub 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)] |
| 73 | pub 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. |
| 80 | struct 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 | |
| 88 | impl 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 | |
| 100 | impl<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)] |
| 1130 | struct ParsedCommand { |
| 1131 | tag: String, |
| 1132 | command: String, |
| 1133 | args: String, |
| 1134 | } |
| 1135 | |
| 1136 | fn 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 | |
| 1148 | struct 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 | |
| 1155 | impl<'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)] |
| 1262 | enum 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)] |
| 1275 | enum Section { |
| 1276 | Whole, |
| 1277 | Header, |
| 1278 | Text, |
| 1279 | HeaderFields(Vec<String>), |
| 1280 | HeaderFieldsNot(Vec<String>), |
| 1281 | } |
| 1282 | |
| 1283 | fn 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(§ion_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 | |
| 1348 | fn 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 | |
| 1375 | fn 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 | |
| 1391 | fn 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 | |
| 1406 | fn 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 | |
| 1416 | fn 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 | |
| 1424 | fn 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 | |
| 1456 | fn 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 | |
| 1479 | fn 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 | |
| 1507 | fn 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. |
| 1516 | fn 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 | |
| 1552 | fn 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. |
| 1563 | fn 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 | |
| 1588 | fn 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 | |
| 1598 | async 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 | |
| 1608 | async 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 | |
| 1613 | async 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 | |
| 1618 | async 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 | |
| 1623 | async 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 | |
| 1654 | fn 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. |
| 1661 | fn 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. |
| 1678 | fn 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 | |
| 1683 | fn 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. |
| 1715 | fn 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. |
| 1727 | fn 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. |
| 1750 | fn 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`. |
| 1770 | fn 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. |
| 1787 | enum FetchSegment { |
| 1788 | Text(String), |
| 1789 | Literal { prefix: String, payload: Vec<u8> }, |
| 1790 | } |