diff --git a/protocols/turnloop-mongodb/README.md b/protocols/turnloop-mongodb/README.md index 9a40cb8..c1206ac 100644 --- a/protocols/turnloop-mongodb/README.md +++ b/protocols/turnloop-mongodb/README.md @@ -44,25 +44,49 @@ the core. WASI 0.2 supports A/AAAA through ip-name-lookup; its standard interfac has no SRV/TXT capability, so use explicit resolved seed URIs there. +## Host entropy + +The sans-I/O core reads no OS entropy and no clock, so every random or unique value +comes from the host. Two are easy to miss because nothing fails locally: a weak SCRAM +nonce authenticates anyway, and an insert without `_id` succeeds against a server. + +| Obligation | Where it goes | What to supply | +|---|---|---| +| SCRAM client nonce | `Connection::connected(now, nonce)` | For a connection with credentials, a fresh CSPRNG nonce per connection: at least 16 printable ASCII bytes and no `,` (base64 of 18 random bytes works). `Scram::new` rejects a short or unprintable nonce; it cannot detect a reused or predictable one. Pass `""` without credentials. | +| Document `_id` | each inserted document, before encoding | Every document sent by `insert` (or through `BulkBatcher`) must already carry `_id`. The server otherwise creates one the host never learns, where the Node driver would have generated it locally and returned it as the inserted id. Create one `ObjectIdGenerator::new(random, counter)` per process with 5 random bytes and a random 24-bit counter start, then call `generate(unix_seconds)` for each document. | +| Session identity | `Session::new(id)`, `RetrySession { id, .. }` | A random RFC 4122 UUID per server session. `SessionPool::checkout` returns `None` when a new one is needed. | +| Server selection | `Topology::choose` / `select_reusing` | Two random `u64`s per selection, for an unbiased pick within the latency window. | + +The `turnloop` feature's `asynchronous` client fills the nonce, session identity and +selection entropy from rustls' ring `SecureRandom`. Adding `_id` remains the caller's job +in every mode. + + ## Driving a connection 1. Parse `Options`. For `mongodb+srv`, execute both `resolution_requests()` and provide the answers to `resolve()`. The host opens the resulting address using its transport. 2. Construct `Connection`, call `connected(now, nonce)`. A nonce must be supplied by a - host CSPRNG for authenticated connections. `UpgradeTls` means finish certificate- - validated TLS on the transport and call `tls_established()` before sending Mongo bytes. + host CSPRNG for authenticated connections (see [Host entropy](#host-entropy)). + `UpgradeTls` means finish certificate-validated TLS on the transport and call + `tls_established()` before sending Mongo bytes. 3. Write `transmit()` and acknowledge only bytes actually written with `consume_transmit(n)`. The slice must not be retained across mutating calls; an adapter using completion I/O keeps the connection immobile until write completion. 4. Feed plaintext bytes through `receive()`. It returns the consumed prefix: loop on the remainder. It deliberately stops at a frame boundary. Drain events and release a complete reply before feeding another frame. Partial reads and writes are normal. -5. On `Ready`, submit one command. An accepted token receives one `Reply`, `Failed` or - `Unacknowledged`; rejected submissions receive no completion. `Reply` is a borrowed - view obtained with `reply()`. Consume or copy it, then `release_reply()`. + `accepts_receive()` is the same check `receive()` makes, so a host can ask whether + to read instead of tracking an `expecting_reply` flag of its own. +5. On `Ready`, submit one command. An accepted token is settled exactly once: by + `Failed`, by `Unacknowledged`, or by `release_reply()` after its `Reply` event; + rejected submissions receive no completion. `Reply` is a borrowed view obtained with + `reply()`. Consume or copy it, then `release_reply()`. 6. Pass transport errors to `fail()`. Drive `handle_timeout(now)` at `next_timeout()`. - A timed-out connection must be physically closed by the adapter. Closing cancels an - outstanding command once, followed by `Closed`. A reply already received remains + A timed-out connection must be physically closed by the adapter. `fail()` settles an + outstanding command, or a reply not yet released, with `Failed`: the reply is revoked + and a still-queued `Reply` event for it is withdrawn. `close()` cancels an + outstanding command once, followed by `Closed`, but a reply already received remains readable after close until released. `Instant` is std::time::Instant on native/WASI. Browser Wasm uses `HostInstant`, created @@ -87,11 +111,15 @@ assert_eq!(command.raw().get_str("find").unwrap(), "tasks"); Insert/update/delete use OP_MSG document sequences named documents/updates/deletes. `BulkBatcher` splits borrowed models by negotiated count, BSON and wire size limits; -pass the **final decorated body** when calculating overhead. `BulkResult` preserves -original error indices across batches. It exposes aggregate wire counts; the caller +pass the **final decorated body** when calculating overhead. A failed write can arrive +inside `ok: 1`: a duplicate key is reported in `writeErrors` and an unmet write concern +in `writeConcernError`. `WriteResult::parse` is the verdict — it runs +`Error::from_response` first and fails on either, so no separate check is needed. +`WriteResult::decode` keeps both as data; check `succeeded()` on its result. `BulkResult` +uses `decode` and preserves original error indices across batches. It exposes aggregate wire counts; the caller maps them and inserted/upserted IDs to the desired JS result structure. Missing `_id` -values must be added before sending; `ObjectIdGenerator` accepts host entropy and Unix -seconds and never reads a clock or OS entropy itself. +values must be added before sending with `ObjectIdGenerator`; see +[Host entropy](#host-entropy). `CursorBatch::rows()` yields borrowed RawDocuments without per-row allocation. `Cursor` tracks namespace, server affinity and client limit. `needs_kill()` means issue diff --git a/protocols/turnloop-mongodb/src/auth.rs b/protocols/turnloop-mongodb/src/auth.rs index b7249ce..3a4c2f5 100644 --- a/protocols/turnloop-mongodb/src/auth.rs +++ b/protocols/turnloop-mongodb/src/auth.rs @@ -1,5 +1,6 @@ //! auth/auth.md §§ SCRAM-SHA-1, SCRAM-SHA-256 and mongodb-handshake/handshake.md -//! § Speculative Authentication. Entropy is provided by the host, never acquired here. +//! § Speculative Authentication. Entropy is provided by the host, never acquired here; +//! see the crate's [host entropy](crate#host-entropy) obligations. use crate::{Error, ErrorKind, Result, uri::Credential}; use base64::{Engine, engine::general_purpose::STANDARD}; use bson::{Binary, Document, doc, spec::BinarySubtype}; diff --git a/protocols/turnloop-mongodb/src/command.rs b/protocols/turnloop-mongodb/src/command.rs index 9c09602..95b708f 100644 --- a/protocols/turnloop-mongodb/src/command.rs +++ b/protocols/turnloop-mongodb/src/command.rs @@ -1,6 +1,8 @@ //! crud/crud.md §§ Read/Write Operations, Write Models and Results. //! Builders retain capacity. Inputs and options are borrowed raw BSON; IDs and //! wall-clock ObjectId timestamps are supplied by the host (no ObjectId::new()). +//! Inserted documents need an `_id` before sending; see [`ObjectIdGenerator`] and +//! the crate's [host entropy](crate#host-entropy) obligations. use crate::{Error, ErrorKind, Result, wire::BsonWriter}; use bson::{ Document, @@ -468,7 +470,8 @@ impl<'a> CursorBatch<'a> { .and_then(|v| v.as_document()), }) } - pub fn rows(&self) -> impl Iterator> { + /// The iterator borrows the reply, not this batch, so it may outlive `self`. + pub fn rows(&self) -> impl Iterator> + use<'a> { self.documents.into_iter().map(|v| match v { Ok(RawBsonRef::Document(d)) => Ok(d), _ => Err(Error::protocol("Cursor batch contains non-document")), @@ -533,16 +536,39 @@ impl Cursor { Ok(true) } } +/// Counts and error detail of an insert/update/delete reply. +/// +/// MongoDB reports a failed write *inside* a successful command: a duplicate +/// key answers `ok: 1` with a `writeErrors` array, and an unsatisfied write +/// concern answers `ok: 1` with `writeConcernError`. `ok` alone is therefore no +/// verdict. [`WriteResult::parse`] is the verdict: it runs +/// [`Error::from_response`] first and fails on either field, so a caller never +/// needs to call `from_response` beforehand. [`WriteResult::decode`] keeps them +/// as data for aggregation (as [`BulkResult`] does); check +/// [`WriteResult::succeeded`] on what it returns. #[derive(Debug)] pub struct WriteResult<'a> { pub count: i64, pub modified_count: i64, pub upserted: Option<&'a RawArray>, + /// Never nonempty in a [`WriteResult::parse`] result. pub write_errors: Option<&'a RawArray>, + /// Never present in a [`WriteResult::parse`] result. pub write_concern_error: Option<&'a RawDocument>, } impl<'a> WriteResult<'a> { + /// Succeeds only for a write that fully succeeded. `ok: 0`, a nonempty + /// `writeErrors` (a [`ErrorKind::BulkWrite`] error) and a + /// `writeConcernError` (a [`ErrorKind::Server`] error) all fail with the + /// error [`Error::from_response`] builds, full server response included. pub fn parse(r: &'a RawDocument) -> Result { + Error::from_response(r)?; + Self::decode(r) + } + /// Fails only when the command itself failed (`ok: 0`) or the reply is + /// malformed; per-document and write-concern errors are returned in the + /// fields, so callers must consult [`WriteResult::succeeded`]. + pub fn decode(r: &'a RawDocument) -> Result { if crate::error::number(r, "ok").unwrap_or(0.0) == 0.0 { Error::from_response(r)?; } @@ -562,6 +588,12 @@ impl<'a> WriteResult<'a> { .and_then(|v| v.as_document()), }) } + /// False when any document failed or the write concern was not satisfied. + pub fn succeeded(&self) -> bool { + self.write_errors + .is_none_or(|a| a.into_iter().next().is_none()) + && self.write_concern_error.is_none() + } } fn integer(d: &RawDocument, k: &str) -> Result { match d.get(k).ok().flatten() { @@ -582,7 +614,7 @@ pub struct BulkResult { } impl BulkResult { pub fn accept(&mut self, reply: &RawDocument, offset: i32, ordered: bool) -> Result { - let result = WriteResult::parse(reply)?; + let result = WriteResult::decode(reply)?; self.count += result.count; self.modified_count += result.modified_count; let mut errors = false; @@ -688,6 +720,12 @@ impl<'a> BulkBatcher<'a> { } /// ObjectId generator with host-supplied process entropy and wall-clock seconds. /// BSON ObjectId specification § Generation: 4 timestamp, 5 random, 3 counter bytes. +/// +/// Nothing in this crate adds a missing `_id`: an insert without one still +/// succeeds, but the server assigns an id the host never learns. Create one +/// generator per process from 5 random bytes and a random counter start, and +/// give every inserted document an `_id` from it. This is one of the crate's +/// [host entropy](crate#host-entropy) obligations. pub struct ObjectIdGenerator { random: [u8; 5], counter: u32, diff --git a/protocols/turnloop-mongodb/src/connection.rs b/protocols/turnloop-mongodb/src/connection.rs index 36d8be2..e2996bd 100644 --- a/protocols/turnloop-mongodb/src/connection.rs +++ b/protocols/turnloop-mongodb/src/connection.rs @@ -1,6 +1,9 @@ //! mongodb-handshake/handshake.md §§ Connection Handshake, Speculative Authentication. //! Events are pull-based. A token is accepted only on successful `command`, and -//! yields exactly one Reply, Failed, or Unacknowledged event before reuse. +//! is settled exactly once before reuse: by a Failed or Unacknowledged event, or +//! by `release_reply()` after its Reply event. `fail()` before that release +//! revokes the reply and settles the token with Failed instead; `close()` keeps +//! a received reply readable until it is released. use crate::Instant; use crate::{ Error, ErrorKind, Result, @@ -81,7 +84,9 @@ impl Connection { max_write_batch_size: 100_000, } } - /// `nonce` must come from a host CSPRNG if credentials are configured. + /// `nonce` must come from a host CSPRNG if credentials are configured: at least + /// 16 printable bytes, fresh per connection. See the crate's + /// [host entropy](crate#host-entropy) obligations. pub fn connected(&mut self, now: Instant, nonce: &str) -> Result<()> { if self.state != State::New { return Err(Error::protocol("Connection already started")); @@ -209,17 +214,17 @@ impl Connection { } Ok(()) } + /// Whether `receive` would accept bytes now: a handshake, authentication or + /// command reply is outstanding. This is the exact check `receive` makes, so + /// a host need not shadow the state machine to decide when to read. It is + /// false before `connected`, during a TLS upgrade, when idle, while an + /// unacknowledged write drains, while a reply is unreleased and after close. + pub fn accepts_receive(&self) -> bool { + matches!(self.state, State::Handshake | State::Auth | State::Command) + } /// Feed may consume only a frame prefix; retain and feed the remaining bytes. pub fn receive(&mut self, input: &[u8]) -> Result { - if matches!( - self.state, - State::New - | State::Tls - | State::Ready - | State::UnackSending - | State::Reply - | State::Closed - ) { + if !self.accepts_receive() { return Err(Error::protocol("Connection is not expecting a reply")); } let result = self.receive_inner(input); @@ -524,16 +529,23 @@ impl Connection { )); } } + /// Closes the connection after a transport or protocol failure. An + /// outstanding command, or a reply not yet released, settles with `Failed`. + /// Such a reply is revoked: `reply()` and `release_reply()` then fail, and + /// its `Reply` event is withdrawn if the host has not yet polled it. pub fn fail(&mut self, error: Error) { if self.state == State::Closed { return; } - if self.state != State::Reply { - self.events.push_back(ConnectionEvent::Failed { - token: self.token.take(), - error, - }); + let token = self.token.take(); + if self.state == State::Reply { + self.decoder.clear(); + self.expanded.clear(); + self.events + .retain(|e| !matches!(e, ConnectionEvent::Reply { token: t } if Some(*t) == token)); } + self.events + .push_back(ConnectionEvent::Failed { token, error }); self.state = State::Closed; self.deadline = None; self.tx.clear(); diff --git a/protocols/turnloop-mongodb/src/lib.rs b/protocols/turnloop-mongodb/src/lib.rs index 78e3d2d..0200832 100644 --- a/protocols/turnloop-mongodb/src/lib.rs +++ b/protocols/turnloop-mongodb/src/lib.rs @@ -9,6 +9,19 @@ //! host owns `LocalExecutor` and calls `turn`; adapters await its streams and //! deadline futures. See the crate README and turnloop-io for ownership, streaming //! and cancellation examples. Default features retain the sans-I/O API. +//! +//! # Host entropy +//! Every random or unique value comes from the host; nothing here reads OS +//! entropy or a clock. The obligations, all described in one place in the +//! [README's Host entropy section](https://github.com/PerryTS/turnloop/blob/main/protocols/turnloop-mongodb/README.md#host-entropy): +//! - a CSPRNG SCRAM nonce for [`Connection::connected`] when credentials are set; +//! - an `_id` on every inserted document, from [`command::ObjectIdGenerator`] +//! (the server otherwise assigns one the host never learns); +//! - random session UUIDs for [`session::Session::new`] and +//! [`operation::RetrySession`]; +//! - selection entropy for [`topology::Topology::choose`]. +//! +//! The `turnloop` feature's client supplies all but the `_id`. #![deny(unsafe_op_in_unsafe_fn)] #![forbid(unsafe_code)] diff --git a/protocols/turnloop-mongodb/src/session.rs b/protocols/turnloop-mongodb/src/session.rs index 04bb066..7f5c9e9 100644 --- a/protocols/turnloop-mongodb/src/session.rs +++ b/protocols/turnloop-mongodb/src/session.rs @@ -32,7 +32,8 @@ pub struct Session { pub pinned_server: Option, } impl Session { - /// Supply an RFC4122 UUID generated by the host. No entropy or clock reads here. + /// Supply an RFC4122 UUID generated by the host. No entropy or clock reads here; + /// see the crate's [host entropy](crate#host-entropy) obligations. pub fn new(id: [u8; 16]) -> Self { Self { id, diff --git a/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs b/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs new file mode 100644 index 0000000..7fe8b01 --- /dev/null +++ b/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs @@ -0,0 +1,63 @@ +//! `CursorBatch::rows()` captures only the reply lifetime (issue #69). Under +//! Rust 2024 capture rules an `impl Iterator` return also captures `&self` +//! unless it says `use<'a>`: the iterator then borrows the local batch although +//! its rows borrow only the reply, and `rows` below fails with E0597. The +//! `match` in `first_row` fails the same way in edition-2021 host crates, whose +//! tail-expression temporaries outlive the block's locals. +use turnloop_mongodb::{ + Error, Result, + bson::{ + doc, + raw::{RawDocument, RawDocumentBuf}, + }, + command::CursorBatch, +}; + +/// The tail-position pattern from the issue: the scrutinee borrows `batch`, +/// a local dropped at the end of the block, while the row escapes it. +#[allow( + clippy::manual_map, + clippy::needless_match, + reason = "the issue's tail-position match, spelled out" +)] +fn first_row(reply: &RawDocument) -> Option> { + let batch = CursorBatch::parse(reply).ok()?; + match batch.rows().next() { + Some(row) => Some(row), + None => None, + } +} + +/// Returning the iterator itself requires that it not borrow the local batch. +fn rows(reply: &RawDocument) -> Result>> { + let batch = CursorBatch::parse(reply)?; + Ok(batch.rows()) +} + +fn x(row: Result<&RawDocument>) -> Result { + row?.get_i32("x").map_err(|_| Error::protocol("missing x")) +} + +#[test] +fn rows_outlive_the_batch_that_produced_them() { + let reply = RawDocumentBuf::try_from(&doc! { + "ok": 1, + "cursor": {"id": 0_i64, "ns": "db.items", "firstBatch": [{"x": 1}, {"x": 2}]}, + }) + .expect("reply"); + assert_eq!(x(first_row(&reply).expect("first row")).expect("x"), 1); + let xs: Vec = rows(&reply) + .expect("batch") + .map(x) + .collect::>() + .expect("rows"); + assert_eq!(xs, [1, 2]); + + let empty = RawDocumentBuf::try_from(&doc! { + "ok": 1, + "cursor": {"id": 0_i64, "ns": "db.items", "nextBatch": []}, + }) + .expect("reply"); + assert!(first_row(&empty).is_none()); + assert_eq!(rows(&empty).expect("batch").count(), 0); +} diff --git a/protocols/turnloop-mongodb/tests/protocol.rs b/protocols/turnloop-mongodb/tests/protocol.rs index ec47658..ad9dc11 100644 --- a/protocols/turnloop-mongodb/tests/protocol.rs +++ b/protocols/turnloop-mongodb/tests/protocol.rs @@ -867,3 +867,222 @@ fn compressed_commands_are_independent_zlib_streams() { assert_eq!(decoded_rounds, 64, "every round must be decoded"); assert_eq!(shrank, 64, "every round must actually compress"); } + +/// Writes the pending request and returns a server reply to it. +fn answer(c: &mut Connection, body: &Document) -> Vec { + let req = wire::i32_at(c.transmit(), 4).unwrap(); + let n = c.transmit().len(); + c.consume_transmit(n).unwrap(); + let mut reply = Vec::new(); + wire::encode(&mut reply, 1, req, 0, &raw(body), &[], 1000).unwrap(); + reply +} + +/// `accepts_receive` must be the check `receive` makes: a refused `receive` +/// changes nothing, and an accepted empty one consumes nothing. +fn accepts(c: &mut Connection) -> bool { + let accepts = c.accepts_receive(); + assert_eq!(c.receive(&[]).is_ok(), accepts); + accepts +} + +#[test] +fn accepts_receive_tracks_every_connection_state() { + let body = raw(&doc! {"ping":1,"$db":"admin"}); + // New, then a requested TLS upgrade, then the handshake it releases. + let mut c = Connection::new(Options::parse("mongodb://a/?tls=true").unwrap()); + assert!(!accepts(&mut c)); + c.connected(clock::now(), "").unwrap(); + assert!(matches!(c.poll_event(), Some(ConnectionEvent::UpgradeTls))); + assert!(!accepts(&mut c)); + c.tls_established().unwrap(); + assert!(accepts(&mut c)); + + // Authentication continues to expect replies after the handshake. + let mut c = Connection::new(Options::parse("mongodb://user:pencil@a/").unwrap()); + c.connected(clock::now(), "fyko+d2lbbFgONRv9qkxdawL") + .unwrap(); + assert!(accepts(&mut c)); + let hello = answer(&mut c, &doc! {"ok":1,"maxWireVersion":27}); + feed(&mut c, &hello); + assert!(c.poll_event().is_none(), "saslStart, not Ready, follows"); + assert!(!c.transmit().is_empty()); + assert!(accepts(&mut c)); + + // Ready, an outstanding command, then its unreleased reply. + let mut c = ready(); + assert!(!accepts(&mut c)); + c.command(1, &body, &[], clock::now()).unwrap(); + assert!(accepts(&mut c)); + let reply = answer(&mut c, &doc! {"ok":1}); + feed(&mut c, &reply); + assert!(matches!( + c.poll_event(), + Some(ConnectionEvent::Reply { token: 1 }) + )); + assert!(!accepts(&mut c)); + c.release_reply().unwrap(); + assert!(!accepts(&mut c)); + assert!(c.is_ready()); + + // An unacknowledged write expects no reply while or after it drains. + let unack = raw(&doc! {"insert":"x","writeConcern":{"w":0},"$db":"test"}); + c.command(2, &unack, &[], clock::now()).unwrap(); + assert!(!accepts(&mut c)); + let n = c.transmit().len(); + c.consume_transmit(n).unwrap(); + assert!(matches!( + c.poll_event(), + Some(ConnectionEvent::Unacknowledged { token: 2 }) + )); + assert!(!accepts(&mut c)); + + // Closed, including when a command was outstanding. + c.command(3, &body, &[], clock::now()).unwrap(); + assert!(accepts(&mut c)); + c.close(); + assert!(!accepts(&mut c)); + assert!(matches!( + c.poll_event(), + Some(ConnectionEvent::Failed { token: Some(3), .. }) + )); + assert!(matches!(c.poll_event(), Some(ConnectionEvent::Closed))); + assert!(c.poll_event().is_none()); +} + +#[test] +fn write_result_verdict_fails_on_write_errors_despite_ok() { + use turnloop_mongodb::command::{BulkResult, WriteResult}; + // A duplicate key is `ok: 1` with the failure in writeErrors. + let duplicate = raw(&doc! {"ok":1,"n":0,"writeErrors":[ + {"index":0,"code":11000,"codeName":"DuplicateKey","errmsg":"E11000 duplicate key"} + ]}); + let e = WriteResult::parse(&duplicate).unwrap_err(); + assert_eq!(e.kind, ErrorKind::BulkWrite); + assert_eq!(e.code, Some(11000)); + assert_eq!(e.message, "E11000 duplicate key"); + assert!( + e.response + .as_ref() + .unwrap() + .get_array("writeErrors") + .is_ok() + ); + let decoded = WriteResult::decode(&duplicate).unwrap(); + assert!(!decoded.succeeded()); + assert_eq!(decoded.count, 0); + assert_eq!(decoded.write_errors.unwrap().into_iter().count(), 1); + + // The write applied, but the requested durability was not confirmed. + let concern = raw(&doc! {"ok":1,"n":1,"writeConcernError": + {"code":64,"codeName":"WriteConcernFailed","errmsg":"waiting for replication timed out"} + }); + let e = WriteResult::parse(&concern).unwrap_err(); + assert_eq!(e.kind, ErrorKind::Server); + assert_eq!(e.code, Some(64)); + let decoded = WriteResult::decode(&concern).unwrap(); + assert!(!decoded.succeeded()); + assert_eq!(decoded.count, 1); + + // Clean writes, including an explicitly empty writeErrors, succeed. + for clean in [ + raw(&doc! {"ok":1,"n":2,"nModified":1}), + raw(&doc! {"ok":1.0,"n":2,"nModified":1,"writeErrors":[]}), + ] { + let parsed = WriteResult::parse(&clean).unwrap(); + assert!(parsed.succeeded()); + assert_eq!((parsed.count, parsed.modified_count), (2, 1)); + assert!(WriteResult::decode(&clean).unwrap().succeeded()); + } + + // A failed command is an error either way. + let failed = raw(&doc! {"ok":0,"code":13,"errmsg":"unauthorized"}); + assert_eq!(WriteResult::parse(&failed).unwrap_err().code, Some(13)); + assert_eq!(WriteResult::decode(&failed).unwrap_err().code, Some(13)); + + // Aggregation still sees per-document errors instead of an early return. + let mut bulk = BulkResult::default(); + assert!(!bulk.accept(&duplicate, 5, true).unwrap()); + assert!(bulk.accept(&concern, 6, false).unwrap()); + assert_eq!(bulk.count, 1); + assert_eq!(bulk.write_errors[0].get_i32("index").unwrap(), 5); + assert_eq!(bulk.write_concern_errors[0].get_i32("code").unwrap(), 64); +} + +/// Drains the queue and returns the settlement of `token` plus whether the +/// terminal `Closed` came last, asserting no event for it appears twice. +fn settlements(c: &mut Connection, token: u64) -> (Vec<&'static str>, bool) { + let mut seen = Vec::new(); + let mut closed_last = false; + while let Some(event) = c.poll_event() { + closed_last = matches!(event, ConnectionEvent::Closed); + match event { + ConnectionEvent::Reply { token: t } if t == token => seen.push("reply"), + ConnectionEvent::Failed { + token: Some(t), + error, + } if t == token => { + assert_eq!(error.kind, ErrorKind::Network); + seen.push("failed"); + } + ConnectionEvent::Closed => {} + e => panic!("unexpected event {e:?}"), + } + } + (seen, closed_last) +} + +#[test] +fn fail_settles_an_unreleased_reply_exactly_once() { + let body = raw(&doc! {"ping":1,"$db":"admin"}); + let reset = || Error::new(ErrorKind::Network, "connection reset"); + + // The reply arrived but its event was never polled: the queued Reply is + // withdrawn so the token's only settlement is Failed. + let mut c = ready(); + c.command(90, &body, &[], clock::now()).unwrap(); + let reply = answer(&mut c, &doc! {"ok":1}); + feed(&mut c, &reply); + c.fail(reset()); + assert_eq!(settlements(&mut c, 90), (vec!["failed"], true)); + assert!(c.reply().is_err()); + assert!(c.release_reply().is_err()); + + // The host took the Reply but had not released it: Failed settles it. + let mut c = ready(); + c.command(91, &body, &[], clock::now()).unwrap(); + let reply = answer(&mut c, &doc! {"ok":1}); + feed(&mut c, &reply); + assert!(matches!( + c.poll_event(), + Some(ConnectionEvent::Reply { token: 91 }) + )); + assert_eq!(c.reply().unwrap().get_i32("ok").unwrap(), 1); + c.fail(reset()); + assert_eq!(settlements(&mut c, 91), (vec!["failed"], true)); + assert!(c.reply().is_err()); + assert!(c.release_reply().is_err()); + // Nothing further is settled by a second failure or a close. + c.fail(reset()); + c.close(); + assert!(c.poll_event().is_none()); + + // A released reply was settled by the release; failing afterwards reports + // only the connection failure, for no token. + let mut c = ready(); + c.command(92, &body, &[], clock::now()).unwrap(); + let reply = answer(&mut c, &doc! {"ok":1}); + feed(&mut c, &reply); + assert!(matches!( + c.poll_event(), + Some(ConnectionEvent::Reply { token: 92 }) + )); + c.release_reply().unwrap(); + c.fail(reset()); + assert!(matches!( + c.poll_event(), + Some(ConnectionEvent::Failed { token: None, .. }) + )); + assert!(matches!(c.poll_event(), Some(ConnectionEvent::Closed))); + assert!(c.poll_event().is_none()); +} diff --git a/protocols/turnloop-smtp/README.md b/protocols/turnloop-smtp/README.md index 2a75310..857ce75 100644 --- a/protocols/turnloop-smtp/README.md +++ b/protocols/turnloop-smtp/README.md @@ -24,7 +24,10 @@ final mailbox delivery; the core never retries a message automatically. TLS options correspond to secure/requireTLS/ignoreTLS: Implicit for secure=true (or host's port-465 default), Required, None for ignoreTLS, otherwise -Opportunistic. STARTTLS is followed by a fresh EHLO and capability parsing. AUTH +Opportunistic. Required never reaches `Ready` in the clear: if the server does not +offer STARTTLS or refuses it, the session fails with code `ETLS` before any AUTH is +sent, so hosts need no `Ready` guard of their own. STARTTLS is followed by a fresh +EHLO and capability parsing. AUTH PLAIN, LOGIN and XOAUTH2 are selectable; OAuth token acquisition/refresh is host policy. No CRAM-MD5. PIPELINING sends MAIL and all RCPT commands together, drains all responses and sends DATA only with an accepted recipient. SIZE uses @@ -41,8 +44,21 @@ are rejected. MIME building materializes owned representations; this is separate from the transport hot path. The transport reuses TX/RX, response, SASL and DATA buffers; warmed success allocates only the returned accepted-vector, recipient string and response string (verified by counting allocator). Maximum response -buffer defaults to 64 KiB. Message size is limited by server SIZE when supplied; -there is no streaming body API yet, so hosts must bound their MIME inputs. +buffer defaults to 64 KiB. Message size is limited by server SIZE when supplied. + +`send` is the one-shot convenience: it copies the whole message into DATA storage. +For large messages, `start_send(token, envelope, message_id, StreamBody { size, +eight_bit }, now)` runs the same MAIL/RCPT/DATA exchange, then emits `BodyReady` +once the server answers 354. Pass the content in any number of `send_chunk(bytes, +now)` calls and end it with `finish_body(now)`; the single `Sent`/`Failed` follows +as for `send`. Draining `output()` between chunks keeps memory bounded by one +chunk. `DataEncoder` normalizes line endings and dot-stuffs incrementally, so a +CRLF pair or a CRLF.CRLF split across chunks encodes exactly as in one piece. +`size` is the declared SIZE parameter (omitted when `None`) and `eight_bit` +declares BODY=8BITMIME up front, since both are sent before any content. Content +cannot be retracted mid-DATA, so undeclared 8-bit bytes or a server reply while +the body is open fail the message and close the transport instead of sending a +terminator. Error fields follow nodemailer response/responseCode/command/code conventions. AUTH failures include mechanism and server response in the message. Full exact diff --git a/protocols/turnloop-smtp/src/lib.rs b/protocols/turnloop-smtp/src/lib.rs index c03b7ad..bbf8f48 100644 --- a/protocols/turnloop-smtp/src/lib.rs +++ b/protocols/turnloop-smtp/src/lib.rs @@ -18,9 +18,16 @@ use std::{ #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Tls { + /// Never attempt STARTTLS (nodemailer `ignoreTLS`). None, + /// Upgrade with STARTTLS when the server offers it, otherwise stay in the clear. Opportunistic, + /// Upgrade with STARTTLS or fail. A server that does not offer STARTTLS, or + /// refuses it, yields `Failed` with code `ETLS` before any AUTH is sent, + /// followed by `CloseTransport` and `Closed`; `Ready` is never reached in + /// the clear (nodemailer `requireTLS`). Required, + /// TLS from the first byte (nodemailer `secure`). Implicit, } #[derive(Clone)] @@ -163,6 +170,11 @@ pub enum Event { token: u64, info: SendInfo, }, + /// A streamed message's DATA was accepted (354): write its content with + /// [`Connection::send_chunk`] and end it with [`Connection::finish_body`]. + BodyReady { + token: u64, + }, Failed { token: Option, error: Error, @@ -209,6 +221,9 @@ pub struct Connection { auth_raw: Vec, auth_encoded: String, body: Vec, + streaming: Option, + encoder: DataEncoder, + body_open: bool, events: VecDeque, deadline: Option, token: Option, @@ -239,6 +254,9 @@ impl Connection { auth_raw: Vec::with_capacity(256), auth_encoded: String::with_capacity(512), body: Vec::with_capacity(8192), + streaming: None, + encoder: DataEncoder::new(), + body_open: false, events: VecDeque::with_capacity(16), deadline: None, token: None, @@ -497,10 +515,17 @@ impl Connection { } } State::Data if code == 354 => { - self.tx.extend_from_slice(&self.body); self.state = State::Body; + if self.streaming.is_some() { + self.encoder = DataEncoder::new(); + self.body_open = true; + let token = self.token.unwrap(); + self.events.push_back(Event::BodyReady { token }); + } else { + self.tx.extend_from_slice(&self.body); + } } - State::Body if code == 250 => { + State::Body if code == 250 && !self.body_open => { let token = self.token.take().unwrap(); let info = SendInfo { response: response.into(), @@ -589,12 +614,31 @@ impl Connection { self.tx.extend_from_slice(self.auth_encoded.as_bytes()); self.tx.extend_from_slice(b"\r\n"); } + /// Whether the configured TLS policy forbids a plaintext session. + fn tls_mandatory(&self) -> bool { + matches!(self.config.tls, Tls::Required | Tls::Implicit) + } + /// The only transition to `Ready`. `after_ehlo` already refuses to AUTH in + /// the clear under `Tls::Required`; this re-checks at the transition itself so + /// no path can hand the host a plaintext session its policy forbids. fn ready(&mut self) { + if self.tls_mandatory() && !self.secure { + self.fail(error( + "ETLS", + "STARTTLS", + "TLS is required but the SMTP session was never encrypted", + None, + "", + )); + self.close_transport(); + return; + } self.state = State::Ready; self.events.push_back(Event::Ready); } /// Move the result's inherent envelope/ID ownership into the core. `content` /// may be lettre's formatted MIME bytes or another already-built message. + /// For a message too large to hold twice in memory, use [`Self::start_send`]. pub fn send( &mut self, token: u64, @@ -603,6 +647,90 @@ impl Connection { content: &[u8], now: Instant, ) -> Result<(), Error> { + let eight_bit = !content.is_ascii(); + self.check_send(&envelope, eight_bit)?; + self.body.clear(); + let size = encode_data(content, &mut self.body); + self.begin( + token, + envelope, + message_id, + Some(size as u64), + eight_bit, + now, + )?; + self.streaming = None; + Ok(()) + } + /// Starts a message whose content is streamed instead of passed whole. + /// MAIL/RCPT/DATA proceed as for [`Self::send`]; once the server accepts + /// DATA the core emits [`Event::BodyReady`], after which the host passes the + /// content in any number of [`Self::send_chunk`] calls and then calls + /// [`Self::finish_body`]. Line-ending normalization and dot-stuffing carry + /// across chunk boundaries, so any split yields the bytes `send` would. + /// The result is the same single `Sent`/`Failed` for `token`. + pub fn start_send( + &mut self, + token: u64, + envelope: Envelope, + message_id: String, + body: StreamBody, + now: Instant, + ) -> Result<(), Error> { + self.check_send(&envelope, body.eight_bit)?; + self.begin(token, envelope, message_id, body.size, body.eight_bit, now)?; + self.streaming = Some(body); + Ok(()) + } + /// True between [`Event::BodyReady`] and [`Self::finish_body`]. + pub fn can_send_body(&self) -> bool { + self.body_open && self.state == State::Body + } + /// Encodes one piece of a streamed message into `output()`. Drain `output()` + /// between chunks to keep memory bounded by the chunk size. 8-bit bytes in + /// a body not declared `eight_bit` cannot be retracted mid-DATA, so they + /// fail the transaction and close the transport. + pub fn send_chunk(&mut self, chunk: &[u8], now: Instant) -> Result<(), Error> { + self.check_body()?; + if !chunk.is_ascii() && !self.streaming.is_some_and(|b| b.eight_bit) { + let e = error( + "EMESSAGE", + "DATA", + "8-bit content in a body not declared eight_bit", + None, + "", + ); + self.fail(e.clone()); + self.close_transport(); + return Err(e); + } + self.encoder.encode(chunk, &mut self.tx); + self.arm(now); + Ok(()) + } + /// Ends a streamed message: terminates its last line if needed and writes + /// the `.` terminator. The server's reply settles the token. + pub fn finish_body(&mut self, now: Instant) -> Result<(), Error> { + self.check_body()?; + self.encoder.finish(&mut self.tx); + self.body_open = false; + self.arm(now); + Ok(()) + } + fn check_body(&self) -> Result<(), Error> { + if self.can_send_body() { + Ok(()) + } else { + Err(error( + "ESTATE", + "DATA", + "No streamed SMTP message body is open", + None, + "", + )) + } + } + fn check_send(&self, envelope: &Envelope, eight_bit: bool) -> Result<(), Error> { if self.state != State::Ready { return Err(error( "ESTATE", @@ -631,9 +759,7 @@ impl Connection { "", )); } - let utf8 = !envelope.from.is_ascii() || envelope.to.iter().any(|s| !s.is_ascii()); - let eight_bit = !content.is_ascii(); - if utf8 && !self.capabilities.smtp_utf8 { + if utf8(envelope) && !self.capabilities.smtp_utf8 { return Err(error( "EENVELOPE", "MAIL FROM", @@ -651,13 +777,23 @@ impl Connection { "", )); } - self.body.clear(); - let size = encode_data(content, &mut self.body); - if self - .capabilities - .size - .is_some_and(|limit| limit != 0 && size as u64 > limit) - { + Ok(()) + } + /// Emits MAIL FROM (and pipelined RCPT TO) once `check_send` has passed. + fn begin( + &mut self, + token: u64, + envelope: Envelope, + message_id: String, + size: Option, + eight_bit: bool, + now: Instant, + ) -> Result<(), Error> { + if size.is_some_and(|size| { + self.capabilities + .size + .is_some_and(|limit| limit != 0 && size > limit) + }) { return Err(error( "EMESSAGE", "MAIL FROM", @@ -670,14 +806,17 @@ impl Connection { self.rejected.clear(); self.mail_error = None; self.recipient_index = 0; + self.body_open = false; write!(self.tx, "MAIL FROM:<{}>", envelope.from).unwrap(); - if self.capabilities.size_supported { + if let Some(size) = size + && self.capabilities.size_supported + { write!(self.tx, " SIZE={size}").unwrap(); } if eight_bit { self.tx.extend_from_slice(b" BODY=8BITMIME"); } - if utf8 { + if utf8(&envelope) { self.tx.extend_from_slice(b" SMTPUTF8"); } self.tx.extend_from_slice(b"\r\n"); @@ -734,7 +873,8 @@ impl Connection { Some(code), response, ); - if matches!(self.state, State::Data | State::Body) { + // RSET cannot follow a DATA body that was never terminated. + if matches!(self.state, State::Data | State::Body) && !self.body_open { self.fail_and_reset(e); } else { self.fail(e); @@ -856,6 +996,7 @@ impl Connection { return; } self.state = State::Closed; + self.body_open = false; self.deadline = None; self.tx.clear(); self.tx_pos = 0; @@ -863,40 +1004,132 @@ impl Connection { self.events.push_back(Event::Closed); } } -/// Normalize LF, CRLF and bare CR, dot-stuff each line, append the SMTP DATA -/// terminator. Returns SIZE octets (normalized content, excluding transparency). -pub fn encode_data(content: &[u8], out: &mut Vec) -> usize { - let mut start = true; - let mut i = 0; - let mut size = 0; - while i < content.len() { - match content[i] { - b'\r' | b'\n' => { - if content[i] == b'\r' && content.get(i + 1) == Some(&b'\n') { - i += 1; +fn utf8(envelope: &Envelope) -> bool { + !envelope.from.is_ascii() || envelope.to.iter().any(|s| !s.is_ascii()) +} +/// How a streamed message body will be sent; see [`Connection::start_send`]. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct StreamBody { + /// Expected normalized content octets, sent as the SIZE parameter when the + /// server supports it and checked against the server's advertised limit. + /// `None` omits the parameter. The server, not this core, enforces it. + pub size: Option, + /// The body may contain 8-bit bytes: sends BODY=8BITMIME, which the server + /// must support. Without it, a non-ASCII chunk fails the transaction. + pub eight_bit: bool, +} +/// Incremental SMTP DATA encoding: normalizes LF, CRLF and bare CR to CRLF and +/// dot-stuffs each line. State carries across [`DataEncoder::encode`] calls, so +/// a CR/LF pair or a line-leading `.` split between chunks encodes exactly as +/// it would in one piece. +#[derive(Debug, Clone)] +pub struct DataEncoder { + at_line_start: bool, + after_cr: bool, + size: usize, +} +impl Default for DataEncoder { + fn default() -> Self { + Self::new() + } +} +impl DataEncoder { + pub fn new() -> Self { + Self { + at_line_start: true, + after_cr: false, + size: 0, + } + } + pub fn encode(&mut self, content: &[u8], out: &mut Vec) { + for &byte in content { + match byte { + // The LF of a CRLF whose CR already produced the line break. + b'\n' if self.after_cr => self.after_cr = false, + b'\r' | b'\n' => { + out.extend_from_slice(b"\r\n"); + self.size += 2; + self.at_line_start = true; + self.after_cr = byte == b'\r'; } - out.extend_from_slice(b"\r\n"); - size += 2; - start = true; - } - byte => { - if start && byte == b'.' { - out.push(b'.'); + byte => { + if self.at_line_start && byte == b'.' { + out.push(b'.'); + } + out.push(byte); + self.size += 1; + self.at_line_start = false; + self.after_cr = false; } - out.push(byte); - size += 1; - start = false; } } - i += 1; } - if !start { - out.extend_from_slice(b"\r\n"); - size += 2; + /// Terminates an unfinished last line and appends the DATA terminator. + /// Returns SIZE octets (normalized content, excluding transparency dots) + /// and resets the encoder for the next message. + pub fn finish(&mut self, out: &mut Vec) -> usize { + if !self.at_line_start { + out.extend_from_slice(b"\r\n"); + self.size += 2; + } + out.extend_from_slice(b".\r\n"); + let size = self.size; + *self = Self::new(); + size } - out.extend_from_slice(b".\r\n"); - size +} +/// Normalize LF, CRLF and bare CR, dot-stuff each line, append the SMTP DATA +/// terminator. Returns SIZE octets (normalized content, excluding transparency). +pub fn encode_data(content: &[u8], out: &mut Vec) -> usize { + let mut encoder = DataEncoder::new(); + encoder.encode(content, out); + encoder.finish(out) } #[cfg(feature = "turnloop")] pub mod asynchronous; + +#[cfg(test)] +mod tests { + use super::*; + /// Reaching the Ready transition in the clear, by any path, must fail the + /// session under a mandatory TLS policy rather than report `Ready`. + #[test] + fn ready_refuses_plaintext_under_mandatory_tls() { + let now = Instant::now(); + for (tls, allowed) in [ + (Tls::Required, false), + (Tls::Implicit, false), + (Tls::Opportunistic, true), + (Tls::None, true), + ] { + let mut c = Connection::new(Config { + tls, + ..Config::default() + }) + .expect("config"); + c.connected(now).expect("connected"); + c.state = State::Ehlo; + // Skip after_ehlo's STARTTLS decision, as a future path might. + c.authenticate(); + let events: Vec<_> = std::iter::from_fn(|| c.poll_event()).collect(); + let events: Vec<_> = events + .into_iter() + .filter(|e| *e != Event::UpgradeTls) + .collect(); + if allowed { + assert_eq!(events, [Event::Ready], "{tls:?}"); + assert_eq!(c.state(), State::Ready); + } else { + assert_eq!(events.len(), 3, "{tls:?}: {events:?}"); + assert!( + matches!(&events[0], Event::Failed { token: None, error, .. } + if error.code == "ETLS" && error.message.contains("never encrypted")), + "{tls:?}: {events:?}" + ); + assert_eq!(events[1..], [Event::CloseTransport, Event::Closed]); + assert_eq!(c.state(), State::Closed); + } + } + } +} diff --git a/protocols/turnloop-smtp/tests/smtp.rs b/protocols/turnloop-smtp/tests/smtp.rs index ba76b8f..ed5f798 100644 --- a/protocols/turnloop-smtp/tests/smtp.rs +++ b/protocols/turnloop-smtp/tests/smtp.rs @@ -8,7 +8,8 @@ use std::{ time::{Duration, Instant, SystemTime}, }; use turnloop_smtp::{ - Auth, Capabilities, Config, Connection, Envelope, Event, State, Tls, encode_data, + Auth, Capabilities, Config, Connection, DataEncoder, Envelope, Event, State, StreamBody, Tls, + encode_data, message::{self, FileAttachment, Mail}, }; fn discard(c: &mut Connection) { @@ -56,6 +57,233 @@ fn normalization_and_dot_stuffing() { encode_data(b"x\r\n", &mut out); assert_eq!(out, b"x\r\n.\r\n"); } +/// Content whose every hazard can land on a chunk boundary: an embedded +/// CRLF.CRLF (the DATA terminator), leading dots, CRLF/bare CR/bare LF, and a +/// CR immediately followed by a line-leading dot. +const HAZARDS: &[u8] = b".lead\r\nbody\r\n.\r\nafter\r\r\n\n.x\r.y\n..z\r\nlast."; +#[test] +fn incremental_dot_stuffing_is_split_invariant() { + let mut whole = Vec::new(); + let whole_size = encode_data(HAZARDS, &mut whole); + assert!( + whole.windows(7).any(|w| w == b"\r\n..\r\na"), + "the embedded terminator is stuffed" + ); + // Every split into two and three chunks, including empty chunks. + for i in 0..=HAZARDS.len() { + for j in i..=HAZARDS.len() { + let mut encoder = DataEncoder::new(); + let mut out = Vec::new(); + for chunk in [&HAZARDS[..i], &HAZARDS[i..j], &HAZARDS[j..]] { + encoder.encode(chunk, &mut out); + } + assert_eq!(encoder.finish(&mut out), whole_size, "split {i}/{j}"); + assert_eq!(out, whole, "split {i}/{j}"); + } + } + // The terminator split exactly as CRLF | .CRLF and CR | LF.CRLF. + for parts in [ + &[&b"a\r\n"[..], b".\r\nb"][..], + &[b"a\r", b"\n.\r\nb"], + &[b"a\r\n.", b"\r\nb"], + &[b"a", b"\r", b"\n", b".", b"\r", b"\n", b"b"], + ] { + let mut encoder = DataEncoder::new(); + let mut out = Vec::new(); + for part in parts { + encoder.encode(part, &mut out); + } + encoder.finish(&mut out); + assert_eq!(out, b"a\r\n..\r\nb\r\n.\r\n", "{parts:?}"); + } + // Byte-at-a-time, and the encoder resets after finish. + let mut encoder = DataEncoder::new(); + for round in 0..2 { + let mut out = Vec::new(); + for byte in HAZARDS { + encoder.encode(std::slice::from_ref(byte), &mut out); + } + assert_eq!(encoder.finish(&mut out), whole_size, "round {round}"); + assert_eq!(out, whole, "round {round}"); + } +} +/// Drives one message from MAIL to the 354 reply, returning every byte the +/// client wrote; `stream` chooses `start_send` over `send`. +fn to_data(c: &mut Connection, token: u64, stream: Option) -> Vec { + let now = Instant::now(); + match stream { + Some(body) => c.start_send(token, envelope(), "".into(), body, now), + None => c.send(token, envelope(), "".into(), HAZARDS, now), + } + .expect("fixture operation must succeed"); + let mut wire = c.output().to_vec(); + discard(c); + c.receive(b"250 mail\r\n250 ok\r\n550 no such user\r\n", now) + .expect("fixture operation must succeed"); + assert_eq!(c.output(), b"DATA\r\n"); + wire.extend_from_slice(c.output()); + discard(c); + c.receive(b"354 go ahead\r\n", now) + .expect("fixture operation must succeed"); + wire.extend_from_slice(c.output()); + discard(c); + wire +} +fn hazards_body() -> StreamBody { + let mut out = Vec::new(); + StreamBody { + size: Some(encode_data(HAZARDS, &mut out) as u64), + eight_bit: false, + } +} +#[test] +fn streamed_send_writes_the_one_shot_bytes() { + let now = Instant::now(); + let mut oneshot = plain_ready(true); + let expected = to_data(&mut oneshot, 1, None); + assert!(expected.ends_with(b"\r\nlast.\r\n.\r\n")); + for split in [1, 7, 13, 14, 15, 16, HAZARDS.len()] { + let mut c = plain_ready(true); + let mut wire = to_data(&mut c, 2, Some(hazards_body())); + assert_eq!(c.poll_event(), Some(Event::BodyReady { token: 2 })); + assert!(c.can_send_body()); + for chunk in HAZARDS.chunks(split) { + c.send_chunk(chunk, now) + .expect("fixture operation must succeed"); + // Draining between chunks keeps only the current chunk buffered. + wire.extend_from_slice(c.output()); + discard(&mut c); + } + c.finish_body(now).expect("fixture operation must succeed"); + assert!(!c.can_send_body()); + wire.extend_from_slice(c.output()); + discard(&mut c); + assert_eq!(wire, expected, "chunks of {split}"); + assert!(c.send_chunk(b"late", now).is_err()); + assert!(c.finish_body(now).is_err()); + assert!(c.output().is_empty()); + c.receive(b"250 queued\r\n", now) + .expect("fixture operation must succeed"); + assert!(matches!( + c.poll_event(), + Some(Event::Sent { token: 2, info }) + if info.accepted == ["ok@example.test"] && info.rejected.len() == 1 + )); + assert_eq!(c.state(), State::Ready); + assert_eq!(c.poll_event(), None); + } + // SIZE is optional for a stream, and a body may be empty. + let mut c = plain_ready(true); + c.start_send(3, envelope(), "id".into(), StreamBody::default(), now) + .expect("fixture operation must succeed"); + assert_eq!( + &c.output()[..c.output().iter().position(|b| *b == b'\n').expect("line") + 1], + b"MAIL FROM:\r\n" + ); + assert!(c.send_chunk(b"early", now).is_err(), "no body before 354"); + discard(&mut c); + c.receive(b"250 mail\r\n250 ok\r\n250 ok\r\n", now) + .expect("fixture operation must succeed"); + discard(&mut c); + c.receive(b"354 go ahead\r\n", now) + .expect("fixture operation must succeed"); + assert_eq!(c.poll_event(), Some(Event::BodyReady { token: 3 })); + c.finish_body(now).expect("fixture operation must succeed"); + assert_eq!(c.output(), b".\r\n"); +} +#[test] +fn streamed_send_failures_never_terminate_a_partial_body() { + let now = Instant::now(); + // Undeclared 8-bit content cannot be retracted: fail and close, no ".". + let mut c = plain_ready(true); + to_data(&mut c, 4, Some(hazards_body())); + assert_eq!(c.poll_event(), Some(Event::BodyReady { token: 4 })); + c.send_chunk(b"Subject: x\r\n\r\n", now) + .expect("fixture operation must succeed"); + discard(&mut c); + let e = c.send_chunk("ü".as_bytes(), now).unwrap_err(); + assert_eq!(e.code, "EMESSAGE"); + assert!(c.output().is_empty(), "nothing, least of all a terminator"); + assert!( + matches!(c.poll_event(), Some(Event::Failed { token: Some(4), error, .. }) if error.code == "EMESSAGE") + ); + assert_eq!(c.poll_event(), Some(Event::CloseTransport)); + assert_eq!(c.poll_event(), Some(Event::Closed)); + // A declared 8-bit body is sent with BODY=8BITMIME and accepted. + let mut c = plain_ready(true); + let wire = to_data( + &mut c, + 5, + Some(StreamBody { + size: None, + eight_bit: true, + }), + ); + assert!(wire.starts_with(b"MAIL FROM: BODY=8BITMIME\r\n")); + assert_eq!(c.poll_event(), Some(Event::BodyReady { token: 5 })); + c.send_chunk("ü\n".as_bytes(), now) + .expect("fixture operation must succeed"); + c.finish_body(now).expect("fixture operation must succeed"); + assert_eq!(c.output(), "ü\r\n.\r\n".as_bytes()); + // A server reply while the body is open ends the session: RSET would be + // read as message content. + let mut c = plain_ready(true); + to_data(&mut c, 6, Some(hazards_body())); + assert_eq!(c.poll_event(), Some(Event::BodyReady { token: 6 })); + c.send_chunk(b"partial", now) + .expect("fixture operation must succeed"); + discard(&mut c); + c.receive(b"250 premature\r\n", now) + .expect("fixture operation must succeed"); + assert!( + matches!(c.poll_event(), Some(Event::Failed { token: Some(6), error, .. }) if error.code == "EMESSAGE") + ); + assert_eq!(c.poll_event(), Some(Event::CloseTransport)); + assert_eq!(c.poll_event(), Some(Event::Closed)); + assert!(c.output().is_empty()); + assert!(!c.can_send_body()); + // DATA refused: no body is requested, and RSET recovers the session. + let mut c = plain_ready(true); + c.start_send(7, envelope(), "id".into(), hazards_body(), now) + .expect("fixture operation must succeed"); + discard(&mut c); + c.receive(b"250 mail\r\n250 ok\r\n250 ok\r\n", now) + .expect("fixture operation must succeed"); + discard(&mut c); + c.receive(b"554 no\r\n", now) + .expect("fixture operation must succeed"); + assert!(matches!( + c.poll_event(), + Some(Event::Failed { token: Some(7), .. }) + )); + assert_eq!(c.output(), b"RSET\r\n"); + assert!(c.send_chunk(b"x", now).is_err()); + // Declarations are checked before any envelope byte is written. + let mut c = plain_ready(false); + let unsupported = StreamBody { + size: None, + eight_bit: true, + }; + assert_eq!( + c.start_send(8, envelope(), "id".into(), unsupported, now) + .unwrap_err() + .code, + "EMESSAGE" + ); + let mut c = plain_ready(true); + let too_big = StreamBody { + size: Some(10_001), + eight_bit: false, + }; + assert!( + c.start_send(9, envelope(), "id".into(), too_big, now) + .unwrap_err() + .message + .contains("SIZE") + ); + assert!(c.output().is_empty()); + assert_eq!(c.state(), State::Ready); +} #[test] fn ehlo_fallback_required_tls_and_deadline() { let now = Instant::now(); @@ -92,6 +320,89 @@ fn ehlo_fallback_required_tls_and_deadline() { assert_eq!(c.poll_event(), Some(Event::Closed)); assert_eq!(c.poll_event(), None); } +/// Drains every queued event. +fn events(c: &mut Connection) -> Vec { + std::iter::from_fn(|| c.poll_event()).collect() +} +#[test] +fn required_tls_never_reaches_ready_in_the_clear() { + let now = Instant::now(); + let required = || { + let mut c = Connection::new(Config { + tls: Tls::Required, + auth: Some(Auth::Plain { + user: "user".into(), + password: "secret".into(), + }), + ..Config::default() + }) + .expect("fixture operation must succeed"); + c.connected(now).expect("fixture operation must succeed"); + c.receive(b"220 hi\r\n", now) + .expect("fixture operation must succeed"); + assert_eq!(c.output(), b"EHLO [127.0.0.1]\r\n"); + discard(&mut c); + c + }; + let refused = |events: &[Event]| { + assert_eq!(events.len(), 3, "{events:?}"); + assert!( + matches!(&events[0], Event::Failed { token: None, error, .. } + if error.code == "ETLS" && error.command == "STARTTLS" + && error.message.contains("STARTTLS")), + "{events:?}" + ); + assert_eq!(events[1..], [Event::CloseTransport, Event::Closed]); + }; + // EHLO without STARTTLS: fail before any AUTH exchange leaves in the clear. + let mut c = required(); + c.receive(b"250-hi\r\n250 AUTH PLAIN\r\n", now) + .expect("fixture operation must succeed"); + refused(&events(&mut c)); + assert!(c.output().is_empty(), "no AUTH may follow"); + assert_eq!(c.state(), State::Closed); + // STARTTLS advertised but refused by the server. + let mut c = required(); + c.receive(b"250-hi\r\n250-STARTTLS\r\n250 AUTH PLAIN\r\n", now) + .expect("fixture operation must succeed"); + assert_eq!(c.output(), b"STARTTLS\r\n"); + discard(&mut c); + c.receive(b"454 TLS not available\r\n", now) + .expect("fixture operation must succeed"); + let drained = events(&mut c); + assert!( + matches!(&drained[0], Event::Failed { error, .. } + if error.code == "ETLS" && error.response_code == Some(454)), + "{drained:?}" + ); + assert!(!drained.contains(&Event::Ready)); + assert!(c.output().is_empty()); + // STARTTLS accepted: Ready follows only after the host's upgrade. + let mut c = required(); + c.receive(b"250-hi\r\n250-STARTTLS\r\n250 AUTH PLAIN\r\n", now) + .expect("fixture operation must succeed"); + discard(&mut c); + c.receive(b"220 go ahead\r\n", now) + .expect("fixture operation must succeed"); + assert_eq!(events(&mut c), [Event::UpgradeTls]); + assert_eq!(c.state(), State::Tls); + c.tls_established(now) + .expect("fixture operation must succeed"); + discard(&mut c); + c.receive(b"250-hi\r\n250 AUTH PLAIN\r\n", now) + .expect("fixture operation must succeed"); + assert!(c.output().starts_with(b"AUTH PLAIN ")); + discard(&mut c); + c.receive(b"235 ok\r\n", now) + .expect("fixture operation must succeed"); + assert_eq!(events(&mut c), [Event::Ready]); + // Opportunistic TLS, by contrast, proceeds in the clear. + let mut c = Connection::new(Config::default()).expect("fixture operation must succeed"); + c.connected(now).expect("fixture operation must succeed"); + c.receive(b"220 hi\r\n250 hi\r\n", now) + .expect("fixture operation must succeed"); + assert_eq!(events(&mut c), [Event::Ready]); +} #[test] fn malformed_multiline_and_injection_are_rejected() { let now = Instant::now(); @@ -507,6 +818,7 @@ fn test_socket(mode: Mode) { ); failed = true; } + Event::BodyReady { .. } => panic!("one-shot send streams no body"), Event::CloseTransport => {} Event::Closed => break, }