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(); } }