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:
1254
src/commands.zig
1254
src/commands.zig
File diff suppressed because it is too large
Load Diff
979
src/cursor.zig
Normal file
979
src/cursor.zig
Normal 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());
|
||||
}
|
||||
127
src/db.zig
127
src/db.zig
@@ -22,6 +22,7 @@ const bson = @import("bson.zig");
|
||||
const storage = @import("storage.zig");
|
||||
const index = @import("index.zig");
|
||||
const pgr = @import("pager.zig");
|
||||
const cursor = @import("cursor.zig");
|
||||
// Always active, including in the default ReleaseFast build -- see assert.zig
|
||||
// for why std.debug.assert is the wrong tool for these invariants.
|
||||
const assert = @import("assert.zig").assert;
|
||||
@@ -97,8 +98,21 @@ pub const Collection = struct {
|
||||
/// replaces the old serialization-guarded docs-map fast path for
|
||||
/// integer/string/etc. _id lookups.
|
||||
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 = .{
|
||||
.doc_count = 0,
|
||||
.pager = pager,
|
||||
@@ -110,6 +124,7 @@ pub const Collection = struct {
|
||||
.hold = .{},
|
||||
.indexes = .empty,
|
||||
.id_index = undefined,
|
||||
.layout_epoch = layout_epoch,
|
||||
};
|
||||
const keys = [_]index.IndexKey{.{ .path = "_id", .descending = false }};
|
||||
// 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`.
|
||||
live_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.
|
||||
/// 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
|
||||
@@ -324,6 +348,12 @@ pub const Engine = struct {
|
||||
/// command reads it while still holding the write lock.
|
||||
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 {
|
||||
var log = try storage.Log.open(gpa, io, path);
|
||||
errdefer log.close();
|
||||
@@ -346,8 +376,10 @@ pub const Engine = struct {
|
||||
.dbs = .empty,
|
||||
.seq = 0,
|
||||
.compact_threshold = 16 * 1024 * 1024,
|
||||
.cursors = try default_cursor_store(gpa, io),
|
||||
};
|
||||
errdefer {
|
||||
engine.cursors.deinit();
|
||||
engine.pager.deinit();
|
||||
engine.dbs.deinit(gpa);
|
||||
}
|
||||
@@ -393,6 +425,17 @@ pub const Engine = struct {
|
||||
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 {
|
||||
var db_it = self.dbs.iterator();
|
||||
while (db_it.next()) |db_entry| {
|
||||
@@ -400,6 +443,11 @@ pub const Engine = struct {
|
||||
self.gpa.free(db_entry.key_ptr.*);
|
||||
}
|
||||
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.gpa.destroy(self.pager);
|
||||
self.log.close();
|
||||
@@ -936,6 +984,11 @@ pub const Engine = struct {
|
||||
const removed = db.collections.fetchRemove(coll_name) orelse return false;
|
||||
self.free_collection(removed.value);
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -943,6 +996,7 @@ pub const Engine = struct {
|
||||
var removed = self.dbs.fetchRemove(db_name) orelse return false;
|
||||
self.free_db(&removed.value);
|
||||
self.gpa.free(removed.key);
|
||||
_ = self.cursors.kill_namespace(self.io, db_name, null);
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -1156,7 +1210,8 @@ pub const Engine = struct {
|
||||
errdefer self.gpa.free(coll_key);
|
||||
const new_coll = try self.gpa.create(Collection);
|
||||
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);
|
||||
try db.collections.put(self.gpa, coll_key, 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 (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(
|
||||
@@ -3886,3 +3948,64 @@ const Reader = struct {
|
||||
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();
|
||||
}
|
||||
|
||||
463
src/index.zig
463
src/index.zig
@@ -106,7 +106,9 @@ const page_size = 4096;
|
||||
/// Bytes of node payload: a 32-byte header plus the slotted region.
|
||||
const page_data = page_size - 32;
|
||||
/// 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
|
||||
/// least its own size. Bounds the split scratch.
|
||||
const max_slots = page_data / slot_size;
|
||||
@@ -234,6 +236,19 @@ pub const Index = struct {
|
||||
depth: u32,
|
||||
/// Total entries, maintained incrementally.
|
||||
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.
|
||||
scratch: [page_data]u8,
|
||||
/// Promoted-key scratch: inline keys being propagated up a split are
|
||||
@@ -268,6 +283,7 @@ pub const Index = struct {
|
||||
.leaf_count = 0,
|
||||
.depth = 0,
|
||||
.entry_count = 0,
|
||||
.epoch = 0,
|
||||
.scratch = undefined,
|
||||
.promo = undefined,
|
||||
};
|
||||
@@ -614,6 +630,10 @@ pub const Index = struct {
|
||||
self.depth = 0;
|
||||
self.entry_count = 0;
|
||||
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.
|
||||
@@ -777,6 +797,14 @@ pub const Index = struct {
|
||||
leaf: 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 {
|
||||
const ix = self.ix;
|
||||
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.
|
||||
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 {
|
||||
const ix = self.ix;
|
||||
while (self.leaf != 0) {
|
||||
@@ -919,6 +954,174 @@ pub const Index = struct {
|
||||
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 ------------------------------------------------------
|
||||
|
||||
/// The canonical spec document bytes
|
||||
@@ -1710,6 +1913,9 @@ pub const Index = struct {
|
||||
self.first_leaf = self.root;
|
||||
self.leaf_count = 1;
|
||||
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
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
/// 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);
|
||||
}
|
||||
|
||||
@@ -10,6 +10,7 @@ pub const db = @import("db.zig");
|
||||
pub const query = @import("query.zig");
|
||||
pub const update = @import("update.zig");
|
||||
pub const index = @import("index.zig");
|
||||
pub const cursor = @import("cursor.zig");
|
||||
pub const pager = @import("pager.zig");
|
||||
|
||||
test {
|
||||
@@ -23,5 +24,6 @@ test {
|
||||
_ = @import("query.zig");
|
||||
_ = @import("update.zig");
|
||||
_ = @import("index.zig");
|
||||
_ = @import("cursor.zig");
|
||||
_ = @import("pager.zig");
|
||||
}
|
||||
|
||||
142
src/main.zig
142
src/main.zig
@@ -10,6 +10,18 @@ const usage =
|
||||
\\ --db <path> database file (default multiforadb.log)
|
||||
\\ --ttl-sweep-secs <n>
|
||||
\\ 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>
|
||||
\\ minimum log bytes between compactions; suffixes k/m/g
|
||||
\\ (default 16m). The actual trigger also scales with the
|
||||
@@ -44,34 +56,77 @@ fn parse_size_suffix(v: []const u8) ?u64 {
|
||||
return n * mult;
|
||||
}
|
||||
|
||||
pub fn main(init: std.process.Init) !void {
|
||||
var port: u16 = 27017;
|
||||
var bind_ip: []const u8 = "127.0.0.1";
|
||||
var db_path: []const u8 = "multiforadb.log";
|
||||
var ttl_sweep_secs: i64 = 60;
|
||||
var compact_threshold: u64 = 16 * 1024 * 1024;
|
||||
/// Everything the CLI can set. Parsed apart from `main` so the option table has
|
||||
/// room to grow without main outgrowing the 70-line limit.
|
||||
const Options = struct {
|
||||
port: u16 = 27017,
|
||||
bind_ip: []const u8 = "127.0.0.1",
|
||||
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);
|
||||
defer it.deinit();
|
||||
_ = it.next(); // program name
|
||||
while (it.next()) |arg| {
|
||||
if (std.mem.eql(u8, arg, "--port")) {
|
||||
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});
|
||||
return error.InvalidPort;
|
||||
};
|
||||
} 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")) {
|
||||
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")) {
|
||||
const v = it.next() orelse return error.MissingValue;
|
||||
// i64 is the width std.Io.Duration.fromSeconds takes, so the
|
||||
// value reaches the sweeper without a cast; negatives are the
|
||||
// only thing parseInt would otherwise let through.
|
||||
ttl_sweep_secs = std.fmt.parseInt(i64, v, 10) catch -1;
|
||||
if (ttl_sweep_secs < 0) {
|
||||
// i64 is the width std.Io.Duration.fromSeconds takes, so the value
|
||||
// reaches the sweeper without a cast; negatives are the only thing
|
||||
// parseInt would otherwise let through.
|
||||
o.ttl_sweep_secs = std.fmt.parseInt(i64, v, 10) catch -1;
|
||||
if (o.ttl_sweep_secs < 0) {
|
||||
std.debug.print("multiforadb: invalid ttl sweep interval '{s}'\n", .{v});
|
||||
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", .{});
|
||||
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")) {
|
||||
try std.Io.File.writeStreamingAll(.stdout(), init.io, usage);
|
||||
return;
|
||||
return null;
|
||||
} else if (try parse_cursor_flag(arg, &it, &o)) {
|
||||
// Handled: one of the cursor-registry flags.
|
||||
} else {
|
||||
std.debug.print("multiforadb: unknown option '{s}'\n{s}", .{ arg, usage });
|
||||
return error.UnknownOption;
|
||||
}
|
||||
}
|
||||
|
||||
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();
|
||||
engine.compact_threshold = compact_threshold;
|
||||
std.debug.print("multiforadb: opened database '{s}' (compact threshold {d})\n", .{
|
||||
db_path,
|
||||
compact_threshold,
|
||||
});
|
||||
|
||||
var server = mongo.server.Server{
|
||||
.gpa = init.gpa,
|
||||
.port = port,
|
||||
.bind_ip = bind_ip,
|
||||
.oid_gen = oid_gen,
|
||||
.connection_counter = .init(1),
|
||||
.engine = &engine,
|
||||
.start_time = std.Io.Timestamp.now(init.io, .real),
|
||||
.ttl_sweep_secs = ttl_sweep_secs,
|
||||
};
|
||||
try server.run();
|
||||
return o;
|
||||
}
|
||||
|
||||
/// The cursor-registry flags, grouped so `parse_args` stays one flat table of
|
||||
/// options. Returns whether `arg` was one of them; consumes its value if so.
|
||||
fn parse_cursor_flag(arg: []const u8, it: *std.process.Args.Iterator, o: *Options) !bool {
|
||||
if (std.mem.eql(u8, arg, "--cursor-timeout-ms")) {
|
||||
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) {
|
||||
std.debug.print("multiforadb: invalid cursor timeout '{s}'\n", .{v});
|
||||
return error.InvalidCursorTimeout;
|
||||
}
|
||||
} else if (std.mem.eql(u8, arg, "--cursor-sweep-secs")) {
|
||||
const v = it.next() orelse return error.MissingValue;
|
||||
o.cursor_sweep_secs = std.fmt.parseInt(i64, v, 10) catch -1;
|
||||
if (o.cursor_sweep_secs < 0) {
|
||||
std.debug.print("multiforadb: invalid cursor sweep interval '{s}'\n", .{v});
|
||||
return error.InvalidCursorSweepSecs;
|
||||
}
|
||||
} 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;
|
||||
}
|
||||
|
||||
@@ -20,6 +20,9 @@ pub const Server = struct {
|
||||
/// because that is what std.Io.Duration.fromSeconds takes — the CLI
|
||||
/// rejects negatives.
|
||||
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 {
|
||||
// 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 `group.cancel` above stops it with everything else.
|
||||
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) {
|
||||
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`.
|
||||
fn handle_connection(io: std.Io, stream: std.Io.net.Stream, server: *Server) error{Canceled}!void {
|
||||
handle_connection_inner(io, stream, server) catch {};
|
||||
|
||||
41
src/wire.zig
41
src/wire.zig
@@ -299,9 +299,20 @@ fn begin_message(
|
||||
}
|
||||
|
||||
/// 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 {
|
||||
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);
|
||||
}
|
||||
|
||||
@@ -411,3 +422,31 @@ test "parse OP_QUERY handshake" {
|
||||
try testing.expectEqual(@as(i32, op_code_query), msg.op_code);
|
||||
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);
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user