oxedyne/fe2o3/fe2o3_net/src/email/file.rs
3.1 KiB, 14 runs
created by r1870400018:565, 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 | use crate::email::msg::EmailMessage; |
| 2 | |
| 3 | use oxedyne_fe2o3_core::prelude::*; |
| 4 | |
| 5 | use std::path::Path; |
| 6 | |
| 7 | use tokio::{ |
| 8 | fs::File, |
| 9 | io::{AsyncBufReadExt, BufReader}, |
| 10 | }; |
| 11 | |
| 12 | |
| 13 | pub struct MboxEmailIterator { |
| 14 | reader: BufReader<File>, |
| 15 | buffer: Vec<u8>, |
| 16 | } |
| 17 | |
| 18 | impl MboxEmailIterator { |
| 19 | pub async fn new<P: AsRef<Path>>(path: P) -> Outcome<Self> { |
| 20 | let file = match tokio::fs::File::open(path).await { |
| 21 | Ok(file) => file, |
| 22 | Err(e) => { |
| 23 | return Err(err!(e, "While opening mbox file."; IO, File, Read)); |
| 24 | } |
| 25 | }; |
| 26 | let reader = BufReader::new(file); |
| 27 | Ok(Self { |
| 28 | reader, |
| 29 | buffer: Vec::new(), |
| 30 | }) |
| 31 | } |
| 32 | |
| 33 | async fn read_line(&mut self) -> Outcome<Option<Vec<u8>>> { |
| 34 | self.buffer.clear(); |
| 35 | let result = self.reader.read_until(b'\n', &mut self.buffer).await; |
| 36 | let bytes_read = res!(result, IO, File, Read); |
| 37 | if bytes_read == 0 { |
| 38 | Ok(None) |
| 39 | } else { |
| 40 | if self.buffer.ends_with(&[b'\r', b'\n']) { |
| 41 | self.buffer.truncate(self.buffer.len() - 2); |
| 42 | } else if self.buffer.ends_with(&[b'\n']) { |
| 43 | self.buffer.truncate(self.buffer.len() - 1); |
| 44 | } |
| 45 | Ok(Some(self.buffer.clone())) |
| 46 | } |
| 47 | } |
| 48 | |
| 49 | async fn read_until_next_from(&mut self) -> Outcome<Option<Vec<u8>>> { |
| 50 | let mut email_content = Vec::new(); |
| 51 | let mut line_num: usize = 0; |
| 52 | loop { |
| 53 | match self.read_line().await { |
| 54 | Ok(Some(line)) => { |
| 55 | debug!("{:05} line={}", line_num, String::from_utf8_lossy(&line)); |
| 56 | line_num += 1; |
| 57 | if line.starts_with(b"From ") { |
| 58 | if !email_content.is_empty() { |
| 59 | break; |
| 60 | } |
| 61 | } |
| 62 | email_content.extend_from_slice(&line); |
| 63 | email_content.extend_from_slice(&[b'\n']); |
| 64 | } |
| 65 | Ok(None) => { |
| 66 | if email_content.is_empty() { |
| 67 | return Ok(None); |
| 68 | } else { |
| 69 | break; |
| 70 | } |
| 71 | } |
| 72 | Err(e) => return Err(e), |
| 73 | } |
| 74 | } |
| 75 | Ok(Some(email_content)) |
| 76 | } |
| 77 | } |
| 78 | |
| 79 | impl Iterator for MboxEmailIterator { |
| 80 | type Item = Outcome<EmailMessage>; |
| 81 | |
| 82 | fn next(&mut self) -> Option<Self::Item> { |
| 83 | let runtime = match tokio::runtime::Runtime::new() { |
| 84 | Ok(runtime) => runtime, |
| 85 | Err(_) => return None, |
| 86 | }; |
| 87 | |
| 88 | runtime.block_on(async { |
| 89 | match self.read_until_next_from().await { |
| 90 | Ok(Some(content)) => { |
| 91 | let cursor = std::io::Cursor::new(content); |
| 92 | let mut reader = tokio::io::BufReader::new(cursor); |
| 93 | match EmailMessage::read(&mut reader).await { |
| 94 | Ok(email) => Some(Ok(email)), |
| 95 | Err(e) => Some(Err(e)), |
| 96 | } |
| 97 | } |
| 98 | Ok(None) => None, |
| 99 | Err(e) => Some(Err(e)), |
| 100 | } |
| 101 | }) |
| 102 | } |
| 103 | } |