cursors: server-side cursors for find, aggregate and the listing commands

Every reply came back in a single batch with `cursor.id = 0`, `getMore` was a
stub answering an empty `nextBatch` on the literal namespace `test.$cmd`, and
nothing read `batchSize`. That caps the useful collection size at what fits in
one 48 MiB message, which is the opposite of the tens-of-GB target and the
reason M0 made whole-index scans stream: the streaming candidate generator
existed with no consumer that could suspend.

## What a cursor is allowed to remember

A cursor holds no lock between requests, so everything it saves has to survive
arbitrary concurrent mutation. Nothing here is a pointer, and the two things
that look like stable addresses are not: `reset_tree` re-creates node ids 0 and
1 as different nodes, and `rebuild_collection` moves every document. Three
sources, chosen by query shape, each with a different memory contract:

- **stream** -- an index-ordered walk resumed from a `(key, off)` anchor plus a
  `(leaf, slot)` hint. O(key) state, so this is what lets a cursor walk a
  collection larger than memory. Survives a rebuild, because a repack changes
  no key.
- **offsets** -- the matched slab offsets a narrowed plan already materialized,
  8 bytes each. Killed by a rebuild with `QueryPlanKilled`, because those
  offsets now name unrelated bytes.
- **buffered** -- canonical BSON copies, for a sort no index provides and for
  aggregate/listing output. Depends on nothing, which is what lets a listing
  hold a cursor over a `$cmd.*` namespace no collection backs.

`Collection.layout_epoch` and `Index.epoch` are the invalidation tokens, both
checked as error returns rather than assertions since a client reaches them by
keeping a cursor open across maintenance.

## Resume

`resume_forward`/`resume_reverse` are O(1) while the hint holds and fall back to
an exact-order band walk bounded by `resume_walk_max`. Without the hint, `seek`
lands at the *start* of an equal-key band, so `sort({status: 1})` over three
distinct values across 10M documents would cost ~5e10 comparisons to drain.

Two hazards found by draining a collection while writing to it, neither
predictable from reading the code:

- A deleted anchor must resume at its *band position*, or the rest of an
  equal-key band is silently dropped -- most of the collection on a
  low-cardinality index. Hence `band_index`.
- On a **unique** index a same-key entry can only be the anchor rewritten, so
  resuming at it returned updated documents twice. Observed as duplicate `_id`s
  while updating underneath a drain.

## Protocol

Measured against mongod 8.3.7 rather than recalled, which corrected three
assumptions: a bare `getMore` does *not* inherit the find's `batchSize` (4998 of
5000 documents come back), a namespace mismatch is `Unauthorized` (13) not
`CursorNotFound`, and `CursorInUse` is 143 not 12051.
`internalQueryFindCommandBatchSize` is 101, `cursorTimeoutMillis` 600000,
`clientCursorMonitorFrequencySecs` 4.

The rule everything follows is **never look ahead**: a batch that met its target
leaves the cursor open even when the source is in fact exhausted, so four
documents at `batchSize: 2` take three commands. `limit` acts as an EOF source,
which is what makes `batchSize == limit` close in one round trip. `skip` is
consumed once. `batchSize: 0` returns an empty batch with a live cursor.

Cursor ids are `(nonce << 20) | slot`, always positive. The nonce is not
decoration: without it a recycled slot serves one client another's documents.
Cursors are not connection-pinned, since the driver spec allows a `getMore` on
any connection to the same server; they end at exhaustion, `killCursors`, or the
idle sweep (a second monitor fiber, separate from the TTL one because the
cadences differ by an order of magnitude and a TTL failure must not stop
reclamation). The registry is fixed-capacity and evicts the least recently used
cursor, whose client sees the same 43 an idle timeout gives.

Fixed alongside, because cursors are what expose them:

- `listCollections` reported `"<db>."` with an *empty* collection part, which
  makes the driver throw client-side -- so it would have broken the moment its
  cursor stopped being id 0. Now `<db>.$cmd.listCollections`, as mongod uses.
- `count` ignored `skip` and `limit` entirely.
- `wire.end_message` now bounds a reply by the 48 MiB we advertise rather than
  by `maxInt(u32)`; a reply past what we told the client to expect is not a large
  reply, it is a desynchronized connection.
- Two `codeName` strings were wrong: 72 is `InvalidOptions` (MongoDB has no
  `InvalidArgument`), and 40324 reports as `Location40324`.

## Verification

Unit 160/160 in ReleaseFast and ReleaseSafe; `tests/e2e/e2e7.js` adds 86 cursor
checks across five phases (batching/lifecycle/errors, streaming across churn,
aggregate+listings+count, expiry+capacity, restart) and is self-contained
because cursor behaviour is only observable with non-default flags. No
regressions: e2e 49, e2e3 16, e2e4 17, e2e2 2, e2e6 72. Spec 168 pass / 124
fail, +5 against the previous scorecard.

Mutation-checked, per the repo's second ground rule: `hint_slot + 1`, the
`band_index` off-by-one, both epoch bumps, the id nonce, the at-least-one-
document rule, and `stream_shape` returning null each turn the intended test
red. One claim was withdrawn rather than kept -- swapping `std.mem.order` for
`cmp_prefix` in the band walk changes nothing observable, so the comment now
says so instead of asserting a check that does not hold.
This commit is contained in:
2026-08-04 14:54:27 +03:00
parent cd88e1a4d1
commit f2844e7894
15 changed files with 3489 additions and 176 deletions

View File

@@ -13,9 +13,11 @@ indexes, and TTL/unique/sparse/compound index support.
**Forward plan**: the project's direction — a full-fledged embedded, **Forward plan**: the project's direction — a full-fledged embedded,
tens-of-GB, maximally MongoDB-compatible database — is decided and written tens-of-GB, maximally MongoDB-compatible database — is decided and written
in **[PLAN.md](PLAN.md)**. The current milestone is **M0 (mmap + WAL in **[PLAN.md](PLAN.md)**. M0 (mmap + WAL storage foundation) has landed;
storage foundation)**. Before starting any work, read PLAN.md; its decision the current milestone is **M1 (cursors + wire polish)**, whose cursor work is
record (D1-D9) and ground rules are binding. done — see `src/cursor.zig` and `tests/e2e/e2e7.js`. Before starting any
work, read PLAN.md; its decision record (D1-D9) and ground rules are
binding.
## Read first, in order ## Read first, in order
@@ -71,6 +73,7 @@ node tests/e2e/e2e2.js crash-b # restart, verify all 50 survived
node tests/e2e/e2e3.js # secondary indexes node tests/e2e/e2e3.js # secondary indexes
node tests/e2e/e2e4.js # TTL indexes (server must run --ttl-sweep-secs 1) node tests/e2e/e2e4.js # TTL indexes (server must run --ttl-sweep-secs 1)
node tests/e2e/e2e6.js # self-contained full lifecycle (spawns its own server, incl. kill -9) node tests/e2e/e2e6.js # self-contained full lifecycle (spawns its own server, incl. kill -9)
node tests/e2e/e2e7.js # self-contained cursors (spawns its own servers; needs no server running)
``` ```
Which suites to run for a given change: Which suites to run for a given change:
@@ -78,6 +81,7 @@ Which suites to run for a given change:
- anything touching the write path or log format → the crash pair - anything touching the write path or log format → the crash pair
(e2e2 crash-a/b) and e2e6 (e2e2 crash-a/b) and e2e6
- anything touching indexes → e2e3.js and e2e4.js - anything touching indexes → e2e3.js and e2e4.js
- anything touching cursors, batching or the reply size → e2e7.js
- everything → all of the above - everything → all of the above
`tests/e2e/README.md` has the full matrix, ports, and harness docs `tests/e2e/README.md` has the full matrix, ports, and harness docs
@@ -180,7 +184,8 @@ src/server.zig TCP accept loop, per-connection handlers, TTL sweep monitor
src/db.zig engine: db → collection → _id → document maps, slab storage src/db.zig engine: db → collection → _id → document maps, slab storage
src/storage.zig append-only log: blocks, LZ4, XxHash3, replay, compaction src/storage.zig append-only log: blocks, LZ4, XxHash3, replay, compaction
src/query.zig filter matcher, regex engine, sort, projection src/query.zig filter matcher, regex engine, sort, projection
src/index.zig B+tree indexes: entries, search, query planner src/index.zig B+tree indexes: entries, search, query planner, scan resume
src/cursor.zig server-side cursor state: registry, batch policy, expiry
src/update.zig update operators with dot-path navigation src/update.zig update operators with dot-path navigation
src/main.zig CLI: --port, --bind, --db, --ttl-sweep-secs, --compact-threshold src/main.zig CLI: --port, --bind, --db, --ttl-sweep-secs, --compact-threshold
``` ```
@@ -193,6 +198,8 @@ src/main.zig CLI: --port, --bind, --db, --ttl-sweep-secs, --compact-threshol
baseline. baseline.
3. Commit scorecard and benchmark results with each milestone (PLAN D9) so 3. Commit scorecard and benchmark results with each milestone (PLAN D9) so
progress stays verifiable across sessions. progress stays verifiable across sessions.
4. Deferred designs (cursors, aggregation, transactions, change streams, 4. Deferred designs (aggregation, transactions, change streams, C API) are
C API) are deliberately *not* specified yet — grill the design with the deliberately *not* specified yet — grill the design with the human before
human before implementing (PLAN section 6). implementing (PLAN section 6). Cursors are no longer among them: the
design was settled and implemented in M1, and `src/cursor.zig`'s module
comment is where it is written down.

44
PLAN.md
View File

@@ -635,9 +635,47 @@ it.
## 6. Deferred designs (grill each at its milestone) ## 6. Deferred designs (grill each at its milestone)
- **M1 cursors**: cursor id allocation, idle expiration, batchSize - **M1 cursors** — *settled and implemented.* The design lives in
semantics, getMore against a lagging/compactable engine, cursor state `src/cursor.zig`'s module comment; the decisions it records, and how each
lifecycle across compaction. was reached:
- **Cursor ids** are `(nonce << 20) | slot`, always positive, never 0. The
nonce is not decoration: without it a recycled slot serves one client
another's documents, which is the worst failure this feature could have.
- **batchSize semantics** were *measured against mongod 8.3.7*, not
recalled, and three assumptions were wrong: a bare `getMore` does **not**
inherit the find's batchSize (4998 of 5000 documents come back), a
namespace mismatch is `Unauthorized` (13) rather than `CursorNotFound`,
and `CursorInUse` is 143 rather than the 12051 an earlier note claimed.
`internalQueryFindCommandBatchSize` is 101, `cursorTimeoutMillis`
600000, `clientCursorMonitorFrequencySecs` 4.
- **Never look ahead**: a batch that met its target leaves the cursor open
even when the source is in fact exhausted, so four documents at
`batchSize: 2` take three commands. The pinned suites assert that count.
- **Idle expiration** is a second monitor fiber, separate from the TTL one:
the cadences differ by an order of magnitude, and a TTL sweep failure
must not stop cursors being reclaimed. The registry is fixed-capacity and
evicts the least recently used cursor, whose client sees the same
`CursorNotFound` an idle timeout gives.
- **Against a lagging/compactable engine**, what survives depends on what
the cursor remembers, so the check is per-source: a repack changes no
key, so a streaming cursor resumes; slab offsets all move, so an offsets
cursor is killed with `QueryPlanKilled`; a snapshot needs no collection at
all. `Collection.layout_epoch` and `Index.epoch` are the tokens.
- **Resume** anchors on `(key, off)` plus a position hint gated on
`Index.epoch`, with an exact-order band walk bounded by
`resume_walk_max`. Two hazards found while implementing: a deleted anchor
must resume at its band position or the rest of an equal-key band is
silently dropped, and on a *unique* index a same-key entry can only be
the anchor rewritten — resuming at it returned updated documents twice,
caught by draining a collection being updated underneath.
Still open in M1: the doc-level free list, sessions plumbing (`lsid`
accepted), and command-monitoring (`expectEvents`) in the spec runner.
**A prerequisite the free list must honour**, recorded here while it is
still being designed: *an offset that was ever a record start must remain a
record start.* `doc_bytes` reads a `u32` length prefix in place, so an
offset landing mid-record after a re-split is a garbage-length read rather
than a wrong answer — and an offsets cursor holds exactly such offsets.
- **M2 aggregation**: stage/expression tiers, which spec-test files are - **M2 aggregation**: stage/expression tiers, which spec-test files are
the gate, whether $lookup/$unwind/facet make the first cut. the gate, whether $lookup/$unwind/facet make the first cut.
- **M4 transactions**: snapshot isolation over mmap (COW vs undo), read - **M4 transactions**: snapshot isolation over mmap (COW vs undo), read

View File

@@ -12,8 +12,9 @@ maximally MongoDB-compatible database — its decision record, milestones
and gates live in [PLAN.md](PLAN.md). Milestone 0 (mmap + WAL storage and gates live in [PLAN.md](PLAN.md). Milestone 0 (mmap + WAL storage
foundation) has landed; its measured gate results are in foundation) has landed; its measured gate results are in
[`tests/e2e/results/m0-gates.txt`](tests/e2e/results/m0-gates.txt). [`tests/e2e/results/m0-gates.txt`](tests/e2e/results/m0-gates.txt).
Milestone 1 (cursors, and the doc-level free list the churn gate showed is Milestone 1 is in progress: server-side cursors have landed (see
needed) is next. **Cursors** below); the doc-level free list the churn gate showed is needed
is still open.
## Quick start ## Quick start
@@ -145,10 +146,43 @@ whole pass, so the interval is the tuning knob: expiry is never more
precise than `--ttl-sweep-secs`, and a very large TTL index wants a precise than `--ttl-sweep-secs`, and a very large TTL index wants a
longer one. longer one.
## Cursors
`find`, `aggregate`, `listCollections` and `listIndexes` return real cursor
ids, and `getMore`/`killCursors` work. Batching follows MongoDB: a first
batch of 101 documents unless `batchSize` says otherwise, a `getMore` with
no `batchSize` bounded only by the 16 MiB batch cap, `batchSize: 0` as an
empty batch with a live cursor, and `limit` honoured across batches. Every
default here was measured against a real `mongod` rather than assumed.
A cursor holds no lock between requests, so what it remembers has to survive
arbitrary concurrent writes. Three shapes, picked by the query:
| query | what the cursor keeps |
| --- | --- |
| a whole-index walk (`find({})`, or a sort an index provides) | the last key and offset it yielded — O(key), whatever the collection size |
| a narrowed index plan | the matching offsets, 8 bytes each |
| a sort no index provides, or aggregate/listing output | a snapshot of the remaining documents |
The first is what lets a cursor walk a collection larger than memory. It
also survives a compaction, because a repack changes no key; the offsets
form cannot, and says so with `QueryPlanKilled` rather than returning
documents from the wrong place.
Cursors are not pinned to the connection that created them, so a `getMore`
may arrive on any connection — which is what the driver specification
allows. They are reclaimed when exhausted, when killed, or after
`--cursor-timeout-ms` idle (default 600000, MongoDB's own
`cursorTimeoutMillis`); `--max-open-cursors` bounds the registry and evicts
the least recently used cursor at capacity, whose client then sees the same
`CursorNotFound` an idle timeout gives.
## Not (yet) implemented ## Not (yet) implemented
- Authentication (SCRAM) — run without credentials - Authentication (SCRAM) — run without credentials
- Real cursors (all results are returned in one batch, cursor id 0) - Tailable/awaitData cursors, which need capped collections; a tailable
`find` is rejected, exactly as MongoDB rejects one on a non-capped
collection
- Transactions, change streams, replicasets - Transactions, change streams, replicasets
- Compression (OP_COMPRESSED) - Compression (OP_COMPRESSED)
- `collMod`, so an index's `expireAfterSeconds` cannot be changed in - `collMod`, so an index's `expireAfterSeconds` cannot be changed in

File diff suppressed because it is too large Load Diff

979
src/cursor.zig Normal file
View File

@@ -0,0 +1,979 @@
//! Server-side cursor state: what a `find`/`aggregate` leaves behind so a later
//! `getMore` can carry on, and the fixed-capacity registry that holds it.
//!
//! This module is deliberately *pure*: it owns state and policy, never
//! execution. It does not import `db.zig` or `commands.zig`, so `db.Engine` can
//! own a `Store` with no import cycle, and the batch policy below is testable
//! with no engine, no socket and no allocator. Filling a batch stays in
//! `commands.zig`, which already owns orchestration the way `index.zig` owns
//! planning.
//!
//! **The one rule the whole batching protocol follows: never look ahead.** A
//! batch ends either because it reached its target -- and the cursor stays open
//! -- or because the source reported EOF, and then the cursor closes with
//! `id: 0` in that same reply. A batch that reached its target leaves the cursor
//! open *even when the source happens to be exhausted*. So four documents at
//! `batchSize: 2` need a third command answering `nextBatch: []` with `id: 0`;
//! that empty terminal batch is correct, not a bug, and the pinned spec suites
//! assert exactly that command count.
//!
//! ## What a cursor is allowed to remember
//!
//! A cursor holds no lock between requests, so everything it saves must survive
//! arbitrary concurrent mutation. Nothing here is a pointer, and the two things
//! that look like stable addresses are not:
//!
//! - A tree position `(leaf, slot)` is invalidated by `Index.reset_tree`,
//! which clears the node table so ids 0 and 1 become a live but *unrelated*
//! root and leaf. Guarded by `Stream.index_epoch`.
//! - A slab offset is invalidated by `rebuild_collection`, which moves every
//! document. Guarded by `layout_epoch`.
//!
//! Both are checked as error returns rather than assertions, because a client
//! can reach either one by keeping a cursor open across maintenance.
const std = @import("std");
const bson = @import("bson.zig");
const index = @import("index.zig");
// Always active, including in the default ReleaseFast build -- see assert.zig.
const assert = @import("assert.zig").assert;
const assert_msg = @import("assert.zig").assert_msg;
// ---------------------------------------------------------------------------
// Bounds
// ---------------------------------------------------------------------------
/// Longest index key a `.stream` cursor will anchor on. Tied to the B+tree's
/// own "this record is normal" threshold rather than picked: a key past it has
/// already spilled to the overflow slab, so the tree itself considers it
/// exceptional. Also what makes the anchor a fixed inline array instead of an
/// allocation.
///
/// Unbounded, this is a memory denial of service and not a subtle one:
/// `bson.encode_key` escapes NULs, so a 16 MB string doubles, and a compound
/// index may carry 32 of them.
pub const anchor_key_max: usize = 1024;
comptime {
// The bound is only defensible if it really is the tree's spill threshold.
// If the page size or the spill fraction ever changes, this fails to
// compile rather than silently becoming an arbitrary number.
std.debug.assert(anchor_key_max == index.inline_limit);
}
pub const ns_db_max: usize = 64;
pub const ns_coll_max: usize = 192;
pub const index_name_max: usize = 128;
/// Documents in a first batch when the client named no `batchSize`. MongoDB's
/// own default (`internalQueryFindCommandBatchSize`).
pub const default_first_batch: u32 = 101;
/// Cap on a batch's document payload: `maxBsonObjectSize`, which is also what
/// leaves room for the reply envelope inside the 48 MiB message limit.
pub const batch_bytes_max: u64 = 16 * 1024 * 1024;
/// Idle milliseconds before the sweep reaps a cursor. MongoDB's
/// `cursorTimeoutMillis`.
pub const default_idle_timeout_ms: i64 = 10 * 60 * 1000;
/// Slots in the registry unless configured otherwise.
pub const default_capacity: u32 = 4096;
/// Low bits of a cursor id that address its slot; the rest is the nonce.
const slot_bits: u6 = 20;
const slot_mask: u64 = (@as(u64, 1) << slot_bits) - 1;
/// Nonce width, leaving the sign bit clear so every id is a positive i64.
const nonce_mask: u64 = (@as(u64, 1) << (63 - slot_bits)) - 1;
pub const max_capacity: u32 = @intCast(slot_mask);
// ---------------------------------------------------------------------------
// Cursor state
// ---------------------------------------------------------------------------
/// A namespace, by value. A cursor cannot hold a `*Collection`: `drop` frees
/// it, and the pointer would dangle exactly the way the M0 notes on
/// heap-allocating collections describe.
pub const Ns = struct {
db: []const u8,
coll: []const u8,
};
/// Where the remaining documents come from.
pub const Source = union(enum) {
/// An index-ordered scan, resumed from a value-typed anchor. O(key)
/// memory, so this is the shape that lets a cursor walk a collection far
/// larger than memory -- the reason M0 made whole-index scans stream.
stream: Stream,
/// Matched slab offsets, 8 bytes each, which `scan_sorted` has already
/// materialized for a narrowed plan.
///
/// Safe against an offset the coming document free list has recycled,
/// because every batch re-applies the full filter -- the index invariant.
/// A recycled offset is therefore either rejected or resolves to a
/// document that genuinely matches. It needs one guarantee from the free
/// list, recorded in PLAN: an offset that was ever a record start must
/// stay a record start, since `doc_bytes` reads a length prefix in place.
offsets: struct { items: []u64, next: u32 = 0 },
/// Canonical BSON bytes owned by the cursor's arena, for results with no
/// stable backing store to point at: a sort no index provides, and
/// aggregate/listCollections/listIndexes output.
buffered: struct { docs: []const []const u8, next: u32 = 0 },
};
/// A resumable index scan. Every field is a value; nothing here is a pointer
/// into the tree, the slab or the request that created it.
pub const Stream = struct {
/// Empty means the implicit `_id_` index. Re-resolved by name on every
/// `getMore`, so a `dropIndexes` cannot leave a dangling `*Index`.
index_name_buf: [index_name_max]u8 = undefined,
index_name_len: u8 = 0,
/// Bumped by `reset_tree`/`replace_root_with_leaf`; if it moved, the hint
/// below addresses a different tree and must not be trusted.
index_epoch: u64 = 0,
backward: bool = false,
anchor_buf: [anchor_key_max]u8 = undefined,
anchor_len: u16 = 0,
anchor_off: u64 = 0,
/// Entries sharing the anchor's key that this cursor has already yielded.
/// Without it, an anchor whose document was deleted between batches would
/// resume past the entire equal-key band -- on a three-value index that is
/// millions of documents silently missing.
band_index: u64 = 0,
/// Last known position of the anchor. A hint, never trusted without
/// re-reading the entry there: it turns resume from a walk down the
/// equal-key band into O(1), which is what keeps a low-cardinality
/// `sort({status: 1})` from costing O(band) per batch.
hint_leaf: u32 = 0,
hint_slot: u32 = 0,
pub fn index_name(self: *const Stream) []const u8 {
return self.index_name_buf[0..self.index_name_len];
}
pub fn anchor_key(self: *const Stream) []const u8 {
return self.anchor_buf[0..self.anchor_len];
}
/// Whether anything has been yielded yet. Derived rather than stored: an
/// encoded index key always begins with `bson.encode_key`'s rank byte, so it
/// is never empty, and a separate `started` flag would be a second field that
/// has to agree with this one.
///
/// It can legitimately be false on a live cursor: `batchSize: 0` returns an
/// empty first batch without consuming anything, and such a cursor starts at
/// `iter()`/`iter_reverse()` rather than resuming.
pub fn started(self: *const Stream) bool {
return self.anchor_len > 0;
}
/// Record the entry just yielded as the point to resume after.
///
/// Asserts the anchor advances in scan order. This is the single check most
/// likely to catch a resume bug: going backwards duplicates documents,
/// standing still makes `getMore` loop forever, and both are far easier to
/// see here than in a client's result set. Equal keys are legal (a
/// duplicate band), which is exactly why `band_index` also has to move.
pub fn advance(self: *Stream, key: []const u8, off: u64, leaf: u32, slot: u32) void {
assert(key.len <= anchor_key_max);
// What makes `started()` derivable, so pin it here rather than trust it.
assert_msg(key.len > 0, "an encoded index key is never empty");
if (self.started()) {
const order = std.mem.order(u8, key, self.anchor_key());
if (self.backward) {
assert_msg(order != .gt, "a reverse cursor's anchor moved forward");
} else {
assert_msg(order != .lt, "a forward cursor's anchor moved backward");
}
if (order == .eq) {
assert_msg(
off != self.anchor_off or self.band_index > 0,
"a cursor re-anchored on the entry it just yielded",
);
// Still inside the anchor's band, so the position within it has
// to move or a resume could not tell the two entries apart.
self.band_index += 1;
} else {
self.band_index = 0;
}
}
@memcpy(self.anchor_buf[0..key.len], key);
self.anchor_len = @intCast(key.len);
self.anchor_off = off;
self.hint_leaf = leaf;
self.hint_slot = slot;
}
};
pub const Cursor = struct {
/// Positive and never 0: `id: 0` is "no cursor" on the wire.
id: i64,
ns_db_buf: [ns_db_max]u8 = undefined,
ns_db_len: u8 = 0,
ns_coll_buf: [ns_coll_max]u8 = undefined,
ns_coll_len: u8 = 0,
/// Bumped when a rebuild moves documents, so a saved offset or anchor
/// offset is stale. Also the drop detector.
layout_epoch: u64 = 0,
/// Serialized so they outlive the request that parsed them: a parsed
/// `[]bson.Pair` points into the per-request message arena, and the reply
/// arena is reset on every request.
filter_bytes: []const u8 = &.{},
proj_bytes: []const u8 = &.{},
/// Documents still owed across all remaining batches; null is unbounded.
/// Reaching 0 is an EOF *source*, which is what closes the cursor in the
/// very batch that exhausts the limit rather than one round trip later.
/// Optional rather than "0 means unbounded" precisely because 0 has to keep
/// its literal meaning here.
remaining_limit: ?u64 = null,
/// The client's `batchSize`, reused when a `getMore` names none.
batch_size: ?u32 = null,
/// Exempt from the idle sweep. Still killable by `killCursors` and by
/// eviction -- a fixed-capacity registry cannot promise "never expires".
no_timeout: bool = false,
/// A request is using this cursor right now. Concurrent use is rejected
/// rather than queued: queueing lets one client turn a single cursor into a
/// connection-count denial of service.
pinned: bool = false,
/// `killCursors` arrived while pinned; the in-flight request frees it.
kill_requested: bool = false,
last_use_ms: i64 = 0,
arena: std.heap.ArenaAllocator,
source: Source,
pub fn ns(self: *const Cursor) Ns {
return .{
.db = self.ns_db_buf[0..self.ns_db_len],
.coll = self.ns_coll_buf[0..self.ns_coll_len],
};
}
pub fn ns_matches(self: *const Cursor, other: Ns) bool {
const own = self.ns();
return std.mem.eql(u8, own.db, other.db) and std.mem.eql(u8, own.coll, other.coll);
}
};
/// Everything a caller must decide before a cursor can exist. Grouped so
/// `open` cannot be called with an argument silently in the wrong position.
pub const OpenSpec = struct {
ns: Ns,
layout_epoch: u64,
filter_bytes: []const u8 = &.{},
proj_bytes: []const u8 = &.{},
remaining_limit: ?u64 = null,
batch_size: ?u32 = null,
no_timeout: bool = false,
source: Source,
};
pub const OpenError = error{
/// The namespace does not fit the fixed buffers. Callers degrade to a
/// single batch rather than failing the query.
NameTooLong,
/// Every slot is pinned by an in-flight request.
TooManyCursors,
OutOfMemory,
/// Taking the store mutex was cancelled (shutdown).
Canceled,
};
pub const PinError = error{
CursorNotFound,
/// The id exists but belongs to another namespace. Distinct from
/// `CursorNotFound` because mongod answers this with `Unauthorized` (13),
/// not 43, and leaves the cursor alive -- the request is wrong, not the
/// cursor.
CursorNamespaceMismatch,
CursorInUse,
Canceled,
};
/// Owned storage for a namespace copied out of the store, so an error message
/// can name a cursor's namespace without holding the store's mutex or a pointer
/// into its slots.
pub const NsBuf = struct {
db_buf: [ns_db_max]u8 = undefined,
db_len: u8 = 0,
coll_buf: [ns_coll_max]u8 = undefined,
coll_len: u8 = 0,
pub fn ns(self: *const NsBuf) Ns {
return .{ .db = self.db_buf[0..self.db_len], .coll = self.coll_buf[0..self.coll_len] };
}
fn set(self: *NsBuf, from: Ns) void {
@memcpy(self.db_buf[0..from.db.len], from.db);
self.db_len = @intCast(from.db.len);
@memcpy(self.coll_buf[0..from.coll.len], from.coll);
self.coll_len = @intCast(from.coll.len);
}
};
pub const KillOutcome = enum { killed, not_found };
// ---------------------------------------------------------------------------
// The registry
// ---------------------------------------------------------------------------
pub const Store = struct {
/// Guards every field below. A **leaf** lock: no other lock -- catalog,
/// collection, log -- is ever acquired while it is held, so it cannot
/// participate in a cycle. In particular a `getMore` copies what it needs
/// out, releases this, and only then iterates under the collection lock;
/// otherwise the reaper would block behind a full scan.
mutex: std.Io.Mutex = .init,
/// Boxed, not inline. A `Cursor` inlines its anchor and namespace buffers and
/// so is ~1.5 KiB; at the default capacity an inline table would be 6.2 MiB
/// allocated and zeroed at *every* `Engine.open` -- paid by every embedded
/// user and by all ~47 engine opens in the unit suite, to hold zero cursors.
/// A pointer table is 32 KiB and the cursor itself is allocated when one
/// actually exists, which is also when its arena is created anyway.
slots: []?*Cursor,
/// Mixed into every id so ids are not guessable across processes, and so a
/// reused slot rejects the previous id exactly. Without this a stale
/// `getMore` can address a recycled slot and read another client's cursor.
nonce: u64,
live: u32 = 0,
idle_timeout_ms: i64 = default_idle_timeout_ms,
gpa: std.mem.Allocator,
pub fn init(
gpa: std.mem.Allocator,
io: std.Io,
capacity: u32,
idle_timeout_ms: i64,
) !Store {
assert(capacity > 0 and capacity <= max_capacity);
var seed: [8]u8 = undefined;
io.random(&seed);
const slots = try gpa.alloc(?*Cursor, capacity);
@memset(slots, null);
return .{
.slots = slots,
// A zero nonce would make the first slot's id equal to its index,
// and slot 0's id would be 0 -- which means "no cursor".
.nonce = std.mem.readInt(u64, &seed, .little) | 1,
.idle_timeout_ms = idle_timeout_ms,
.gpa = gpa,
};
}
pub fn deinit(self: *Store) void {
for (self.slots) |maybe| {
if (maybe) |c| destroy_cursor(self.gpa, c);
}
self.gpa.free(self.slots);
self.slots = &.{};
}
/// Free a cursor: its arena first, then the box the slot pointed at.
fn destroy_cursor(gpa: std.mem.Allocator, c: *Cursor) void {
c.arena.deinit();
gpa.destroy(c);
}
fn slot_of(id: i64) usize {
return @intCast(@as(u64, @bitCast(id)) & slot_mask);
}
/// Build the id for `slot` at the store's current nonce, then advance the
/// nonce so the next cursor in this slot gets a different id.
fn mint(self: *Store, slot: usize) i64 {
const n = self.nonce & nonce_mask;
self.nonce +%= 1;
const raw = (n << slot_bits) | @as(u64, @intCast(slot));
const id: i64 = @intCast(raw & ~(@as(u64, 1) << 63));
// Both properties are load-bearing on the wire and in lookup.
assert_msg(id > 0, "a cursor id must be a positive int64");
assert_msg(slot_of(id) == slot, "a cursor id must address its own slot");
return id;
}
/// Register a cursor and return its id, or null when the caller should
/// answer in a single batch instead (`NameTooLong` is not worth failing a
/// query over -- the degradation is exactly today's behaviour).
///
/// Takes ownership of `spec.source` and of the arena backing it.
pub fn open(
self: *Store,
io: std.Io,
now_ms: i64,
arena: std.heap.ArenaAllocator,
spec: OpenSpec,
) OpenError!i64 {
if (spec.ns.db.len > ns_db_max or spec.ns.coll.len > ns_coll_max) {
return error.NameTooLong;
}
try self.mutex.lock(io);
defer self.mutex.unlock(io);
const slot = self.free_slot(now_ms) orelse return error.TooManyCursors;
assert(self.slots[slot] == null);
const c = try self.gpa.create(Cursor);
errdefer self.gpa.destroy(c);
c.* = .{
.id = self.mint(slot),
.layout_epoch = spec.layout_epoch,
.filter_bytes = spec.filter_bytes,
.proj_bytes = spec.proj_bytes,
.remaining_limit = spec.remaining_limit,
.batch_size = spec.batch_size,
.no_timeout = spec.no_timeout,
.last_use_ms = now_ms,
.arena = arena,
.source = spec.source,
};
@memcpy(c.ns_db_buf[0..spec.ns.db.len], spec.ns.db);
c.ns_db_len = @intCast(spec.ns.db.len);
@memcpy(c.ns_coll_buf[0..spec.ns.coll.len], spec.ns.coll);
c.ns_coll_len = @intCast(spec.ns.coll.len);
self.slots[slot] = c;
self.live += 1;
return c.id;
}
/// An empty slot: a genuinely free one, else the least-recently-used
/// unpinned cursor. Evicting is legal and cheap to reason about, because
/// the victim's client gets `CursorNotFound` on its next `getMore` -- the
/// same answer an idle timeout gives, which every driver already handles.
/// Caller holds the mutex.
fn free_slot(self: *Store, now_ms: i64) ?usize {
var lru: ?usize = null;
var lru_ms: i64 = std.math.maxInt(i64);
for (self.slots, 0..) |maybe, i| {
const c = maybe orelse return i;
// Reap on the way past, so a store that has gone quiet does not
// wait for the sweep tick to reclaim what already expired.
if (self.expired(c, now_ms)) {
self.destroy(i);
return i;
}
if (c.pinned) continue;
if (c.last_use_ms < lru_ms) {
lru_ms = c.last_use_ms;
lru = i;
}
}
if (lru) |i| {
self.destroy(i);
return i;
}
return null;
}
/// Caller holds the mutex.
fn expired(self: *const Store, c: *const Cursor, now_ms: i64) bool {
if (c.pinned or c.no_timeout or self.idle_timeout_ms <= 0) return false;
return now_ms -| c.last_use_ms >= self.idle_timeout_ms;
}
/// Caller holds the mutex.
fn destroy(self: *Store, slot: usize) void {
const c = self.slots[slot] orelse return;
assert_msg(!c.pinned, "a pinned cursor must not be destroyed under its user");
destroy_cursor(self.gpa, c);
self.slots[slot] = null;
self.live -= 1;
}
/// Claim a cursor for one request. The returned pointer is stable only
/// until `release`, and only the pinning request may touch it.
///
/// The namespace check is not cosmetic: `dispatch` locks the collection
/// named in the *message*, so a `getMore` quoting one cursor's id and
/// another collection's name would otherwise iterate the first collection's
/// index while holding the second collection's lock. A mismatch leaves the
/// cursor alive -- it is the request that is wrong, not the cursor.
pub fn pin(self: *Store, io: std.Io, id: i64, ns: Ns, now_ms: i64) PinError!*Cursor {
try self.mutex.lock(io);
defer self.mutex.unlock(io);
if (id <= 0) return error.CursorNotFound;
const slot = slot_of(id);
if (slot >= self.slots.len) return error.CursorNotFound;
const c = self.slots[slot] orelse return error.CursorNotFound;
// Compare the whole id, not just the slot: this is what makes a
// recycled slot reject its predecessor's id.
if (c.id != id) return error.CursorNotFound;
if (self.expired(c, now_ms)) {
self.destroy(slot);
return error.CursorNotFound;
}
if (!c.ns_matches(ns)) return error.CursorNamespaceMismatch;
if (c.pinned) return error.CursorInUse;
c.pinned = true;
c.last_use_ms = now_ms;
return c;
}
/// The namespace a live cursor belongs to, copied out. Only used to build
/// the namespace-mismatch error message, so a second lock acquisition on an
/// error path is the right trade for not threading an out-parameter through
/// the success path.
pub fn ns_of(self: *Store, io: std.Io, id: i64, out: *NsBuf) bool {
self.mutex.lock(io) catch return false;
defer self.mutex.unlock(io);
if (id <= 0) return false;
const slot = slot_of(id);
if (slot >= self.slots.len) return false;
const c = self.slots[slot] orelse return false;
if (c.id != id) return false;
out.set(c.ns());
return true;
}
/// Hand a pinned cursor back. `exhausted` destroys it, and so does a
/// `killCursors` that arrived while it was pinned.
pub fn release(self: *Store, io: std.Io, c: *Cursor, now_ms: i64, exhausted: bool) void {
self.mutex.lock(io) catch {
// Cancellation while returning a cursor would otherwise leave it
// pinned forever, unreachable and un-reapable. Unpinning without
// the lock is the lesser evil: the field is only ever written by
// the one request that owns the pin.
c.pinned = false;
return;
};
defer self.mutex.unlock(io);
assert_msg(c.pinned, "released a cursor that was not pinned");
c.pinned = false;
c.last_use_ms = now_ms;
if (exhausted or c.kill_requested) {
const slot = slot_of(c.id);
assert(self.slots[slot].? == c);
self.destroy(slot);
}
}
/// `killCursors` for one id. A pinned cursor is marked and reported killed:
/// the client's intent is satisfied, and the in-flight request frees it on
/// release. Storage is never freed under a running request.
pub fn kill(self: *Store, io: std.Io, id: i64, ns: Ns) KillOutcome {
self.mutex.lock(io) catch return .not_found;
defer self.mutex.unlock(io);
if (id <= 0) return .not_found;
const slot = slot_of(id);
if (slot >= self.slots.len) return .not_found;
const c = self.slots[slot] orelse return .not_found;
if (c.id != id or !c.ns_matches(ns)) return .not_found;
if (c.pinned) {
c.kill_requested = true;
return .killed;
}
self.destroy(slot);
return .killed;
}
/// Kill every cursor on a namespace. Called when the collection or its
/// database is dropped: a later `getMore` would fail anyway, since the
/// cursor holds names rather than a pointer, but reaping here frees the
/// slots at once and keeps the open-cursor metric honest.
pub fn kill_namespace(
self: *Store,
io: std.Io,
db_name: []const u8,
coll_name: ?[]const u8,
) u32 {
self.mutex.lock(io) catch return 0;
defer self.mutex.unlock(io);
var n: u32 = 0;
for (self.slots, 0..) |maybe, i| {
const c = maybe orelse continue;
const own = c.ns();
if (!std.mem.eql(u8, own.db, db_name)) continue;
if (coll_name) |name| {
if (!std.mem.eql(u8, own.coll, name)) continue;
}
if (c.pinned) {
c.kill_requested = true;
} else {
self.destroy(i);
}
n += 1;
}
return n;
}
/// Reap idle cursors. Returns how many went.
pub fn sweep(self: *Store, io: std.Io, now_ms: i64) u32 {
self.mutex.lock(io) catch return 0;
defer self.mutex.unlock(io);
if (self.live == 0) return 0;
var n: u32 = 0;
for (self.slots, 0..) |maybe, i| {
const c = maybe orelse continue;
if (!self.expired(c, now_ms)) continue;
self.destroy(i);
n += 1;
}
return n;
}
};
// ---------------------------------------------------------------------------
// Batch policy
// ---------------------------------------------------------------------------
/// What `offer` decided about one document.
pub const Offered = enum {
appended,
/// The batch is full. **Not** EOF: the cursor stays open, and this document
/// has not been consumed -- the caller must hand it to the next batch.
batch_full,
};
/// Accumulates one batch and owns the two limits that end it.
///
/// Split out from the emit path so the whole policy is a pure function of
/// `(target, emitted, bytes, size)` and can be unit-tested without an engine,
/// a socket or an allocator. The subtle parts are all here: a target of 0 is
/// unbounded (a `getMore` naming no `batchSize`), a `batchSize: 0` first batch
/// is a target that is *reached immediately*, and the byte cap must still let
/// the first document through or an oversized document would wedge the cursor
/// forever, returning empty batches with no progress.
pub const BatchBuilder = struct {
/// Documents wanted; null means no document target, fill to the byte cap.
/// Optional rather than "0 means unbounded" because `batchSize: 0` is a real
/// request for an empty batch, and conflating the two returned the whole
/// collection where mongod returns nothing.
target: ?u32,
bytes_max: u64 = batch_bytes_max,
emitted: u32 = 0,
bytes: u64 = 0,
pub fn init(target: ?u32) BatchBuilder {
return .{ .target = target };
}
/// Whether the batch has already met its document target, checked before
/// pulling from the source so a full batch never consumes a document it
/// cannot carry.
pub fn full(self: *const BatchBuilder) bool {
const t = self.target orelse return false;
return self.emitted >= t;
}
/// Account for a document of `size` serialized bytes.
pub fn offer(self: *BatchBuilder, size: u64) Offered {
assert_msg(!self.full(), "offered a document to a batch that was already full");
// The at-least-one rule: an empty batch takes the document whatever it
// measures. Stored documents cannot exceed the cap (inserts enforce
// 16 MiB), so this only arises for a generated one.
if (self.emitted > 0 and self.bytes + size > self.bytes_max) return .batch_full;
self.emitted += 1;
self.bytes += size;
return .appended;
}
};
/// The document target for a batch: the client's `batchSize` if it named one,
/// otherwise 101 for a first batch and *no* document target for a `getMore`.
///
/// Both defaults are measured against mongod 8.3.7 rather than assumed.
/// `internalQueryFindCommandBatchSize` reports 101, and a `getMore` carrying no
/// `batchSize` after a `find` with `batchSize: 2` returns 4998 of 5000
/// documents -- so a bare `getMore` is bounded by bytes alone and does *not*
/// inherit the `batchSize` the cursor was created with.
pub fn batch_target(batch_size: ?u32, first: bool) ?u32 {
if (batch_size) |n| return n;
return if (first) default_first_batch else null;
}
// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
const testing = std.testing;
/// A Store for tests, with a threaded Io so the mutex is real.
const TestStore = struct {
threaded: std.Io.Threaded,
store: Store,
fn init(capacity: u32, idle_timeout_ms: i64) !TestStore {
var self: TestStore = undefined;
self.threaded = std.Io.Threaded.init(testing.allocator, .{});
self.store = try Store.init(
testing.allocator,
self.threaded.io(),
capacity,
idle_timeout_ms,
);
return self;
}
fn io(self: *TestStore) std.Io {
return self.threaded.io();
}
fn deinit(self: *TestStore) void {
self.store.deinit();
self.threaded.deinit();
}
fn open_one(self: *TestStore, coll: []const u8, now_ms: i64) !i64 {
const arena = std.heap.ArenaAllocator.init(testing.allocator);
return self.store.open(self.io(), now_ms, arena, .{
.ns = .{ .db = "t", .coll = coll },
.layout_epoch = 0,
.source = .{ .buffered = .{ .docs = &.{} } },
});
}
};
test "cursor ids are positive, address their slot, and never repeat" {
var ts = try TestStore.init(4, default_idle_timeout_ms);
defer ts.deinit();
var seen: [16]i64 = undefined;
for (0..16) |i| {
const id = try ts.open_one("c", 0);
try testing.expect(id > 0);
// Freeing the slot immediately means the next open reuses it, which is
// exactly the case the nonce has to survive.
const killed = ts.store.kill(ts.io(), id, .{ .db = "t", .coll = "c" });
try testing.expectEqual(KillOutcome.killed, killed);
seen[i] = id;
}
for (seen, 0..) |a, i| {
for (seen[i + 1 ..]) |b| try testing.expect(a != b);
}
}
test "a recycled slot rejects the id it used to hold" {
// The guard that keeps one client from reading another's cursor. Without
// the nonce in the id, `slot_of(stale) == slot_of(fresh)` and the stale
// getMore would be served the new cursor's documents.
var ts = try TestStore.init(1, default_idle_timeout_ms);
defer ts.deinit();
const ns = Ns{ .db = "t", .coll = "c" };
const stale = try ts.open_one("c", 0);
try testing.expectEqual(KillOutcome.killed, ts.store.kill(ts.io(), stale, ns));
const fresh = try ts.open_one("c", 0);
try testing.expectEqual(Store.slot_of(stale), Store.slot_of(fresh));
try testing.expect(stale != fresh);
try testing.expectError(error.CursorNotFound, ts.store.pin(ts.io(), stale, ns, 0));
_ = try ts.store.pin(ts.io(), fresh, ns, 0);
}
test "pin rejects a wrong namespace and a second holder, and leaves the cursor alive" {
var ts = try TestStore.init(4, default_idle_timeout_ms);
defer ts.deinit();
const ns = Ns{ .db = "t", .coll = "c" };
const id = try ts.open_one("c", 0);
// A wrong namespace must not kill the cursor: the request is wrong, not
// the cursor, and the client is allowed to retry correctly. Reported apart
// from CursorNotFound because mongod answers it with Unauthorized (13).
const wrong_coll = Ns{ .db = "t", .coll = "other" };
const wrong_db = Ns{ .db = "other", .coll = "c" };
const mismatch = error.CursorNamespaceMismatch;
try testing.expectError(mismatch, ts.store.pin(ts.io(), id, wrong_coll, 0));
try testing.expectError(mismatch, ts.store.pin(ts.io(), id, wrong_db, 0));
var found: NsBuf = .{};
try testing.expect(ts.store.ns_of(ts.io(), id, &found));
try testing.expectEqualStrings("t", found.ns().db);
try testing.expectEqualStrings("c", found.ns().coll);
const c = try ts.store.pin(ts.io(), id, ns, 0);
try testing.expectError(error.CursorInUse, ts.store.pin(ts.io(), id, ns, 0));
ts.store.release(ts.io(), c, 1, false);
// Released, so it can be pinned again.
const again = try ts.store.pin(ts.io(), id, ns, 2);
ts.store.release(ts.io(), again, 3, true);
try testing.expectError(error.CursorNotFound, ts.store.pin(ts.io(), id, ns, 4));
}
test "a full store evicts the least recently used unpinned cursor" {
var ts = try TestStore.init(3, default_idle_timeout_ms);
defer ts.deinit();
const ns = Ns{ .db = "t", .coll = "c" };
const a = try ts.open_one("c", 100);
const b = try ts.open_one("c", 200);
const c = try ts.open_one("c", 300);
// Touch `a` so `b` becomes the least recently used.
const pinned_a = try ts.store.pin(ts.io(), a, ns, 400);
ts.store.release(ts.io(), pinned_a, 400, false);
const d = try ts.open_one("c", 500);
try testing.expectError(error.CursorNotFound, ts.store.pin(ts.io(), b, ns, 500));
for ([_]i64{ a, c, d }) |id| {
const live = try ts.store.pin(ts.io(), id, ns, 500);
ts.store.release(ts.io(), live, 500, false);
}
}
test "a store whose every slot is pinned refuses rather than evicting" {
var ts = try TestStore.init(2, default_idle_timeout_ms);
defer ts.deinit();
const ns = Ns{ .db = "t", .coll = "c" };
const a = try ts.open_one("c", 0);
const b = try ts.open_one("c", 0);
_ = try ts.store.pin(ts.io(), a, ns, 0);
_ = try ts.store.pin(ts.io(), b, ns, 0);
try testing.expectError(error.TooManyCursors, ts.open_one("c", 0));
}
test "the sweep reaps idle cursors and spares noCursorTimeout" {
var ts = try TestStore.init(4, 1000);
defer ts.deinit();
const ns = Ns{ .db = "t", .coll = "c" };
const perishable = try ts.open_one("c", 0);
const arena = std.heap.ArenaAllocator.init(testing.allocator);
const immortal = try ts.store.open(ts.io(), 0, arena, .{
.ns = ns,
.layout_epoch = 0,
.no_timeout = true,
.source = .{ .buffered = .{ .docs = &.{} } },
});
// Just short of the timeout: nothing goes.
try testing.expectEqual(@as(u32, 0), ts.store.sweep(ts.io(), 999));
try testing.expectEqual(@as(u32, 1), ts.store.sweep(ts.io(), 1000));
try testing.expectError(error.CursorNotFound, ts.store.pin(ts.io(), perishable, ns, 1000));
// The exempt one survives an interval it would otherwise have died in...
try testing.expectEqual(@as(u32, 0), ts.store.sweep(ts.io(), 100_000));
const live = try ts.store.pin(ts.io(), immortal, ns, 100_000);
ts.store.release(ts.io(), live, 100_000, false);
// ...but is still killable explicitly.
try testing.expectEqual(KillOutcome.killed, ts.store.kill(ts.io(), immortal, ns));
}
test "killCursors reports a pinned cursor killed and frees it on release" {
var ts = try TestStore.init(4, default_idle_timeout_ms);
defer ts.deinit();
const ns = Ns{ .db = "t", .coll = "c" };
const id = try ts.open_one("c", 0);
const c = try ts.store.pin(ts.io(), id, ns, 0);
try testing.expectEqual(KillOutcome.killed, ts.store.kill(ts.io(), id, ns));
// Still pinned, so its storage must not have been freed under the request.
try testing.expect(c.kill_requested);
// Not exhausted, but the pending kill wins.
ts.store.release(ts.io(), c, 1, false);
try testing.expectError(error.CursorNotFound, ts.store.pin(ts.io(), id, ns, 2));
try testing.expectEqual(KillOutcome.not_found, ts.store.kill(ts.io(), id, ns));
}
test "kill_namespace reaps a collection's cursors and leaves the rest" {
var ts = try TestStore.init(8, default_idle_timeout_ms);
defer ts.deinit();
const doomed = try ts.open_one("doomed", 0);
const spared = try ts.open_one("spared", 0);
try testing.expectEqual(@as(u32, 1), ts.store.kill_namespace(ts.io(), "t", "doomed"));
const doomed_ns = Ns{ .db = "t", .coll = "doomed" };
try testing.expectError(error.CursorNotFound, ts.store.pin(ts.io(), doomed, doomed_ns, 0));
const live = try ts.store.pin(ts.io(), spared, .{ .db = "t", .coll = "spared" }, 0);
ts.store.release(ts.io(), live, 0, false);
// Whole-database form.
try testing.expectEqual(@as(u32, 1), ts.store.kill_namespace(ts.io(), "t", null));
try testing.expectEqual(@as(u32, 0), ts.store.live);
}
test "a namespace too long for the fixed buffers declines rather than failing" {
var ts = try TestStore.init(2, default_idle_timeout_ms);
defer ts.deinit();
const long = "c" ** (ns_coll_max + 1);
try testing.expectError(error.NameTooLong, ts.open_one(long, 0));
}
test "batch_target: 101 for a first batch, unbounded for a getMore, honoured when given" {
try testing.expectEqual(@as(?u32, default_first_batch), batch_target(null, true));
// A bare getMore has no document target at all. Measured against mongod:
// it does not inherit the batchSize the cursor was created with.
try testing.expectEqual(@as(?u32, null), batch_target(null, false));
try testing.expectEqual(@as(?u32, 7), batch_target(7, true));
// batchSize: 0 is a real target of zero -- an empty first batch with a live
// cursor, which drivers use to obtain a cursor cheaply. It must NOT read as
// "unbounded": conflating the two returns the whole collection where mongod
// returns nothing, which is exactly the bug this optional prevents.
try testing.expectEqual(@as(?u32, 0), batch_target(0, true));
var zero = BatchBuilder.init(batch_target(0, true));
try testing.expect(zero.full());
var bare = BatchBuilder.init(batch_target(null, false));
try testing.expect(!bare.full());
}
test "BatchBuilder stops at its document target" {
var b = BatchBuilder.init(2);
try testing.expect(!b.full());
try testing.expectEqual(Offered.appended, b.offer(10));
try testing.expect(!b.full());
try testing.expectEqual(Offered.appended, b.offer(10));
try testing.expect(b.full());
try testing.expectEqual(@as(u32, 2), b.emitted);
}
test "BatchBuilder: a null target is unbounded by documents" {
var b = BatchBuilder.init(null);
for (0..5000) |_| {
try testing.expect(!b.full());
try testing.expectEqual(Offered.appended, b.offer(1));
}
try testing.expect(!b.full());
}
test "BatchBuilder stops on bytes, but always takes at least one document" {
// Hitting the byte cap must not read as EOF, or the cursor would close and
// silently drop the rest of the result.
var b = BatchBuilder.init(null);
b.bytes_max = 100;
try testing.expectEqual(Offered.appended, b.offer(60));
try testing.expectEqual(Offered.batch_full, b.offer(60));
// The refused document was not accounted for, so the caller can hand it to
// the next batch.
try testing.expectEqual(@as(u32, 1), b.emitted);
try testing.expectEqual(@as(u64, 60), b.bytes);
// An oversized document on an empty batch goes through anyway: refusing it
// would wedge the cursor, returning empty batches and never progressing.
var solo = BatchBuilder.init(null);
solo.bytes_max = 100;
try testing.expectEqual(Offered.appended, solo.offer(1_000_000));
try testing.expectEqual(@as(u32, 1), solo.emitted);
try testing.expectEqual(Offered.batch_full, solo.offer(1));
}
test "Stream.advance records the anchor and counts an equal-key band" {
var s = Stream{};
try testing.expect(!s.started());
s.advance("aaa", 10, 3, 4);
try testing.expect(s.started());
try testing.expectEqualStrings("aaa", s.anchor_key());
try testing.expectEqual(@as(u64, 10), s.anchor_off);
try testing.expectEqual(@as(u32, 3), s.hint_leaf);
try testing.expectEqual(@as(u32, 4), s.hint_slot);
try testing.expectEqual(@as(u64, 0), s.band_index);
// Same key, different document: still inside the band, so the position
// within it has to advance or a resume could not tell them apart.
s.advance("aaa", 11, 3, 5);
try testing.expectEqual(@as(u64, 1), s.band_index);
s.advance("aaa", 12, 3, 6);
try testing.expectEqual(@as(u64, 2), s.band_index);
// A new key ends the band.
s.advance("bbb", 13, 3, 7);
try testing.expectEqual(@as(u64, 0), s.band_index);
try testing.expectEqualStrings("bbb", s.anchor_key());
}
test "Stream.advance accepts a reverse cursor moving down" {
var s = Stream{ .backward = true };
s.advance("ccc", 1, 1, 5);
s.advance("bbb", 2, 1, 4);
s.advance("aaa", 3, 1, 3);
try testing.expectEqualStrings("aaa", s.anchor_key());
}

View File

@@ -22,6 +22,7 @@ const bson = @import("bson.zig");
const storage = @import("storage.zig"); const storage = @import("storage.zig");
const index = @import("index.zig"); const index = @import("index.zig");
const pgr = @import("pager.zig"); const pgr = @import("pager.zig");
const cursor = @import("cursor.zig");
// Always active, including in the default ReleaseFast build -- see assert.zig // Always active, including in the default ReleaseFast build -- see assert.zig
// for why std.debug.assert is the wrong tool for these invariants. // for why std.debug.assert is the wrong tool for these invariants.
const assert = @import("assert.zig").assert; const assert = @import("assert.zig").assert;
@@ -97,8 +98,21 @@ pub const Collection = struct {
/// replaces the old serialization-guarded docs-map fast path for /// replaces the old serialization-guarded docs-map fast path for
/// integer/string/etc. _id lookups. /// integer/string/etc. _id lookups.
id_index: index.Index, id_index: index.Index,
/// Identity-and-layout token for open cursors. Drawn from
/// `Engine.layout_epoch_seq`, so it is unique across the engine's life and
/// bumped again by every rebuild.
///
/// It answers two questions a cursor cannot answer any other way. A rebuild
/// moves every document, so a saved slab offset (or a saved index anchor's
/// offset) is stale -- and the keys surviving unchanged makes that *worse*,
/// because a lookup then succeeds and quietly resolves to the wrong bytes.
/// And a cursor holds namespace *strings*, not a `*Collection`, so a
/// drop-and-recreate under the same name would otherwise be invisible to it;
/// drawing from an engine-wide sequence rather than starting each collection
/// at zero is what makes the recreated one compare unequal.
layout_epoch: u64,
fn init(gpa: std.mem.Allocator, pager: *pgr.Pager) !Collection { fn init(gpa: std.mem.Allocator, pager: *pgr.Pager, layout_epoch: u64) !Collection {
var self: Collection = .{ var self: Collection = .{
.doc_count = 0, .doc_count = 0,
.pager = pager, .pager = pager,
@@ -110,6 +124,7 @@ pub const Collection = struct {
.hold = .{}, .hold = .{},
.indexes = .empty, .indexes = .empty,
.id_index = undefined, .id_index = undefined,
.layout_epoch = layout_epoch,
}; };
const keys = [_]index.IndexKey{.{ .path = "_id", .descending = false }}; const keys = [_]index.IndexKey{.{ .path = "_id", .descending = false }};
// unique: the tree, not the docs map, is what enforces _id uniqueness // unique: the tree, not the docs map, is what enforces _id uniqueness
@@ -290,6 +305,15 @@ pub const Engine = struct {
/// rewrite is worth doing — see `note_compact`. /// rewrite is worth doing — see `note_compact`.
live_docs: u64 = 0, live_docs: u64 = 0,
dead_docs: u64 = 0, dead_docs: u64 = 0,
/// Hands out `Collection.layout_epoch` values. Monotonic and never reset, so
/// no two collection instances -- including a drop followed by a recreate
/// under the same name -- ever share one.
layout_epoch_seq: u64 = 0,
/// Open cursors. Lives on the engine rather than the server because the C
/// API seam (PLAN D1) lists cursor iteration, and because the unit tests
/// build an Engine with no server at all. Its mutex is a leaf: see
/// `cursor.Store`.
cursors: cursor.Store,
/// The same question in bytes, about the *data file* rather than the log. /// The same question in bytes, about the *data file* rather than the log.
/// Once a checkpoint truncates the log, the log no longer holds the garbage /// Once a checkpoint truncates the log, the log no longer holds the garbage
/// -- the doc slab does, and only a rebuild reclaims it. These are what /// -- the doc slab does, and only a rebuild reclaims it. These are what
@@ -324,6 +348,12 @@ pub const Engine = struct {
/// command reads it while still holding the write lock. /// command reads it while still holding the write lock.
dup_index: ?[]const u8 = null, dup_index: ?[]const u8 = null,
/// The registry an embedded caller gets without configuring anything; the
/// CLI replaces it through `reconfigure_cursors`.
fn default_cursor_store(gpa: std.mem.Allocator, io: std.Io) !cursor.Store {
return cursor.Store.init(gpa, io, cursor.default_capacity, cursor.default_idle_timeout_ms);
}
pub fn open(gpa: std.mem.Allocator, io: std.Io, path: []const u8) !Engine { pub fn open(gpa: std.mem.Allocator, io: std.Io, path: []const u8) !Engine {
var log = try storage.Log.open(gpa, io, path); var log = try storage.Log.open(gpa, io, path);
errdefer log.close(); errdefer log.close();
@@ -346,8 +376,10 @@ pub const Engine = struct {
.dbs = .empty, .dbs = .empty,
.seq = 0, .seq = 0,
.compact_threshold = 16 * 1024 * 1024, .compact_threshold = 16 * 1024 * 1024,
.cursors = try default_cursor_store(gpa, io),
}; };
errdefer { errdefer {
engine.cursors.deinit();
engine.pager.deinit(); engine.pager.deinit();
engine.dbs.deinit(gpa); engine.dbs.deinit(gpa);
} }
@@ -393,6 +425,17 @@ pub const Engine = struct {
return engine; return engine;
} }
/// Replace the cursor registry with one of a different shape. Only legal
/// before the server starts accepting connections, because it drops every
/// cursor -- asserted rather than left to the comment, since the method is
/// public and a later caller would otherwise get silent data loss.
pub fn reconfigure_cursors(self: *Engine, capacity: u32, idle_timeout_ms: i64) !void {
assert_msg(self.cursors.live == 0, "reconfigured the cursor registry with cursors open");
const fresh = try cursor.Store.init(self.gpa, self.io, capacity, idle_timeout_ms);
self.cursors.deinit();
self.cursors = fresh;
}
pub fn deinit(self: *Engine) void { pub fn deinit(self: *Engine) void {
var db_it = self.dbs.iterator(); var db_it = self.dbs.iterator();
while (db_it.next()) |db_entry| { while (db_it.next()) |db_entry| {
@@ -400,6 +443,11 @@ pub const Engine = struct {
self.gpa.free(db_entry.key_ptr.*); self.gpa.free(db_entry.key_ptr.*);
} }
self.dbs.deinit(self.gpa); self.dbs.deinit(self.gpa);
// Before the pager: a cursor's arena is its own, but freeing cursors
// first keeps the teardown order the same as the construction order
// reversed, which is the only order that stays obviously correct as
// cursors grow to hold more.
self.cursors.deinit();
self.pager.deinit(); self.pager.deinit();
self.gpa.destroy(self.pager); self.gpa.destroy(self.pager);
self.log.close(); self.log.close();
@@ -936,6 +984,11 @@ pub const Engine = struct {
const removed = db.collections.fetchRemove(coll_name) orelse return false; const removed = db.collections.fetchRemove(coll_name) orelse return false;
self.free_collection(removed.value); self.free_collection(removed.value);
self.gpa.free(removed.key); self.gpa.free(removed.key);
// A cursor on this namespace is already safe -- it holds names, so its
// next getMore finds nothing to lock -- but reaping here frees the slots
// now instead of at the idle timeout, and keeps the open-cursor metric
// describing cursors that can still return something.
_ = self.cursors.kill_namespace(self.io, db_name, coll_name);
return true; return true;
} }
@@ -943,6 +996,7 @@ pub const Engine = struct {
var removed = self.dbs.fetchRemove(db_name) orelse return false; var removed = self.dbs.fetchRemove(db_name) orelse return false;
self.free_db(&removed.value); self.free_db(&removed.value);
self.gpa.free(removed.key); self.gpa.free(removed.key);
_ = self.cursors.kill_namespace(self.io, db_name, null);
return true; return true;
} }
@@ -1156,7 +1210,8 @@ pub const Engine = struct {
errdefer self.gpa.free(coll_key); errdefer self.gpa.free(coll_key);
const new_coll = try self.gpa.create(Collection); const new_coll = try self.gpa.create(Collection);
errdefer self.gpa.destroy(new_coll); errdefer self.gpa.destroy(new_coll);
new_coll.* = try Collection.init(self.gpa, self.pager); self.layout_epoch_seq += 1;
new_coll.* = try Collection.init(self.gpa, self.pager, self.layout_epoch_seq);
errdefer new_coll.id_index.deinit(self.gpa); errdefer new_coll.id_index.deinit(self.gpa);
try db.collections.put(self.gpa, coll_key, new_coll); try db.collections.put(self.gpa, coll_key, new_coll);
return new_coll; return new_coll;
@@ -1389,6 +1444,13 @@ pub const Engine = struct {
for (coll.indexes.items) |ix| try self.repack_index(coll, ix, moved.items); for (coll.indexes.items) |ix| try self.repack_index(coll, ix, moved.items);
for (old_extents) |e| try self.pager.free_pages(e.first, e.pages); for (old_extents) |e| try self.pager.free_pages(e.first, e.pages);
// Every document has moved, so every offset an open cursor is holding
// now names different bytes. Bumped last, after the rebuild can no
// longer fail: a cursor invalidated by a rebuild that then errored out
// would have been invalidated for nothing.
self.layout_epoch_seq += 1;
coll.layout_epoch = self.layout_epoch_seq;
} }
fn repack_index( fn repack_index(
@@ -3886,3 +3948,64 @@ const Reader = struct {
return self.take(n); return self.take(n);
} }
}; };
test "the epochs that invalidate a cursor move exactly when they must" {
// Three separate promises, each one load-bearing for an open cursor:
//
// - a rebuild moves every document, so a saved slab offset is stale;
// - a drop-and-recreate under the same name is a different collection,
// which a cursor holding only namespace strings cannot otherwise see;
// - `Index.reset_tree` re-creates node ids 0 and 1 as different nodes, so a
// saved (leaf, slot) position becomes valid-and-wrong rather than absent.
//
// A cursor's whole safety story is these three bumps, so assert them here
// rather than inferring them from cursor behaviour later.
var threaded: std.Io.Threaded = .init_single_threaded;
defer threaded.deinit();
var env = test_env(&threaded);
const io = env.io;
const gpa = testing.allocator;
var tmp = try TmpLog.init(gpa);
defer tmp.deinit(gpa);
var engine = try Engine.open(gpa, io, tmp.path);
defer engine.deinit();
try engine.lock();
var i: i32 = 0;
while (i < 40) : (i += 1) {
var d = try make_doc(gpa, i, "payload");
defer d.deinit();
try engine.insert("app", "c", &d, &env.gen);
}
const before = engine.get_collection("app", "c").?.layout_epoch;
engine.unlock();
try testing.expect(before != 0);
// A rebuild moves documents, so the epoch must move with them.
try engine.compact();
try engine.lock();
const after_rebuild = engine.get_collection("app", "c").?.layout_epoch;
engine.unlock();
try testing.expect(after_rebuild != before);
// A recreated collection must not be mistaken for the one that was
// dropped. Starting each collection's epoch at zero would fail here.
try engine.lock();
try testing.expect(try engine.drop_collection("app", "c"));
var fresh_doc = try make_doc(gpa, 1, "fresh");
defer fresh_doc.deinit();
try engine.insert("app", "c", &fresh_doc, &env.gen);
const after_recreate = engine.get_collection("app", "c").?.layout_epoch;
engine.unlock();
try testing.expect(after_recreate != after_rebuild);
try testing.expect(after_recreate != before);
// And the index-level token, which guards the position hint.
try engine.lock();
const coll = engine.get_collection("app", "c").?;
const index_before = coll.id_index.epoch;
try coll.id_index.reset_tree(gpa);
try testing.expect(coll.id_index.epoch != index_before);
engine.unlock();
}

View File

@@ -106,7 +106,9 @@ const page_size = 4096;
/// Bytes of node payload: a 32-byte header plus the slotted region. /// Bytes of node payload: a 32-byte header plus the slotted region.
const page_data = page_size - 32; const page_data = page_size - 32;
/// Records longer than a quarter of a node spill to the overflow slab. /// Records longer than a quarter of a node spill to the overflow slab.
const inline_limit = page_size / 4; /// Public because it is also the bound a resumable cursor anchors within: a key
/// past it has already spilled, so the tree itself treats it as exceptional.
pub const inline_limit = page_size / 4;
/// Upper bound on the slots one node can hold, since every slot costs at /// Upper bound on the slots one node can hold, since every slot costs at
/// least its own size. Bounds the split scratch. /// least its own size. Bounds the split scratch.
const max_slots = page_data / slot_size; const max_slots = page_data / slot_size;
@@ -234,6 +236,19 @@ pub const Index = struct {
depth: u32, depth: u32,
/// Total entries, maintained incrementally. /// Total entries, maintained incrementally.
entry_count: usize, entry_count: usize,
/// Bumped whenever a node id stops meaning what it meant, which is the one
/// thing that makes a saved `(leaf, slot)` position dangerous rather than
/// merely stale. Node ids are otherwise append-only (`alloc_node`, and
/// `drop_child` abandons a page without recycling its id), and `page()`
/// resolves ids through `node_pages`, so copy-on-write and checkpoints move
/// pages without disturbing ids. Only `reset_tree` and
/// `replace_root_with_leaf` reuse an id for different contents.
///
/// A resumable cursor keeps a position hint to avoid walking an equal-key
/// band on every `getMore`; it must compare this first. Without it the hint
/// would address a live but unrelated leaf after a compaction and the cursor
/// would iterate a tree that no longer exists.
epoch: u64,
/// Repack scratch: any single node's record bytes fit here. /// Repack scratch: any single node's record bytes fit here.
scratch: [page_data]u8, scratch: [page_data]u8,
/// Promoted-key scratch: inline keys being propagated up a split are /// Promoted-key scratch: inline keys being propagated up a split are
@@ -268,6 +283,7 @@ pub const Index = struct {
.leaf_count = 0, .leaf_count = 0,
.depth = 0, .depth = 0,
.entry_count = 0, .entry_count = 0,
.epoch = 0,
.scratch = undefined, .scratch = undefined,
.promo = undefined, .promo = undefined,
}; };
@@ -614,6 +630,10 @@ pub const Index = struct {
self.depth = 0; self.depth = 0;
self.entry_count = 0; self.entry_count = 0;
self.multikey = false; self.multikey = false;
// Node ids 0 and 1 were just re-created as different nodes, so every
// position anyone saved into the old tree now points somewhere valid
// and wrong. This is the bump that tells them apart.
self.epoch += 1;
} }
/// Remove every entry for `id`, in one pass over the leaves. Infallible. /// Remove every entry for `id`, in one pass over the leaves. Infallible.
@@ -777,6 +797,14 @@ pub const Index = struct {
leaf: u32, leaf: u32,
slot: u32, slot: u32,
/// `next`, plus the position of the entry it yielded. Exact because
/// `next` leaves `leaf` alone on the call that yields and has already
/// incremented `slot` past the entry.
pub fn positioned(self: *Iter) ?Positioned {
const e = self.next() orelse return null;
return .{ .key = e.key, .off = e.off, .leaf = self.leaf, .slot = self.slot - 1 };
}
pub fn next(self: *Iter) ?EntryRef { pub fn next(self: *Iter) ?EntryRef {
const ix = self.ix; const ix = self.ix;
while (self.leaf != 0) { while (self.leaf != 0) {
@@ -872,6 +900,13 @@ pub const Index = struct {
/// One past the slot to yield next, so 0 means this leaf is done. /// One past the slot to yield next, so 0 means this leaf is done.
slot: u32, slot: u32,
/// As `Iter.positioned`, but `RevIter.next` decrements *onto* the entry
/// it yields, so the slot needs no adjustment.
pub fn positioned(self: *RevIter) ?Positioned {
const e = self.next() orelse return null;
return .{ .key = e.key, .off = e.off, .leaf = self.leaf, .slot = self.slot };
}
pub fn next(self: *RevIter) ?EntryRef { pub fn next(self: *RevIter) ?EntryRef {
const ix = self.ix; const ix = self.ix;
while (self.leaf != 0) { while (self.leaf != 0) {
@@ -919,6 +954,174 @@ pub const Index = struct {
return .{ .ix = self, .leaf = b.leaf, .slot = b.slot }; return .{ .ix = self, .leaf = b.leaf, .slot = b.slot };
} }
// -- resuming an interrupted scan ---------------------------------------
/// Entries a resume will walk past before giving up and reporting `capped`.
/// A bound rather than a hope: `seek` lands at the *start* of an equal-key
/// band, so without one a key with millions of duplicates would make every
/// batch cost O(band) and a full drain quadratic.
pub const resume_walk_max: u32 = 1 << 16;
/// One entry, with enough of its position to resume after it next time.
pub const Positioned = struct {
key: []const u8,
off: u64,
leaf: u32,
slot: u32,
};
/// A resumed forward walk. `capped` means the anchor could not be located
/// within `resume_walk_max` steps, so the position is not trustworthy and
/// the caller must fail rather than return documents from the wrong place.
pub const Resumed = struct { it: Iter, capped: bool = false };
pub const ResumedRev = struct { it: RevIter, capped: bool = false };
/// Does `(leaf, slot)` still hold exactly `(key, off)`?
///
/// A hint is never believed, only checked, and the checks are ordered so the
/// cheap structural ones run first: `off_of` asserts `is_leaf` with
/// `std.debug.assert`, which in ReleaseFast is a promise to the optimizer
/// rather than a check, so `is_leaf` must be tested for real beforehand.
///
/// Node ids are append-only, so a stale id is always in bounds; what makes a
/// hint dangerous rather than merely wrong is `reset_tree` re-creating ids 0
/// and 1 as different nodes, and `Index.epoch` is what the caller compares
/// for that.
fn hint_holds(self: *const Index, leaf: u32, slot: u32, key: []const u8, off: u64) bool {
if (leaf == 0 or leaf >= self.node_pages.items.len) return false;
const node = self.page(leaf);
if (node.is_leaf != 1) return false;
if (slot >= node.count) return false;
if (self.off_of(leaf, slot) != off) return false;
return std.mem.eql(u8, self.key_of(leaf, slot), key);
}
/// Locate the anchor `(key, off)` by walking its equal-key band.
///
/// Comparison is `std.mem.order`, not `cmp_prefix`, because the band is
/// defined as the entries whose key is byte-equal to the anchor's and prefix
/// semantics would call `"ab"` and `"abc"` equal. In fairness the two happen
/// to agree on where this function resumes -- the fallback is positional, and
/// the first out-of-band entry is the same entry either way -- so this is a
/// clarity choice, not a bug fix; an attempted mutation to `cmp_prefix` does
/// not change any observable result. What does matter is that `lower_bound`
/// uses prefix semantics and therefore errs *before* the band, never past it,
/// so the walk cannot start beyond the anchor and skip it.
///
/// When the anchor is gone, "gone" turns out to mean two different things and
/// they want opposite answers:
///
/// - **Deleted.** A sibling has moved up into the anchor's band position, and
/// that sibling has not been returned yet. Resume *at* band position
/// `band_index`. Resuming after the whole band instead would silently drop
/// every remaining member, which on a three-value index is most of the
/// collection.
/// - **Updated.** The document was rewritten, so its key is unchanged but its
/// offset moved. The entry at the anchor's band position *is* the anchor,
/// already returned. Resume *after* it.
///
/// The index cannot tell these apart in general -- both look like "same key,
/// different offset". On a **unique** index it can: two entries cannot share a
/// key, so a same-key entry is necessarily the same document, hence the update
/// case, hence resume past the band. That covers `_id_` and so every unsorted
/// scan and every `_id` sort, which is where an update-during-drain otherwise
/// returns a document twice -- observed as duplicate `_id`s draining a
/// collection that was being updated underneath.
///
/// On a non-unique index the positional fallback stands, so an updated document
/// may come back a second time. That is legal: MongoDB documents that a
/// non-snapshot cursor may return a document more than once if an intervening
/// write moves it.
fn band_resume(self: *const Index, key: []const u8, off: u64, band_index: u64) Resumed {
var it = self.seek(key);
var fallback: ?Iter = null;
var pos: u64 = 0;
var steps: u32 = 0;
while (steps < resume_walk_max) : (steps += 1) {
// The iterator state that would yield the entry we are about to
// look at, i.e. "resume *at* this entry".
const before = it;
const e = it.next() orelse break;
if (std.mem.order(u8, e.key, key) != .eq) {
// Past the band. Prefer the fallback if the band held one.
return .{ .it = fallback orelse before };
}
if (e.off == off) return .{ .it = it }; // resume just after the anchor
// On a unique index a same-key entry can only be the anchor itself,
// rewritten, so there is no sibling to fall back to.
if (!self.unique and pos == band_index and fallback == null) fallback = before;
pos += 1;
}
if (steps == resume_walk_max) return .{ .it = it, .capped = true };
return .{ .it = fallback orelse it };
}
/// A forward walk positioned just after `(key, off)`.
///
/// O(1) whenever the hint still holds, which is the case unless something
/// wrote to that exact leaf between batches. The band walk is the fallback,
/// and it is what the walk bound exists to contain.
pub fn resume_forward(
self: *const Index,
key: []const u8,
off: u64,
band_index: u64,
hint_leaf: u32,
hint_slot: u32,
hint_trusted: bool,
) Resumed {
if (hint_trusted and self.hint_holds(hint_leaf, hint_slot, key, off)) {
return .{ .it = .{ .ix = self, .leaf = hint_leaf, .slot = hint_slot + 1 } };
}
return self.band_resume(key, off, band_index);
}
/// A reverse walk positioned just before `(key, off)` in key order, i.e. at
/// the next entry a descending scan owes.
///
/// `RevIter` decrements before yielding, so slot `s` yields `s - 1` -- the
/// entry immediately below the anchor -- and crosses into `prev` when the
/// anchor sat at slot 0.
///
/// Known limitation, and it is a deliberate trade. When the anchor is gone
/// *and* it had duplicates, this resumes below the whole band rather than at
/// the anchor's position within it, so the band's remaining members are not
/// returned. Placing a reverse fallback exactly would need the band's length,
/// which is only known after walking it, hence a second walk on a path that
/// requires a descending scan over a duplicate-heavy index whose anchor was
/// deleted mid-cursor. Forward resumes -- every unsorted scan and every
/// ascending sort -- use `band_index` and have no such gap.
pub fn resume_reverse(
self: *const Index,
key: []const u8,
off: u64,
hint_leaf: u32,
hint_slot: u32,
hint_trusted: bool,
) ResumedRev {
if (hint_trusted and self.hint_holds(hint_leaf, hint_slot, key, off)) {
return .{ .it = .{ .ix = self, .leaf = hint_leaf, .slot = hint_slot } };
}
// Find the anchor by walking forward, then turn around on it.
var it = self.seek(key);
var steps: u32 = 0;
while (steps < resume_walk_max) : (steps += 1) {
const e = it.next() orelse break;
if (std.mem.order(u8, e.key, key) != .eq) break; // past the band
if (e.off == off) {
// `it` has already stepped past the anchor, so the anchor sat at
// `it.slot - 1` and a RevIter there yields the entry below it.
return .{ .it = .{ .ix = self, .leaf = it.leaf, .slot = it.slot - 1 } };
}
}
if (steps == resume_walk_max) {
return .{ .it = .{ .ix = self, .leaf = 0, .slot = 0 }, .capped = true };
}
// Anchor gone: resume below the band.
const b = self.lower_bound(key);
return .{ .it = .{ .ix = self, .leaf = b.leaf, .slot = b.slot } };
}
// -- serialization ------------------------------------------------------ // -- serialization ------------------------------------------------------
/// The canonical spec document bytes /// The canonical spec document bytes
@@ -1710,6 +1913,9 @@ pub const Index = struct {
self.first_leaf = self.root; self.first_leaf = self.root;
self.leaf_count = 1; self.leaf_count = 1;
self.depth = 0; self.depth = 0;
// The root's id is unchanged but it is a leaf now, so a saved position
// that named it as an internal node describes a different tree shape.
self.epoch += 1;
} }
/// Pack the sorted staging array into a fresh tree: leaves filled in /// Pack the sorted staging array into a fresh tree: leaves filled in
@@ -3672,3 +3878,258 @@ test "planner picks eq run, ranges, and bails on sparse null" {
try testing.expect((try plan(gpa, null, &.{&ix}, &f, &.{})) == null); try testing.expect((try plan(gpa, null, &.{&ix}, &f, &.{})) == null);
} }
} }
/// Drain `ix` by resuming every `stride` entries, the way a cursor with that
/// batch size would, and return the offsets in the order they came out.
fn drain_resuming(
gpa: std.mem.Allocator,
ix: *const Index,
stride: u32,
backward: bool,
out: *std.ArrayListUnmanaged(u64),
) !void {
var started = false;
var anchor: std.ArrayListUnmanaged(u8) = .empty;
defer anchor.deinit(gpa);
var anchor_off: u64 = 0;
var band_index: u64 = 0;
var hint_leaf: u32 = 0;
var hint_slot: u32 = 0;
while (true) {
// One "batch": open a walk where the last one stopped.
var fwd: Index.Iter = undefined;
var rev: Index.RevIter = undefined;
if (!started) {
if (backward) rev = ix.iter_reverse() else fwd = ix.iter();
} else if (backward) {
const r = ix.resume_reverse(anchor.items, anchor_off, hint_leaf, hint_slot, true);
try testing.expect(!r.capped);
rev = r.it;
} else {
const r = ix.resume_forward(
anchor.items,
anchor_off,
band_index,
hint_leaf,
hint_slot,
true,
);
try testing.expect(!r.capped);
fwd = r.it;
}
var n: u32 = 0;
while (n < stride) : (n += 1) {
const e = if (backward)
rev.positioned()
else
fwd.positioned();
const got = e orelse return;
try out.append(gpa, got.off);
if (started and std.mem.eql(u8, anchor.items, got.key)) {
band_index += 1;
} else {
band_index = 0;
}
anchor.clearRetainingCapacity();
try anchor.appendSlice(gpa, got.key);
anchor_off = got.off;
hint_leaf = got.leaf;
hint_slot = got.slot;
started = true;
}
}
}
test "a resumed walk yields exactly what an uninterrupted one does" {
// The property the whole streaming cursor rests on: stopping and restarting
// a scan changes nothing about what it returns, in either direction, at any
// batch size, including one entry at a time.
//
// Mutation-checked: changing `resume_forward`'s `hint_slot + 1` to
// `hint_slot` makes every batch boundary repeat an entry, and this test goes
// red. (A third mutation was tried and rejected as meaningless: swapping
// `std.mem.order` for `cmp_prefix` inside `band_resume` changes nothing
// observable, so no test can catch it -- see the note there.)
//
// Note this test always resumes from a *valid* hint, since nothing mutates
// the tree between its batches. The band walk is covered by the two tests
// below, which invalidate the hint on purpose.
const gpa = testing.allocator;
// Three corpora, each hard for a different reason: distinct keys spanning
// several leaves and an interior level; a low-cardinality index whose bands
// span leaves; and *variable-length string keys in prefix relationships*
// ("a" < "ab" < "abc"), which is the only shape that can tell `std.mem.order`
// apart from `cmp_prefix` -- fixed-width integer keys never differ, so an
// integer-only corpus cannot catch that mistake at all.
const Shape = enum { distinct, duplicates, prefixes };
for ([_]Shape{ .distinct, .duplicates, .prefixes }) |shape| {
var ix = try simple_index(gpa, test_pager(), &.{"a"}, false, false);
defer ix.deinit(gpa);
const n = 400;
var key_buf: [40]u8 = undefined;
for (0..n) |i| {
const value: bson.Value = switch (shape) {
.distinct => .{ .int32 = @intCast(i + 1) },
.duplicates => .{ .int32 = @intCast(i % 3) },
// Every key is a prefix of the next in its group of eight, so
// each band start is also a proper prefix of later keys.
.prefixes => blk: {
const written = try std.fmt.bufPrint(&key_buf, "k{d}", .{i / 8});
const depth = (i % 8) + 1;
@memset(key_buf[written.len .. written.len + depth], 'x');
break :blk .{ .string = key_buf[0 .. written.len + depth] };
},
};
const d = try bytes_of(gpa, &.{
.{ .key = "_id", .value = .{ .int32 = @intCast(i + 1) } },
.{ .key = "a", .value = value },
});
defer gpa.free(d);
_ = try ix.add_doc(gpa, d, @intCast(i + 1), false);
}
try testing.expect(ix.depth >= 1);
for ([_]bool{ false, true }) |backward| {
var whole: std.ArrayListUnmanaged(u64) = .empty;
defer whole.deinit(gpa);
if (backward) {
var it = ix.iter_reverse();
while (it.next()) |e| try whole.append(gpa, e.off);
} else {
var it = ix.iter();
while (it.next()) |e| try whole.append(gpa, e.off);
}
try testing.expectEqual(@as(usize, n), whole.items.len);
for ([_]u32{ 1, 2, 7, 101, 399, 400, 1000 }) |stride| {
var resumed: std.ArrayListUnmanaged(u64) = .empty;
defer resumed.deinit(gpa);
try drain_resuming(gpa, &ix, stride, backward, &resumed);
try testing.expectEqualSlices(u64, whole.items, resumed.items);
}
}
}
}
test "a resume survives a split between batches" {
// A cursor holds no lock, so the tree it comes back to is not the tree it
// left. Inserting mid-drain moves entries between leaves and invalidates the
// position hint, which is exactly what the anchor is for.
const gpa = testing.allocator;
var ix = try simple_index(gpa, test_pager(), &.{"a"}, false, false);
defer ix.deinit(gpa);
const n = 200;
for (0..n) |i| {
// Even keys only, so the inserts below land between existing entries.
const d = try bytes_of(gpa, &.{
.{ .key = "_id", .value = .{ .int32 = @intCast(i + 1) } },
.{ .key = "a", .value = .{ .int32 = @intCast((i + 1) * 2) } },
});
defer gpa.free(d);
_ = try ix.add_doc(gpa, d, @intCast(i + 1), false);
}
var seen: std.ArrayListUnmanaged(u64) = .empty;
defer seen.deinit(gpa);
var anchor: std.ArrayListUnmanaged(u8) = .empty;
defer anchor.deinit(gpa);
var anchor_off: u64 = 0;
var hint_leaf: u32 = 0;
var hint_slot: u32 = 0;
var started = false;
var next_id: i32 = 10_000;
while (true) {
var it = if (!started) ix.iter() else blk: {
const r = ix.resume_forward(anchor.items, anchor_off, 0, hint_leaf, hint_slot, true);
try testing.expect(!r.capped);
break :blk r.it;
};
var n_in_batch: u32 = 0;
while (n_in_batch < 5) : (n_in_batch += 1) {
const got = it.positioned() orelse break;
try seen.append(gpa, got.off);
anchor.clearRetainingCapacity();
try anchor.appendSlice(gpa, got.key);
anchor_off = got.off;
hint_leaf = got.leaf;
hint_slot = got.slot;
started = true;
}
if (n_in_batch < 5) break;
// Between batches, insert odd keys across the whole range: guaranteed to
// split leaves and to appear both before and after the anchor.
for (0..20) |k| {
next_id += 1;
const d = try bytes_of(gpa, &.{
.{ .key = "_id", .value = .{ .int32 = next_id } },
.{ .key = "a", .value = .{ .int32 = @intCast(k * 19 + 1) } },
});
defer gpa.free(d);
_ = try ix.add_doc(gpa, d, @intCast(next_id), false);
}
}
// The 200 originals must each appear exactly once. Entries inserted behind
// the cursor may or may not show up -- that is ordinary non-snapshot cursor
// behaviour -- but nothing may be duplicated or lost.
var originals: u32 = 0;
var counts = std.AutoHashMap(u64, u32).init(gpa);
defer counts.deinit();
for (seen.items) |off| {
const e = try counts.getOrPutValue(off, 0);
e.value_ptr.* += 1;
try testing.expectEqual(@as(u32, 1), e.value_ptr.*); // no duplicates
if (off <= n) originals += 1;
}
try testing.expectEqual(@as(u32, n), originals);
}
test "a resume whose anchor was deleted keeps the rest of its band" {
// The failure this guards against is silent and large: with the anchor gone,
// resuming after the whole equal-key band drops every remaining member, and
// on a low-cardinality index that is most of the collection.
//
// Mutation-checked: `pos == band_index + 1` in `band_resume` shifts the
// resume by one entry and this test goes red.
const gpa = testing.allocator;
var ix = try simple_index(gpa, test_pager(), &.{"a"}, false, false);
defer ix.deinit(gpa);
// One key, 50 documents: a single band.
for (0..50) |i| {
const d = try bytes_of(gpa, &.{
.{ .key = "_id", .value = .{ .int32 = @intCast(i + 1) } },
.{ .key = "a", .value = .{ .int32 = 7 } },
});
defer gpa.free(d);
_ = try ix.add_doc(gpa, d, @intCast(i + 1), false);
}
// Yield three, then delete the third -- the anchor itself.
var it = ix.iter();
var third: Index.Positioned = undefined;
var band_index: u64 = 0;
for (0..3) |i| {
third = it.positioned().?;
if (i > 0) band_index += 1;
}
const anchor_key = try gpa.dupe(u8, third.key);
defer gpa.free(anchor_key);
ix.remove_off(third.off);
const r = ix.resume_forward(anchor_key, third.off, band_index, third.leaf, third.slot, true);
try testing.expect(!r.capped);
var rest: u32 = 0;
var walk = r.it;
while (walk.next()) |_| rest += 1;
// 50 inserted, 1 deleted, 2 already returned before the anchor: 47 left.
try testing.expectEqual(@as(u32, 47), rest);
}

View File

@@ -10,6 +10,7 @@ pub const db = @import("db.zig");
pub const query = @import("query.zig"); pub const query = @import("query.zig");
pub const update = @import("update.zig"); pub const update = @import("update.zig");
pub const index = @import("index.zig"); pub const index = @import("index.zig");
pub const cursor = @import("cursor.zig");
pub const pager = @import("pager.zig"); pub const pager = @import("pager.zig");
test { test {
@@ -23,5 +24,6 @@ test {
_ = @import("query.zig"); _ = @import("query.zig");
_ = @import("update.zig"); _ = @import("update.zig");
_ = @import("index.zig"); _ = @import("index.zig");
_ = @import("cursor.zig");
_ = @import("pager.zig"); _ = @import("pager.zig");
} }

View File

@@ -10,6 +10,18 @@ const usage =
\\ --db <path> database file (default multiforadb.log) \\ --db <path> database file (default multiforadb.log)
\\ --ttl-sweep-secs <n> \\ --ttl-sweep-secs <n>
\\ seconds between TTL index sweeps (default 60, 0 disables) \\ seconds between TTL index sweeps (default 60, 0 disables)
\\ --cursor-timeout-ms <n>
\\ idle milliseconds before an open cursor is reaped
\\ (default 600000, matching MongoDB's cursorTimeoutMillis;
\\ 0 disables expiry)
\\ --cursor-sweep-secs <n>
\\ seconds between idle-cursor sweeps (default 4, matching
\\ MongoDB's clientCursorMonitorFrequencySecs; 0 disables)
\\ --max-open-cursors <n>
\\ cursor registry capacity (default 4096). At capacity the
\\ least recently used cursor is evicted, and its client
\\ sees CursorNotFound -- the same answer an idle timeout
\\ gives, which every driver already handles.
\\ --compact-threshold <bytes> \\ --compact-threshold <bytes>
\\ minimum log bytes between compactions; suffixes k/m/g \\ minimum log bytes between compactions; suffixes k/m/g
\\ (default 16m). The actual trigger also scales with the \\ (default 16m). The actual trigger also scales with the
@@ -44,34 +56,77 @@ fn parse_size_suffix(v: []const u8) ?u64 {
return n * mult; return n * mult;
} }
pub fn main(init: std.process.Init) !void { /// Everything the CLI can set. Parsed apart from `main` so the option table has
var port: u16 = 27017; /// room to grow without main outgrowing the 70-line limit.
var bind_ip: []const u8 = "127.0.0.1"; const Options = struct {
var db_path: []const u8 = "multiforadb.log"; port: u16 = 27017,
var ttl_sweep_secs: i64 = 60; bind_ip: []const u8 = "127.0.0.1",
var compact_threshold: u64 = 16 * 1024 * 1024; db_path: []const u8 = "multiforadb.log",
ttl_sweep_secs: i64 = 60,
compact_threshold: u64 = 16 * 1024 * 1024,
cursor_timeout_ms: i64 = mongo.cursor.default_idle_timeout_ms,
cursor_sweep_secs: i64 = 4,
max_open_cursors: u32 = mongo.cursor.default_capacity,
/// Set when --help was given: print usage and exit without opening anything.
help: bool = false,
};
pub fn main(init: std.process.Init) !void {
const opts = try parse_args(init) orelse {
try std.Io.File.writeStreamingAll(.stdout(), init.io, usage);
return;
};
const oid_gen = mongo.bson.ObjectIdGen.init(init.io);
var engine = try mongo.db.Engine.open(init.gpa, init.io, opts.db_path);
defer engine.deinit();
engine.compact_threshold = opts.compact_threshold;
// The engine builds its registry with the defaults so an embedded caller
// needs no configuration; the CLI replaces it when asked for something else.
try engine.reconfigure_cursors(opts.max_open_cursors, opts.cursor_timeout_ms);
std.debug.print("multiforadb: opened database '{s}' (compact threshold {d})\n", .{
opts.db_path,
opts.compact_threshold,
});
var server = mongo.server.Server{
.gpa = init.gpa,
.port = opts.port,
.bind_ip = opts.bind_ip,
.oid_gen = oid_gen,
.connection_counter = .init(1),
.engine = &engine,
.start_time = std.Io.Timestamp.now(init.io, .real),
.ttl_sweep_secs = opts.ttl_sweep_secs,
.cursor_sweep_secs = opts.cursor_sweep_secs,
};
try server.run();
}
/// Null means --help: the caller prints usage and exits.
fn parse_args(init: std.process.Init) !?Options {
var o = Options{};
var it = std.process.Args.Iterator.init(init.minimal.args); var it = std.process.Args.Iterator.init(init.minimal.args);
defer it.deinit(); defer it.deinit();
_ = it.next(); // program name _ = it.next(); // program name
while (it.next()) |arg| { while (it.next()) |arg| {
if (std.mem.eql(u8, arg, "--port")) { if (std.mem.eql(u8, arg, "--port")) {
const v = it.next() orelse return error.MissingValue; const v = it.next() orelse return error.MissingValue;
port = std.fmt.parseInt(u16, v, 10) catch { o.port = std.fmt.parseInt(u16, v, 10) catch {
std.debug.print("multiforadb: invalid port '{s}'\n", .{v}); std.debug.print("multiforadb: invalid port '{s}'\n", .{v});
return error.InvalidPort; return error.InvalidPort;
}; };
} else if (std.mem.eql(u8, arg, "--bind")) { } else if (std.mem.eql(u8, arg, "--bind")) {
bind_ip = it.next() orelse return error.MissingValue; o.bind_ip = it.next() orelse return error.MissingValue;
} else if (std.mem.eql(u8, arg, "--db")) { } else if (std.mem.eql(u8, arg, "--db")) {
db_path = it.next() orelse return error.MissingValue; o.db_path = it.next() orelse return error.MissingValue;
} else if (std.mem.eql(u8, arg, "--ttl-sweep-secs")) { } else if (std.mem.eql(u8, arg, "--ttl-sweep-secs")) {
const v = it.next() orelse return error.MissingValue; const v = it.next() orelse return error.MissingValue;
// i64 is the width std.Io.Duration.fromSeconds takes, so the // i64 is the width std.Io.Duration.fromSeconds takes, so the value
// value reaches the sweeper without a cast; negatives are the // reaches the sweeper without a cast; negatives are the only thing
// only thing parseInt would otherwise let through. // parseInt would otherwise let through.
ttl_sweep_secs = std.fmt.parseInt(i64, v, 10) catch -1; o.ttl_sweep_secs = std.fmt.parseInt(i64, v, 10) catch -1;
if (ttl_sweep_secs < 0) { if (o.ttl_sweep_secs < 0) {
std.debug.print("multiforadb: invalid ttl sweep interval '{s}'\n", .{v}); std.debug.print("multiforadb: invalid ttl sweep interval '{s}'\n", .{v});
return error.InvalidTtlSweepSecs; return error.InvalidTtlSweepSecs;
} }
@@ -85,34 +140,45 @@ pub fn main(init: std.process.Init) !void {
std.debug.print("multiforadb: compact threshold must be at least 1m\n", .{}); std.debug.print("multiforadb: compact threshold must be at least 1m\n", .{});
return error.InvalidCompactThreshold; return error.InvalidCompactThreshold;
} }
compact_threshold = parsed; o.compact_threshold = parsed;
} else if (std.mem.eql(u8, arg, "--help") or std.mem.eql(u8, arg, "-h")) { } else if (std.mem.eql(u8, arg, "--help") or std.mem.eql(u8, arg, "-h")) {
try std.Io.File.writeStreamingAll(.stdout(), init.io, usage); return null;
return; } else if (try parse_cursor_flag(arg, &it, &o)) {
// Handled: one of the cursor-registry flags.
} else { } else {
std.debug.print("multiforadb: unknown option '{s}'\n{s}", .{ arg, usage }); std.debug.print("multiforadb: unknown option '{s}'\n{s}", .{ arg, usage });
return error.UnknownOption; return error.UnknownOption;
} }
} }
return o;
const oid_gen = mongo.bson.ObjectIdGen.init(init.io); }
var engine = try mongo.db.Engine.open(init.gpa, init.io, db_path);
defer engine.deinit(); /// The cursor-registry flags, grouped so `parse_args` stays one flat table of
engine.compact_threshold = compact_threshold; /// options. Returns whether `arg` was one of them; consumes its value if so.
std.debug.print("multiforadb: opened database '{s}' (compact threshold {d})\n", .{ fn parse_cursor_flag(arg: []const u8, it: *std.process.Args.Iterator, o: *Options) !bool {
db_path, if (std.mem.eql(u8, arg, "--cursor-timeout-ms")) {
compact_threshold, const v = it.next() orelse return error.MissingValue;
}); o.cursor_timeout_ms = std.fmt.parseInt(i64, v, 10) catch -1;
if (o.cursor_timeout_ms < 0) {
var server = mongo.server.Server{ std.debug.print("multiforadb: invalid cursor timeout '{s}'\n", .{v});
.gpa = init.gpa, return error.InvalidCursorTimeout;
.port = port, }
.bind_ip = bind_ip, } else if (std.mem.eql(u8, arg, "--cursor-sweep-secs")) {
.oid_gen = oid_gen, const v = it.next() orelse return error.MissingValue;
.connection_counter = .init(1), o.cursor_sweep_secs = std.fmt.parseInt(i64, v, 10) catch -1;
.engine = &engine, if (o.cursor_sweep_secs < 0) {
.start_time = std.Io.Timestamp.now(init.io, .real), std.debug.print("multiforadb: invalid cursor sweep interval '{s}'\n", .{v});
.ttl_sweep_secs = ttl_sweep_secs, return error.InvalidCursorSweepSecs;
}; }
try server.run(); } else if (std.mem.eql(u8, arg, "--max-open-cursors")) {
const v = it.next() orelse return error.MissingValue;
o.max_open_cursors = std.fmt.parseInt(u32, v, 10) catch 0;
if (o.max_open_cursors == 0 or o.max_open_cursors > mongo.cursor.max_capacity) {
std.debug.print("multiforadb: invalid max open cursors '{s}'\n", .{v});
return error.InvalidMaxOpenCursors;
}
} else {
return false;
}
return true;
} }

View File

@@ -20,6 +20,9 @@ pub const Server = struct {
/// because that is what std.Io.Duration.fromSeconds takes — the CLI /// because that is what std.Io.Duration.fromSeconds takes — the CLI
/// rejects negatives. /// rejects negatives.
ttl_sweep_secs: i64, ttl_sweep_secs: i64,
/// Seconds between idle-cursor sweeps; 0 leaves that monitor unspawned.
/// mongod's own `clientCursorMonitorFrequencySecs` default is 4.
cursor_sweep_secs: i64,
pub fn run(self: *Server) !void { pub fn run(self: *Server) !void {
// Unbounded async limit: connection handlers otherwise fall back to // Unbounded async limit: connection handlers otherwise fall back to
@@ -43,6 +46,12 @@ pub const Server = struct {
// The TTL monitor is just another member of the connection group, so // The TTL monitor is just another member of the connection group, so
// the `group.cancel` above stops it with everything else. // the `group.cancel` above stops it with everything else.
if (self.ttl_sweep_secs > 0) group.async(io, ttl_monitor, .{ io, self }); if (self.ttl_sweep_secs > 0) group.async(io, ttl_monitor, .{ io, self });
// A separate fiber rather than a branch inside ttl_monitor, for two
// reasons: the cadences differ by more than an order of magnitude (4 s
// against 60 s), and a TTL sweep that fails must not stop cursors being
// reclaimed. It is also spawned when TTL sweeping is disabled entirely,
// which is the configuration the spec runner uses.
if (self.cursor_sweep_secs > 0) group.async(io, cursor_monitor, .{ io, self });
while (true) { while (true) {
const stream = listener.accept(io) catch |err| switch (err) { const stream = listener.accept(io) catch |err| switch (err) {
@@ -87,6 +96,20 @@ fn ttl_monitor(io: std.Io, server: *Server) error{Canceled}!void {
} }
} }
/// Reap cursors nobody has touched for `cursor_timeout_ms`, until the group is
/// canceled. Takes no engine lock: a cursor owns its own arena, and the store's
/// mutex is a leaf.
fn cursor_monitor(io: std.Io, server: *Server) error{Canceled}!void {
const interval: std.Io.Duration = .fromSeconds(server.cursor_sweep_secs);
while (true) {
// Sleep first, for the same reason the TTL monitor does: at startup
// there is nothing to reap and the listener wants the CPU.
try std.Io.sleep(io, interval, .awake);
const now_ms = std.Io.Timestamp.now(io, .real).toMilliseconds();
_ = server.engine.cursors.sweep(io, now_ms);
}
}
/// Entry point required by `Group.async`: must return only `error.Canceled`. /// Entry point required by `Group.async`: must return only `error.Canceled`.
fn handle_connection(io: std.Io, stream: std.Io.net.Stream, server: *Server) error{Canceled}!void { fn handle_connection(io: std.Io, stream: std.Io.net.Stream, server: *Server) error{Canceled}!void {
handle_connection_inner(io, stream, server) catch {}; handle_connection_inner(io, stream, server) catch {};

View File

@@ -299,9 +299,20 @@ fn begin_message(
} }
/// Patch in the total length of the message started at `len_pos`. /// Patch in the total length of the message started at `len_pos`.
///
/// The bound is `max_message_size`, the same 48 MiB we advertise to drivers as
/// `maxMessageSizeBytes`, not `maxInt(u32)`. A reply past what we told the
/// client to expect is not a large reply, it is a desynchronized connection:
/// the driver reads the length, refuses or mis-frames it, and every later
/// command on that socket reads the wrong bytes. Failing here turns that into
/// one honest error on the request that caused it.
///
/// Reachable today: nothing caps how many documents a `find` puts in its single
/// batch, so ~3000 documents of 16 KiB clears 48 MB. The cursor batch budget
/// makes it unreachable, which is the point of keeping this as the backstop.
fn end_message(out: *std.ArrayListUnmanaged(u8), len_pos: usize) !void { fn end_message(out: *std.ArrayListUnmanaged(u8), len_pos: usize) !void {
const total = out.items.len - len_pos; const total = out.items.len - len_pos;
if (total > std.math.maxInt(u32)) return error.MessageTooLarge; if (total > max_message_size) return error.MessageTooLarge;
std.mem.writeInt(u32, out.items[len_pos..][0..4], @intCast(total), .little); std.mem.writeInt(u32, out.items[len_pos..][0..4], @intCast(total), .little);
} }
@@ -411,3 +422,31 @@ test "parse OP_QUERY handshake" {
try testing.expectEqual(@as(i32, op_code_query), msg.op_code); try testing.expectEqual(@as(i32, op_code_query), msg.op_code);
try testing.expectEqualStrings("isMaster", msg.command_name()); try testing.expectEqualStrings("isMaster", msg.command_name());
} }
test "a reply past the advertised message size fails to build" {
// The guard exists because exceeding it desynchronizes the connection
// rather than merely making one reply large, so it must be an error return
// and not a truncation. One oversized string is the cheapest way past it
// without allocating 48 MB of documents.
const gpa = testing.allocator;
var reply = Reply.init(gpa);
defer reply.deinit();
const big = try reply.arena_alloc().alloc(u8, max_message_size + 1);
@memset(big, 'x');
try reply.put_ok();
try reply.put("payload", .{ .string = big });
var out: std.ArrayListUnmanaged(u8) = .empty;
defer out.deinit(gpa);
try testing.expectError(error.MessageTooLarge, reply.build(gpa, 1, 1, &out));
// And a reply comfortably inside the bound still builds, so the guard is
// not simply rejecting everything.
var small = Reply.init(gpa);
defer small.deinit();
try small.put_ok();
out.clearRetainingCapacity();
try small.build(gpa, 1, 1, &out);
try testing.expect(out.items.len < max_message_size);
}

View File

@@ -41,6 +41,19 @@ node tests/e2e/e2e6.js # 73 checks, ~15 s, needs no running server
E2E6_PORT=27300 node tests/e2e/e2e6.js # different port if 27220 is taken E2E6_PORT=27300 node tests/e2e/e2e6.js # different port if 27220 is taken
``` ```
`e2e7.js` is the cursor suite, self-contained for a different reason: cursor
behaviour is only observable with non-default flags. It spawns three servers in
turn -- default flags for batching/streaming/aggregate, then
`--cursor-timeout-ms 800 --cursor-sweep-secs 1 --max-open-cursors 4` for idle
expiry and registry capacity, then a restart on the same database to confirm a
cursor does not survive one. Most of it uses raw `runCommand`, because the
driver hides `cursor.id` and that is the thing under test:
```sh
node tests/e2e/e2e7.js # 86 checks, needs no running server
E2E7_PORT=27310 node tests/e2e/e2e7.js # different port if 27230 is taken
```
Rebuild with `zig build` after any change under `src/` before restarting the Rebuild with `zig build` after any change under `src/` before restarting the
server: `zig build test` compiles the test binary only and leaves server: `zig build test` compiles the test binary only and leaves
`zig-out/bin/multiforadb` stale, so the suites keep running against the old `zig-out/bin/multiforadb` stale, so the suites keep running against the old

474
tests/e2e/e2e7.js Normal file
View File

@@ -0,0 +1,474 @@
// E2E part 7: server-side cursors, self-contained.
//
// Spawns its own multiforadb servers, because cursor behaviour is only
// observable with non-default flags (a short idle timeout, a tiny registry) and
// with raw `runCommand` — the driver hides `cursor.id`, which is the thing under
// test.
//
// node tests/e2e/e2e7.js
//
// Env: E2E7_PORT listen port (default 27230)
// MFDB_BIN server binary (default ../../zig-out/bin/multiforadb)
// E2E7_KEEP keep the log files after the run
//
// The one rule most of this file is about: **never look ahead.** A batch ends
// either because it reached its target — cursor stays open — or because the
// source reported EOF, and only then does the cursor close with `id: 0`. So four
// documents at `batchSize: 2` require a third command answering an empty
// `nextBatch` with `id: 0`. That empty terminal batch is correct, and it is what
// real mongod does (measured, not assumed — see checks 5 and 6).
const { MongoClient, Long } = require('mongodb');
const { spawn } = require('child_process');
const fs = require('fs');
const path = require('path');
const PORT = Number(process.env.E2E7_PORT || 27230);
const BIN = process.env.MFDB_BIN || path.resolve(__dirname, '../../zig-out/bin/multiforadb');
const DBFILE = path.resolve(__dirname, '../../.zig-cache/e2e7-cursors.log');
const URL = `mongodb://127.0.0.1:${PORT}`;
const results = [];
function check(name, cond, detail = '') {
results.push({ name, ok: !!cond, detail: String(detail) });
if (!cond) console.error(` x ${name} ${detail}`);
}
function eq(name, got, want) {
const ok = JSON.stringify(got) === JSON.stringify(want);
check(name, ok, ok ? '' : `got ${JSON.stringify(got)} want ${JSON.stringify(want)}`);
}
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
/// The error code a command produces, or 0 when it succeeds.
async function codeOf(fn) {
try {
await fn();
return 0;
} catch (e) {
return e.code === undefined ? -1 : e.code;
}
}
let server = null;
let serverLog = '';
let serverDead = false;
function cleanup() {
if (server && !serverDead) {
try { server.kill('SIGKILL'); } catch {}
}
}
process.on('exit', cleanup);
process.on('SIGINT', () => { cleanup(); process.exit(130); });
process.on('SIGTERM', () => { cleanup(); process.exit(143); });
function startServer(args, fresh = true) {
return new Promise((resolve, reject) => {
if (fresh) {
fs.rmSync(DBFILE, { force: true });
fs.rmSync(DBFILE + '.data', { force: true });
}
serverDead = false;
serverLog = '';
server = spawn(BIN, ['--port', String(PORT), '--db', DBFILE, ...args], {
stdio: ['ignore', 'pipe', 'pipe'],
});
server.stdout.on('data', (d) => (serverLog += d));
server.stderr.on('data', (d) => (serverLog += d));
server.on('error', (e) => reject(new Error(`cannot start ${BIN}: ${e.message}`)));
server.on('exit', (code, sig) => {
serverDead = true;
if (code !== null && sig === null) serverLog += `\n[child exited rc=${code}]`;
});
const deadline = Date.now() + 15000;
(async () => {
while (Date.now() < deadline) {
if (serverDead) {
reject(new Error(`server exited during start (port ${PORT} busy?)\n${serverLog}`));
return;
}
const c = new MongoClient(URL, { serverSelectionTimeoutMS: 1000 });
try {
await c.connect();
await c.db('admin').command({ ping: 1 });
await c.close();
return resolve();
} catch {
try { await c.close(); } catch {}
await sleep(100);
}
}
reject(new Error(`server did not come up on :${PORT}\n${serverLog}`));
})();
});
}
async function stopServer(sig = 'SIGTERM') {
if (!server) return;
const exited = new Promise((r) => server.once('exit', r));
server.kill(sig);
await Promise.race([exited, sleep(5000)]);
serverDead = true;
server = null;
}
// ---------------------------------------------------------------------------
// Phase A — batching, lifecycle and errors, on default flags
// ---------------------------------------------------------------------------
async function phaseA(db) {
const col = db.collection('c');
await col.deleteMany({});
await col.insertMany([...Array(250)].map((_, i) => ({ _id: i + 1, x: i })));
// 1. A cursor is a real cursor: nonzero id, a namespace with both parts.
let r = await db.command({ find: 'c', filter: {}, batchSize: 2 });
eq('1 batchSize 2 returns 2', r.cursor.firstBatch.length, 2);
check('1 cursor id is nonzero', r.cursor.id > 0, r.cursor.id);
eq('1 ns is db.coll', r.cursor.ns, 'e2e7.c');
// 2. The default first batch is 101, MongoDB's own
// internalQueryFindCommandBatchSize.
eq('2 default first batch is 101', (await db.command({ find: 'c', filter: {} })).cursor.firstBatch.length, 101);
// 3. A getMore naming no batchSize is bounded by bytes, not by the batchSize
// the cursor was created with. Measured against mongod 8.3.7: find with
// batchSize 2 then a bare getMore returns 4998 of 5000 documents.
const id3 = (await db.command({ find: 'c', filter: {}, batchSize: 2 })).cursor.id;
let g = await db.command({ getMore: id3, collection: 'c' });
eq('3 a bare getMore drains the rest', g.cursor.nextBatch.length, 248);
eq('3 and closes at EOF', String(g.cursor.id), '0');
// 4. A getMore's batchSize applies to that batch only.
const id4 = (await db.command({ find: 'c', filter: {}, batchSize: 2 })).cursor.id;
g = await db.command({ getMore: id4, collection: 'c', batchSize: 3 });
eq('4 getMore batchSize 3', g.cursor.nextBatch.map((d) => d._id), [3, 4, 5]);
check('4 still open', g.cursor.id > 0);
// 5. No look-ahead: 4 documents at batchSize 2 needs a third command whose
// nextBatch is empty. Closing on "batch full and source dry" would break
// the command counts the pinned spec suites assert.
const four = db.collection('four');
await four.deleteMany({});
await four.insertMany([1, 2, 3, 4].map((i) => ({ _id: i })));
r = await db.command({ find: 'four', filter: {}, batchSize: 2 });
g = await db.command({ getMore: r.cursor.id, collection: 'four', batchSize: 2 });
eq('5 second full batch is 2 documents', g.cursor.nextBatch.length, 2);
check('5 and leaves the cursor open', g.cursor.id > 0);
g = await db.command({ getMore: g.cursor.id, collection: 'four', batchSize: 2 });
eq('5 terminal batch is empty and closed', [g.cursor.nextBatch.length, String(g.cursor.id)], [0, '0']);
// 6. limit is an EOF source, so the batch that exhausts it also closes the
// cursor. This is why the driver sends batchSize = limit + 1.
r = await db.command({ find: 'c', filter: {}, sort: { _id: 1 }, limit: 4, batchSize: 5 });
eq('6 limit 4 batchSize 5 closes in one reply', String(r.cursor.id), '0');
eq('6 and returns exactly the limit', r.cursor.firstBatch.map((d) => d._id), [1, 2, 3, 4]);
r = await db.command({ find: 'c', filter: {}, sort: { _id: 1 }, limit: 4, batchSize: 2 });
check('6 limit 4 batchSize 2 stays open', r.cursor.id > 0);
g = await db.command({ getMore: r.cursor.id, collection: 'c', batchSize: 2 });
eq('6 the batch reaching the limit closes', String(g.cursor.id), '0');
eq('6 across batches, limit still honoured', g.cursor.nextBatch.map((d) => d._id), [3, 4]);
// 7. skip is consumed once, at creation, and never re-applied.
r = await db.command({ find: 'c', filter: {}, skip: 20, batchSize: 3 });
eq('7 skip 20 starts at 21', r.cursor.firstBatch.map((d) => d._id), [21, 22, 23]);
g = await db.command({ getMore: r.cursor.id, collection: 'c', batchSize: 3 });
eq('7 skip not re-applied on getMore', g.cursor.nextBatch.map((d) => d._id), [24, 25, 26]);
// 8. batchSize 0 is a real request for an empty batch with a live cursor, not
// "unbounded". Drivers use it to obtain a cursor cheaply. Nothing may be
// consumed.
r = await db.command({ find: 'c', filter: {}, batchSize: 0 });
eq('8 batchSize 0 returns nothing', r.cursor.firstBatch.length, 0);
check('8 but a live cursor', r.cursor.id > 0);
g = await db.command({ getMore: r.cursor.id, collection: 'c', batchSize: 1 });
eq('8 nothing was consumed', g.cursor.nextBatch.map((d) => d._id), [1]);
// 9. singleBatch, and its wire-legacy form, a negative limit.
eq('9 singleBatch closes', String((await db.command({ find: 'c', filter: {}, batchSize: 2, singleBatch: true })).cursor.id), '0');
r = await db.command({ find: 'c', filter: {}, limit: -3 });
eq('9 negative limit is one batch', [r.cursor.firstBatch.length, String(r.cursor.id)], [3, '0']);
// 10. A cursor is not pinned to the connection that created it: the driver
// spec allows a getMore from any connection to the same server.
const other = new MongoClient(URL);
await other.connect();
const shared = (await db.command({ find: 'c', filter: {}, batchSize: 2 })).cursor.id;
eq('10 getMore from another connection', await codeOf(() => other.db('e2e7').command({ getMore: shared, collection: 'c', batchSize: 2 })), 0);
await other.close();
// 11. A getMore naming the wrong collection is Unauthorized (13), not 43, and
// leaves the cursor alive — the request is wrong, not the cursor.
// Measured against mongod, which answers exactly this code.
const live = (await db.command({ find: 'c', filter: {}, batchSize: 2 })).cursor.id;
eq('11 wrong collection is Unauthorized 13', await codeOf(() => db.command({ getMore: live, collection: 'four' })), 13);
eq('11 the cursor survived it', (await db.command({ getMore: live, collection: 'c', batchSize: 1 })).cursor.nextBatch.length, 1);
// 12. killCursors: all four arrays, and the right partitioning.
let k = await db.command({ killCursors: 'c', cursors: [live] });
eq('12 a live cursor is killed', k.cursorsKilled.map(String), [String(live)]);
check('12 all four arrays present', ['cursorsKilled', 'cursorsNotFound', 'cursorsAlive', 'cursorsUnknown'].every((f) => Array.isArray(k[f])), Object.keys(k).join(','));
k = await db.command({ killCursors: 'c', cursors: [live] });
eq('12 killing it twice reports notFound', k.cursorsNotFound.map(String), [String(live)]);
const other_ns = (await db.command({ find: 'four', filter: {}, batchSize: 1 })).cursor.id;
k = await db.command({ killCursors: 'c', cursors: [other_ns] });
eq('12 a wrong-namespace id reports notFound', k.cursorsNotFound.map(String), [String(other_ns)]);
eq('12 and that cursor still lives', await codeOf(() => db.command({ getMore: other_ns, collection: 'four', batchSize: 1 })), 0);
// 13. Malformed and unknown ids.
eq('13 getMore after kill is 43', await codeOf(() => db.command({ getMore: live, collection: 'c' })), 43);
eq('13 id 0 is 43', await codeOf(() => db.command({ getMore: Long.fromNumber(0), collection: 'c' })), 43);
eq('13 an unknown id is 43', await codeOf(() => db.command({ getMore: Long.fromString('987654321'), collection: 'c' })), 43);
eq('13 a non-numeric id is TypeMismatch 14', await codeOf(() => db.command({ getMore: 'nope', collection: 'c' })), 14);
eq('13 a missing collection is BadValue 2', await codeOf(() => db.command({ getMore: Long.fromNumber(1) })), 2);
// 14. tailable is refused, which is parity: mongod rejects it on a non-capped
// collection and this engine has none. Ignoring it would make a driver's
// tail loop exit, which the application reads as data loss.
eq('14 tailable is BadValue 2', await codeOf(() => db.command({ find: 'c', filter: {}, tailable: true })), 2);
eq('14 awaitData alone is BadValue 2', await codeOf(() => db.command({ find: 'c', filter: {}, awaitData: true })), 2);
// 15. The driver's own iteration, which is the point of all of the above.
const ids = (await col.find({}).batchSize(7).toArray()).map((d) => d._id);
eq('15 driver drains 250 at batchSize 7', ids.length, 250);
eq('15 no duplicates and no gaps', [new Set(ids).size, Math.min(...ids), Math.max(...ids)], [250, 1, 250]);
eq('15 sort+skip+limit unchanged', (await col.find({ _id: { $gt: 2 } }, { sort: { _id: 1 }, skip: 2, limit: 2 }).toArray()).map((d) => d._id), [5, 6]);
// A sort no index provides must materialize; it still has to drain correctly.
eq('15 an unindexed sort drains in order', (await col.find({}, { sort: { x: -1 } }).batchSize(10).toArray()).map((d) => d.x)[0], 249);
}
// ---------------------------------------------------------------------------
// Phase B — streaming cursors: resume across writes, and what survives a rebuild
// ---------------------------------------------------------------------------
async function phaseB(db) {
const col = db.collection('s');
await col.deleteMany({});
await col.insertMany([...Array(300)].map((_, i) => ({ _id: i + 1, a: i % 5, pad: 'q'.repeat(200) })));
// 16. A whole-index walk holds a key, not a list, so it resumes across writes
// that move documents. The bug this caught: an update rewrites a document
// to a new offset, and resuming by band position returned it twice.
let r = await db.command({ find: 's', filter: {}, batchSize: 10 });
const seen = new Set(r.cursor.firstBatch.map((d) => d._id));
let dupes = 0;
let id = r.cursor.id;
let rounds = 0;
let errored = 0;
while (String(id) !== '0' && rounds++ < 200) {
// Churn between every batch: updates rewrite documents, which both moves
// them in the slab and can split leaves.
await col.updateMany({ _id: { $lt: 60 } }, { $inc: { n: 1 } });
let g;
try {
g = await db.command({ getMore: id, collection: 's', batchSize: 10 });
} catch (e) {
errored = e.code;
break;
}
for (const d of g.cursor.nextBatch) {
if (seen.has(d._id)) dupes++;
seen.add(d._id);
}
id = g.cursor.id;
}
eq('16 draining across churn did not error', errored, 0);
eq('16 no document came back twice', dupes, 0);
eq('16 every document was returned', seen.size, 300);
check('16 and nothing outside the collection', [...seen].every((v) => v >= 1 && v <= 300));
// 17. Both directions stream, over the _id_ index and a secondary one.
eq('17 ascending _id sort drains', (await col.find({}, { sort: { _id: 1 } }).batchSize(9).toArray()).length, 300);
const desc = (await col.find({}, { sort: { _id: -1 } }).batchSize(9).toArray()).map((d) => d._id);
eq('17 descending drains in order', [desc.length, desc[0], desc[299]], [300, 300, 1]);
await col.createIndex({ a: 1 });
const bya = await col.find({}, { sort: { a: 1 } }).batchSize(11).toArray();
eq('17 a secondary-index sort drains', bya.length, 300);
check('17 and in the index order', bya.every((d, i) => i === 0 || bya[i - 1].a <= d.a));
// 18. Dropping the index a stream is following cannot be resumed — the walk
// has nothing left to walk. That must be a clean error, not garbage.
await col.createIndex({ b: 1 });
const onb = (await db.command({ find: 's', filter: {}, sort: { b: 1 }, batchSize: 3 })).cursor.id;
await col.dropIndex('b_1');
eq('18 dropping the streamed index is 175', await codeOf(() => db.command({ getMore: onb, collection: 's', batchSize: 3 })), 175);
// 19. Dropping the collection kills every kind of cursor.
const doomed = (await db.command({ find: 's', filter: {}, batchSize: 3 })).cursor.id;
await col.drop();
const dc = await codeOf(() => db.command({ getMore: doomed, collection: 's', batchSize: 3 }));
check('19 dropping the collection kills the cursor', dc === 175 || dc === 43, dc);
}
// ---------------------------------------------------------------------------
// Phase C — aggregate, the listing commands, and count
// ---------------------------------------------------------------------------
async function phaseC(db) {
const col = db.collection('g');
await col.deleteMany({});
await col.insertMany([...Array(250)].map((_, i) => ({ _id: i + 1, g: i % 40, v: i })));
// 20. aggregate batches through cursor.batchSize; a bare cursor is the default.
let r = await db.command({ aggregate: 'g', pipeline: [], cursor: { batchSize: 3 } });
eq('20 aggregate batchSize 3', r.cursor.firstBatch.length, 3);
check('20 aggregate cursor is real', r.cursor.id > 0);
eq('20 aggregate ns', r.cursor.ns, 'e2e7.g');
const g20 = await db.command({ getMore: r.cursor.id, collection: 'g', batchSize: 5 });
eq('20 aggregate getMore continues', g20.cursor.nextBatch.map((d) => d._id), [4, 5, 6, 7, 8]);
eq('20 bare cursor defaults to 101', (await db.command({ aggregate: 'g', pipeline: [], cursor: {} })).cursor.firstBatch.length, 101);
eq('20 driver aggregate drains', (await col.aggregate([], { batchSize: 7 }).toArray()).length, 250);
const groups = await col.aggregate([{ $group: { _id: '$g', n: { $sum: 1 } } }], { batchSize: 6 }).toArray();
eq('20 $group drains across batches', [groups.length, groups.reduce((a, x) => a + x.n, 0)], [40, 250]);
eq('20 $count stage', await col.aggregate([{ $count: 'total' }]).toArray(), [{ total: 250 }]);
// 21. listCollections' namespace. It used to be "<db>." with an empty
// collection part, and the driver throws client-side on a namespace like
// that — so the moment the cursor stopped being id 0 it would have broken.
for (let i = 0; i < 12; i++) await db.createCollection('k' + i);
r = await db.command({ listCollections: 1, cursor: { batchSize: 4 } });
eq('21 listCollections ns has a collection part', r.cursor.ns, 'e2e7.$cmd.listCollections');
eq('21 listCollections honours batchSize', r.cursor.firstBatch.length, 4);
check('21 listCollections cursor is real', r.cursor.id > 0);
const g21 = await db.command({ getMore: r.cursor.id, collection: '$cmd.listCollections', batchSize: 100 });
check('21 its getMore works', g21.cursor.nextBatch.length >= 8, g21.cursor.nextBatch.length);
const listed = await db.listCollections({}, { batchSize: 3 }).toArray();
check('21 driver listCollections drains', listed.length >= 13, listed.length);
// 22. listIndexes.
await col.createIndex({ v: 1 });
await col.createIndex({ g: 1 });
await col.createIndex({ v: -1, g: 1 });
r = await db.command({ listIndexes: 'g', cursor: { batchSize: 2 } });
eq('22 listIndexes honours batchSize', r.cursor.firstBatch.length, 2);
eq('22 listIndexes ns', r.cursor.ns, 'e2e7.g');
eq('22 driver listIndexes drains', (await col.listIndexes({ batchSize: 1 }).toArray()).length, 4);
// 23. count honoured neither skip nor limit before, which made
// countDocuments(f, {limit}) a silent wrong answer.
eq('23 count plain', (await db.command({ count: 'g' })).n, 250);
eq('23 count limit', (await db.command({ count: 'g', limit: 10 })).n, 10);
eq('23 count skip', (await db.command({ count: 'g', skip: 240 })).n, 10);
eq('23 count skip and limit', (await db.command({ count: 'g', skip: 245, limit: 10 })).n, 5);
eq('23 count skip past the end', (await db.command({ count: 'g', skip: 1000 })).n, 0);
eq('23 count with a query and limit', (await db.command({ count: 'g', query: { g: 0 }, limit: 3 })).n, 3);
eq('23 driver countDocuments limit', await col.countDocuments({}, { limit: 7 }), 7);
// 24. A batch is capped by bytes as well as by documents, so a large-document
// result splits instead of building a reply past the advertised message
// size. 40 documents of ~1 MiB cannot all fit one 16 MiB batch.
const big = db.collection('big');
await big.deleteMany({});
const pad = 'p'.repeat(1024 * 1024 - 64);
for (let i = 0; i < 40; i++) await big.insertOne({ _id: i + 1, pad });
r = await db.command({ find: 'big', filter: {}, batchSize: 40 });
check('24 the byte cap split the batch', r.cursor.firstBatch.length >= 1 && r.cursor.firstBatch.length <= 16, r.cursor.firstBatch.length);
check('24 and left the cursor open', r.cursor.id > 0);
eq('24 the whole result still drains', (await big.find({}).batchSize(40).toArray()).length, 40);
}
// ---------------------------------------------------------------------------
// Phase D — expiry, capacity, and what a restart does
// ---------------------------------------------------------------------------
async function phaseD(db) {
const col = db.collection('e');
await col.deleteMany({});
await col.insertMany([...Array(50)].map((_, i) => ({ _id: i + 1 })));
// 25. An idle cursor is reaped; noCursorTimeout exempts one from that.
const perishable = (await db.command({ find: 'e', filter: {}, batchSize: 2 })).cursor.id;
const immortal = (await db.command({ find: 'e', filter: {}, batchSize: 2, noCursorTimeout: true })).cursor.id;
await sleep(300);
eq('25 before the timeout it is alive', await codeOf(() => db.command({ getMore: perishable, collection: 'e', batchSize: 1 })), 0);
await sleep(2500);
eq('25 an idle cursor is reaped', await codeOf(() => db.command({ getMore: perishable, collection: 'e', batchSize: 1 })), 43);
eq('25 noCursorTimeout survives', await codeOf(() => db.command({ getMore: immortal, collection: 'e', batchSize: 1 })), 0);
const k = await db.command({ killCursors: 'e', cursors: [immortal] });
eq('25 but is still killable', k.cursorsKilled.map(String), [String(immortal)]);
// 26. A full registry evicts the least recently used cursor rather than
// refusing the new one. The victim sees the same 43 an idle timeout gives,
// which every driver already handles.
const ids = [];
for (let i = 0; i < 5; i++) ids.push((await db.command({ find: 'e', filter: {}, batchSize: 1 })).cursor.id);
eq('26 the oldest was evicted', await codeOf(() => db.command({ getMore: ids[0], collection: 'e', batchSize: 1 })), 43);
const alive = [];
for (const id of ids.slice(1)) alive.push(await codeOf(() => db.command({ getMore: id, collection: 'e', batchSize: 1 })));
eq('26 the newest four are alive', alive, [0, 0, 0, 0]);
}
async function phaseE(db, staleId) {
// 27. Cursors do not survive a restart, and a stale id must be a clean 43 —
// not a hang, and not an empty batch claiming the result ended.
eq('27 a cursor from before the restart is 43', await codeOf(() => db.command({ getMore: staleId, collection: 'e', batchSize: 1 })), 43);
const fresh = await db.command({ find: 'e', filter: {}, batchSize: 2 });
check('27 and new cursors work after a restart', fresh.cursor.id > 0);
}
async function main() {
// The same guard e2e6.js and big.js carry: without it a missing binary
// surfaces as a generic spawn error instead of saying what to do about it.
if (!fs.existsSync(BIN)) {
console.error(`E2E7_FAIL server binary not found: ${BIN}\n run: zig build`);
process.exit(1);
}
// Phase A-C on default cursor flags.
await startServer(['--ttl-sweep-secs', '0', '--compact-threshold', '1m'], true);
let client = new MongoClient(URL);
await client.connect();
let db = client.db('e2e7');
console.log('phase A: batching, lifecycle, errors');
await phaseA(db);
console.log('phase B: streaming cursors across writes');
await phaseB(db);
console.log('phase C: aggregate, listings, count');
await phaseC(db);
await client.close();
await stopServer('SIGTERM');
// Phase D needs a short timeout and a tiny registry.
console.log('phase D: idle expiry and registry capacity');
await startServer(
['--ttl-sweep-secs', '0', '--cursor-timeout-ms', '800', '--cursor-sweep-secs', '1', '--max-open-cursors', '4'],
true,
);
client = new MongoClient(URL);
await client.connect();
db = client.db('e2e7');
await phaseD(db);
const staleId = (await db.command({ find: 'e', filter: {}, batchSize: 1 })).cursor.id;
await client.close();
await stopServer('SIGTERM');
// Phase E: the same database, a new process.
console.log('phase E: a cursor does not survive a restart');
await startServer(['--ttl-sweep-secs', '0'], false);
client = new MongoClient(URL);
await client.connect();
await phaseE(client.db('e2e7'), staleId);
await client.close();
if (process.env.E2E7_KEEP !== '1') {
fs.rmSync(DBFILE, { force: true });
fs.rmSync(DBFILE + '.data', { force: true });
}
await stopServer('SIGTERM');
const failed = results.filter((r) => !r.ok);
console.log(`\n${results.length - failed.length}/${results.length} checks passed`);
if (failed.length) {
console.log('FAILED:', failed.map((f) => f.name).join(', '));
console.log('--- server log tail ---');
console.log(serverLog.split('\n').slice(-30).join('\n'));
process.exit(1);
}
console.log('E2E7_OK');
}
main().catch((e) => {
console.error('E2E7_FAIL', e);
console.log('--- server log tail ---');
console.log(serverLog.split('\n').slice(-40).join('\n'));
process.exit(1);
});

View File

@@ -314,11 +314,38 @@ class Unsupported extends Error {}
// Argument keys the spec passes positionally rather than as driver options. // Argument keys the spec passes positionally rather than as driver options.
const POSITIONAL = new Set(['filter', 'document', 'documents', 'update', 'replacement', 'pipeline', 'fieldName', 'models', 'requests', 'keys', 'name', 'indexes', 'command', 'session', 'entity', 'to']); const POSITIONAL = new Set(['filter', 'document', 'documents', 'update', 'replacement', 'pipeline', 'fieldName', 'models', 'requests', 'keys', 'name', 'indexes', 'command', 'session', 'entity', 'to']);
// Driver options the driver only honours as a JavaScript number. The suites are
// parsed with `EJSON.parse(text, {relaxed: false})` so that `$numberLong` and
// friends keep their exact BSON type in *data* -- but that also turns a plain
// JSON `2` in an *option* into a BSON Int32 object, and the driver gates every
// one of these on `typeof options.skip === 'number'`
// (node_modules/mongodb/lib/operations/find.js:68-95). A BSON wrapper therefore
// failed the check and the option was dropped on the floor: `skip`, `limit` and
// `batchSize` never reached the wire at all, and three find.json cases failed
// with the *unclipped* match count while the engine was applying both correctly.
// Read as an engine bug for a whole milestone. Coerce by name, not by shape:
// unwrapping every numeric-looking value would rewrite the wire type of the
// `comment` and `hint` values that other suites assert on.
const NUMERIC_OPTIONS = new Set([
'skip',
'limit',
'batchSize',
'maxTimeMS',
'maxAwaitTimeMS',
'expireAfterSeconds',
]);
function numeric_option(v) {
if (v === null || typeof v !== 'object' || typeof v.valueOf !== 'function') return v;
const n = v.valueOf();
return typeof n === 'number' ? n : v;
}
function options(args, drop = []) { function options(args, drop = []) {
const o = {}; const o = {};
for (const [k, v] of Object.entries(args || {})) { for (const [k, v] of Object.entries(args || {})) {
if (POSITIONAL.has(k) || drop.includes(k)) continue; if (POSITIONAL.has(k) || drop.includes(k)) continue;
o[k] = v; o[k] = NUMERIC_OPTIONS.has(k) ? numeric_option(v) : v;
} }
return Object.keys(o).length ? o : undefined; return Object.keys(o).length ? o : undefined;
} }

View File

@@ -13,7 +13,7 @@
# semantics; ignoring them makes some cases pass that a full runner would # semantics; ignoring them makes some cases pass that a full runner would
# fail, so treat `pass` as an upper bound until M1 wires events up. # fail, so treat `pass` as an upper bound until M1 wires events up.
total 163 pass 129 fail 195 skip 175 files 0 errored total 168 pass 124 fail 195 skip 175 files 0 errored
# per-file: name pass fail skip # per-file: name pass fail skip
aggregate-allowdiskuse.json 3 0 0 aggregate-allowdiskuse.json 3 0 0
@@ -78,7 +78,7 @@ client-bulkWrite-updateOne-sort.json 0 0 1
count-collation.json 1 0 1 count-collation.json 1 0 1
count-empty.json 2 0 1 count-empty.json 2 0 1
count-rawdata.json 0 0 2 count-rawdata.json 0 0 2
count.json 3 1 3 count.json 4 0 3
countDocuments-comment.json 2 0 1 countDocuments-comment.json 2 0 1
countDocuments-rawdata.json 1 0 1 countDocuments-rawdata.json 1 0 1
create-null-ids.json 0 6 1 create-null-ids.json 0 6 1
@@ -116,8 +116,8 @@ find-collation.json 0 1 0
find-comment.json 1 2 2 find-comment.json 1 2 2
find-let.json 0 1 1 find-let.json 0 1 1
find-rawdata.json 1 0 1 find-rawdata.json 1 0 1
find.json 2 3 0 find.json 5 0 0
findOne.json 1 1 0 findOne.json 2 0 0
findOneAndDelete-collation.json 0 1 0 findOneAndDelete-collation.json 0 1 0
findOneAndDelete-comment.json 2 0 1 findOneAndDelete-comment.json 2 0 1
findOneAndDelete-hint-serverError.json 0 0 2 findOneAndDelete-hint-serverError.json 0 0 2
@@ -292,7 +292,6 @@ count-collation.json SKIP Deprecated count with collation runner: operation coun
count-empty.json SKIP Deprecated count with empty collection runner: operation count count-empty.json SKIP Deprecated count with empty collection runner: operation count
count-rawdata.json SKIP Deprecated count with rawData option needs server >= 8.2.0 count-rawdata.json SKIP Deprecated count with rawData option needs server >= 8.2.0
count-rawdata.json SKIP Deprecated count with rawData option on less than 8.2.0 - ignore argument runner: operation count count-rawdata.json SKIP Deprecated count with rawData option on less than 8.2.0 - ignore argument runner: operation count
count.json FAIL Count documents with skip and limit countDocuments: expected 2, got 3
count.json SKIP Deprecated count without a filter runner: operation count count.json SKIP Deprecated count without a filter runner: operation count
count.json SKIP Deprecated count with a filter runner: operation count count.json SKIP Deprecated count with a filter runner: operation count
count.json SKIP Deprecated count with skip and limit runner: operation count count.json SKIP Deprecated count with skip and limit runner: operation count
@@ -354,10 +353,6 @@ find-comment.json SKIP find with comment does not set comment on getMore - pre 4
find-let.json SKIP Find with let option needs server >= 5.0 find-let.json SKIP Find with let option needs server >= 5.0
find-let.json FAIL Find with let option unsupported (server-side error) find: expected an error, the operation succeeded find-let.json FAIL Find with let option unsupported (server-side error) find: expected an error, the operation succeeded
find-rawdata.json SKIP Find with rawData option needs server >= 8.2.0 find-rawdata.json SKIP Find with rawData option needs server >= 8.2.0
find.json FAIL Find with filter, sort, skip, and limit find: expected 2 elements, got 4
find.json FAIL Find with limit, sort, and batchsize find: expected 4 elements, got 6
find.json FAIL Find with batchSize equal to limit find: expected 4 elements, got 5
findOne.json FAIL FindOne with filter, sort, and skip findOne._id: expected 5, got 3
findOneAndDelete-collation.json FAIL FindOneAndDelete when one document matches with collation findOneAndDelete: expected a document, got null findOneAndDelete-collation.json FAIL FindOneAndDelete when one document matches with collation findOneAndDelete: expected a document, got null
findOneAndDelete-comment.json SKIP findOneAndDelete with comment - pre 4.4 needs server <= 4.2.99 findOneAndDelete-comment.json SKIP findOneAndDelete with comment - pre 4.4 needs server <= 4.2.99
findOneAndDelete-hint-serverError.json SKIP * needs server <= 4.3.3 findOneAndDelete-hint-serverError.json SKIP * needs server <= 4.3.3