//! In-memory database engine backed by the append-only log. Maps //! db -> collection -> _id(serialized) -> owned Document. All mutations are //! logged and synced before they become visible in memory, so a crash never //! loses a committed write. Callers must hold the write lock (`lock`) around //! any command that mutates state, and the read lock (`lock_read`) around //! read-only commands so reads overlap with each other. const std = @import("std"); const bson = @import("bson.zig"); const storage = @import("storage.zig"); pub const Collection = struct { docs: std.StringHashMapUnmanaged(*bson.Document), fn init() Collection { return .{ .docs = .empty }; } }; pub const Db = struct { collections: std.StringHashMapUnmanaged(Collection), }; pub const Engine = struct { gpa: std.mem.Allocator, io: std.Io, // One writer at a time (log append + fsync, map mutation); many // concurrent readers (find/count/aggregate scans). Writer-preferring: // a queued writer blocks new readers rather than starving. rwlock: std.Io.RwLock, log: storage.Log, dbs: std.StringHashMapUnmanaged(Db), seq: u64, compact_threshold: u64, pub fn open(gpa: std.mem.Allocator, io: std.Io, path: []const u8) !Engine { var engine = Engine{ .gpa = gpa, .io = io, .rwlock = .init, .log = try storage.Log.open(gpa, io, path), .dbs = .empty, .seq = 0, .compact_threshold = 16 * 1024 * 1024, }; errdefer { engine.log.close(); engine.dbs.deinit(gpa); } try engine.log.replay(&engine, apply_record); return engine; } pub fn deinit(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| { var doc_it = coll_entry.value_ptr.docs.iterator(); while (doc_it.next()) |doc_entry| { doc_entry.value_ptr.*.deinit(); self.gpa.destroy(doc_entry.value_ptr.*); self.gpa.free(doc_entry.key_ptr.*); } coll_entry.value_ptr.docs.deinit(self.gpa); self.gpa.free(coll_entry.key_ptr.*); } db_entry.value_ptr.collections.deinit(self.gpa); self.gpa.free(db_entry.key_ptr.*); } self.dbs.deinit(self.gpa); self.log.close(); } // -- commands (callers must hold the matching lock) --------------------- /// Exclusive lock: for commands that mutate the engine. pub fn lock(self: *Engine) !void { try self.rwlock.lock(self.io); } pub fn unlock(self: *Engine) void { self.rwlock.unlock(self.io); } /// Shared lock: for read-only commands (find, count, aggregate, list*). /// Multiple readers may hold it simultaneously; writers wait for them. pub fn lock_read(self: *Engine) !void { try self.rwlock.lockShared(self.io); } pub fn unlock_read(self: *Engine) void { self.rwlock.unlockShared(self.io); } /// Insert a document. Fails with error.DuplicateKey if the _id exists. /// Generates an ObjectId _id when absent. pub fn insert(self: *Engine, db_name: []const u8, coll_name: []const u8, doc: *const bson.Document, oid_gen: *bson.ObjectIdGen) !void { const coll = try self.get_or_create_collection(db_name, coll_name); const owned = try self.own_with_id(doc, oid_gen); errdefer { owned.deinit(); self.gpa.destroy(owned); } const id_value = owned.get("_id") orelse unreachable; const id_key = try bson.serialize_value(self.gpa, id_value); defer self.gpa.free(id_key); if (coll.docs.contains(id_key)) return error.DuplicateKey; 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); const key_owned = try self.gpa.dupe(u8, id_key); try coll.docs.put(self.gpa, key_owned, owned); try self.maybe_compact(); } /// Insert or replace a document by _id (upsert without existence check). pub fn replace(self: *Engine, db_name: []const u8, coll_name: []const u8, doc: *const bson.Document, oid_gen: *bson.ObjectIdGen) !void { const coll = try self.get_or_create_collection(db_name, coll_name); const owned = try self.own_with_id(doc, oid_gen); errdefer { owned.deinit(); self.gpa.destroy(owned); } const id_value = owned.get("_id") orelse unreachable; const id_key = try bson.serialize_value(self.gpa, id_value); defer self.gpa.free(id_key); 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); if (coll.docs.fetchRemove(id_key)) |old| { old.value.*.deinit(); self.gpa.destroy(old.value); self.gpa.free(old.key); } const key_owned = try self.gpa.dupe(u8, id_key); try coll.docs.put(self.gpa, key_owned, owned); try self.maybe_compact(); } /// Remove a document by _id. Returns true if it existed. pub fn remove(self: *Engine, db_name: []const u8, coll_name: []const u8, id_key: []const u8) !bool { const db = self.dbs.get(db_name) orelse return false; const coll = db.collections.getPtr(coll_name) orelse return false; const doc = coll.docs.get(id_key) orelse return false; // Log (and sync) the delete before removing it from memory, so the // log always describes at least as much as the in-memory state. const doc_bytes = try serialize_doc(self.gpa, doc); defer self.gpa.free(doc_bytes); self.seq += 1; try self.log.append_delete(db_name, coll_name, doc_bytes, self.seq); const removed = coll.docs.fetchRemove(id_key) orelse unreachable; removed.value.*.deinit(); self.gpa.destroy(removed.value); self.gpa.free(removed.key); return true; } pub fn get_collection(self: *Engine, db_name: []const u8, coll_name: []const u8) ?*Collection { const db = self.dbs.get(db_name) orelse return null; return db.collections.getPtr(coll_name); } pub fn get_doc(self: *Engine, db_name: []const u8, coll_name: []const u8, id_key: []const u8) ?*const bson.Document { const coll = self.get_collection(db_name, coll_name) orelse return null; return coll.docs.get(id_key); } pub fn drop_collection(self: *Engine, db_name: []const u8, coll_name: []const u8) !bool { const db = self.dbs.getPtr(db_name) orelse return false; var removed = db.collections.fetchRemove(coll_name) orelse return false; var doc_it = removed.value.docs.iterator(); while (doc_it.next()) |doc_entry| { doc_entry.value_ptr.*.deinit(); self.gpa.destroy(doc_entry.value_ptr.*); self.gpa.free(doc_entry.key_ptr.*); } removed.value.docs.deinit(self.gpa); self.gpa.free(removed.key); return true; } pub fn drop_database(self: *Engine, db_name: []const u8) !bool { var removed = self.dbs.fetchRemove(db_name) orelse return false; var coll_it = removed.value.collections.iterator(); while (coll_it.next()) |coll_entry| { var docs_it = coll_entry.value_ptr.docs.iterator(); while (docs_it.next()) |doc_entry| { doc_entry.value_ptr.*.deinit(); self.gpa.destroy(doc_entry.value_ptr.*); self.gpa.free(doc_entry.key_ptr.*); } coll_entry.value_ptr.docs.deinit(self.gpa); self.gpa.free(coll_entry.key_ptr.*); } removed.value.collections.deinit(self.gpa); self.gpa.free(removed.key); 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.*); } pub fn collection_names(self: *Engine, db_name: []const u8, out: *std.ArrayListUnmanaged([]const u8)) !void { const db = self.dbs.get(db_name) orelse return; var it = db.collections.iterator(); while (it.next()) |entry| try out.append(self.gpa, entry.key_ptr.*); } // -- internals ----------------------------------------------------------- pub fn get_or_create_collection(self: *Engine, db_name: []const u8, coll_name: []const u8) !*Collection { const db = self.dbs.getPtr(db_name) orelse { const db_key = try self.gpa.dupe(u8, db_name); errdefer self.gpa.free(db_key); try self.dbs.put(self.gpa, db_key, .{ .collections = .empty }); return self.get_or_create_collection(db_name, coll_name); }; if (db.collections.getPtr(coll_name)) |coll| return coll; const coll_key = try self.gpa.dupe(u8, coll_name); errdefer self.gpa.free(coll_key); try db.collections.put(self.gpa, coll_key, Collection.init()); return db.collections.getPtr(coll_name) orelse unreachable; } /// Deep-copy a document into engine-owned storage, prepending a /// generated ObjectId `_id` when absent. fn own_with_id(self: *Engine, doc: *const bson.Document, oid_gen: *bson.ObjectIdGen) !*bson.Document { var pairs: std.ArrayListUnmanaged(bson.Pair) = .empty; defer pairs.deinit(self.gpa); if (doc.get("_id") == null) { const oid = oid_gen.new(self.io); try pairs.append(self.gpa, .{ .key = "_id", .value = .{ .object_id = oid } }); } try pairs.appendSlice(self.gpa, doc.pairs); var out: std.ArrayListUnmanaged(u8) = .empty; defer out.deinit(self.gpa); try bson.write_doc(pairs.items, self.gpa, &out); const owned = try self.gpa.create(bson.Document); errdefer self.gpa.destroy(owned); owned.* = try bson.Document.parse(self.gpa, out.items); return owned; } fn maybe_compact(self: *Engine) !void { if (self.log.log_bytes < self.compact_threshold) return; try self.compact(); } /// Rewrite the log with only live documents, atomically swapping the file. /// Callers must hold the write lock. pub fn compact(self: *Engine) !void { const tmp_path = try std.fmt.allocPrint(self.gpa, "{s}.tmp", .{self.log.path}); defer self.gpa.free(tmp_path); std.Io.Dir.cwd().deleteFile(self.io, tmp_path) catch {}; var new_log = try storage.Log.open(self.gpa, self.io, tmp_path); defer new_log.close(); 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| { 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.*); defer self.gpa.free(doc_bytes); try new_log.append_upsert(db_entry.key_ptr.*, coll_entry.key_ptr.*, doc_bytes, self.seq); } } } const new_end_pos = new_log.end_pos; try std.Io.Dir.renameAbsolute(tmp_path, self.log.path, self.io); // Persist the rename: fsync the parent directory so the new // directory entry survives a power loss right after compaction. const parent = parent_dir(self.log.path); var dir_file = try std.Io.Dir.cwd().openFile(self.io, parent, .{ .mode = .read_only, .allow_directory = true }); defer dir_file.close(self.io); try dir_file.sync(self.io); const old_path = try self.gpa.dupe(u8, self.log.path); self.log.close(); self.log = try storage.Log.open(self.gpa, self.io, old_path); // Log.open starts at end_pos 0 and does not replay; continue appending // where the compacted file actually ends. self.log.end_pos = new_end_pos; self.gpa.free(old_path); } }; fn parent_dir(path: []const u8) []const u8 { const last = std.mem.lastIndexOfScalar(u8, path, '/') orelse return "."; if (last == 0) return "/"; return path[0..last]; } fn serialize_doc(gpa: std.mem.Allocator, doc: *const bson.Document) ![]u8 { var out: std.ArrayListUnmanaged(u8) = .empty; errdefer out.deinit(gpa); try doc.to_bytes(gpa, &out); return out.toOwnedSlice(gpa); } fn apply_record(ctx: *anyopaque, record: storage.Record, doc: *bson.Document) anyerror!void { const self: *Engine = @ptrCast(@alignCast(ctx)); var stored = false; defer if (!stored) { doc.deinit(); self.gpa.destroy(doc); }; const id_value = doc.get("_id") orelse { std.debug.print("mongo-light: log record without _id, skipping\n", .{}); return; }; const id_key = try bson.serialize_value(self.gpa, id_value); defer 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 => { if (coll.docs.fetchRemove(id_key)) |old| { old.value.*.deinit(); self.gpa.destroy(old.value); self.gpa.free(old.key); } const key_owned = try self.gpa.dupe(u8, id_key); try coll.docs.put(self.gpa, key_owned, doc); stored = true; }, storage.record_type_delete => { if (coll.docs.fetchRemove(id_key)) |old| { old.value.*.deinit(); self.gpa.destroy(old.value); self.gpa.free(old.key); } }, else => {}, } } // --------------------------------------------------------------------------- // Tests // --------------------------------------------------------------------------- const testing = std.testing; const TmpLog = struct { tmp: std.testing.TmpDir, path: []u8, fn init(gpa: std.mem.Allocator) !TmpLog { const tmp = std.testing.tmpDir(.{}); const path = try std.fmt.allocPrint(gpa, ".zig-cache/tmp/{s}/test.log", .{tmp.sub_path}); return .{ .tmp = tmp, .path = path }; } fn deinit(self: *TmpLog, gpa: std.mem.Allocator) void { self.tmp.cleanup(); gpa.free(self.path); } }; fn test_env(threaded: *std.Io.Threaded) struct { io: std.Io, gen: bson.ObjectIdGen } { const io = threaded.io(); const gen = bson.ObjectIdGen.init(io); return .{ .io = io, .gen = gen }; } fn make_doc(gpa: std.mem.Allocator, id: i32, name: []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, "name"), .value = .{ .string = try arena.allocator().dupe(u8, name) } }; return .{ .arena = arena, .pairs = pairs }; } test "insert, query, remove" { 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 d1 = try make_doc(gpa, 1, "alice"); defer d1.deinit(); var d2 = try make_doc(gpa, 2, "bob"); defer d2.deinit(); try engine.lock(); try engine.insert("app", "users", &d1, &env.gen); try engine.insert("app", "users", &d2, &env.gen); engine.unlock(); // duplicate key var d3 = try make_doc(gpa, 1, "alice2"); defer d3.deinit(); try engine.lock(); try testing.expectError(error.DuplicateKey, engine.insert("app", "users", &d3, &env.gen)); engine.unlock(); // find by id const id_key = try bson.serialize_value(gpa, bson.Value{ .int32 = 2 }); defer gpa.free(id_key); try engine.lock(); const found = engine.get_doc("app", "users", id_key).?; try testing.expectEqualStrings("bob", found.get("name").?.string); const removed = try engine.remove("app", "users", id_key); try testing.expect(removed); engine.unlock(); } test "reopen replays log" { 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 d1 = try make_doc(gpa, 1, "alice"); defer d1.deinit(); var d2 = try make_doc(gpa, 2, "bob"); defer d2.deinit(); try engine.lock(); try engine.insert("app", "users", &d1, &env.gen); try engine.insert("app", "users", &d2, &env.gen); const id_key2 = try bson.serialize_value(gpa, bson.Value{ .int32 = 2 }); defer gpa.free(id_key2); _ = try engine.remove("app", "users", id_key2); 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", "users", id_key) == null); const id_key1 = try bson.serialize_value(gpa, bson.Value{ .int32 = 1 }); defer gpa.free(id_key1); try testing.expectEqualStrings("alice", engine2.get_doc("app", "users", id_key1).?.get("name").?.string); engine2.unlock(); } test "auto _id generation 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 doc = try make_doc(gpa, 0, "no-id-here"); defer doc.deinit(); // strip _id const stripped = doc.pairs[1..]; var arena = std.heap.ArenaAllocator.init(gpa); defer arena.deinit(); var d2 = try bson.Document.alloc(gpa, try arena.allocator().dupe(bson.Pair, stripped)); defer d2.deinit(); try engine.lock(); try engine.insert("app", "no_ids", &d2, &env.gen); engine.unlock(); } var engine2 = try Engine.open(gpa, io, tmp.path); defer engine2.deinit(); const coll = engine2.get_collection("app", "no_ids").?; var it = coll.docs.iterator(); var count: usize = 0; while (it.next()) |entry| { count += 1; try testing.expect(entry.value_ptr.*.get("_id").?.object_id.len == 12); } try testing.expectEqual(@as(usize, 1), count); } test "compaction rewrites log and keeps data" { 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); engine.compact_threshold = 1; // always compact defer engine.deinit(); var docs: [4]bson.Document = undefined; defer for (&docs) |*d| d.deinit(); try engine.lock(); for (0..4) |i| { docs[i] = try make_doc(gpa, @intCast(i + 1), "user-{d}"); try engine.insert("app", "users", &docs[i], &env.gen); } engine.unlock(); } // Reopen after compaction and keep writing: with the log reopened at // end_pos 0, appends would clobber the compacted records. var engine2 = try Engine.open(gpa, io, tmp.path); defer engine2.deinit(); try engine2.lock(); var extra = try make_doc(gpa, 5, "eve"); defer extra.deinit(); try engine2.insert("app", "users", &extra, &env.gen); engine2.unlock(); var engine3 = try Engine.open(gpa, io, tmp.path); defer engine3.deinit(); try engine3.lock(); for (1..6) |i| { const id_key = try bson.serialize_value(gpa, bson.Value{ .int32 = @intCast(i) }); defer gpa.free(id_key); try testing.expect(engine3.get_doc("app", "users", id_key) != null); } engine3.unlock(); }