Skip to content
Merged
50 changes: 39 additions & 11 deletions protocols/turnloop-mongodb/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
3 changes: 2 additions & 1 deletion protocols/turnloop-mongodb/src/auth.rs
Original file line number Diff line number Diff line change
@@ -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};
Expand Down
42 changes: 40 additions & 2 deletions protocols/turnloop-mongodb/src/command.rs
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -468,7 +470,8 @@ impl<'a> CursorBatch<'a> {
.and_then(|v| v.as_document()),
})
}
pub fn rows(&self) -> impl Iterator<Item = Result<&'a RawDocument>> {
/// The iterator borrows the reply, not this batch, so it may outlive `self`.
pub fn rows(&self) -> impl Iterator<Item = Result<&'a RawDocument>> + use<'a> {
self.documents.into_iter().map(|v| match v {
Ok(RawBsonRef::Document(d)) => Ok(d),
_ => Err(Error::protocol("Cursor batch contains non-document")),
Expand Down Expand Up @@ -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<Self> {
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<Self> {
if crate::error::number(r, "ok").unwrap_or(0.0) == 0.0 {
Error::from_response(r)?;
}
Expand All @@ -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<i64> {
match d.get(k).ok().flatten() {
Expand All @@ -582,7 +614,7 @@ pub struct BulkResult {
}
impl BulkResult {
pub fn accept(&mut self, reply: &RawDocument, offset: i32, ordered: bool) -> Result<bool> {
let result = WriteResult::parse(reply)?;
let result = WriteResult::decode(reply)?;
self.count += result.count;
self.modified_count += result.modified_count;
let mut errors = false;
Expand Down Expand Up @@ -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,
Expand Down
44 changes: 28 additions & 16 deletions protocols/turnloop-mongodb/src/connection.rs
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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"));
Expand Down Expand Up @@ -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<usize> {
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);
Expand Down Expand Up @@ -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();
Expand Down
13 changes: 13 additions & 0 deletions protocols/turnloop-mongodb/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]

Expand Down
3 changes: 2 additions & 1 deletion protocols/turnloop-mongodb/src/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,8 @@ pub struct Session {
pub pinned_server: Option<String>,
}
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,
Expand Down
63 changes: 63 additions & 0 deletions protocols/turnloop-mongodb/tests/cursor_rows_capture.rs
Original file line number Diff line number Diff line change
@@ -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<Result<&RawDocument>> {
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<impl Iterator<Item = Result<&RawDocument>>> {
let batch = CursorBatch::parse(reply)?;
Ok(batch.rows())
}

fn x(row: Result<&RawDocument>) -> Result<i32> {
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<i32> = rows(&reply)
.expect("batch")
.map(x)
.collect::<Result<_>>()
.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);
}
Loading