From 453928fecb02ffb4a93958c3acf751fb88788916 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:30:26 +0200 Subject: [PATCH 1/8] Let CursorBatch::rows outlive the batch that produced it `rows(&self) -> impl Iterator>` is defined in an edition-2024 crate, so its opaque return type captured every lifetime in scope, including the anonymous `&self` borrow. The rows only borrow the reply (`'a`), but the iterator was tied to the `CursorBatch` value: a host that parsed a batch in a helper and returned `batch.rows()` got E0597, and an edition-2021 host hit the same error on `match batch.rows().next()` in tail position, whose temporaries outlive the block's locals. The return type now says `+ use<'a>`, which is what rustc's own diagnostic suggests: the closure copies the `&'a RawArray` out of `self` and holds nothing borrowed from the batch. Callers that used rows() in place are unaffected. tests/cursor_rows_capture.rs compiles both patterns and asserts the rows they return. Without the change the test target fails to build with E0597 at the function returning `batch.rows()`. Closes #69 --- protocols/turnloop-mongodb/src/command.rs | 3 +- .../tests/cursor_rows_capture.rs | 58 +++++++++++++++++++ 2 files changed, 60 insertions(+), 1 deletion(-) create mode 100644 protocols/turnloop-mongodb/tests/cursor_rows_capture.rs diff --git a/protocols/turnloop-mongodb/src/command.rs b/protocols/turnloop-mongodb/src/command.rs index 9c09602..faf2eed 100644 --- a/protocols/turnloop-mongodb/src/command.rs +++ b/protocols/turnloop-mongodb/src/command.rs @@ -468,7 +468,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")), 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..f815a57 --- /dev/null +++ b/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs @@ -0,0 +1,58 @@ +//! `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. +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); +} From 4f44b690f663e210db22d80bcaa40b695d33ad29 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:31:24 +0200 Subject: [PATCH 2/8] Let a host ask a MongoDB connection whether it expects a reply `Connection::receive` refuses bytes unless a handshake, authentication or command reply is outstanding, but the only state a host could observe was `is_ready()`. Every host therefore kept its own `expecting_reply` flag to decide when reading was allowed, a second copy of the state machine free to drift from the real one. `accepts_receive()` now reports exactly that condition, and `receive` itself calls it, so the two cannot disagree. A named query was chosen over exposing the private `State` enum: it answers the host's actual question and leaves the state machine free to grow states without a breaking change. The new test walks a connection through every state (New, TLS upgrade, handshake, SCRAM conversation, Ready, command, unreleased reply, draining unacknowledged write, Closed) and checks at each step that `accepts_receive()` agrees with whether `receive(&[])` is accepted. The README's driving steps mention the query. Closes #70 --- protocols/turnloop-mongodb/README.md | 2 + protocols/turnloop-mongodb/src/connection.rs | 18 ++--- protocols/turnloop-mongodb/tests/protocol.rs | 82 ++++++++++++++++++++ 3 files changed, 93 insertions(+), 9 deletions(-) diff --git a/protocols/turnloop-mongodb/README.md b/protocols/turnloop-mongodb/README.md index 9a40cb8..a5d6b54 100644 --- a/protocols/turnloop-mongodb/README.md +++ b/protocols/turnloop-mongodb/README.md @@ -57,6 +57,8 @@ has no SRV/TXT capability, so use explicit resolved seed URIs there. 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. + `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 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()`. diff --git a/protocols/turnloop-mongodb/src/connection.rs b/protocols/turnloop-mongodb/src/connection.rs index 36d8be2..c3388a9 100644 --- a/protocols/turnloop-mongodb/src/connection.rs +++ b/protocols/turnloop-mongodb/src/connection.rs @@ -209,17 +209,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); diff --git a/protocols/turnloop-mongodb/tests/protocol.rs b/protocols/turnloop-mongodb/tests/protocol.rs index ec47658..b0c82c0 100644 --- a/protocols/turnloop-mongodb/tests/protocol.rs +++ b/protocols/turnloop-mongodb/tests/protocol.rs @@ -867,3 +867,85 @@ 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()); +} From 003cafe506ea54914a30e47b99ef675586c91cd9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:32:26 +0200 Subject: [PATCH 3/8] Make WriteResult::parse fail a write the server reported as failed MongoDB reports per-document failures inside a successful command: a duplicate-key insert answers `ok: 1` with a `writeErrors` array, and an unmet write concern answers `ok: 1` with `writeConcernError`. `WriteResult::parse` only consulted `Error::from_response` when `ok` was 0, so a host that trusted its `Ok` treated a rejected insert as a successful one. The safe order (check `from_response` first) was documented nowhere. The issue offered two fixes: make `parse` itself fail, or document that callers must run `from_response` first. Documentation alone leaves the trap armed, so `parse` now runs `Error::from_response` unconditionally. A nonempty `writeErrors` fails as `ErrorKind::BulkWrite` and a `writeConcernError` as `ErrorKind::Server`, each with the full server response attached, which is also what the Node driver throws. A write-concern error does fail even though the write was applied, because the durability the caller asked for was not confirmed. Bulk aggregation still needs the errors as data, so the old lenient behaviour moves to `WriteResult::decode` (fails only on `ok: 0` or a malformed reply), and `BulkResult::accept` uses it. `WriteResult::succeeded()` gives `decode` callers the verdict. The type docs and README explain which to use and that `parse` needs no prior `from_response`. Behaviour change: `WriteResult::parse` now returns `Err` for `ok: 1` replies with write or write-concern errors. Callers that aggregated those errors from `parse` should switch to `decode`. The new test covers the duplicate-key and write-concern replies, clean replies with and without an empty `writeErrors`, `ok: 0`, and `BulkResult` aggregation with original indices. With `parse` reverted to the old check, it fails on the duplicate-key `unwrap_err`. Closes #72 --- protocols/turnloop-mongodb/README.md | 8 ++- protocols/turnloop-mongodb/src/command.rs | 31 +++++++++- protocols/turnloop-mongodb/tests/protocol.rs | 59 ++++++++++++++++++++ 3 files changed, 95 insertions(+), 3 deletions(-) diff --git a/protocols/turnloop-mongodb/README.md b/protocols/turnloop-mongodb/README.md index a5d6b54..c0c5123 100644 --- a/protocols/turnloop-mongodb/README.md +++ b/protocols/turnloop-mongodb/README.md @@ -89,8 +89,12 @@ 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. diff --git a/protocols/turnloop-mongodb/src/command.rs b/protocols/turnloop-mongodb/src/command.rs index faf2eed..95c4a38 100644 --- a/protocols/turnloop-mongodb/src/command.rs +++ b/protocols/turnloop-mongodb/src/command.rs @@ -534,16 +534,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)?; } @@ -563,6 +586,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() { @@ -583,7 +612,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; diff --git a/protocols/turnloop-mongodb/tests/protocol.rs b/protocols/turnloop-mongodb/tests/protocol.rs index b0c82c0..3fb5d68 100644 --- a/protocols/turnloop-mongodb/tests/protocol.rs +++ b/protocols/turnloop-mongodb/tests/protocol.rs @@ -949,3 +949,62 @@ fn accepts_receive_tracks_every_connection_state() { 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); +} From 397a87842b296f44521e06644299805bdd8a9676 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:33:59 +0200 Subject: [PATCH 4/8] Settle an unreleased MongoDB reply with Failed when the connection fails `Connection::fail()` pushed `Failed` only when no reply was held. With a reply received but not yet released it pushed nothing: the token's reply stayed readable on a connection that had just been declared broken, and a host that settles operations from the event stream never learned that the operation it was still holding had failed. The settlement contract is now explicit and applied on every path: an accepted token is settled exactly once, by `Failed`, by `Unacknowledged`, or by `release_reply()` after its `Reply` event. `fail()` before that release revokes the reply (`reply()` and `release_reply()` then return errors) and pushes `Failed` for the token ahead of `Closed`. If the host has not polled the `Reply` event yet, it is withdrawn from the queue, so that token sees exactly one event. A reply already released was already settled, and a later failure reports no token for it. `close()` is unchanged: it is the orderly path (the async adapter uses it on EOF) and still keeps a received reply readable until released, as the existing close test requires. Revoking in `fail()` rather than preserving is deliberate. `fail()` means the host has declared the transport or protocol broken, and a network error after the request was sent is already the ambiguous may-or-may-not-have-applied outcome the retry rules are built for. Behaviour change: after `fail()`, an unreleased reply can no longer be read. The new test covers a reply whose event was never polled, one whose event was polled but not released, and one already released, asserting the exact event sequence and that the reply accessors fail. Against the previous `fail()` it reports `["reply"]` where `["failed"]` is required. Closes #73 --- protocols/turnloop-mongodb/README.md | 13 ++-- protocols/turnloop-mongodb/src/connection.rs | 22 ++++-- protocols/turnloop-mongodb/tests/protocol.rs | 78 ++++++++++++++++++++ 3 files changed, 102 insertions(+), 11 deletions(-) diff --git a/protocols/turnloop-mongodb/README.md b/protocols/turnloop-mongodb/README.md index c0c5123..29e5679 100644 --- a/protocols/turnloop-mongodb/README.md +++ b/protocols/turnloop-mongodb/README.md @@ -59,12 +59,15 @@ has no SRV/TXT capability, so use explicit resolved seed URIs there. a complete reply before feeding another frame. Partial reads and writes are normal. `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 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()`. +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 diff --git a/protocols/turnloop-mongodb/src/connection.rs b/protocols/turnloop-mongodb/src/connection.rs index c3388a9..dae9a81 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, @@ -524,16 +527,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/tests/protocol.rs b/protocols/turnloop-mongodb/tests/protocol.rs index 3fb5d68..ad9dc11 100644 --- a/protocols/turnloop-mongodb/tests/protocol.rs +++ b/protocols/turnloop-mongodb/tests/protocol.rs @@ -1008,3 +1008,81 @@ fn write_result_verdict_fails_on_write_errors_despite_ok() { 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()); +} From 40c1ad3024a97dd32b30e1e08eac285c69da3adb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:35:11 +0200 Subject: [PATCH 5/8] Document the MongoDB host-entropy obligations in one place The crate needs the host to supply every random or unique value, but those obligations were scattered: the SCRAM nonce was mentioned in the README's driving steps, the `ObjectIdGenerator` requirement lived in a sentence of the command-results section, and session UUIDs and selection entropy only in item docs. The `_id` one in particular fails silently: an insert without `_id` succeeds, and the server assigns an id the host never learns. The README now has a "Host entropy" section, placed before the driving steps where a new integrator starts, listing all four obligations (SCRAM nonce, document `_id`, session identity, server-selection entropy), where each value goes, what to supply, and what goes wrong without it. It also states which ones the `turnloop` feature's client fills on its own; `_id` is not one of them. The crate docs carry a short `# Host entropy` section linking to it, and `Connection::connected`, the `command` and `auth` modules, `ObjectIdGenerator` and `Session::new` link to that section. The two old README mentions now point to it instead of restating it. Documentation only; rustdoc builds with `-D warnings`. Closes #71 --- protocols/turnloop-mongodb/README.md | 27 +++++++++++++++++--- protocols/turnloop-mongodb/src/auth.rs | 3 ++- protocols/turnloop-mongodb/src/command.rs | 8 ++++++ protocols/turnloop-mongodb/src/connection.rs | 4 ++- protocols/turnloop-mongodb/src/lib.rs | 13 ++++++++++ protocols/turnloop-mongodb/src/session.rs | 3 ++- 6 files changed, 51 insertions(+), 7 deletions(-) diff --git a/protocols/turnloop-mongodb/README.md b/protocols/turnloop-mongodb/README.md index 29e5679..c1206ac 100644 --- a/protocols/turnloop-mongodb/README.md +++ b/protocols/turnloop-mongodb/README.md @@ -44,13 +44,32 @@ 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. @@ -99,8 +118,8 @@ in `writeConcernError`. `WriteResult::parse` is the verdict — it runs `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 95c4a38..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, @@ -718,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 dae9a81..e2996bd 100644 --- a/protocols/turnloop-mongodb/src/connection.rs +++ b/protocols/turnloop-mongodb/src/connection.rs @@ -84,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")); 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, From ff7547c69d16763f788d14690bc1142accdbf555 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:37:10 +0200 Subject: [PATCH 6/8] Refuse Ready in the clear under Tls::Required, at the transition itself The issue reports that a `Tls::Required` session reaches `Ready` in plaintext when the server offers no STARTTLS, so hosts such as Perry's P6 lane guard `on_ready` themselves. On current main the reachable paths already refuse: `after_ehlo` fails with `ETLS` when STARTTLS is not advertised (EHLO or HELO fallback), and a refused STARTTLS fails through `unexpected`. None of this was documented, though, and the refusal lived in one branch of the STARTTLS decision rather than at the point the issue names, so any future path to `ready()` that skipped that branch would silently hand out a plaintext session. `ready()` is the only transition to `Ready`, and it now re-checks the policy: under `Tls::Required` (or `Implicit`) without an established TLS layer it emits `Failed` with code `ETLS` ("TLS is required but the SMTP session was never encrypted"), then `CloseTransport` and `Closed`, and never `Ready`. The earlier check in `after_ehlo` stays, because it must also stop AUTH credentials from being sent in the clear, which happens before `Ready`. The `Tls` variants and the README now state the guarantee, so hosts can drop their own guard. Tests: - `tests::ready_refuses_plaintext_under_mandatory_tls` (unit) enters authentication directly, skipping the STARTTLS decision, and asserts `Failed(ETLS)`, `CloseTransport`, `Closed` for Required and Implicit, and `Ready` for Opportunistic and None. With the new check disabled it fails ("Required: [Ready]"). - `required_tls_never_reaches_ready_in_the_clear` pins the wire behaviour: no STARTTLS advertised means `ETLS` with no AUTH bytes emitted; a 454 reply to STARTTLS means `ETLS` and no `Ready`; an accepted STARTTLS reaches `Ready` only after `tls_established` and a fresh EHLO; Opportunistic reaches `Ready` in the clear. This test also passes on the previous code, which already handled these paths. Closes #56 --- protocols/turnloop-smtp/README.md | 5 +- protocols/turnloop-smtp/src/lib.rs | 70 ++++++++++++++++++++++ protocols/turnloop-smtp/tests/smtp.rs | 83 +++++++++++++++++++++++++++ 3 files changed, 157 insertions(+), 1 deletion(-) diff --git a/protocols/turnloop-smtp/README.md b/protocols/turnloop-smtp/README.md index 2a75310..49c457b 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 diff --git a/protocols/turnloop-smtp/src/lib.rs b/protocols/turnloop-smtp/src/lib.rs index c03b7ad..eb3af37 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)] @@ -589,7 +596,25 @@ 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); } @@ -900,3 +925,48 @@ pub fn encode_data(content: &[u8], out: &mut Vec) -> usize { #[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..8b5b69f 100644 --- a/protocols/turnloop-smtp/tests/smtp.rs +++ b/protocols/turnloop-smtp/tests/smtp.rs @@ -92,6 +92,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(); From a648ca6f48c57e67bb9c3cdc5f13aa924c13e7c5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:41:33 +0200 Subject: [PATCH 7/8] Stream SMTP message bodies in chunks `Connection::send` took the whole message as one `&[u8]` and `encode_data` copied it again into DATA storage, so a 25 MB attachment meant two full copies of the body in memory before the first byte reached the socket. A streaming path now sits beside the one-shot `send`, shaped like the HTTP/1 client's `send_body`/`finish_body`: - `start_send(token, envelope, message_id, StreamBody { size, eight_bit }, now)` runs the same envelope checks and MAIL/RCPT/DATA exchange as `send`. - On the server's 354 the core emits the new `Event::BodyReady { token }`. - The host then calls `send_chunk(bytes, now)` any number of times and `finish_body(now)`, and `can_send_body()` reports when that is allowed. - The token settles with the usual single `Sent` or `Failed`. Draining `output()` between chunks keeps memory bounded by one chunk. `send` stays as the convenience and now shares its checks with `start_send`. Encoding is incremental. `DataEncoder` (public) carries two bits of state across calls: whether the next byte starts a line (for dot-stuffing) and whether the previous byte was a CR (so a CR|LF split across chunks collapses to one CRLF). `encode_data` is now `DataEncoder` in one call, so both paths emit identical bytes by construction. SIZE and BODY=8BITMIME go in MAIL FROM before any content exists, so the stream declares them up front. `size` is optional and is checked against the server's advertised limit; `eight_bit` must be supported by the server. Once DATA content has been written it cannot be retracted, and RSET would be read as message text. So an undeclared 8-bit chunk, or any server reply while the body is open, fails the message and closes the transport rather than terminating a partial message. Breaking: `Event` gains the `BodyReady` variant, so exhaustive matches need an arm (the crate's own adapter already uses a wildcard). Tests: - `incremental_dot_stuffing_is_split_invariant` encodes a message full of hazards (an embedded CRLF.CRLF, leading dots, CR/LF/CRLF mixes, CR then a line-leading dot) at every two- and three-way split and byte by byte, and requires output and SIZE identical to `encode_data`. It also covers CRLF.CRLF split as CRLF|.CRLF, CR|LF.CRLF, CRLF.|CRLF and one byte per chunk, which must encode as `..`. - `streamed_send_writes_the_one_shot_bytes` requires the full client transcript of a streamed send, at several chunk sizes, to equal the one-shot transcript byte for byte, followed by the same `Sent`. It also covers an empty body without SIZE and chunks refused before 354 or after finish. - `streamed_send_failures_never_terminate_a_partial_body` covers undeclared 8-bit content (no terminator written; Failed, CloseTransport, Closed), a declared 8-bit body, a premature server reply, a refused DATA (RSET, no BodyReady), and declarations rejected before any envelope byte. Carrying no state across `encode` calls fails the split tests. The allocation gate for the one-shot path is unchanged. Closes #55 --- protocols/turnloop-smtp/README.md | 17 +- protocols/turnloop-smtp/src/lib.rs | 247 +++++++++++++++++++++----- protocols/turnloop-smtp/tests/smtp.rs | 231 +++++++++++++++++++++++- 3 files changed, 450 insertions(+), 45 deletions(-) diff --git a/protocols/turnloop-smtp/README.md b/protocols/turnloop-smtp/README.md index 49c457b..857ce75 100644 --- a/protocols/turnloop-smtp/README.md +++ b/protocols/turnloop-smtp/README.md @@ -44,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 eb3af37..bbf8f48 100644 --- a/protocols/turnloop-smtp/src/lib.rs +++ b/protocols/turnloop-smtp/src/lib.rs @@ -170,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, @@ -216,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, @@ -246,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, @@ -504,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(), @@ -620,6 +638,7 @@ impl Connection { } /// 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, @@ -628,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", @@ -656,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", @@ -676,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", @@ -695,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"); @@ -759,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); @@ -881,6 +996,7 @@ impl Connection { return; } self.state = State::Closed; + self.body_open = false; self.deadline = None; self.tx.clear(); self.tx_pos = 0; @@ -888,39 +1004,86 @@ 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")] diff --git a/protocols/turnloop-smtp/tests/smtp.rs b/protocols/turnloop-smtp/tests/smtp.rs index 8b5b69f..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(); @@ -590,6 +818,7 @@ fn test_socket(mode: Mode) { ); failed = true; } + Event::BodyReady { .. } => panic!("one-shot send streams no body"), Event::CloseTransport => {} Event::Closed => break, } From c7b598e898913c363ad4793dc948623d9dbe944a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:42:27 +0200 Subject: [PATCH 8/8] Keep the CursorBatch::rows capture test clippy-clean The test spells out the tail-position match from #69 on purpose, which clippy's manual_map and needless_match lints flag under -D warnings. Allow both on that one function, with the reason stated. --- protocols/turnloop-mongodb/tests/cursor_rows_capture.rs | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs b/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs index f815a57..7fe8b01 100644 --- a/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs +++ b/protocols/turnloop-mongodb/tests/cursor_rows_capture.rs @@ -15,6 +15,11 @@ use turnloop_mongodb::{ /// 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() {