Oregami
Repositories/oxedyne/fe2o3

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

1use crate::email::msg::EmailMessage;
2
3use oxedyne_fe2o3_core::prelude::*;
4
5use std::path::Path;
6
7use tokio::{
8 fs::File,
9 io::{AsyncBufReadExt, BufReader},
10};
11
12
13pub struct MboxEmailIterator {
14 reader: BufReader<File>,
15 buffer: Vec<u8>,
16}
17
18impl 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
79impl 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}