db: engine maintenance for secondary indexes

Collection gains an indexes list; evict_doc removes entries at the single
document-death chokepoint; upsert does build → check → reserve → log →
evict → publish so entry insertion after the append is infallible and a
rejected unique write never reaches the log. Replay registers empty
indexes from create/drop records (types 3/4) and Engine.open rebuilds them
from live docs. compact re-emits index-create records. Engine gains
create_index/drop_index and dup_index for E11000 naming.
This commit is contained in:
2026-08-02 12:25:03 +03:00
parent a38ddc2f50
commit a733fc1993

View File

@@ -8,12 +8,14 @@
const std = @import("std");
const bson = @import("bson.zig");
const storage = @import("storage.zig");
const index = @import("index.zig");
pub const Collection = struct {
docs: std.StringHashMapUnmanaged(*bson.Document),
indexes: std.ArrayListUnmanaged(index.Index),
fn init() Collection {
return .{ .docs = .empty };
return .{ .docs = .empty, .indexes = .empty };
}
};
@@ -32,6 +34,10 @@ pub const Engine = struct {
dbs: std.StringHashMapUnmanaged(Db),
seq: u64,
compact_threshold: u64,
/// Set to the failing index's own stable name when an upsert is
/// rejected by a unique secondary index (error.DuplicateKeyIndex). The
/// command reads it while still holding the write lock.
dup_index: ?[]const u8 = null,
pub fn open(gpa: std.mem.Allocator, io: std.Io, path: []const u8) !Engine {
var engine = Engine{
@@ -49,6 +55,9 @@ pub const Engine = struct {
}
try engine.log.replay(&engine, apply_record);
// Replay registers empty indexes; build them from the live docs
// once replay completes (order-independent).
try engine.build_all_indexes();
return engine;
}
@@ -62,8 +71,12 @@ pub const Engine = struct {
self.log.close();
}
/// Free every document in a collection along with its owned _id keys.
/// Free every document in a collection along with its owned _id keys
/// and secondary indexes (whose entries alias the documents — freed
/// first).
fn free_collection(self: *Engine, coll: *Collection) void {
for (coll.indexes.items) |*ix| ix.deinit(self.gpa);
coll.indexes.deinit(self.gpa);
var doc_it = coll.docs.iterator();
while (doc_it.next()) |doc_entry| {
doc_entry.value_ptr.*.deinit();
@@ -84,8 +97,12 @@ pub const Engine = struct {
}
/// Drop the document stored under `id_key`, freeing it and its key.
/// No-op when the id is absent.
/// No-op when the id is absent. This is the single chokepoint where a
/// document dies, so index entries are removed here — before
/// old.value.deinit() and gpa.free(old.key) — keeping the entry aliasing
/// (values into the document arena, id into the docs map key) safe.
fn evict_doc(self: *Engine, coll: *Collection, id_key: []const u8) void {
for (coll.indexes.items) |*ix| ix.remove_id(self.gpa, id_key);
const old = coll.docs.fetchRemove(id_key) orelse return;
old.value.*.deinit();
self.gpa.destroy(old.value);
@@ -124,6 +141,13 @@ pub const Engine = struct {
return self.upsert(db_name, coll_name, doc, oid_gen, .replace);
}
/// One document's built entries for one index, tracked so a failure
/// anywhere before the log append frees them all.
const Built = struct {
built: index.BuiltEntries,
ix: *index.Index,
};
/// Shared body of `insert` and `replace`: they differ only in how an
/// existing _id is treated. Logs (and syncs) the new document before it
/// becomes visible in memory.
@@ -147,16 +171,58 @@ pub const Engine = struct {
self.gpa.destroy(owned);
self.gpa.free(id_key);
};
self.dup_index = null;
// 1. Build entries for every index. ParallelArrays escapes here,
// before anything is logged or mutated.
var built_list: std.ArrayListUnmanaged(Built) = .empty;
defer {
for (built_list.items) |*b| b.built.deinit(self.gpa);
built_list.deinit(self.gpa);
}
for (coll.indexes.items) |*ix| {
var built = try ix.build_entries(self.gpa, owned, id_key);
built_list.append(self.gpa, .{ .built = built, .ix = ix }) catch |err| {
built.deinit(self.gpa);
return err;
};
}
// 2. The _id check, mirroring the pre-index behavior.
if (mode == .insert and coll.docs.contains(id_key)) return error.DuplicateKey;
// 3. Unique secondary-index checks; a rejected write never reaches
// the log.
for (built_list.items) |*b| {
if (!b.ix.unique) continue;
b.ix.check_unique(b.built.entries.items, id_key) catch {
self.dup_index = b.ix.name;
return error.DuplicateKeyIndex;
};
}
// 4. Reserve entry capacity — the last fallible step, so the entry
// insertion after the log append is infallible.
for (built_list.items) |*b| {
if (b.built.entries.items.len == 0) continue;
try b.ix.entries.ensureUnusedCapacity(self.gpa, b.built.entries.items.len);
}
// 5. Log (and sync) before anything becomes visible.
const doc_bytes = try serialize_doc(self.gpa, owned);
defer self.gpa.free(doc_bytes);
self.seq += 1;
try self.log.append_upsert(db_name, coll_name, doc_bytes, self.seq);
// 6. Replace drops the old document (and its index entries).
if (mode == .replace) self.evict_doc(coll, id_key);
// 7. Publish the document and its entries.
try coll.docs.put(self.gpa, id_key, owned);
for (built_list.items) |*b| {
if (b.built.multikey) b.ix.multikey = true;
b.ix.insert_entries(self.gpa, &b.built);
}
stored = true;
try self.maybe_compact();
}
@@ -214,6 +280,93 @@ pub const Engine = struct {
return true;
}
/// Build and register a secondary index from a spec document
/// ({key, name, unique?, sparse?}). The create record is written only
/// after the index builds over the existing documents and passes
/// uniqueness, so a rejected create persists nothing. Returns the new
/// index (or the existing one when the spec matches — idempotent).
pub fn create_index(self: *Engine, db_name: []const u8, coll_name: []const u8, spec_doc: *const bson.Document) !*index.Index {
const coll = try self.get_or_create_collection(db_name, coll_name);
var ix = try index.parse_spec(self.gpa, spec_doc);
var committed = false;
errdefer if (!committed) ix.deinit(self.gpa);
for (coll.indexes.items) |*existing| {
if (std.mem.eql(u8, existing.name, ix.name)) {
if (index.Index.spec_equal(existing, &ix)) return existing;
return error.IndexOptionsConflict;
}
}
// Build entries over the existing documents, checking uniqueness as
// we go (the index is not exposed until the end, so mutating it is
// safe). On any failure the built entries are freed and nothing is
// persisted.
var built_list: std.ArrayListUnmanaged(index.BuiltEntries) = .empty;
defer {
for (built_list.items) |*b| b.deinit(self.gpa);
built_list.deinit(self.gpa);
}
var doc_it = coll.docs.iterator();
while (doc_it.next()) |entry| {
var built = try ix.build_entries(self.gpa, entry.value_ptr.*, entry.key_ptr.*);
built_list.append(self.gpa, built) catch |err| {
built.deinit(self.gpa);
return err;
};
if (built.multikey) ix.multikey = true;
if (ix.unique) {
try ix.check_unique(built.entries.items, entry.key_ptr.*);
}
if (built.entries.items.len == 0) continue;
try ix.entries.ensureUnusedCapacity(self.gpa, built.entries.items.len);
ix.insert_entries(self.gpa, &built);
}
// Reserve the collection slot, then persist and publish.
try coll.indexes.ensureUnusedCapacity(self.gpa, 1);
var spec_bytes: std.ArrayListUnmanaged(u8) = .empty;
defer spec_bytes.deinit(self.gpa);
try ix.write_spec(self.gpa, &spec_bytes);
self.seq += 1;
try self.log.append_index_create(db_name, coll_name, spec_bytes.items, self.seq);
coll.indexes.appendAssumeCapacity(ix);
committed = true;
return &coll.indexes.items[coll.indexes.items.len - 1];
}
/// Remove a secondary index by name, persisting a drop record first.
/// Returns false when no such index exists.
pub fn drop_index(self: *Engine, db_name: []const u8, coll_name: []const u8, index_name: []const u8) !bool {
const db = self.dbs.get(db_name) orelse return false;
const coll = db.collections.getPtr(coll_name) orelse return false;
var found = false;
for (coll.indexes.items) |ix| {
if (std.mem.eql(u8, ix.name, index_name)) {
found = true;
break;
}
}
if (!found) return false;
const name_pairs = [_]bson.Pair{.{ .key = "name", .value = .{ .string = index_name } }};
var name_doc: std.ArrayListUnmanaged(u8) = .empty;
defer name_doc.deinit(self.gpa);
try bson.write_doc(&name_pairs, self.gpa, &name_doc);
self.seq += 1;
try self.log.append_index_drop(db_name, coll_name, name_doc.items, self.seq);
var i: usize = 0;
while (i < coll.indexes.items.len) {
if (std.mem.eql(u8, coll.indexes.items[i].name, index_name)) {
var removed = coll.indexes.orderedRemove(i);
removed.deinit(self.gpa);
} else i += 1;
}
return true;
}
pub fn database_names(self: *Engine, out: *std.ArrayListUnmanaged([]const u8)) !void {
var it = self.dbs.iterator();
while (it.next()) |entry| try out.append(self.gpa, entry.key_ptr.*);
@@ -279,6 +432,15 @@ pub const Engine = struct {
while (db_it.next()) |db_entry| {
var coll_it = db_entry.value_ptr.collections.iterator();
while (coll_it.next()) |coll_entry| {
// Re-emit the index definitions first: a compacted log that
// dropped them would resurrect the collections without
// indexes on replay.
for (coll_entry.value_ptr.indexes.items) |*ix| {
var spec_bytes: std.ArrayListUnmanaged(u8) = .empty;
defer spec_bytes.deinit(self.gpa);
try ix.write_spec(self.gpa, &spec_bytes);
try new_log.append_index_create(db_entry.key_ptr.*, coll_entry.key_ptr.*, spec_bytes.items, self.seq);
}
var doc_it = coll_entry.value_ptr.docs.iterator();
while (doc_it.next()) |doc_entry| {
const doc_bytes = try serialize_doc(self.gpa, doc_entry.value_ptr.*);
@@ -305,6 +467,58 @@ pub const Engine = struct {
self.log.end_pos = new_end_pos;
self.gpa.free(old_path);
}
/// Rebuild every empty index from the live documents. Runs after replay
/// completes, so it is order-independent: a create record, the documents
/// it indexes, and any drop record all replay first. A duplicate under a
/// unique index logs a loud warning and keeps the index (still correct
/// as a candidate generator; future writes are still enforced) — the
/// database always opens, leaving dropIndexes as an in-band recovery
/// path.
fn build_all_indexes(self: *Engine) !void {
var db_it = self.dbs.iterator();
while (db_it.next()) |db_entry| {
var coll_it = db_entry.value_ptr.collections.iterator();
while (coll_it.next()) |coll_entry| {
for (coll_entry.value_ptr.indexes.items) |*ix| {
if (ix.entries.items.len > 0) continue; // defensive
var doc_it = coll_entry.value_ptr.docs.iterator();
while (doc_it.next()) |doc_entry| {
var built = ix.build_entries(self.gpa, doc_entry.value_ptr.*, doc_entry.key_ptr.*) catch |err| switch (err) {
error.ParallelArrays => {
std.debug.print("mongo-light: WARNING: index '{s}' cannot index an existing document; entry skipped\n", .{ix.name});
continue;
},
else => return err,
};
defer built.deinit(self.gpa);
if (built.multikey) ix.multikey = true;
if (ix.unique) {
ix.check_unique(built.entries.items, doc_entry.key_ptr.*) catch {
std.debug.print("mongo-light: WARNING: unique index '{s}' has duplicate keys in existing data; duplicates not enforced for existing documents\n", .{ix.name});
};
}
if (built.entries.items.len == 0) continue;
try ix.entries.ensureUnusedCapacity(self.gpa, built.entries.items.len);
ix.insert_entries(self.gpa, &built);
}
}
}
}
}
/// Register an (empty) index from a persisted spec document. A repeated
/// create record for the same name is an idempotent no-op.
fn register_index_from_spec(self: *Engine, coll: *Collection, spec_doc: *const bson.Document) !void {
var ix = try index.parse_spec(self.gpa, spec_doc);
var committed = false;
defer if (!committed) ix.deinit(self.gpa);
for (coll.indexes.items) |existing| {
if (std.mem.eql(u8, existing.name, ix.name)) return;
}
try coll.indexes.append(self.gpa, ix);
committed = true;
}
};
fn parent_dir(path: []const u8) []const u8 {
@@ -327,6 +541,38 @@ fn apply_record(ctx: *anyopaque, record: storage.Record, doc: *bson.Document) an
doc.deinit();
self.gpa.destroy(doc);
};
const coll = self.get_or_create_collection(record.db, record.coll) catch return;
// Index records carry no _id — handle them before the lookup. Replay
// registers indexes empty; Engine.open builds them from the live docs
// after replay completes.
switch (record.type) {
storage.record_type_index_create => {
self.register_index_from_spec(coll, doc) catch |err| {
std.debug.print("mongo-light: index create record failed to apply: {s}\n", .{@errorName(err)});
return;
};
return;
},
storage.record_type_index_drop => {
const name_value = doc.get("name") orelse return;
const name = switch (name_value) {
.string => |s| s,
else => return,
};
var i: usize = 0;
while (i < coll.indexes.items.len) {
if (std.mem.eql(u8, coll.indexes.items[i].name, name)) {
var removed = coll.indexes.orderedRemove(i);
removed.deinit(self.gpa);
} else i += 1;
}
return;
},
else => {},
}
const id_value = doc.get("_id") orelse {
std.debug.print("mongo-light: log record without _id, skipping\n", .{});
return;
@@ -335,8 +581,6 @@ fn apply_record(ctx: *anyopaque, record: storage.Record, doc: *bson.Document) an
var key_owned = false;
defer if (!key_owned) self.gpa.free(id_key);
const coll = self.get_or_create_collection(record.db, record.coll) catch return;
switch (record.type) {
storage.record_type_upsert => {
self.evict_doc(coll, id_key);
@@ -599,3 +843,245 @@ test "concurrent readers and writers on a threaded Io" {
try testing.expect(engine.get_doc("app", "users", id_key) != null);
}
}
// -- index tests -----------------------------------------------------------
/// A spec document for a single-path index, built by serializing and
/// re-parsing so the pairs are arena-owned.
fn index_spec(gpa: std.mem.Allocator, path: []const u8, name: []const u8, unique: bool, sparse: bool) !bson.Document {
var out: std.ArrayListUnmanaged(u8) = .empty;
defer out.deinit(gpa);
const pairs = [_]bson.Pair{
.{ .key = "key", .value = .{ .doc = &.{.{ .key = path, .value = .{ .int32 = 1 } }} } },
.{ .key = "name", .value = .{ .string = name } },
.{ .key = "unique", .value = .{ .bool = unique } },
.{ .key = "sparse", .value = .{ .bool = sparse } },
};
try bson.write_doc(&pairs, gpa, &out);
return bson.Document.parse(gpa, out.items);
}
/// Number of entries the named index has for a single-value equality key.
fn index_count(gpa: std.mem.Allocator, engine: *Engine, db_name: []const u8, coll_name: []const u8, name: []const u8, key_value: bson.Value) !usize {
const coll = engine.get_collection(db_name, coll_name) orelse return 0;
for (coll.indexes.items) |*ix| {
if (std.mem.eql(u8, ix.name, name)) {
var out: std.ArrayListUnmanaged([]const u8) = .empty;
defer out.deinit(gpa);
try ix.lookup_eq(gpa, &.{key_value}, &out);
return out.items.len;
}
}
return 0;
}
fn make_user(gpa: std.mem.Allocator, id: i32, email: []const u8) !bson.Document {
var arena = std.heap.ArenaAllocator.init(gpa);
errdefer arena.deinit();
const pairs = try arena.allocator().alloc(bson.Pair, 2);
pairs[0] = .{ .key = try arena.allocator().dupe(u8, "_id"), .value = .{ .int32 = id } };
pairs[1] = .{ .key = try arena.allocator().dupe(u8, "email"), .value = .{ .string = try arena.allocator().dupe(u8, email) } };
return .{ .arena = arena, .pairs = pairs };
}
test "unique index enforced on insert, replace, and upsert-conflict" {
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();
var spec = try index_spec(gpa, "email", "email_1", true, false);
defer spec.deinit();
try engine.lock();
_ = try engine.create_index("app", "users", &spec);
var d1 = try make_user(gpa, 1, "a@x.io");
defer d1.deinit();
try engine.insert("app", "users", &d1, &env.gen);
// A second doc with the same email is rejected and never logged.
var d2 = try make_user(gpa, 2, "a@x.io");
defer d2.deinit();
try testing.expectError(error.DuplicateKeyIndex, engine.insert("app", "users", &d2, &env.gen));
try testing.expectEqualStrings("email_1", engine.dup_index.?);
// A replace that keeps its own email is fine (own entries excluded).
var d1b = try make_user(gpa, 1, "a@x.io");
defer d1b.deinit();
try engine.replace("app", "users", &d1b, &env.gen);
try testing.expectEqual(@as(usize, 1), try index_count(gpa, &engine, "app", "users", "email_1", .{ .string = "a@x.io" }));
// An update that would collide is rejected.
var d2b = try make_user(gpa, 2, "a@x.io");
defer d2b.deinit();
try testing.expectError(error.DuplicateKeyIndex, engine.replace("app", "users", &d2b, &env.gen));
// A different email still inserts.
var d3 = try make_user(gpa, 3, "b@x.io");
defer d3.deinit();
try engine.insert("app", "users", &d3, &env.gen);
engine.unlock();
}
test "index maintained across update and delete" {
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();
var spec = try index_spec(gpa, "a", "a_1", false, false);
defer spec.deinit();
try engine.lock();
_ = try engine.create_index("app", "items", &spec);
var d1 = try doc_with_a(gpa, 1, 10);
defer d1.deinit();
var d2 = try doc_with_a(gpa, 2, 20);
defer d2.deinit();
try engine.insert("app", "items", &d1, &env.gen);
try engine.insert("app", "items", &d2, &env.gen);
try testing.expectEqual(@as(usize, 1), try index_count(gpa, &engine, "app", "items", "a_1", .{ .int32 = 20 }));
// Replace doc 1 with a new value: old entry gone, new entry present.
var d1b = try doc_with_a(gpa, 1, 30);
defer d1b.deinit();
try engine.replace("app", "items", &d1b, &env.gen);
try testing.expectEqual(@as(usize, 0), try index_count(gpa, &engine, "app", "items", "a_1", .{ .int32 = 10 }));
try testing.expectEqual(@as(usize, 1), try index_count(gpa, &engine, "app", "items", "a_1", .{ .int32 = 30 }));
// Delete doc 2: its entry is removed.
_ = try engine.remove_by_id("app", "items", .{ .int32 = 2 });
try testing.expectEqual(@as(usize, 0), try index_count(gpa, &engine, "app", "items", "a_1", .{ .int32 = 20 }));
engine.unlock();
}
/// A document with an integer `a` field (on top of _id + name).
fn doc_with_a(gpa: std.mem.Allocator, id: i32, a: i32) !bson.Document {
var arena = std.heap.ArenaAllocator.init(gpa);
errdefer arena.deinit();
const pairs = try arena.allocator().alloc(bson.Pair, 3);
pairs[0] = .{ .key = try arena.allocator().dupe(u8, "_id"), .value = .{ .int32 = id } };
pairs[1] = .{ .key = try arena.allocator().dupe(u8, "name"), .value = .{ .string = try arena.allocator().dupe(u8, "x") } };
pairs[2] = .{ .key = try arena.allocator().dupe(u8, "a"), .value = .{ .int32 = a } };
return .{ .arena = arena, .pairs = pairs };
}
test "index survives reopen and compaction" {
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();
engine.compact_threshold = 1; // every write compacts
var spec = try index_spec(gpa, "email", "email_1", false, false);
defer spec.deinit();
try engine.lock();
_ = try engine.create_index("app", "users", &spec);
var d1 = try make_user(gpa, 1, "a@x.io");
defer d1.deinit();
var d2 = try make_user(gpa, 2, "b@x.io");
defer d2.deinit();
try engine.insert("app", "users", &d1, &env.gen);
try engine.insert("app", "users", &d2, &env.gen);
engine.unlock();
}
// Reopen: the index (rebuilt from the compacted log) still finds docs.
var engine2 = try Engine.open(gpa, io, tmp.path);
defer engine2.deinit();
try engine2.lock();
try testing.expectEqual(@as(usize, 1), try index_count(gpa, &engine2, "app", "users", "email_1", .{ .string = "b@x.io" }));
engine2.unlock();
}
test "index drop survives reopen" {
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();
var spec = try index_spec(gpa, "email", "email_1", false, false);
defer spec.deinit();
try engine.lock();
_ = try engine.create_index("app", "users", &spec);
var d1 = try make_user(gpa, 1, "a@x.io");
defer d1.deinit();
try engine.insert("app", "users", &d1, &env.gen);
try testing.expect(try engine.drop_index("app", "users", "email_1"));
engine.unlock();
}
var engine2 = try Engine.open(gpa, io, tmp.path);
defer engine2.deinit();
try engine2.lock();
try testing.expectEqual(@as(usize, 0), engine2.get_collection("app", "users").?.indexes.items.len);
engine2.unlock();
}
test "drop_collection frees indexes; log without index records replays" {
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();
var spec = try index_spec(gpa, "email", "email_1", false, false);
defer spec.deinit();
try engine.lock();
_ = try engine.create_index("app", "users", &spec);
var d1 = try make_user(gpa, 1, "a@x.io");
defer d1.deinit();
try engine.insert("app", "users", &d1, &env.gen);
// Dropped in memory; free_collection releases the index memory
// (verified by testing.allocator at engine.deinit).
try testing.expect(try engine.drop_collection("app", "users"));
try testing.expect(engine.get_collection("app", "users") == null);
// A log that only ever contained plain upserts replays fine.
var d2 = try make_doc(gpa, 2, "bob");
defer d2.deinit();
try engine.insert("app", "plain", &d2, &env.gen);
engine.unlock();
}
var engine2 = try Engine.open(gpa, io, tmp.path);
defer engine2.deinit();
try engine2.lock();
const id_key = try bson.serialize_value(gpa, bson.Value{ .int32 = 2 });
defer gpa.free(id_key);
try testing.expect(engine2.get_doc("app", "plain", id_key) != null);
// Pre-existing limitation (documented in the README): drop_collection
// writes no log record, so the collection and its index resurrect.
const users = engine2.get_collection("app", "users").?;
try testing.expectEqual(@as(usize, 1), users.indexes.items.len);
try testing.expectEqual(@as(usize, 1), try index_count(gpa, &engine2, "app", "users", "email_1", .{ .string = "a@x.io" }));
engine2.unlock();
}