Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_net/src/smtp/msg.rs

5.7 KiB, 15 runs

created by r1870400018:601, 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::{
2 conc::AsyncReadIterator,
3 constant,
4 email::msg::EmailMessage,
5 smtp::cmd::SmtpCommand,
6};
7
8use oxedyne_fe2o3_core::{
9 prelude::*,
10 count::ErrorWhen,
11};
12
13use std::{
14 future::Future,
15 pin::Pin,
16};
17
18use tokio::{
19 io::{
20 AsyncRead,
21 //AsyncReadExt,
22 AsyncBufRead,
23 AsyncBufReadExt,
24 //AsyncWriteExt,
25 },
26};
27
28
29new_type!(SmtpMessage, Vec<SmtpCommand>, Debug, Default);
30
31impl SmtpMessage {
32
33 pub async fn read<R: AsyncRead + AsyncBufRead + Unpin>(
34 stream: &mut R,
35 _buffer: &mut Vec<u8>,
36 )
37 -> Outcome<Option<Self>>
38 {
39 let mut message = SmtpMessage::default();
40
41 let mut safety = ErrorWhen::new(constant::READ_LOOP_SAFETY_LIMIT);
42 loop {
43 res!(safety.inc());
44 let mut line = Vec::new();
45 let result = Self::read_line(stream, &mut line).await;
46 let byts_read = res!(result);
47
48 if byts_read == 0 {
49 if message.is_empty() {
50 return Ok(None);
51 } else {
52 break;
53 }
54 }
55
56 let line = String::from_utf8_lossy(&line);
57 let line = line.trim_end();
58
59 if line.is_empty() {
60 continue;
61 }
62
63 // Handle multiline responses.
64 if line.starts_with('-') {
65 // Remove the '-' prefix and append the line to the previous command.
66 let mut last_command = match message.pop() {
67 Some(command) => command,
68 None => return Err(err!(
69 "Unexpected end of message";
70 Unexpected, Input, Missing)),
71 };
72 let multiline_response = line[1..].trim();
73 match last_command {
74 SmtpCommand::Response(ref _code, ref mut msg) => {
75 msg.push('\n');
76 msg.push_str(multiline_response);
77 }
78 _ => return Err(err!(
79 "Unexpected multiline response: {}", line;
80 Invalid, Input)),
81 }
82 message.push(last_command);
83 continue;
84 }
85
86 // Parse the SMTP command
87 match SmtpCommand::from_str(line) {
88 Ok(command) => match command {
89 SmtpCommand::Data => {
90 let result = EmailMessage::read(stream).await;
91 let email = res!(result);
92 message.push(SmtpCommand::Email(email));
93 }
94 SmtpCommand::Quit => {
95 message.push(SmtpCommand::Quit);
96 break;
97 }
98 _ => {
99 message.push(command);
100 }
101 },
102 Err(e) => return Err(err!(e,
103 "Invalid SMTP command: {}", line;
104 Invalid, Input)),
105 }
106 }
107
108 Ok(Some(message))
109 }
110
111 pub async fn read_line<R: AsyncRead + AsyncBufRead + Unpin>(
112 stream: &mut R,
113 bfr: &mut Vec<u8>,
114 )
115 -> Outcome<usize>
116 {
117 let mut tmp_bfr = Vec::new();
118 let mut total_byts_read = 0;
119
120 let mut safety = ErrorWhen::new(constant::READ_LOOP_SAFETY_LIMIT);
121 loop {
122 debug!("safety={:?}", safety);
123 res!(safety.inc());
124 debug!("safety={:?}", safety);
125 let result = stream.read_until(b'\n', &mut tmp_bfr).await;
126 let byts_read = res!(result, IO, Network, Read);
127
128 if byts_read == 0 {
129 break;
130 }
131
132 // Ignore fragment consisting only of '\n'.
133 if byts_read > 1 {
134 if tmp_bfr[byts_read - 2] == b'\r' {
135 // We've found the \r\n to end the line, remove it.
136 if byts_read > 2 {
137 let byts = &tmp_bfr[..byts_read - 2];
138 bfr.extend_from_slice(byts);
139 total_byts_read += byts_read - 2;
140 }
141 break
142 } else {
143 // We only found an \n ending, continue accumulating.
144 let byts = &tmp_bfr[..byts_read];
145 bfr.extend_from_slice(byts);
146 total_byts_read += byts_read;
147 }
148 }
149 }
150 Ok(total_byts_read)
151 }
152
153}
154
155pub struct SmtpMessageReader<
156 'a,
157 R: AsyncRead + AsyncBufRead + Unpin + Send
158> {
159 stream: Pin<&'a mut R>,
160 buffer: Vec<u8>,
161}
162
163impl<
164 'a,
165 R: AsyncRead + AsyncBufRead + Unpin + Send
166>
167 SmtpMessageReader<'a, R>
168{
169 pub fn new(stream: Pin<&'a mut R>) -> Self {
170 Self {
171 stream,
172 buffer: Vec::new(),
173 }
174 }
175}
176
177impl<
178 'a,
179 R: AsyncRead + AsyncBufRead + Unpin + Send
180>
181 AsyncReadIterator for SmtpMessageReader<'a, R>
182{
183 type Item = Outcome<SmtpMessage>;
184
185 fn next<'b>(&'b mut self) -> Pin<Box<dyn Future<Output = Option<Self::Item>> + Send + 'b>> {
186 let mut stream = self.stream.as_mut();
187 let buffer = &mut self.buffer;
188
189 Box::pin(async move {
190 let result = SmtpMessage::read::<_>(
191 &mut stream.as_mut(),
192 buffer,
193 )
194 .await;
195
196 match result {
197 Ok(Some(message)) => Some(Ok(message)),
198 Ok(None) => None,
199 Err(e) => Some(Err(e)),
200 }
201 })
202 }
203}
204
205//#[derive(Debug)]
206//pub struct SmtpSession {
207// pub client_ip: IpAddr,
208// pub client_hostname: Fqdn,
209// pub messages: Vec<SmtpCommand>,
210//}