From 938379d6d7f03747c8d7aeac4be534c9e20cfd05 Mon Sep 17 00:00:00 2001 From: Mroik Date: Mon, 29 Jun 2026 18:19:41 +0200 Subject: Implement SMTP client This is used to forward requests to the main MTA to send emails to the subscribers of the mailing list. Signed-off-by: Mroik --- src/smtp_client.rs | 195 +++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 195 insertions(+) create mode 100644 src/smtp_client.rs (limited to 'src/smtp_client.rs') diff --git a/src/smtp_client.rs b/src/smtp_client.rs new file mode 100644 index 0000000..3debe7c --- /dev/null +++ b/src/smtp_client.rs @@ -0,0 +1,195 @@ +use std::net::SocketAddr; + +use anyhow::{Result, anyhow}; +use tokio::{ + io::{AsyncBufReadExt, AsyncWriteExt, BufReader}, + net::{TcpSocket, TcpStream}, +}; + +use crate::model::User; + +const CLIENT_NAME: &str = "poul-list"; + +pub struct SmtpClientConnection { + pub stream: TcpStream, +} + +impl Drop for SmtpClientConnection { + fn drop(&mut self) { + let _ = self.stream.try_write(b"QUIT"); + } +} + +impl SmtpClientConnection { + /// Returns the hostname of the SMTP server + async fn init(&mut self) -> Result { + let mut buffer = String::new(); + let (r, mut w) = self.stream.split(); + let mut reader = BufReader::new(r); + + reader.read_line(&mut buffer).await?; + if !buffer.starts_with("220") { + return Err(anyhow!( + "server '{}' is not ready (not 220)", + self.stream.peer_addr()? + )); + } + + w.write_all(format!("HELO {}\r\n", CLIENT_NAME).as_bytes()) + .await?; + + buffer.clear(); + reader.read_line(&mut buffer).await?; + if !buffer.starts_with("250") { + return Err(anyhow!( + "could not start SMTP connection with '{}' (not 250 on HELO)", + self.stream.peer_addr()? + )); + } + + let mut ris = String::new(); + match buffer.split(' ').nth(1) { + Some(x) => ris.push_str(x), + None => (), + } + Ok(ris) + } + + /// Returns an error if it could not reset the buffers. In that the caller should just drop the + /// connection as the server is in an unreliable state. + async fn reset(&mut self) -> Result> { + let (r, mut w) = self.stream.split(); + let mut reader = BufReader::new(r); + let mut buffer = String::new(); + w.write_all(b"RSET\r\n").await?; + reader.read_line(&mut buffer).await?; + if !buffer.starts_with("250") { + return Ok(Some(anyhow!( + "could not reset buffers with '{}'. Connection should be terminated\r\n", + self.stream.peer_addr()? + ))); + } + Ok(None) + } + + async fn send_email(&mut self, data: &str, recipients: &[User], sender: &str) -> Result<()> { + let (r, mut w) = self.stream.split(); + let mut reader = BufReader::new(r); + let mut buffer = String::new(); + + w.write_all(format!("MAIL FROM:<{}>\r\n", sender).as_bytes()) + .await?; + reader.read_line(&mut buffer).await?; + if !buffer.starts_with("250") { + self.reset().await?; + return Err(anyhow!( + "could not initiate a mail transfer with '{}' (not 250)", + self.stream.peer_addr()? + )); + } + + for recp in recipients { + w.write_all(format!("RCPT TO:<{}>\r\n", recp.email).as_bytes()) + .await?; + buffer.clear(); + reader.read_line(&mut buffer).await?; + if !buffer.starts_with("250") { + self.reset().await?; + return Err(anyhow!( + "'{}' rejected RCTP command (not 250)", + self.stream.peer_addr()? + )); + } + } + + w.write_all(b"DATA\r\n").await?; + buffer.clear(); + reader.read_line(&mut buffer).await?; + if !buffer.starts_with("354") { + self.reset().await?; + return Err(anyhow!( + "'{}' not ready to receive DATA (not 354)", + self.stream.peer_addr()? + )); + } + + w.write_all(data.as_bytes()).await?; + if !data.ends_with("\r\n") { + w.write_all(b"\r\n").await?; + } + w.write_all(b".\r\n").await?; + + buffer.clear(); + reader.read_line(&mut buffer).await?; + if !buffer.starts_with("250") { + self.reset().await?; + return Err(anyhow!( + "'{}' refused to queue (not 250)", + self.stream.peer_addr()? + )); + } + + Ok(()) + } +} + +pub async fn connect(server: &str) -> Result { + let remote_peer: SocketAddr = server.parse()?; + let sock = if remote_peer.is_ipv4() { + TcpSocket::new_v4()? + } else { + TcpSocket::new_v6()? + }; + + let mut conn = SmtpClientConnection { + stream: sock.connect(remote_peer).await?, + }; + conn.init().await?; + Ok(conn) +} + +#[cfg(test)] +mod tests { + use tokio::{spawn, sync::mpsc::channel}; + + use crate::{model::User, smtp_client::connect, smtp_server::SmtpServer}; + + #[tokio::test] + async fn send_email() { + let (tx, mut rx) = channel(5); + let mut server = SmtpServer::new([127, 0, 0, 1], 9090, "testerino") + .await + .expect("init server"); + let j = spawn(async move { + server.run(tx).await.expect("start server"); + }); + + let mut connection = connect("127.0.0.1:9090").await.expect("connect to server"); + let lines = [ + "Date: 29 Jun 2026 15:18:58 -0000", + "From: Mroik ", + "To: testerino@example.org", + "Subject: This is a test", + "", + "Please ignore this, thanks.", + "", + "Mroik", + ]; + + connection + .send_email( + lines.join("\r\n").as_str(), + &[User { + name: Some(String::from("testerino")), + email: String::from("testerino@example.org"), + }], + "mroik@delayed.space", + ) + .await + .expect("send email"); + + let message = rx.recv().await.expect("retrieve from channel"); + assert_eq!(message.data.trim(), lines.join("\r\n").trim()); + j.abort(); + } +} -- cgit v1.3.1