tickets
All repositories: gitoria
63.0 KB
// mpackdb DB engine — Zig port of mpackdb v1.0.7 (MPackDB.js / IndexManager.js)// semantics. File-level compatible with the JS implementation (D12)://// <base>.mpack append-only records: [4-byte LE total size][msgpack map]// <base>.meta.json { nextId?, deleted:[offsets], version, schema? }// <base>.<field>.txt sorted index lines "key,offset,length\n"// <base>.idxstate.json { coveredBytes }// <base>.lock lock file containing the holder's pid (O_EXCL protocol,// stale takeover after staleLockTimeout ms)//// Divergences from the JS implementation (documented in the 066 report):// - the in-memory index is the COMPLETE entry set (JS keeps on-disk + delta);// persist() writes the full set instead of read-merge-write. Sequential// cross-process use (close one side, open the other) behaves identically.// - string index keys sort bytewise (JS: String.localeCompare, ICU collation).// Equality lookups are unaffected; range ordering can differ for// mixed-case/accented keys.// - null field values are not indexed (JS indexes them under the key "null").// - external data-file replacement while open (inode change) is not detected.// - a fresh `*id` table's first id is 1 (JS: 0) — ticket #4, see FIRST_ID.const std = @import("std");const linux = std.os.linux;const msgpack = @import("msgpack.zig");pub const PkType = enum(u8) { number = 0, uuid = 1, string = 2 };/// A `*id` table's first auto-increment id (ticket #4, the creator's ruling:/// "start at 1"). Only a table WITHOUT a stored counter starts here: a/// .meta.json that carries `nextId` — every existing table, including one the/// JS mpackdb wrote — keeps counting from it. The file format is unchanged;/// the JS mpackdb 1.0.7 starts a fresh table at 0, which is the one divergence.pub const FIRST_ID: f64 = 1;pub const IdxType = enum(u8) { lexical = 0, numeric = 1 };pub const Error = error{NoDbFile,NoPrimaryKey,/// mission 077 (074 GAP 10): a query named a field with no index. There is/// no index to binary-search and mpackdb does not silently full-scan, so/// the answer used to be an empty result — indistinguishable from "no such/// record". It is a mistake, and it says so now.NoSuchIndex,DuplicateKey,/// ticket #29: update()'s record must carry the primary key. Without it/// the old record was deleted, the new one written without a key, and the/// table left corrupt (fetch null, a unique index still finding it).RecordLacksPrimaryKey,LockTimeout,CorruptRecord,IoError,OutOfMemory,InvalidFormat,Truncated,MapTooLarge,};pub const Key = union(enum) {num: f64,str: []const u8,};pub const Entry = struct {key: Key,off: u64,len: u64,};const FieldIndex = struct {typ: IdxType,entries: std.ArrayList(Entry) = .empty, // sorted by (key, off)};pub const Loc = struct { off: u64, len: u64 };pub const Record = struct {value: msgpack.Value,off: u64,len: u64,};// =========================================================================// Low-level file IO (raw linux syscalls — plugins avoid std.Io plumbing)// =========================================================================const fio = struct {fn openz(alloc: std.mem.Allocator, path: []const u8, flags: linux.O, mode: u32) ?i32 {const path_z = alloc.dupeZ(u8, path) catch return null;defer alloc.free(path_z);const rc = linux.open(path_z, flags, mode);const fd: i32 = @bitCast(@as(u32, @truncate(rc)));if (fd < 0) return null;return fd;}fn openRead(alloc: std.mem.Allocator, path: []const u8) ?i32 {return openz(alloc, path, .{}, 0);}fn openAppend(alloc: std.mem.Allocator, path: []const u8) ?i32 {return openz(alloc, path, .{ .ACCMODE = .WRONLY, .CREAT = true, .APPEND = true }, 0o644);}fn openTrunc(alloc: std.mem.Allocator, path: []const u8) ?i32 {return openz(alloc, path, .{ .ACCMODE = .WRONLY, .CREAT = true, .TRUNC = true }, 0o644);}/// O_CREAT|O_EXCL — returns null when the file already exists.fn openExcl(alloc: std.mem.Allocator, path: []const u8) ?i32 {return openz(alloc, path, .{ .ACCMODE = .WRONLY, .CREAT = true, .EXCL = true }, 0o644);}fn close(fd: i32) void {_ = linux.close(fd);}fn writeAll(fd: i32, bytes: []const u8) bool {var written: usize = 0;while (written < bytes.len) {const rc = linux.write(fd, bytes.ptr + written, bytes.len - written);if (@as(isize, @bitCast(rc)) <= 0) return false;written += rc;}return true;}fn pread(fd: i32, buf: []u8, offset: u64) ?usize {var got: usize = 0;while (got < buf.len) {const rc = linux.pread(fd, buf.ptr + got, buf.len - got, @intCast(offset + got));const n: isize = @bitCast(rc);if (n < 0) return null;if (n == 0) break;got += @intCast(n);}return got;}fn readAll(alloc: std.mem.Allocator, path: []const u8) ?[]u8 {const fd = openRead(alloc, path) orelse return null;defer close(fd);var content: std.ArrayList(u8) = .empty;var buf: [65536]u8 = undefined;while (true) {const rc = linux.read(fd, &buf, buf.len);if (@as(isize, @bitCast(rc)) <= 0) break;content.appendSlice(alloc, buf[0..rc]) catch return null;}return content.toOwnedSlice(alloc) catch null;}fn fileSize(alloc: std.mem.Allocator, path: []const u8) ?u64 {const path_z = alloc.dupeZ(u8, path) catch return null;defer alloc.free(path_z);var stx: linux.Statx = undefined;const rc = linux.statx(linux.AT.FDCWD, path_z, 0, linux.STATX.BASIC_STATS, &stx);if (rc != 0) return null;return stx.size;}fn mtimeMs(alloc: std.mem.Allocator, path: []const u8) ?i64 {const path_z = alloc.dupeZ(u8, path) catch return null;defer alloc.free(path_z);var stx: linux.Statx = undefined;const rc = linux.statx(linux.AT.FDCWD, path_z, 0, linux.STATX.BASIC_STATS, &stx);if (rc != 0) return null;return @as(i64, stx.mtime.sec) * 1000 + @divTrunc(@as(i64, stx.mtime.nsec), 1_000_000);}fn rename(alloc: std.mem.Allocator, old_path: []const u8, new_path: []const u8) bool {const old_z = alloc.dupeZ(u8, old_path) catch return false;defer alloc.free(old_z);const new_z = alloc.dupeZ(u8, new_path) catch return false;defer alloc.free(new_z);return linux.renameat(linux.AT.FDCWD, old_z, linux.AT.FDCWD, new_z) == 0;}fn unlink(alloc: std.mem.Allocator, path: []const u8) bool {const path_z = alloc.dupeZ(u8, path) catch return false;defer alloc.free(path_z);return linux.unlinkat(linux.AT.FDCWD, path_z, 0) == 0;}fn mkdirAll(alloc: std.mem.Allocator, dir_path: []const u8) void {if (dir_path.len == 0 or std.mem.eql(u8, dir_path, ".")) return;var i: usize = 1;while (i <= dir_path.len) : (i += 1) {if (i == dir_path.len or dir_path[i] == '/') {const part = alloc.dupeZ(u8, dir_path[0..i]) catch return;defer alloc.free(part);_ = linux.mkdirat(linux.AT.FDCWD, part, 0o755);}}}fn sleepMs(ms: u64) void {var req: linux.timespec = .{ .sec = @intCast(ms / 1000), .nsec = @intCast((ms % 1000) * 1_000_000) };var rem: linux.timespec = undefined;_ = linux.nanosleep(&req, &rem);}fn nowMs() i64 {var ts: linux.timespec = undefined;_ = linux.clock_gettime(linux.CLOCK.REALTIME, &ts);return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);}};// =========================================================================// JS-compatible helpers// =========================================================================/// Format an f64 the way JS String(n) does for the values mpackdb produces:/// integral values print without a decimal point. (Exotic floats may diverge/// from V8's shortest-round-trip formatting; index keys are ints/strings in/// practice.)pub fn jsNumFmt(buf: []u8, v: f64) []const u8 {if (v == @floor(v) and @abs(v) <= 9007199254740992.0) {const i: i64 = @intFromFloat(v);return std.fmt.bufPrint(buf, "{d}", .{i}) catch buf[0..0];}return std.fmt.bufPrint(buf, "{d}", .{v}) catch buf[0..0];}/// JS parseInt(s, 10) semantics: optional sign, leading digits, ignore rest./// Returns null for NaN.fn jsParseInt(s: []const u8) ?f64 {var i: usize = 0;while (i < s.len and (s[i] == ' ' or s[i] == '\t')) i += 1;var sign: f64 = 1;if (i < s.len and (s[i] == '+' or s[i] == '-')) {if (s[i] == '-') sign = -1;i += 1;}var got = false;var v: f64 = 0;while (i < s.len and s[i] >= '0' and s[i] <= '9') : (i += 1) {v = v * 10 + @as(f64, @floatFromInt(s[i] - '0'));got = true;}if (!got) return null;return sign * v;}/// _compareKeys: numbers numerically; otherwise stringified byte compare/// (JS uses localeCompare — see divergence note in the header).pub fn cmpKeys(a: Key, b: Key) i32 {if (a == .num and b == .num) {if (a.num < b.num) return -1;if (a.num > b.num) return 1;return 0;}var buf_a: [32]u8 = undefined;var buf_b: [32]u8 = undefined;const sa = if (a == .str) a.str else jsNumFmt(&buf_a, a.num);const sb = if (b == .str) b.str else jsNumFmt(&buf_b, b.num);return switch (std.mem.order(u8, sa, sb)) {.lt => -1,.eq => 0,.gt => 1,};}fn entryLess(_: void, a: Entry, b: Entry) bool {const c = cmpKeys(a.key, b.key);if (c != 0) return c < 0;return a.off < b.off;}/// mpack.js uuid(): 9-char base36 ms timestamp (padStart '0') + 3 base36 chars.pub fn genUuid(buf: *[12]u8) []const u8 {const digits = "0123456789abcdefghijklmnopqrstuvwxyz";const t: u64 = @intCast(fio.nowMs());var tmp: [16]u8 = undefined;var n: usize = 0;var v = t;while (v > 0) : (v /= 36) {tmp[n] = digits[@intCast(v % 36)];n += 1;}var i: usize = 0;while (i < 9) : (i += 1) {buf[8 - i] = if (i < n) tmp[i] else '0';}if (!rng_init) {var ts: linux.timespec = undefined;_ = linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts);rng = std.Random.DefaultPrng.init(@bitCast(@as(i64, ts.sec) *% 1_000_000_000 +% ts.nsec));rng_init = true;}var r = rng.random();for (0..3) |j| {buf[9 + j] = digits[r.uintLessThan(u8, 36)];}return buf[0..12];}var rng: std.Random.DefaultPrng = .{ .s = undefined };var rng_init: bool = false;// =========================================================================// The database// =========================================================================pub const OpenOptions = struct {primary_key: ?[]const u8 = null, // with */@/! prefixes, like the JS APIindexes: []const []const u8 = &.{}, // with prefixescompact: bool = true,stale_lock_timeout_ms: i64 = 30000,debug: bool = false,};/// The identity of a table on disk (ticket #110): dirname + basename with any/// extension stripped, exactly the computation `Db.open` uses to turn a given/// path into `<base>.mpack` — so two spellings of the same file (`x.db`, `x`)/// share one key. The plugin ABI layer (mpackdb.zig) keys its process-wide/// handle registry on this, so a table opened twice in one process — from any/// module instance or realm — shares the one live `Db` instead of each open/// racing the other's file state.pub fn tableKey(alloc: std.mem.Allocator, db_file: []const u8) Error![]u8 {const dir = std.fs.path.dirname(db_file) orelse ".";var base = std.fs.path.basename(db_file);if (std.mem.lastIndexOfScalar(u8, base, '.')) |dot| {if (dot > 0) base = base[0..dot];}return std.fmt.allocPrint(alloc, "{s}/{s}", .{ dir, base }) catch Error.OutOfMemory;}pub const Db = struct {arena_state: std.heap.ArenaAllocator,alloc: std.mem.Allocator, // arena — freed wholesale on close// THE BUFFERS OF ONE OPERATION ARE NOT THE HANDLE'S (hybriel#126): the meta read// back and rendered anew, the index files, the paths. On the arena they stayed for// the handle's life, and the meta grows with every tombstone, so each update leaked// more than the last (notes: 2000 updates of one note, +100 MB). They are freed now.tmp: std.mem.Allocator,// pathsdata_path: []const u8,meta_path: []const u8,lock_path: []const u8,idxstate_path: []const u8,// schemapk: ?[]const u8 = null,pk_type: PkType = .string,index_fields: std.ArrayList([]const u8) = .empty, // pk first (when set), like JSunique_fields: std.ArrayList([]const u8) = .empty,idx: std.StringArrayHashMapUnmanaged(FieldIndex) = .empty,// metanext_id: f64 = FIRST_ID,deleted: std.ArrayList(u64) = .empty,version: u64 = 0,covered_bytes: u64 = 0,compact_on_open: bool = true,stale_lock_timeout_ms: i64 = 30000,debug: bool = false,// last DuplicateKey detail for the plugin surfacelast_error_buf: [256]u8 = undefined,last_error: []const u8 = "",pub fn open(gpa: std.mem.Allocator, db_file: []const u8, opts: OpenOptions) Error!*Db {if (db_file.len == 0) return Error.NoDbFile;const self = gpa.create(Db) catch return Error.OutOfMemory;self.* = .{.arena_state = std.heap.ArenaAllocator.init(gpa),.alloc = undefined,.tmp = gpa,.data_path = undefined,.meta_path = undefined,.lock_path = undefined,.idxstate_path = undefined,};self.alloc = self.arena_state.allocator();errdefer {self.arena_state.deinit();gpa.destroy(self);}const a = self.alloc;self.compact_on_open = opts.compact;self.stale_lock_timeout_ms = opts.stale_lock_timeout_ms;self.debug = opts.debug;// dirname / basename (extension stripped, like JS basename(f, extname(f)))const dir = std.fs.path.dirname(db_file) orelse ".";var base = std.fs.path.basename(db_file);if (std.mem.lastIndexOfScalar(u8, base, '.')) |dot| {if (dot > 0) base = base[0..dot];}fio.mkdirAll(a, dir);self.data_path = std.fmt.allocPrint(a, "{s}/{s}.mpack", .{ dir, base }) catch return Error.OutOfMemory;self.meta_path = std.fmt.allocPrint(a, "{s}/{s}.meta.json", .{ dir, base }) catch return Error.OutOfMemory;self.lock_path = std.fmt.allocPrint(a, "{s}/{s}.lock", .{ dir, base }) catch return Error.OutOfMemory;self.idxstate_path = std.fmt.allocPrint(a, "{s}/{s}.idxstate.json", .{ dir, base }) catch return Error.OutOfMemory;// ---- parse primary key ----if (opts.primary_key) |pk_raw| {if (pk_raw.len > 0) {if (pk_raw[0] == '*') {self.pk = a.dupe(u8, pk_raw[1..]) catch return Error.OutOfMemory;self.pk_type = .number;} else if (pk_raw[0] == '@') {self.pk = a.dupe(u8, pk_raw[1..]) catch return Error.OutOfMemory;self.pk_type = .uuid;} else {self.pk = a.dupe(u8, pk_raw) catch return Error.OutOfMemory;self.pk_type = .string;}}}if (self.pk) |p| self.unique_fields.append(a, p) catch return Error.OutOfMemory;// ---- parse indexes (pk first, then configured; dedupe) ----if (self.pk) |p| {self.index_fields.append(a, p) catch return Error.OutOfMemory;const t: IdxType = if (self.pk_type == .number) .numeric else .lexical;self.idx.put(a, p, .{ .typ = t }) catch return Error.OutOfMemory;}for (opts.indexes) |raw| {var clean = raw;var is_unique = false;if (clean.len > 0 and clean[0] == '!') {is_unique = true;clean = clean[1..];}var typ: IdxType = .lexical;if (clean.len > 0 and clean[0] == '*') {typ = .numeric;clean = clean[1..];} else if (clean.len > 0 and clean[0] == '@') {clean = clean[1..];}if (self.idx.contains(clean)) continue;const owned = a.dupe(u8, clean) catch return Error.OutOfMemory;self.index_fields.append(a, owned) catch return Error.OutOfMemory;self.idx.put(a, owned, .{ .typ = typ }) catch return Error.OutOfMemory;if (is_unique) self.unique_fields.append(a, owned) catch return Error.OutOfMemory;}// ---- load meta ----self.loadMeta();// ---- compact on init (JS default; also creates an empty data file) ----// AN OPEN WITH NOTHING TO CHANGE WRITES NOTHING (ticket #21): a store// with no tombstones has nothing to compact, and rewriting it anyway// made opening a backup change it.var did_compact = false;const nothing_to_compact = self.deleted.items.len == 0 andfio.fileSize(self.tmp, self.data_path) != null;if (self.compact_on_open and !nothing_to_compact) {try self.acquireLock();const cr = self.compactLocked();self.releaseLock();try cr;did_compact = true;}// ---- persist schema into meta (JS: init-schema under lock).// compactLocked already wrote meta (with schema); skip the extra bump.if (!did_compact and (self.pk != null or self.index_fields.items.len > 0) and!self.metaOnDiskIsCurrent()){try self.acquireLock();const mr = self.persistMetaLocked();self.releaseLock();try mr;}// ---- indexes (JS: (re)build under the lock; force after compaction) ----if (self.index_fields.items.len > 0) {try self.acquireLock();const ir = self.initIndexes(did_compact);self.releaseLock();try ir;}return self;}pub fn close(self: *Db) void {// persist indexes + idxstate (JS: IndexManager.close → persist under lock)if (self.index_fields.items.len > 0) {if (self.acquireLock()) {self.persistIndexesLocked() catch {};self.releaseLock();} else |_| {}}const gpa = self.arena_state.child_allocator;self.arena_state.deinit();gpa.destroy(self);}fn dbg(self: *Db, comptime fmt: []const u8, args: anytype) void {if (self.debug) std.debug.print("[mpackdb] " ++ fmt ++ "\n", args);}// =====================================================================// Locking (JS _acquireFileLock protocol)// =====================================================================fn acquireLock(self: *Db) Error!void {var retries: u32 = 0;while (true) {if (fio.openExcl(self.tmp, self.lock_path)) |fd| {var pid_buf: [16]u8 = undefined;const pid_s = std.fmt.bufPrint(&pid_buf, "{d}", .{linux.getpid()}) catch "0";_ = fio.writeAll(fd, pid_s);fio.close(fd);return;}// stale lock takeoverif (self.stale_lock_timeout_ms > 0) {if (fio.mtimeMs(self.tmp, self.lock_path)) |mt| {if (fio.nowMs() - mt > self.stale_lock_timeout_ms) {_ = fio.unlink(self.tmp, self.lock_path);continue;}}}if (retries > 480) return Error.LockTimeout; // ~12s at 25msfio.sleepMs(25);retries += 1;}}fn releaseLock(self: *Db) void {_ = fio.unlink(self.tmp, self.lock_path);}// =====================================================================// Meta (.meta.json)// =====================================================================fn appendFmt(self: *Db, list: *std.ArrayList(u8), comptime fmt: []const u8, args: anytype) Error!void {const s = std.fmt.allocPrint(self.tmp, fmt, args) catch return Error.OutOfMemory;defer self.tmp.free(s);list.appendSlice(self.tmp, s) catch return Error.OutOfMemory;}fn loadMeta(self: *Db) void {const content = fio.readAll(self.tmp, self.meta_path) orelse return;defer self.tmp.free(content);self.applyMetaJson(content);}fn applyMetaJson(self: *Db, content: []const u8) void {const parsed = std.json.parseFromSlice(std.json.Value, self.tmp, content, .{}) catch return;defer parsed.deinit();if (parsed.value != .object) return;const obj = parsed.value.object;self.next_id = FIRST_ID;self.deleted.clearRetainingCapacity();self.version = 0;if (obj.get("nextId")) |v| {self.next_id = switch (v) {.integer => |i| @floatFromInt(i),.float => |f| f,else => FIRST_ID,};}if (obj.get("version")) |v| {if (v == .integer) self.version = @intCast(@max(v.integer, 0));}if (obj.get("deleted")) |v| {if (v == .array) {for (v.array.items) |it| {const off: u64 = switch (it) {.integer => |i| @intCast(@max(i, 0)),.float => |f| @intFromFloat(@max(f, 0)),else => continue,};self.deleted.append(self.alloc, off) catch {};}}}}/// Pick up other processes' persisted state (JS refresh()): adopt disk meta/// when its version is newer; index appended tail records.pub fn refresh(self: *Db) void {if (fio.readAll(self.tmp, self.meta_path)) |content| {defer self.tmp.free(content);const parsed = std.json.parseFromSlice(std.json.Value, self.tmp, content, .{}) catch return;defer parsed.deinit();if (parsed.value == .object) {var disk_version: u64 = 0;if (parsed.value.object.get("version")) |v| {if (v == .integer) disk_version = @intCast(@max(v.integer, 0));}if (disk_version > self.version) self.applyMetaJson(content);}}// catchUp: index records another process appendedif (self.index_fields.items.len > 0) {const size = fio.fileSize(self.tmp, self.data_path) orelse 0;if (size > self.covered_bytes) self.catchUp(size);}}fn persistMetaLocked(self: *Db) Error!void {self.version += 1;var out: std.ArrayList(u8) = .empty;defer out.deinit(self.tmp);try self.renderMeta(&out);// atomic tmp + rename (JS: `${metaPath}.${pid}.tmp`)var tmp_buf: [512]u8 = undefined;const tmp_path = std.fmt.bufPrint(&tmp_buf, "{s}.{d}.tmp", .{ self.meta_path, linux.getpid() }) catch return Error.IoError;const fd = fio.openTrunc(self.tmp, tmp_path) orelse return Error.IoError;const ok = fio.writeAll(fd, out.items);fio.close(fd);if (!ok) return Error.IoError;if (!fio.rename(self.tmp, tmp_path, self.meta_path)) return Error.IoError;}/// Does .meta.json already say, byte for byte, what this handle would write/// at its current version? Then an open has nothing to persist (ticket #21).fn metaOnDiskIsCurrent(self: *Db) bool {const disk = fio.readAll(self.tmp, self.meta_path) orelse return false;defer self.tmp.free(disk);var out: std.ArrayList(u8) = .empty;defer out.deinit(self.tmp);self.renderMeta(&out) catch return false;return std.mem.eql(u8, disk, out.items);}fn renderMeta(self: *Db, out: *std.ArrayList(u8)) Error!void {try self.appendFmt(out, "{{", .{});if (self.pk != null and self.pk_type == .number) {var nbuf: [32]u8 = undefined;try self.appendFmt(out, "\"nextId\":{s},", .{jsNumFmt(&nbuf, self.next_id)});}try self.appendFmt(out, "\"deleted\":[", .{});for (self.deleted.items, 0..) |off, i| {if (i > 0) try self.appendFmt(out, ",", .{});try self.appendFmt(out, "{d}", .{off});}try self.appendFmt(out, "],\"version\":{d}", .{self.version});if (self.pk != null or self.index_fields.items.len > 0) {try self.appendFmt(out, ",\"schema\":{{", .{});var first = true;if (self.pk) |p| {const prefix: []const u8 = switch (self.pk_type) {.number => "*",.uuid => "@",.string => "",};try self.appendFmt(out, "\"primaryKey\":\"{s}{s}\"", .{ prefix, p });first = false;}if (!first) try self.appendFmt(out, ",", .{});try self.appendFmt(out, "\"indexes\":[", .{});var n: usize = 0;for (self.index_fields.items) |field| {if (self.pk != null and std.mem.eql(u8, field, self.pk.?)) continue;if (n > 0) try self.appendFmt(out, ",", .{});const uni = for (self.unique_fields.items) |u| {if (std.mem.eql(u8, u, field)) break true;} else false;const numeric = (self.idx.get(field) orelse FieldIndex{ .typ = .lexical }).typ == .numeric;try self.appendFmt(out, "\"{s}{s}{s}\"", .{if (uni) "!" else "",if (numeric) "*" else "",field,});n += 1;}try self.appendFmt(out, "]}}", .{});}try self.appendFmt(out, "}}", .{});}// =====================================================================// Data file scanning// =====================================================================fn deletedSet(self: *Db, alloc: std.mem.Allocator) std.AutoHashMapUnmanaged(u64, void) {var set: std.AutoHashMapUnmanaged(u64, void) = .empty;for (self.deleted.items) |off| set.put(alloc, off, {}) catch {};return set;}/// Sequential scan yielding all live records. Caller supplies an arena for/// decoded values. Used for rebuilds, compaction and unindexed finds.pub const Scanner = struct {db: *Db,fd: i32 = -1,offset: u64 = 0,size: u64 = 0,skip_deleted: bool = true,deleted_set: std.AutoHashMapUnmanaged(u64, void) = .empty,scratch: std.mem.Allocator,pub fn init(db: *Db, scratch: std.mem.Allocator, skip_deleted: bool) Scanner {var s = Scanner{ .db = db, .scratch = scratch, .skip_deleted = skip_deleted };s.size = fio.fileSize(scratch, db.data_path) orelse 0;if (s.size > 0) {s.fd = fio.openRead(scratch, db.data_path) orelse -1;}if (skip_deleted) s.deleted_set = db.deletedSet(scratch);return s;}pub fn deinit(self: *Scanner) void {if (self.fd >= 0) fio.close(self.fd);self.fd = -1;}/// Returns the next record (decoded into `arena`) or null at EOF.pub fn next(self: *Scanner, arena: std.mem.Allocator) Error!?Record {while (true) {if (self.fd < 0 or self.offset + 4 > self.size) return null;var hdr: [4]u8 = undefined;const got = fio.pread(self.fd, &hdr, self.offset) orelse return Error.IoError;if (got < 4) return null;const rec_size = std.mem.readInt(u32, &hdr, .little);if (rec_size <= 4) return Error.CorruptRecord;if (self.offset + rec_size > self.size) return null; // incomplete tailconst off = self.offset;self.offset += rec_size;if (self.skip_deleted and self.deleted_set.contains(off)) continue;const buf = arena.alloc(u8, rec_size - 4) catch return Error.OutOfMemory;const got2 = fio.pread(self.fd, buf, off + 4) orelse return Error.IoError;if (got2 < rec_size - 4) return Error.Truncated;const d = msgpack.decode(arena, buf) catch return Error.CorruptRecord;return Record{ .value = d.value, .off = off, .len = rec_size };}}/// Like next() but without decoding — yields the raw framed bytes.pub fn nextRaw(self: *Scanner, arena: std.mem.Allocator) Error!?struct { bytes: []u8, off: u64 } {while (true) {if (self.fd < 0 or self.offset + 4 > self.size) return null;var hdr: [4]u8 = undefined;const got = fio.pread(self.fd, &hdr, self.offset) orelse return Error.IoError;if (got < 4) return null;const rec_size = std.mem.readInt(u32, &hdr, .little);if (rec_size <= 4) return Error.CorruptRecord;if (self.offset + rec_size > self.size) return null;const off = self.offset;self.offset += rec_size;if (self.skip_deleted and self.deleted_set.contains(off)) continue;const buf = arena.alloc(u8, rec_size) catch return Error.OutOfMemory;const got2 = fio.pread(self.fd, buf, off) orelse return Error.IoError;if (got2 < rec_size) return Error.Truncated;return .{ .bytes = buf, .off = off };}}};/// Read + decode one record by location.pub fn readAt(self: *Db, arena: std.mem.Allocator, loc: Loc) Error!msgpack.Value {const fd = fio.openRead(self.tmp, self.data_path) orelse return Error.IoError;defer fio.close(fd);if (loc.len <= 4) return Error.CorruptRecord;const buf = arena.alloc(u8, loc.len - 4) catch return Error.OutOfMemory;const got = fio.pread(fd, buf, loc.off + 4) orelse return Error.IoError;if (got < buf.len) return Error.Truncated;const d = msgpack.decode(arena, buf) catch return Error.CorruptRecord;return d.value;}// =====================================================================// Index management// =====================================================================fn keyFromValue(v: msgpack.Value) ?Key {return switch (v) {.number => |n| Key{ .num = n },.str => |s| Key{ .str = s },else => null, // null/bool/objects are not indexed (see header note)};}/// Duplicate a key's string into the db arena so it outlives the op arena.fn ownKey(self: *Db, k: Key) Error!Key {return switch (k) {.num => k,.str => |s| Key{ .str = self.alloc.dupe(u8, s) catch return Error.OutOfMemory },};}fn lowerBound(entries: []const Entry, key: Key) usize {var lo: usize = 0;var hi: usize = entries.len;while (lo < hi) {const mid = lo + (hi - lo) / 2;if (cmpKeys(entries[mid].key, key) < 0) {lo = mid + 1;} else {hi = mid;}}return lo;}/// Binary search: all entries whose key equals `key` (appended to `out`).pub fn indexGet(self: *Db, field: []const u8, key: Key, alloc: std.mem.Allocator, out: *std.ArrayList(Entry)) Error!void {const fi = self.idx.getPtr(field) orelse return Error.NoSuchIndex;const parsed = self.parseKeyForField(field, key);const entries = fi.entries.items;var i = lowerBound(entries, parsed);while (i < entries.len and cmpKeys(entries[i].key, parsed) == 0) : (i += 1) {out.append(alloc, entries[i]) catch return Error.OutOfMemory;}}/// Range scan [from, to] (either side optional), ascending.pub fn indexRange(self: *Db, field: []const u8, from: ?Key, to: ?Key, alloc: std.mem.Allocator, out: *std.ArrayList(Entry)) Error!void {const fi = self.idx.getPtr(field) orelse return Error.NoSuchIndex;const entries = fi.entries.items;var i: usize = if (from) |f| lowerBound(entries, self.parseKeyForField(field, f)) else 0;const to_key: ?Key = if (to) |t| self.parseKeyForField(field, t) else null;while (i < entries.len) : (i += 1) {if (to_key) |t| {if (cmpKeys(entries[i].key, t) > 0) break;}out.append(alloc, entries[i]) catch return Error.OutOfMemory;}}/// _parseKey semantics: numeric index fields parse string keys via/// parseInt; on NaN the key stays a string.fn parseKeyForField(self: *Db, field: []const u8, key: Key) Key {const fi = self.idx.get(field) orelse return key;if (fi.typ == .numeric and key == .str) {if (jsParseInt(key.str)) |n| return Key{ .num = n };}return key;}fn indexInsertEntry(self: *Db, field: []const u8, key: Key, loc: Loc) Error!void {const fi = self.idx.getPtr(field) orelse return;const owned = try self.ownKey(key);const e = Entry{ .key = owned, .off = loc.off, .len = loc.len };// insert at sorted positionconst pos = blk: {var lo: usize = 0;var hi: usize = fi.entries.items.len;while (lo < hi) {const mid = lo + (hi - lo) / 2;if (entryLess({}, fi.entries.items[mid], e)) {lo = mid + 1;} else {hi = mid;}}break :blk lo;};fi.entries.insert(self.alloc, pos, e) catch return Error.OutOfMemory;}/// Index a record's fields at loc.fn indexInsertRecord(self: *Db, rec: msgpack.Value, loc: Loc) Error!void {for (self.index_fields.items) |field| {const v = rec.get(field) orelse continue;const key = keyFromValue(v) orelse continue;try self.indexInsertEntry(field, key, loc);}self.covered_bytes = @max(self.covered_bytes, loc.off + loc.len);}/// Remove all entries pointing at offset `off` (all fields).fn indexRemoveOffset(self: *Db, off: u64) void {var it = self.idx.iterator();while (it.next()) |kv| {const list = &kv.value_ptr.entries;var i: usize = 0;while (i < list.items.len) {if (list.items[i].off == off) {_ = list.orderedRemove(i);} else {i += 1;}}}}fn indexPath(self: *Db, buf: []u8, field: []const u8) []const u8 {// <dir>/<base>.<field>.txt — derive from idxstate path (…/base.idxstate.json)const prefix = self.idxstate_path[0 .. self.idxstate_path.len - "idxstate.json".len];return std.fmt.bufPrint(buf, "{s}{s}.txt", .{ prefix, field }) catch buf[0..0];}fn initIndexes(self: *Db, force_rebuild: bool) Error!void {var need_rebuild = force_rebuild;const data_size = fio.fileSize(self.tmp, self.data_path) orelse 0;if (!need_rebuild) {for (self.index_fields.items) |field| {var pbuf: [512]u8 = undefined;const p = self.indexPath(&pbuf, field);const isize_ = fio.fileSize(self.tmp, p);if (isize_ == null or isize_.? == 0) {if (data_size > 0) {need_rebuild = true;break;}}}}if (!need_rebuild) {// coveredBytes from idxstate — no state file means rebuild (JS)if (self.readIdxState()) |cb| {self.covered_bytes = cb;} else {need_rebuild = true;}}if (need_rebuild) {try self.rebuildIndexes();return;}// Load index files into memoryfor (self.index_fields.items) |field| {var pbuf: [512]u8 = undefined;const p = self.indexPath(&pbuf, field);const content = fio.readAll(self.tmp, p) orelse continue;defer self.tmp.free(content);self.loadIndexLines(field, content);}// sort (files are sorted by JS localeCompare; re-sort under our order)var it = self.idx.iterator();while (it.next()) |kv| {std.sort.pdq(Entry, kv.value_ptr.entries.items, {}, entryLess);}if (data_size > self.covered_bytes) self.catchUp(data_size);}fn loadIndexLines(self: *Db, field: []const u8, content: []const u8) void {const fi = self.idx.getPtr(field) orelse return;var lines = std.mem.splitScalar(u8, content, '\n');while (lines.next()) |line| {if (line.len == 0) continue;// key = up to FIRST comma (JS split(',')[0]); then offset, lengthconst c1 = std.mem.indexOfScalar(u8, line, ',') orelse continue;const rest = line[c1 + 1 ..];const c2 = std.mem.indexOfScalar(u8, rest, ',') orelse continue;const off_s = rest[0..c2];const len_s = rest[c2 + 1 ..];const off = std.fmt.parseInt(u64, off_s, 10) catch continue;const len = std.fmt.parseInt(u64, len_s, 10) catch continue;if (len == 0) continue;const key_s = line[0..c1];var key: Key = undefined;if (fi.typ == .numeric) {key = if (jsParseInt(key_s)) |n| Key{ .num = n } else Key{ .str = self.alloc.dupe(u8, key_s) catch return };} else {key = Key{ .str = self.alloc.dupe(u8, key_s) catch return };}fi.entries.append(self.alloc, .{ .key = key, .off = off, .len = len }) catch return;}}fn rebuildIndexes(self: *Db) Error!void {var it0 = self.idx.iterator();while (it0.next()) |kv| kv.value_ptr.entries.clearRetainingCapacity();var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);defer scratch_state.deinit();const scratch = scratch_state.allocator();// JS rebuild scans ALL records (no tombstone filter — deleted offsets// simply get filtered at read time). Mirror that.var scanner = Scanner.init(self, scratch, false);defer scanner.deinit();while (try scanner.next(scratch)) |rec| {try self.indexInsertRecord(rec.value, .{ .off = rec.off, .len = rec.len });}self.covered_bytes = fio.fileSize(self.tmp, self.data_path) orelse 0;try self.writeIndexFiles();try self.writeIdxState();}fn catchUp(self: *Db, file_size: u64) void {var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);defer scratch_state.deinit();const scratch = scratch_state.allocator();var scanner = Scanner.init(self, scratch, false);defer scanner.deinit();scanner.offset = self.covered_bytes;while (true) {const maybe = scanner.next(scratch) catch break;const rec = maybe orelse break;self.indexInsertRecord(rec.value, .{ .off = rec.off, .len = rec.len }) catch break;}self.covered_bytes = @max(self.covered_bytes, file_size);}/// A file that already holds these exact bytes is left alone — its mtime/// included — so closing a handle that changed nothing writes nothing.fn fileHolds(alloc: std.mem.Allocator, path: []const u8, bytes: []const u8) bool {const disk = fio.readAll(alloc, path) orelse return false;defer alloc.free(disk);return std.mem.eql(u8, disk, bytes);}fn writeIndexFiles(self: *Db) Error!void {for (self.index_fields.items) |field| {const fi = self.idx.getPtr(field) orelse continue;var out: std.ArrayList(u8) = .empty;defer out.deinit(self.tmp);for (fi.entries.items) |e| {var kbuf: [32]u8 = undefined;const ks = switch (e.key) {.num => |n| jsNumFmt(&kbuf, n),.str => |s| s,};try self.appendFmt(&out, "{s},{d},{d}\n", .{ ks, e.off, e.len });}var pbuf: [512]u8 = undefined;const p = self.indexPath(&pbuf, field);if (fileHolds(self.tmp, p, out.items)) continue;var tbuf: [512]u8 = undefined;const tmp = std.fmt.bufPrint(&tbuf, "{s}.{d}.tmp", .{ p, linux.getpid() }) catch return Error.IoError;const fd = fio.openTrunc(self.tmp, tmp) orelse return Error.IoError;const ok = fio.writeAll(fd, out.items);fio.close(fd);if (!ok) return Error.IoError;if (!fio.rename(self.tmp, tmp, p)) return Error.IoError;}}fn readIdxState(self: *Db) ?u64 {const content = fio.readAll(self.tmp, self.idxstate_path) orelse return null;defer self.tmp.free(content);const parsed = std.json.parseFromSlice(std.json.Value, self.tmp, content, .{}) catch return null;defer parsed.deinit();if (parsed.value != .object) return null;const v = parsed.value.object.get("coveredBytes") orelse return null;return switch (v) {.integer => |i| @intCast(@max(i, 0)),.float => |f| @intFromFloat(@max(f, 0)),else => null,};}fn writeIdxState(self: *Db) Error!void {var buf: [128]u8 = undefined;const json = std.fmt.bufPrint(&buf, "{{\"coveredBytes\":{d}}}", .{self.covered_bytes}) catch return Error.IoError;if (fileHolds(self.tmp, self.idxstate_path, json)) return;var tbuf: [512]u8 = undefined;const tmp = std.fmt.bufPrint(&tbuf, "{s}.{d}.tmp", .{ self.idxstate_path, linux.getpid() }) catch return Error.IoError;const fd = fio.openTrunc(self.tmp, tmp) orelse return Error.IoError;const ok = fio.writeAll(fd, json);fio.close(fd);if (!ok) return Error.IoError;if (!fio.rename(self.tmp, tmp, self.idxstate_path)) return Error.IoError;}fn persistIndexesLocked(self: *Db) Error!void {try self.writeIndexFiles();// coverage can only grow (JS: max of ours and on-disk state)if (self.readIdxState()) |cb| self.covered_bytes = @max(self.covered_bytes, cb);try self.writeIdxState();}// =====================================================================// Operations// =====================================================================/// A mutable record under construction (op-arena entries).pub const MutableRecord = struct {entries: std.ArrayList(msgpack.Entry) = .empty,pub fn get(self: *const MutableRecord, key: []const u8) ?msgpack.Value {for (self.entries.items) |e| {if (e.key == .str and std.mem.eql(u8, e.key.str, key)) return e.value;}return null;}pub fn set(self: *MutableRecord, alloc: std.mem.Allocator, key: []const u8, v: msgpack.Value) Error!void {for (self.entries.items) |*e| {if (e.key == .str and std.mem.eql(u8, e.key.str, key)) {e.value = v;return;}}self.entries.append(alloc, .{ .key = .{ .str = key }, .value = v }) catch return Error.OutOfMemory;}pub fn toValue(self: *const MutableRecord) msgpack.Value {return .{ .map = self.entries.items };}};pub const InsertResult = union(enum) {pk_num: f64,pk_str: []const u8, // op-arenarecord: msgpack.Value,};/// insert() — auto primary key, unique checks, append, index./// `record` must be a map value; `arena` is the op arena (record memory).pub fn insert(self: *Db, arena: std.mem.Allocator, record: msgpack.Value, skip_pk: bool) Error!InsertResult {if (record != .map) return Error.CorruptRecord;try self.acquireLock();defer self.releaseLock();self.refresh();return self.insertLocked(arena, record, skip_pk);}fn insertLocked(self: *Db, arena: std.mem.Allocator, record: msgpack.Value, skip_pk: bool) Error!InsertResult {// copy into mutable formvar rec = MutableRecord{};for (record.map) |e| rec.entries.append(arena, e) catch return Error.OutOfMemory;// _hasPrimaryKeyValue: present and not null/undefinedvar has_pk_value = false;if (self.pk) |p| {if (rec.get(p)) |v| {has_pk_value = (v != .nil and v != .undef);}}var auto_gen = false;if (self.pk != null and !skip_pk and !has_pk_value) {if (self.pk_type == .number) {rec.set(arena, self.pk.?, .{ .number = self.next_id }) catch return Error.OutOfMemory;self.next_id += 1;auto_gen = true;} else if (self.pk_type == .uuid) {// A GENERATED ID IS CHECKED, NOT TRUSTED (ticket #113): the format is a// millisecond stamp + 3 base36 chars, so a bulk put collides within a// millisecond (two of 7800 in a measured run) and the unique check below// skips generated keys — two live rows then shared one pk and the next// update() deleted both and re-inserted one. Draw again until it is free.var ubuf: [12]u8 = undefined;var tries: usize = 0;while (true) : (tries += 1) {const u = genUuid(&ubuf);var taken: std.ArrayList(Loc) = .empty;defer taken.deinit(arena);try self.findPkLocs(arena, .{ .str = u }, &taken);if (taken.items.len == 0 or tries >= 64) break;}const owned = arena.dupe(u8, ubuf[0..12]) catch return Error.OutOfMemory;rec.set(arena, self.pk.?, .{ .str = owned }) catch return Error.OutOfMemory;auto_gen = true;}}// unique constraints (skip auto-generated pk)if (self.index_fields.items.len > 0 and self.unique_fields.items.len > 0) {var del_set = self.deletedSet(arena);for (self.unique_fields.items) |field| {if (auto_gen and self.pk != null and std.mem.eql(u8, field, self.pk.?)) continue;const v = rec.get(field) orelse continue;if (v == .undef) continue;const key = keyFromValue(v) orelse continue;var hits: std.ArrayList(Entry) = .empty;defer hits.deinit(arena);try self.indexGet(field, key, arena, &hits);var live: usize = 0;for (hits.items) |h| {if (!del_set.contains(h.off)) live += 1;}if (live > 0) {var kbuf: [32]u8 = undefined;const ks = switch (key) {.num => |n| jsNumFmt(&kbuf, n),.str => |s| s,};self.last_error = std.fmt.bufPrint(&self.last_error_buf, "Duplicate key: {s}={s}", .{ field, ks }) catch "Duplicate key";return Error.DuplicateKey;}}}// serialize + appendconst framed = msgpack.serialize(arena, rec.toValue()) catch return Error.OutOfMemory;const offset = fio.fileSize(self.tmp, self.data_path) orelse 0;const fd = fio.openAppend(self.tmp, self.data_path) orelse return Error.IoError;const ok = fio.writeAll(fd, framed);fio.close(fd);if (!ok) return Error.IoError;const loc = Loc{ .off = offset, .len = framed.len };if (self.index_fields.items.len > 0) {try self.indexInsertRecord(rec.toValue(), loc);}if (self.pk != null and self.pk_type == .number and !skip_pk and !has_pk_value) {try self.persistMetaLocked();}if (self.pk) |p| {const v = rec.get(p) orelse return Error.CorruptRecord;return switch (v) {.number => |n| InsertResult{ .pk_num = n },.str => |s| InsertResult{ .pk_str = s },else => InsertResult{ .record = rec.toValue() },};}return InsertResult{ .record = rec.toValue() };}/// Locate live records by primary key (index-backed when available).pub fn findPkLocs(self: *Db, arena: std.mem.Allocator, key: Key, out: *std.ArrayList(Loc)) Error!void {if (self.pk == null) return Error.NoPrimaryKey;self.refresh();var del_set = self.deletedSet(arena);if (self.index_fields.items.len > 0) {var hits: std.ArrayList(Entry) = .empty;defer hits.deinit(arena);try self.indexGet(self.pk.?, key, arena, &hits);for (hits.items) |h| {if (del_set.contains(h.off)) continue;out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;}return;}// no indexes: full scan comparing the pk fieldvar scanner = Scanner.init(self, arena, true);defer scanner.deinit();while (try scanner.next(arena)) |rec| {const v = rec.value.get(self.pk.?) orelse continue;const k = keyFromValue(v) orelse continue;if (cmpKeys(self.parseKeyForField(self.pk.?, k), self.parseKeyForField(self.pk.?, key)) == 0) {out.append(arena, .{ .off = rec.off, .len = rec.len }) catch return Error.OutOfMemory;}}}/// Locate live records by secondary index equality.pub fn findIndexLocs(self: *Db, arena: std.mem.Allocator, field: []const u8, key: Key, out: *std.ArrayList(Loc)) Error!void {self.refresh();var del_set = self.deletedSet(arena);var hits: std.ArrayList(Entry) = .empty;defer hits.deinit(arena);try self.indexGet(field, key, arena, &hits);for (hits.items) |h| {if (del_set.contains(h.off)) continue;out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;}}/// Locate live records by index range [from, to] ascending.pub fn findRangeLocs(self: *Db, arena: std.mem.Allocator, field: []const u8, from: ?Key, to: ?Key, out: *std.ArrayList(Loc)) Error!void {self.refresh();var del_set = self.deletedSet(arena);var hits: std.ArrayList(Entry) = .empty;defer hits.deinit(arena);try self.indexRange(field, from, to, arena, &hits);var seen: std.AutoHashMapUnmanaged(u64, void) = .empty;for (hits.items) |h| {if (del_set.contains(h.off)) continue;if (seen.contains(h.off)) continue;seen.put(arena, h.off, {}) catch {};out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;}}/// delete by primary key. Returns the deleted records (decoded into arena).pub fn deletePk(self: *Db, arena: std.mem.Allocator, key: Key, out_records: *std.ArrayList(msgpack.Value)) Error!usize {try self.acquireLock();defer self.releaseLock();self.refresh();return self.deletePkLocked(arena, key, out_records);}fn deletePkLocked(self: *Db, arena: std.mem.Allocator, key: Key, out_records: *std.ArrayList(msgpack.Value)) Error!usize {var locs: std.ArrayList(Loc) = .empty;defer locs.deinit(arena);try self.findPkLocs(arena, key, &locs);for (locs.items) |loc| {const rec = try self.readAt(arena, loc);out_records.append(arena, rec) catch return Error.OutOfMemory;self.deleted.append(self.alloc, loc.off) catch return Error.OutOfMemory;self.indexRemoveOffset(loc.off);}try self.persistMetaLocked();return locs.items.len;}/// update by primary key: delete + insert(skip_pk) — JS update() semantics./// `new_record` must already carry the primary key (the JS callback/// contract: the record keeps its pk unless the caller removes it).pub fn updatePk(self: *Db, arena: std.mem.Allocator, key: Key, new_record: msgpack.Value) Error!usize {if (self.pk) |p| {const carried = for (new_record.map) |e| {if (e.key == .str and std.mem.eql(u8, e.key.str, p)) break e.value != .nil and e.value != .undef;} else false;if (!carried) return Error.RecordLacksPrimaryKey;}try self.acquireLock();defer self.releaseLock();self.refresh();var old: std.ArrayList(msgpack.Value) = .empty;defer old.deinit(arena);const n = try self.deletePkLocked(arena, key, &old);if (n == 0) return 0;var i: usize = 0;while (i < n) : (i += 1) {_ = try self.insertLocked(arena, new_record, true);}return n;}/// compact() — rewrite the data file without tombstones, rebuild indexes.pub fn compact(self: *Db) Error!void {try self.acquireLock();defer self.releaseLock();self.refresh();try self.compactLocked();if (self.index_fields.items.len > 0) {try self.rebuildIndexes();}}fn compactLocked(self: *Db) Error!void {var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);defer scratch_state.deinit();const scratch = scratch_state.allocator();var tbuf: [512]u8 = undefined;const tmp = std.fmt.bufPrint(&tbuf, "{s}.tmp", .{self.data_path}) catch return Error.IoError;const out_fd = fio.openTrunc(self.tmp, tmp) orelse return Error.IoError;var scanner = Scanner.init(self, scratch, true);var ok = true;while (true) {const maybe = scanner.nextRaw(scratch) catch {ok = false;break;};const raw = maybe orelse break;if (!fio.writeAll(out_fd, raw.bytes)) {ok = false;break;}}scanner.deinit();fio.close(out_fd);if (!ok) return Error.IoError;if (!fio.rename(self.tmp, tmp, self.data_path)) return Error.IoError;self.deleted.clearRetainingCapacity();try self.persistMetaLocked();}pub fn persistNow(self: *Db) Error!void {try self.acquireLock();defer self.releaseLock();self.refresh();try self.persistIndexesLocked();}};// =========================================================================// Tests// =========================================================================const testing = std.testing;fn tmpBase(buf: []u8, comptime name: []const u8) []const u8 {return std.fmt.bufPrint(buf, "/tmp/mpackdb-zigtest-{d}-" ++ name, .{linux.getpid()}) catch unreachable;}fn cleanup(alloc: std.mem.Allocator, base: []const u8) void {var buf: [512]u8 = undefined;const suffixes = [_][]const u8{ ".mpack", ".meta.json", ".idxstate.json", ".lock", ".id.txt", ".email.txt", ".age.txt", ".uuid.txt" };for (suffixes) |suffix| {const p = std.fmt.bufPrint(&buf, "{s}{s}", .{ base, suffix }) catch continue;_ = fio.unlink(alloc, p);}}fn strKey(s: []const u8) Key {return .{ .str = s };}fn numKey(n: f64) Key {return .{ .num = n };}fn makeUser(arena: std.mem.Allocator, name: []const u8, email: []const u8, age: f64) !msgpack.Value {const entries = try arena.alloc(msgpack.Entry, 3);entries[0] = .{ .key = .{ .str = "name" }, .value = .{ .str = name } };entries[1] = .{ .key = .{ .str = "email" }, .value = .{ .str = email } };entries[2] = .{ .key = .{ .str = "age" }, .value = .{ .number = age } };return .{ .map = entries };}test "engine: insert/find/delete/update with numeric pk + indexes" {var base_buf: [128]u8 = undefined;const base = tmpBase(&base_buf, "crud");cleanup(testing.allocator, base);defer cleanup(testing.allocator, base);var arena_state = std.heap.ArenaAllocator.init(testing.allocator);defer arena_state.deinit();const arena = arena_state.allocator();const db = try Db.open(testing.allocator, base, .{.primary_key = "*id",.indexes = &.{ "email", "*age" },});// insert three — a fresh table counts from 1 (ticket #4)const r1 = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 30), false);try testing.expectEqual(@as(f64, 1), r1.pk_num);const r2 = try db.insert(arena, try makeUser(arena, "Bob", "[email protected]", 25), false);try testing.expectEqual(@as(f64, 2), r2.pk_num);_ = try db.insert(arena, try makeUser(arena, "Carol", "[email protected]", 35), false);// find by pkvar locs: std.ArrayList(Loc) = .empty;try db.findPkLocs(arena, numKey(2), &locs);try testing.expectEqual(@as(usize, 1), locs.items.len);const bob = try db.readAt(arena, locs.items[0]);try testing.expectEqualStrings("Bob", bob.get("name").?.str);// find by secondary indexvar locs2: std.ArrayList(Loc) = .empty;try db.findIndexLocs(arena, "email", strKey("[email protected]"), &locs2);try testing.expectEqual(@as(usize, 1), locs2.items.len);// range on numeric index: age 26..40 → Alice(30), Carol(35)var locs3: std.ArrayList(Loc) = .empty;try db.findRangeLocs(arena, "age", numKey(26), numKey(40), &locs3);try testing.expectEqual(@as(usize, 2), locs3.items.len);// update Bob's agevar bob_new = Db.MutableRecord{};for (bob.map) |e| try bob_new.entries.append(arena, e);try bob_new.set(arena, "age", .{ .number = 26 });const updated = try db.updatePk(arena, numKey(2), bob_new.toValue());try testing.expectEqual(@as(usize, 1), updated);var locs4: std.ArrayList(Loc) = .empty;try db.findPkLocs(arena, numKey(2), &locs4);try testing.expectEqual(@as(usize, 1), locs4.items.len);const bob2 = try db.readAt(arena, locs4.items[0]);try testing.expectEqual(@as(f64, 26), bob2.get("age").?.number);// delete Alicevar deleted_recs: std.ArrayList(msgpack.Value) = .empty;const dn = try db.deletePk(arena, numKey(1), &deleted_recs);try testing.expectEqual(@as(usize, 1), dn);try testing.expectEqualStrings("Alice", deleted_recs.items[0].get("name").?.str);var locs5: std.ArrayList(Loc) = .empty;try db.findPkLocs(arena, numKey(1), &locs5);try testing.expectEqual(@as(usize, 0), locs5.items.len);db.close();// reopen (compaction drops the tombstone) and verify persistenceconst db2 = try Db.open(testing.allocator, base, .{.primary_key = "*id",.indexes = &.{ "email", "*age" },});defer db2.close();var locs6: std.ArrayList(Loc) = .empty;try db2.findPkLocs(arena, numKey(2), &locs6);try testing.expectEqual(@as(usize, 1), locs6.items.len);const bob3 = try db2.readAt(arena, locs6.items[0]);try testing.expectEqual(@as(f64, 26), bob3.get("age").?.number);// nextId continues after reopenconst r4 = try db2.insert(arena, try makeUser(arena, "Dan", "[email protected]", 40), false);try testing.expectEqual(@as(f64, 4), r4.pk_num);}test "engine: a stored counter wins over FIRST_ID (ticket #4)" {var base_buf: [128]u8 = undefined;const base = tmpBase(&base_buf, "counter");cleanup(testing.allocator, base);defer cleanup(testing.allocator, base);var arena_state = std.heap.ArenaAllocator.init(testing.allocator);defer arena_state.deinit();const arena = arena_state.allocator();// A table the JS mpackdb created and never wrote to: its meta says 0, and// an existing table keeps its own counter.var mbuf: [160]u8 = undefined;const meta_path = try std.fmt.bufPrint(&mbuf, "{s}.meta.json", .{base});const fd = fio.openTrunc(testing.allocator, meta_path).?;try testing.expect(fio.writeAll(fd, "{\"nextId\":0,\"deleted\":[],\"version\":1}"));fio.close(fd);const db = try Db.open(testing.allocator, base, .{ .primary_key = "*id" });defer db.close();const r = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 1), false);try testing.expectEqual(@as(f64, 0), r.pk_num);}test "engine: unique index rejects duplicates" {var base_buf: [128]u8 = undefined;const base = tmpBase(&base_buf, "uniq");cleanup(testing.allocator, base);defer cleanup(testing.allocator, base);var arena_state = std.heap.ArenaAllocator.init(testing.allocator);defer arena_state.deinit();const arena = arena_state.allocator();const db = try Db.open(testing.allocator, base, .{.primary_key = "*id",.indexes = &.{"!email"},});defer db.close();_ = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 30), false);const dup = db.insert(arena, try makeUser(arena, "Evil", "[email protected]", 31), false);try testing.expectError(Error.DuplicateKey, dup);try testing.expect(std.mem.indexOf(u8, db.last_error, "[email protected]") != null);}test "engine: uuid primary key" {var base_buf: [128]u8 = undefined;const base = tmpBase(&base_buf, "uuid");cleanup(testing.allocator, base);defer cleanup(testing.allocator, base);var arena_state = std.heap.ArenaAllocator.init(testing.allocator);defer arena_state.deinit();const arena = arena_state.allocator();const db = try Db.open(testing.allocator, base, .{ .primary_key = "@uuid" });defer db.close();const r = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 1), false);try testing.expectEqual(@as(usize, 12), r.pk_str.len);var locs: std.ArrayList(Loc) = .empty;try db.findPkLocs(arena, strKey(r.pk_str), &locs);try testing.expectEqual(@as(usize, 1), locs.items.len);}test "jsNumFmt / jsParseInt / uuid shape" {var buf: [32]u8 = undefined;try testing.expectEqualStrings("42", jsNumFmt(&buf, 42));try testing.expectEqualStrings("-7", jsNumFmt(&buf, -7));try testing.expectEqualStrings("1.5", jsNumFmt(&buf, 1.5));try testing.expectEqual(@as(f64, 1), jsParseInt("1.5").?);try testing.expectEqual(@as(f64, -12), jsParseInt("-12abc").?);try testing.expect(jsParseInt("abc") == null);var ubuf: [12]u8 = undefined;const u = genUuid(&ubuf);try testing.expectEqual(@as(usize, 12), u.len);for (u) |ch| try testing.expect((ch >= '0' and ch <= '9') or (ch >= 'a' and ch <= 'z'));}
Branches
- mainmain branch
Latest commits
- 38f9d10ftickets: Hybriel master 06617221 (plugin allocators 3a781359 + 413f60e4, mpackdb 2cb7ae5e, http1 773de63e); gate 249/0, connect 60/0mre
- d3db6139tickets: Hybriel master 190aa11d (fc838894 GC correctness, #126 closure scopes, #127); gate 249/0, connect 60/0mre
- bce182e3tickets: Hybriel master 7eea0d32 (#126 memory, #48 lambda copies its argument); migrate.hl lambdas take &logmre
- 4137be0fantcolony#40: mission references point to the moved missionsmre
- 9bfba36aantcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmre
- c7bd2645tickets: Hybriel master 73267707 (#122 fixed); compactNow workaround removed (#110 covered)mre
- 2ab91ee9tickets: gate checks rows appear once (session sync); re-vendor to ff51cf46 stopped on hybriel#122, stays 837fe120mre
- e01c2b1dtickets#24: installable app (manifest, service worker, offline list), own icon; gate waits for the hello's pongmre
- 752fbb7fdeploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
- 38bdd5e4deploy.sh: never send .git or .gitignore to Byrodinmre
- f12fa1bcState of 2026-09-27, before the move to gitoriamre