diff --git a/src/db.zig b/src/db.zig index 826b21a..853f150 100644 --- a/src/db.zig +++ b/src/db.zig @@ -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(); +}