gitoriaLog in with ident

notes

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Main branchmain2149e902notes mission 002 (4/4): code order — README file map + same-output test, STATUS, LOG, report; tests/letcount.py, tests/realdata-baseline.mjs, tests/realdata-compare.pymremain/plugins/mpackdb/engine.zig

62.7 KB

  1. // mpackdb DB engine — Zig port of mpackdb v1.0.7 (MPackDB.js / IndexManager.js)
  2. // semantics. File-level compatible with the JS implementation (D12):
  3. //
  4. // <base>.mpack append-only records: [4-byte LE total size][msgpack map]
  5. // <base>.meta.json { nextId?, deleted:[offsets], version, schema? }
  6. // <base>.<field>.txt sorted index lines "key,offset,length\n"
  7. // <base>.idxstate.json { coveredBytes }
  8. // <base>.lock lock file containing the holder's pid (O_EXCL protocol,
  9. // stale takeover after staleLockTimeout ms)
  10. //
  11. // Divergences from the JS implementation (documented in the 066 report):
  12. // - the in-memory index is the COMPLETE entry set (JS keeps on-disk + delta);
  13. // persist() writes the full set instead of read-merge-write. Sequential
  14. // cross-process use (close one side, open the other) behaves identically.
  15. // - string index keys sort bytewise (JS: String.localeCompare, ICU collation).
  16. // Equality lookups are unaffected; range ordering can differ for
  17. // mixed-case/accented keys.
  18. // - null field values are not indexed (JS indexes them under the key "null").
  19. // - external data-file replacement while open (inode change) is not detected.
  20. // - a fresh `*id` table's first id is 1 (JS: 0) — ticket #4, see FIRST_ID.
  21. const std = @import("std");
  22. const linux = std.os.linux;
  23. const msgpack = @import("msgpack.zig");
  24. pub const PkType = enum(u8) { number = 0, uuid = 1, string = 2 };
  25. /// A `*id` table's first auto-increment id (ticket #4, the creator's ruling:
  26. /// "start at 1"). Only a table WITHOUT a stored counter starts here: a
  27. /// .meta.json that carries `nextId` — every existing table, including one the
  28. /// JS mpackdb wrote — keeps counting from it. The file format is unchanged;
  29. /// the JS mpackdb 1.0.7 starts a fresh table at 0, which is the one divergence.
  30. pub const FIRST_ID: f64 = 1;
  31. pub const IdxType = enum(u8) { lexical = 0, numeric = 1 };
  32. pub const Error = error{
  33. NoDbFile,
  34. NoPrimaryKey,
  35. /// mission 077 (074 GAP 10): a query named a field with no index. There is
  36. /// no index to binary-search and mpackdb does not silently full-scan, so
  37. /// the answer used to be an empty result — indistinguishable from "no such
  38. /// record". It is a mistake, and it says so now.
  39. NoSuchIndex,
  40. DuplicateKey,
  41. /// ticket #29: update()'s record must carry the primary key. Without it
  42. /// the old record was deleted, the new one written without a key, and the
  43. /// table left corrupt (fetch null, a unique index still finding it).
  44. RecordLacksPrimaryKey,
  45. LockTimeout,
  46. CorruptRecord,
  47. IoError,
  48. OutOfMemory,
  49. InvalidFormat,
  50. Truncated,
  51. MapTooLarge,
  52. };
  53. pub const Key = union(enum) {
  54. num: f64,
  55. str: []const u8,
  56. };
  57. pub const Entry = struct {
  58. key: Key,
  59. off: u64,
  60. len: u64,
  61. };
  62. const FieldIndex = struct {
  63. typ: IdxType,
  64. entries: std.ArrayList(Entry) = .empty, // sorted by (key, off)
  65. };
  66. pub const Loc = struct { off: u64, len: u64 };
  67. pub const Record = struct {
  68. value: msgpack.Value,
  69. off: u64,
  70. len: u64,
  71. };
  72. // =========================================================================
  73. // Low-level file IO (raw linux syscalls — plugins avoid std.Io plumbing)
  74. // =========================================================================
  75. const fio = struct {
  76. fn openz(alloc: std.mem.Allocator, path: []const u8, flags: linux.O, mode: u32) ?i32 {
  77. const path_z = alloc.dupeZ(u8, path) catch return null;
  78. defer alloc.free(path_z);
  79. const rc = linux.open(path_z, flags, mode);
  80. const fd: i32 = @bitCast(@as(u32, @truncate(rc)));
  81. if (fd < 0) return null;
  82. return fd;
  83. }
  84. fn openRead(alloc: std.mem.Allocator, path: []const u8) ?i32 {
  85. return openz(alloc, path, .{}, 0);
  86. }
  87. fn openAppend(alloc: std.mem.Allocator, path: []const u8) ?i32 {
  88. return openz(alloc, path, .{ .ACCMODE = .WRONLY, .CREAT = true, .APPEND = true }, 0o644);
  89. }
  90. fn openTrunc(alloc: std.mem.Allocator, path: []const u8) ?i32 {
  91. return openz(alloc, path, .{ .ACCMODE = .WRONLY, .CREAT = true, .TRUNC = true }, 0o644);
  92. }
  93. /// O_CREAT|O_EXCL — returns null when the file already exists.
  94. fn openExcl(alloc: std.mem.Allocator, path: []const u8) ?i32 {
  95. return openz(alloc, path, .{ .ACCMODE = .WRONLY, .CREAT = true, .EXCL = true }, 0o644);
  96. }
  97. fn close(fd: i32) void {
  98. _ = linux.close(fd);
  99. }
  100. fn writeAll(fd: i32, bytes: []const u8) bool {
  101. var written: usize = 0;
  102. while (written < bytes.len) {
  103. const rc = linux.write(fd, bytes.ptr + written, bytes.len - written);
  104. if (@as(isize, @bitCast(rc)) <= 0) return false;
  105. written += rc;
  106. }
  107. return true;
  108. }
  109. fn pread(fd: i32, buf: []u8, offset: u64) ?usize {
  110. var got: usize = 0;
  111. while (got < buf.len) {
  112. const rc = linux.pread(fd, buf.ptr + got, buf.len - got, @intCast(offset + got));
  113. const n: isize = @bitCast(rc);
  114. if (n < 0) return null;
  115. if (n == 0) break;
  116. got += @intCast(n);
  117. }
  118. return got;
  119. }
  120. fn readAll(alloc: std.mem.Allocator, path: []const u8) ?[]u8 {
  121. const fd = openRead(alloc, path) orelse return null;
  122. defer close(fd);
  123. var content: std.ArrayList(u8) = .empty;
  124. var buf: [65536]u8 = undefined;
  125. while (true) {
  126. const rc = linux.read(fd, &buf, buf.len);
  127. if (@as(isize, @bitCast(rc)) <= 0) break;
  128. content.appendSlice(alloc, buf[0..rc]) catch return null;
  129. }
  130. return content.toOwnedSlice(alloc) catch null;
  131. }
  132. fn fileSize(alloc: std.mem.Allocator, path: []const u8) ?u64 {
  133. const path_z = alloc.dupeZ(u8, path) catch return null;
  134. defer alloc.free(path_z);
  135. var stx: linux.Statx = undefined;
  136. const rc = linux.statx(linux.AT.FDCWD, path_z, 0, linux.STATX.BASIC_STATS, &stx);
  137. if (rc != 0) return null;
  138. return stx.size;
  139. }
  140. fn mtimeMs(alloc: std.mem.Allocator, path: []const u8) ?i64 {
  141. const path_z = alloc.dupeZ(u8, path) catch return null;
  142. defer alloc.free(path_z);
  143. var stx: linux.Statx = undefined;
  144. const rc = linux.statx(linux.AT.FDCWD, path_z, 0, linux.STATX.BASIC_STATS, &stx);
  145. if (rc != 0) return null;
  146. return @as(i64, stx.mtime.sec) * 1000 + @divTrunc(@as(i64, stx.mtime.nsec), 1_000_000);
  147. }
  148. fn rename(alloc: std.mem.Allocator, old_path: []const u8, new_path: []const u8) bool {
  149. const old_z = alloc.dupeZ(u8, old_path) catch return false;
  150. defer alloc.free(old_z);
  151. const new_z = alloc.dupeZ(u8, new_path) catch return false;
  152. defer alloc.free(new_z);
  153. return linux.renameat(linux.AT.FDCWD, old_z, linux.AT.FDCWD, new_z) == 0;
  154. }
  155. fn unlink(alloc: std.mem.Allocator, path: []const u8) bool {
  156. const path_z = alloc.dupeZ(u8, path) catch return false;
  157. defer alloc.free(path_z);
  158. return linux.unlinkat(linux.AT.FDCWD, path_z, 0) == 0;
  159. }
  160. fn mkdirAll(alloc: std.mem.Allocator, dir_path: []const u8) void {
  161. if (dir_path.len == 0 or std.mem.eql(u8, dir_path, ".")) return;
  162. var i: usize = 1;
  163. while (i <= dir_path.len) : (i += 1) {
  164. if (i == dir_path.len or dir_path[i] == '/') {
  165. const part = alloc.dupeZ(u8, dir_path[0..i]) catch return;
  166. defer alloc.free(part);
  167. _ = linux.mkdirat(linux.AT.FDCWD, part, 0o755);
  168. }
  169. }
  170. }
  171. fn sleepMs(ms: u64) void {
  172. var req: linux.timespec = .{ .sec = @intCast(ms / 1000), .nsec = @intCast((ms % 1000) * 1_000_000) };
  173. var rem: linux.timespec = undefined;
  174. _ = linux.nanosleep(&req, &rem);
  175. }
  176. fn nowMs() i64 {
  177. var ts: linux.timespec = undefined;
  178. _ = linux.clock_gettime(linux.CLOCK.REALTIME, &ts);
  179. return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);
  180. }
  181. };
  182. // =========================================================================
  183. // JS-compatible helpers
  184. // =========================================================================
  185. /// Format an f64 the way JS String(n) does for the values mpackdb produces:
  186. /// integral values print without a decimal point. (Exotic floats may diverge
  187. /// from V8's shortest-round-trip formatting; index keys are ints/strings in
  188. /// practice.)
  189. pub fn jsNumFmt(buf: []u8, v: f64) []const u8 {
  190. if (v == @floor(v) and @abs(v) <= 9007199254740992.0) {
  191. const i: i64 = @intFromFloat(v);
  192. return std.fmt.bufPrint(buf, "{d}", .{i}) catch buf[0..0];
  193. }
  194. return std.fmt.bufPrint(buf, "{d}", .{v}) catch buf[0..0];
  195. }
  196. /// JS parseInt(s, 10) semantics: optional sign, leading digits, ignore rest.
  197. /// Returns null for NaN.
  198. fn jsParseInt(s: []const u8) ?f64 {
  199. var i: usize = 0;
  200. while (i < s.len and (s[i] == ' ' or s[i] == '\t')) i += 1;
  201. var sign: f64 = 1;
  202. if (i < s.len and (s[i] == '+' or s[i] == '-')) {
  203. if (s[i] == '-') sign = -1;
  204. i += 1;
  205. }
  206. var got = false;
  207. var v: f64 = 0;
  208. while (i < s.len and s[i] >= '0' and s[i] <= '9') : (i += 1) {
  209. v = v * 10 + @as(f64, @floatFromInt(s[i] - '0'));
  210. got = true;
  211. }
  212. if (!got) return null;
  213. return sign * v;
  214. }
  215. /// _compareKeys: numbers numerically; otherwise stringified byte compare
  216. /// (JS uses localeCompare — see divergence note in the header).
  217. pub fn cmpKeys(a: Key, b: Key) i32 {
  218. if (a == .num and b == .num) {
  219. if (a.num < b.num) return -1;
  220. if (a.num > b.num) return 1;
  221. return 0;
  222. }
  223. var buf_a: [32]u8 = undefined;
  224. var buf_b: [32]u8 = undefined;
  225. const sa = if (a == .str) a.str else jsNumFmt(&buf_a, a.num);
  226. const sb = if (b == .str) b.str else jsNumFmt(&buf_b, b.num);
  227. return switch (std.mem.order(u8, sa, sb)) {
  228. .lt => -1,
  229. .eq => 0,
  230. .gt => 1,
  231. };
  232. }
  233. fn entryLess(_: void, a: Entry, b: Entry) bool {
  234. const c = cmpKeys(a.key, b.key);
  235. if (c != 0) return c < 0;
  236. return a.off < b.off;
  237. }
  238. /// mpack.js uuid(): 9-char base36 ms timestamp (padStart '0') + 3 base36 chars.
  239. pub fn genUuid(buf: *[12]u8) []const u8 {
  240. const digits = "0123456789abcdefghijklmnopqrstuvwxyz";
  241. const t: u64 = @intCast(fio.nowMs());
  242. var tmp: [16]u8 = undefined;
  243. var n: usize = 0;
  244. var v = t;
  245. while (v > 0) : (v /= 36) {
  246. tmp[n] = digits[@intCast(v % 36)];
  247. n += 1;
  248. }
  249. var i: usize = 0;
  250. while (i < 9) : (i += 1) {
  251. buf[8 - i] = if (i < n) tmp[i] else '0';
  252. }
  253. if (!rng_init) {
  254. var ts: linux.timespec = undefined;
  255. _ = linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts);
  256. rng = std.Random.DefaultPrng.init(@bitCast(@as(i64, ts.sec) *% 1_000_000_000 +% ts.nsec));
  257. rng_init = true;
  258. }
  259. var r = rng.random();
  260. for (0..3) |j| {
  261. buf[9 + j] = digits[r.uintLessThan(u8, 36)];
  262. }
  263. return buf[0..12];
  264. }
  265. var rng: std.Random.DefaultPrng = .{ .s = undefined };
  266. var rng_init: bool = false;
  267. // =========================================================================
  268. // The database
  269. // =========================================================================
  270. pub const OpenOptions = struct {
  271. primary_key: ?[]const u8 = null, // with */@/! prefixes, like the JS API
  272. indexes: []const []const u8 = &.{}, // with prefixes
  273. compact: bool = true,
  274. stale_lock_timeout_ms: i64 = 30000,
  275. debug: bool = false,
  276. };
  277. /// The identity of a table on disk (ticket #110): dirname + basename with any
  278. /// extension stripped, exactly the computation `Db.open` uses to turn a given
  279. /// path into `<base>.mpack` — so two spellings of the same file (`x.db`, `x`)
  280. /// share one key. The plugin ABI layer (mpackdb.zig) keys its process-wide
  281. /// handle registry on this, so a table opened twice in one process — from any
  282. /// module instance or realm — shares the one live `Db` instead of each open
  283. /// racing the other's file state.
  284. pub fn tableKey(alloc: std.mem.Allocator, db_file: []const u8) Error![]u8 {
  285. const dir = std.fs.path.dirname(db_file) orelse ".";
  286. var base = std.fs.path.basename(db_file);
  287. if (std.mem.lastIndexOfScalar(u8, base, '.')) |dot| {
  288. if (dot > 0) base = base[0..dot];
  289. }
  290. return std.fmt.allocPrint(alloc, "{s}/{s}", .{ dir, base }) catch Error.OutOfMemory;
  291. }
  292. pub const Db = struct {
  293. arena_state: std.heap.ArenaAllocator,
  294. alloc: std.mem.Allocator, // arena — freed wholesale on close
  295. // paths
  296. data_path: []const u8,
  297. meta_path: []const u8,
  298. lock_path: []const u8,
  299. idxstate_path: []const u8,
  300. // schema
  301. pk: ?[]const u8 = null,
  302. pk_type: PkType = .string,
  303. index_fields: std.ArrayList([]const u8) = .empty, // pk first (when set), like JS
  304. unique_fields: std.ArrayList([]const u8) = .empty,
  305. idx: std.StringArrayHashMapUnmanaged(FieldIndex) = .empty,
  306. // meta
  307. next_id: f64 = FIRST_ID,
  308. deleted: std.ArrayList(u64) = .empty,
  309. version: u64 = 0,
  310. covered_bytes: u64 = 0,
  311. compact_on_open: bool = true,
  312. stale_lock_timeout_ms: i64 = 30000,
  313. debug: bool = false,
  314. // last DuplicateKey detail for the plugin surface
  315. last_error_buf: [256]u8 = undefined,
  316. last_error: []const u8 = "",
  317. pub fn open(gpa: std.mem.Allocator, db_file: []const u8, opts: OpenOptions) Error!*Db {
  318. if (db_file.len == 0) return Error.NoDbFile;
  319. const self = gpa.create(Db) catch return Error.OutOfMemory;
  320. self.* = .{
  321. .arena_state = std.heap.ArenaAllocator.init(gpa),
  322. .alloc = undefined,
  323. .data_path = undefined,
  324. .meta_path = undefined,
  325. .lock_path = undefined,
  326. .idxstate_path = undefined,
  327. };
  328. self.alloc = self.arena_state.allocator();
  329. errdefer {
  330. self.arena_state.deinit();
  331. gpa.destroy(self);
  332. }
  333. const a = self.alloc;
  334. self.compact_on_open = opts.compact;
  335. self.stale_lock_timeout_ms = opts.stale_lock_timeout_ms;
  336. self.debug = opts.debug;
  337. // dirname / basename (extension stripped, like JS basename(f, extname(f)))
  338. const dir = std.fs.path.dirname(db_file) orelse ".";
  339. var base = std.fs.path.basename(db_file);
  340. if (std.mem.lastIndexOfScalar(u8, base, '.')) |dot| {
  341. if (dot > 0) base = base[0..dot];
  342. }
  343. fio.mkdirAll(a, dir);
  344. self.data_path = std.fmt.allocPrint(a, "{s}/{s}.mpack", .{ dir, base }) catch return Error.OutOfMemory;
  345. self.meta_path = std.fmt.allocPrint(a, "{s}/{s}.meta.json", .{ dir, base }) catch return Error.OutOfMemory;
  346. self.lock_path = std.fmt.allocPrint(a, "{s}/{s}.lock", .{ dir, base }) catch return Error.OutOfMemory;
  347. self.idxstate_path = std.fmt.allocPrint(a, "{s}/{s}.idxstate.json", .{ dir, base }) catch return Error.OutOfMemory;
  348. // ---- parse primary key ----
  349. if (opts.primary_key) |pk_raw| {
  350. if (pk_raw.len > 0) {
  351. if (pk_raw[0] == '*') {
  352. self.pk = a.dupe(u8, pk_raw[1..]) catch return Error.OutOfMemory;
  353. self.pk_type = .number;
  354. } else if (pk_raw[0] == '@') {
  355. self.pk = a.dupe(u8, pk_raw[1..]) catch return Error.OutOfMemory;
  356. self.pk_type = .uuid;
  357. } else {
  358. self.pk = a.dupe(u8, pk_raw) catch return Error.OutOfMemory;
  359. self.pk_type = .string;
  360. }
  361. }
  362. }
  363. if (self.pk) |p| self.unique_fields.append(a, p) catch return Error.OutOfMemory;
  364. // ---- parse indexes (pk first, then configured; dedupe) ----
  365. if (self.pk) |p| {
  366. self.index_fields.append(a, p) catch return Error.OutOfMemory;
  367. const t: IdxType = if (self.pk_type == .number) .numeric else .lexical;
  368. self.idx.put(a, p, .{ .typ = t }) catch return Error.OutOfMemory;
  369. }
  370. for (opts.indexes) |raw| {
  371. var clean = raw;
  372. var is_unique = false;
  373. if (clean.len > 0 and clean[0] == '!') {
  374. is_unique = true;
  375. clean = clean[1..];
  376. }
  377. var typ: IdxType = .lexical;
  378. if (clean.len > 0 and clean[0] == '*') {
  379. typ = .numeric;
  380. clean = clean[1..];
  381. } else if (clean.len > 0 and clean[0] == '@') {
  382. clean = clean[1..];
  383. }
  384. if (self.idx.contains(clean)) continue;
  385. const owned = a.dupe(u8, clean) catch return Error.OutOfMemory;
  386. self.index_fields.append(a, owned) catch return Error.OutOfMemory;
  387. self.idx.put(a, owned, .{ .typ = typ }) catch return Error.OutOfMemory;
  388. if (is_unique) self.unique_fields.append(a, owned) catch return Error.OutOfMemory;
  389. }
  390. // ---- load meta ----
  391. self.loadMeta();
  392. // ---- compact on init (JS default; also creates an empty data file) ----
  393. // AN OPEN WITH NOTHING TO CHANGE WRITES NOTHING (ticket #21): a store
  394. // with no tombstones has nothing to compact, and rewriting it anyway
  395. // made opening a backup change it.
  396. var did_compact = false;
  397. const nothing_to_compact = self.deleted.items.len == 0 and
  398. fio.fileSize(self.alloc, self.data_path) != null;
  399. if (self.compact_on_open and !nothing_to_compact) {
  400. try self.acquireLock();
  401. const cr = self.compactLocked();
  402. self.releaseLock();
  403. try cr;
  404. did_compact = true;
  405. }
  406. // ---- persist schema into meta (JS: init-schema under lock).
  407. // compactLocked already wrote meta (with schema); skip the extra bump.
  408. if (!did_compact and (self.pk != null or self.index_fields.items.len > 0) and
  409. !self.metaOnDiskIsCurrent())
  410. {
  411. try self.acquireLock();
  412. const mr = self.persistMetaLocked();
  413. self.releaseLock();
  414. try mr;
  415. }
  416. // ---- indexes (JS: (re)build under the lock; force after compaction) ----
  417. if (self.index_fields.items.len > 0) {
  418. try self.acquireLock();
  419. const ir = self.initIndexes(did_compact);
  420. self.releaseLock();
  421. try ir;
  422. }
  423. return self;
  424. }
  425. pub fn close(self: *Db) void {
  426. // persist indexes + idxstate (JS: IndexManager.close → persist under lock)
  427. if (self.index_fields.items.len > 0) {
  428. if (self.acquireLock()) {
  429. self.persistIndexesLocked() catch {};
  430. self.releaseLock();
  431. } else |_| {}
  432. }
  433. const gpa = self.arena_state.child_allocator;
  434. self.arena_state.deinit();
  435. gpa.destroy(self);
  436. }
  437. fn dbg(self: *Db, comptime fmt: []const u8, args: anytype) void {
  438. if (self.debug) std.debug.print("[mpackdb] " ++ fmt ++ "\n", args);
  439. }
  440. // =====================================================================
  441. // Locking (JS _acquireFileLock protocol)
  442. // =====================================================================
  443. fn acquireLock(self: *Db) Error!void {
  444. var retries: u32 = 0;
  445. while (true) {
  446. if (fio.openExcl(self.alloc, self.lock_path)) |fd| {
  447. var pid_buf: [16]u8 = undefined;
  448. const pid_s = std.fmt.bufPrint(&pid_buf, "{d}", .{linux.getpid()}) catch "0";
  449. _ = fio.writeAll(fd, pid_s);
  450. fio.close(fd);
  451. return;
  452. }
  453. // stale lock takeover
  454. if (self.stale_lock_timeout_ms > 0) {
  455. if (fio.mtimeMs(self.alloc, self.lock_path)) |mt| {
  456. if (fio.nowMs() - mt > self.stale_lock_timeout_ms) {
  457. _ = fio.unlink(self.alloc, self.lock_path);
  458. continue;
  459. }
  460. }
  461. }
  462. if (retries > 480) return Error.LockTimeout; // ~12s at 25ms
  463. fio.sleepMs(25);
  464. retries += 1;
  465. }
  466. }
  467. fn releaseLock(self: *Db) void {
  468. _ = fio.unlink(self.alloc, self.lock_path);
  469. }
  470. // =====================================================================
  471. // Meta (.meta.json)
  472. // =====================================================================
  473. fn appendFmt(self: *Db, list: *std.ArrayList(u8), comptime fmt: []const u8, args: anytype) Error!void {
  474. const s = std.fmt.allocPrint(self.alloc, fmt, args) catch return Error.OutOfMemory;
  475. defer self.alloc.free(s);
  476. list.appendSlice(self.alloc, s) catch return Error.OutOfMemory;
  477. }
  478. fn loadMeta(self: *Db) void {
  479. const content = fio.readAll(self.alloc, self.meta_path) orelse return;
  480. self.applyMetaJson(content);
  481. }
  482. fn applyMetaJson(self: *Db, content: []const u8) void {
  483. const parsed = std.json.parseFromSlice(std.json.Value, self.alloc, content, .{}) catch return;
  484. defer parsed.deinit();
  485. if (parsed.value != .object) return;
  486. const obj = parsed.value.object;
  487. self.next_id = FIRST_ID;
  488. self.deleted.clearRetainingCapacity();
  489. self.version = 0;
  490. if (obj.get("nextId")) |v| {
  491. self.next_id = switch (v) {
  492. .integer => |i| @floatFromInt(i),
  493. .float => |f| f,
  494. else => FIRST_ID,
  495. };
  496. }
  497. if (obj.get("version")) |v| {
  498. if (v == .integer) self.version = @intCast(@max(v.integer, 0));
  499. }
  500. if (obj.get("deleted")) |v| {
  501. if (v == .array) {
  502. for (v.array.items) |it| {
  503. const off: u64 = switch (it) {
  504. .integer => |i| @intCast(@max(i, 0)),
  505. .float => |f| @intFromFloat(@max(f, 0)),
  506. else => continue,
  507. };
  508. self.deleted.append(self.alloc, off) catch {};
  509. }
  510. }
  511. }
  512. }
  513. /// Pick up other processes' persisted state (JS refresh()): adopt disk meta
  514. /// when its version is newer; index appended tail records.
  515. pub fn refresh(self: *Db) void {
  516. if (fio.readAll(self.alloc, self.meta_path)) |content| {
  517. defer self.alloc.free(content);
  518. const parsed = std.json.parseFromSlice(std.json.Value, self.alloc, content, .{}) catch return;
  519. defer parsed.deinit();
  520. if (parsed.value == .object) {
  521. var disk_version: u64 = 0;
  522. if (parsed.value.object.get("version")) |v| {
  523. if (v == .integer) disk_version = @intCast(@max(v.integer, 0));
  524. }
  525. if (disk_version > self.version) self.applyMetaJson(content);
  526. }
  527. }
  528. // catchUp: index records another process appended
  529. if (self.index_fields.items.len > 0) {
  530. const size = fio.fileSize(self.alloc, self.data_path) orelse 0;
  531. if (size > self.covered_bytes) self.catchUp(size);
  532. }
  533. }
  534. fn persistMetaLocked(self: *Db) Error!void {
  535. self.version += 1;
  536. var out: std.ArrayList(u8) = .empty;
  537. defer out.deinit(self.alloc);
  538. try self.renderMeta(&out);
  539. // atomic tmp + rename (JS: `${metaPath}.${pid}.tmp`)
  540. var tmp_buf: [512]u8 = undefined;
  541. const tmp_path = std.fmt.bufPrint(&tmp_buf, "{s}.{d}.tmp", .{ self.meta_path, linux.getpid() }) catch return Error.IoError;
  542. const fd = fio.openTrunc(self.alloc, tmp_path) orelse return Error.IoError;
  543. const ok = fio.writeAll(fd, out.items);
  544. fio.close(fd);
  545. if (!ok) return Error.IoError;
  546. if (!fio.rename(self.alloc, tmp_path, self.meta_path)) return Error.IoError;
  547. }
  548. /// Does .meta.json already say, byte for byte, what this handle would write
  549. /// at its current version? Then an open has nothing to persist (ticket #21).
  550. fn metaOnDiskIsCurrent(self: *Db) bool {
  551. const disk = fio.readAll(self.alloc, self.meta_path) orelse return false;
  552. defer self.alloc.free(disk);
  553. var out: std.ArrayList(u8) = .empty;
  554. defer out.deinit(self.alloc);
  555. self.renderMeta(&out) catch return false;
  556. return std.mem.eql(u8, disk, out.items);
  557. }
  558. fn renderMeta(self: *Db, out: *std.ArrayList(u8)) Error!void {
  559. try self.appendFmt(out, "{{", .{});
  560. if (self.pk != null and self.pk_type == .number) {
  561. var nbuf: [32]u8 = undefined;
  562. try self.appendFmt(out, "\"nextId\":{s},", .{jsNumFmt(&nbuf, self.next_id)});
  563. }
  564. try self.appendFmt(out, "\"deleted\":[", .{});
  565. for (self.deleted.items, 0..) |off, i| {
  566. if (i > 0) try self.appendFmt(out, ",", .{});
  567. try self.appendFmt(out, "{d}", .{off});
  568. }
  569. try self.appendFmt(out, "],\"version\":{d}", .{self.version});
  570. if (self.pk != null or self.index_fields.items.len > 0) {
  571. try self.appendFmt(out, ",\"schema\":{{", .{});
  572. var first = true;
  573. if (self.pk) |p| {
  574. const prefix: []const u8 = switch (self.pk_type) {
  575. .number => "*",
  576. .uuid => "@",
  577. .string => "",
  578. };
  579. try self.appendFmt(out, "\"primaryKey\":\"{s}{s}\"", .{ prefix, p });
  580. first = false;
  581. }
  582. if (!first) try self.appendFmt(out, ",", .{});
  583. try self.appendFmt(out, "\"indexes\":[", .{});
  584. var n: usize = 0;
  585. for (self.index_fields.items) |field| {
  586. if (self.pk != null and std.mem.eql(u8, field, self.pk.?)) continue;
  587. if (n > 0) try self.appendFmt(out, ",", .{});
  588. const uni = for (self.unique_fields.items) |u| {
  589. if (std.mem.eql(u8, u, field)) break true;
  590. } else false;
  591. const numeric = (self.idx.get(field) orelse FieldIndex{ .typ = .lexical }).typ == .numeric;
  592. try self.appendFmt(out, "\"{s}{s}{s}\"", .{
  593. if (uni) "!" else "",
  594. if (numeric) "*" else "",
  595. field,
  596. });
  597. n += 1;
  598. }
  599. try self.appendFmt(out, "]}}", .{});
  600. }
  601. try self.appendFmt(out, "}}", .{});
  602. }
  603. // =====================================================================
  604. // Data file scanning
  605. // =====================================================================
  606. fn deletedSet(self: *Db, alloc: std.mem.Allocator) std.AutoHashMapUnmanaged(u64, void) {
  607. var set: std.AutoHashMapUnmanaged(u64, void) = .empty;
  608. for (self.deleted.items) |off| set.put(alloc, off, {}) catch {};
  609. return set;
  610. }
  611. /// Sequential scan yielding all live records. Caller supplies an arena for
  612. /// decoded values. Used for rebuilds, compaction and unindexed finds.
  613. pub const Scanner = struct {
  614. db: *Db,
  615. fd: i32 = -1,
  616. offset: u64 = 0,
  617. size: u64 = 0,
  618. skip_deleted: bool = true,
  619. deleted_set: std.AutoHashMapUnmanaged(u64, void) = .empty,
  620. scratch: std.mem.Allocator,
  621. pub fn init(db: *Db, scratch: std.mem.Allocator, skip_deleted: bool) Scanner {
  622. var s = Scanner{ .db = db, .scratch = scratch, .skip_deleted = skip_deleted };
  623. s.size = fio.fileSize(scratch, db.data_path) orelse 0;
  624. if (s.size > 0) {
  625. s.fd = fio.openRead(scratch, db.data_path) orelse -1;
  626. }
  627. if (skip_deleted) s.deleted_set = db.deletedSet(scratch);
  628. return s;
  629. }
  630. pub fn deinit(self: *Scanner) void {
  631. if (self.fd >= 0) fio.close(self.fd);
  632. self.fd = -1;
  633. }
  634. /// Returns the next record (decoded into `arena`) or null at EOF.
  635. pub fn next(self: *Scanner, arena: std.mem.Allocator) Error!?Record {
  636. while (true) {
  637. if (self.fd < 0 or self.offset + 4 > self.size) return null;
  638. var hdr: [4]u8 = undefined;
  639. const got = fio.pread(self.fd, &hdr, self.offset) orelse return Error.IoError;
  640. if (got < 4) return null;
  641. const rec_size = std.mem.readInt(u32, &hdr, .little);
  642. if (rec_size <= 4) return Error.CorruptRecord;
  643. if (self.offset + rec_size > self.size) return null; // incomplete tail
  644. const off = self.offset;
  645. self.offset += rec_size;
  646. if (self.skip_deleted and self.deleted_set.contains(off)) continue;
  647. const buf = arena.alloc(u8, rec_size - 4) catch return Error.OutOfMemory;
  648. const got2 = fio.pread(self.fd, buf, off + 4) orelse return Error.IoError;
  649. if (got2 < rec_size - 4) return Error.Truncated;
  650. const d = msgpack.decode(arena, buf) catch return Error.CorruptRecord;
  651. return Record{ .value = d.value, .off = off, .len = rec_size };
  652. }
  653. }
  654. /// Like next() but without decoding — yields the raw framed bytes.
  655. pub fn nextRaw(self: *Scanner, arena: std.mem.Allocator) Error!?struct { bytes: []u8, off: u64 } {
  656. while (true) {
  657. if (self.fd < 0 or self.offset + 4 > self.size) return null;
  658. var hdr: [4]u8 = undefined;
  659. const got = fio.pread(self.fd, &hdr, self.offset) orelse return Error.IoError;
  660. if (got < 4) return null;
  661. const rec_size = std.mem.readInt(u32, &hdr, .little);
  662. if (rec_size <= 4) return Error.CorruptRecord;
  663. if (self.offset + rec_size > self.size) return null;
  664. const off = self.offset;
  665. self.offset += rec_size;
  666. if (self.skip_deleted and self.deleted_set.contains(off)) continue;
  667. const buf = arena.alloc(u8, rec_size) catch return Error.OutOfMemory;
  668. const got2 = fio.pread(self.fd, buf, off) orelse return Error.IoError;
  669. if (got2 < rec_size) return Error.Truncated;
  670. return .{ .bytes = buf, .off = off };
  671. }
  672. }
  673. };
  674. /// Read + decode one record by location.
  675. pub fn readAt(self: *Db, arena: std.mem.Allocator, loc: Loc) Error!msgpack.Value {
  676. const fd = fio.openRead(self.alloc, self.data_path) orelse return Error.IoError;
  677. defer fio.close(fd);
  678. if (loc.len <= 4) return Error.CorruptRecord;
  679. const buf = arena.alloc(u8, loc.len - 4) catch return Error.OutOfMemory;
  680. const got = fio.pread(fd, buf, loc.off + 4) orelse return Error.IoError;
  681. if (got < buf.len) return Error.Truncated;
  682. const d = msgpack.decode(arena, buf) catch return Error.CorruptRecord;
  683. return d.value;
  684. }
  685. // =====================================================================
  686. // Index management
  687. // =====================================================================
  688. fn keyFromValue(v: msgpack.Value) ?Key {
  689. return switch (v) {
  690. .number => |n| Key{ .num = n },
  691. .str => |s| Key{ .str = s },
  692. else => null, // null/bool/objects are not indexed (see header note)
  693. };
  694. }
  695. /// Duplicate a key's string into the db arena so it outlives the op arena.
  696. fn ownKey(self: *Db, k: Key) Error!Key {
  697. return switch (k) {
  698. .num => k,
  699. .str => |s| Key{ .str = self.alloc.dupe(u8, s) catch return Error.OutOfMemory },
  700. };
  701. }
  702. fn lowerBound(entries: []const Entry, key: Key) usize {
  703. var lo: usize = 0;
  704. var hi: usize = entries.len;
  705. while (lo < hi) {
  706. const mid = lo + (hi - lo) / 2;
  707. if (cmpKeys(entries[mid].key, key) < 0) {
  708. lo = mid + 1;
  709. } else {
  710. hi = mid;
  711. }
  712. }
  713. return lo;
  714. }
  715. /// Binary search: all entries whose key equals `key` (appended to `out`).
  716. pub fn indexGet(self: *Db, field: []const u8, key: Key, alloc: std.mem.Allocator, out: *std.ArrayList(Entry)) Error!void {
  717. const fi = self.idx.getPtr(field) orelse return Error.NoSuchIndex;
  718. const parsed = self.parseKeyForField(field, key);
  719. const entries = fi.entries.items;
  720. var i = lowerBound(entries, parsed);
  721. while (i < entries.len and cmpKeys(entries[i].key, parsed) == 0) : (i += 1) {
  722. out.append(alloc, entries[i]) catch return Error.OutOfMemory;
  723. }
  724. }
  725. /// Range scan [from, to] (either side optional), ascending.
  726. pub fn indexRange(self: *Db, field: []const u8, from: ?Key, to: ?Key, alloc: std.mem.Allocator, out: *std.ArrayList(Entry)) Error!void {
  727. const fi = self.idx.getPtr(field) orelse return Error.NoSuchIndex;
  728. const entries = fi.entries.items;
  729. var i: usize = if (from) |f| lowerBound(entries, self.parseKeyForField(field, f)) else 0;
  730. const to_key: ?Key = if (to) |t| self.parseKeyForField(field, t) else null;
  731. while (i < entries.len) : (i += 1) {
  732. if (to_key) |t| {
  733. if (cmpKeys(entries[i].key, t) > 0) break;
  734. }
  735. out.append(alloc, entries[i]) catch return Error.OutOfMemory;
  736. }
  737. }
  738. /// _parseKey semantics: numeric index fields parse string keys via
  739. /// parseInt; on NaN the key stays a string.
  740. fn parseKeyForField(self: *Db, field: []const u8, key: Key) Key {
  741. const fi = self.idx.get(field) orelse return key;
  742. if (fi.typ == .numeric and key == .str) {
  743. if (jsParseInt(key.str)) |n| return Key{ .num = n };
  744. }
  745. return key;
  746. }
  747. fn indexInsertEntry(self: *Db, field: []const u8, key: Key, loc: Loc) Error!void {
  748. const fi = self.idx.getPtr(field) orelse return;
  749. const owned = try self.ownKey(key);
  750. const e = Entry{ .key = owned, .off = loc.off, .len = loc.len };
  751. // insert at sorted position
  752. const pos = blk: {
  753. var lo: usize = 0;
  754. var hi: usize = fi.entries.items.len;
  755. while (lo < hi) {
  756. const mid = lo + (hi - lo) / 2;
  757. if (entryLess({}, fi.entries.items[mid], e)) {
  758. lo = mid + 1;
  759. } else {
  760. hi = mid;
  761. }
  762. }
  763. break :blk lo;
  764. };
  765. fi.entries.insert(self.alloc, pos, e) catch return Error.OutOfMemory;
  766. }
  767. /// Index a record's fields at loc.
  768. fn indexInsertRecord(self: *Db, rec: msgpack.Value, loc: Loc) Error!void {
  769. for (self.index_fields.items) |field| {
  770. const v = rec.get(field) orelse continue;
  771. const key = keyFromValue(v) orelse continue;
  772. try self.indexInsertEntry(field, key, loc);
  773. }
  774. self.covered_bytes = @max(self.covered_bytes, loc.off + loc.len);
  775. }
  776. /// Remove all entries pointing at offset `off` (all fields).
  777. fn indexRemoveOffset(self: *Db, off: u64) void {
  778. var it = self.idx.iterator();
  779. while (it.next()) |kv| {
  780. const list = &kv.value_ptr.entries;
  781. var i: usize = 0;
  782. while (i < list.items.len) {
  783. if (list.items[i].off == off) {
  784. _ = list.orderedRemove(i);
  785. } else {
  786. i += 1;
  787. }
  788. }
  789. }
  790. }
  791. fn indexPath(self: *Db, buf: []u8, field: []const u8) []const u8 {
  792. // <dir>/<base>.<field>.txt — derive from idxstate path (…/base.idxstate.json)
  793. const prefix = self.idxstate_path[0 .. self.idxstate_path.len - "idxstate.json".len];
  794. return std.fmt.bufPrint(buf, "{s}{s}.txt", .{ prefix, field }) catch buf[0..0];
  795. }
  796. fn initIndexes(self: *Db, force_rebuild: bool) Error!void {
  797. var need_rebuild = force_rebuild;
  798. const data_size = fio.fileSize(self.alloc, self.data_path) orelse 0;
  799. if (!need_rebuild) {
  800. for (self.index_fields.items) |field| {
  801. var pbuf: [512]u8 = undefined;
  802. const p = self.indexPath(&pbuf, field);
  803. const isize_ = fio.fileSize(self.alloc, p);
  804. if (isize_ == null or isize_.? == 0) {
  805. if (data_size > 0) {
  806. need_rebuild = true;
  807. break;
  808. }
  809. }
  810. }
  811. }
  812. if (!need_rebuild) {
  813. // coveredBytes from idxstate — no state file means rebuild (JS)
  814. if (self.readIdxState()) |cb| {
  815. self.covered_bytes = cb;
  816. } else {
  817. need_rebuild = true;
  818. }
  819. }
  820. if (need_rebuild) {
  821. try self.rebuildIndexes();
  822. return;
  823. }
  824. // Load index files into memory
  825. for (self.index_fields.items) |field| {
  826. var pbuf: [512]u8 = undefined;
  827. const p = self.indexPath(&pbuf, field);
  828. const content = fio.readAll(self.alloc, p) orelse continue;
  829. defer self.alloc.free(content);
  830. self.loadIndexLines(field, content);
  831. }
  832. // sort (files are sorted by JS localeCompare; re-sort under our order)
  833. var it = self.idx.iterator();
  834. while (it.next()) |kv| {
  835. std.sort.pdq(Entry, kv.value_ptr.entries.items, {}, entryLess);
  836. }
  837. if (data_size > self.covered_bytes) self.catchUp(data_size);
  838. }
  839. fn loadIndexLines(self: *Db, field: []const u8, content: []const u8) void {
  840. const fi = self.idx.getPtr(field) orelse return;
  841. var lines = std.mem.splitScalar(u8, content, '\n');
  842. while (lines.next()) |line| {
  843. if (line.len == 0) continue;
  844. // key = up to FIRST comma (JS split(',')[0]); then offset, length
  845. const c1 = std.mem.indexOfScalar(u8, line, ',') orelse continue;
  846. const rest = line[c1 + 1 ..];
  847. const c2 = std.mem.indexOfScalar(u8, rest, ',') orelse continue;
  848. const off_s = rest[0..c2];
  849. const len_s = rest[c2 + 1 ..];
  850. const off = std.fmt.parseInt(u64, off_s, 10) catch continue;
  851. const len = std.fmt.parseInt(u64, len_s, 10) catch continue;
  852. if (len == 0) continue;
  853. const key_s = line[0..c1];
  854. var key: Key = undefined;
  855. if (fi.typ == .numeric) {
  856. key = if (jsParseInt(key_s)) |n| Key{ .num = n } else Key{ .str = self.alloc.dupe(u8, key_s) catch return };
  857. } else {
  858. key = Key{ .str = self.alloc.dupe(u8, key_s) catch return };
  859. }
  860. fi.entries.append(self.alloc, .{ .key = key, .off = off, .len = len }) catch return;
  861. }
  862. }
  863. fn rebuildIndexes(self: *Db) Error!void {
  864. var it0 = self.idx.iterator();
  865. while (it0.next()) |kv| kv.value_ptr.entries.clearRetainingCapacity();
  866. var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);
  867. defer scratch_state.deinit();
  868. const scratch = scratch_state.allocator();
  869. // JS rebuild scans ALL records (no tombstone filter — deleted offsets
  870. // simply get filtered at read time). Mirror that.
  871. var scanner = Scanner.init(self, scratch, false);
  872. defer scanner.deinit();
  873. while (try scanner.next(scratch)) |rec| {
  874. try self.indexInsertRecord(rec.value, .{ .off = rec.off, .len = rec.len });
  875. }
  876. self.covered_bytes = fio.fileSize(self.alloc, self.data_path) orelse 0;
  877. try self.writeIndexFiles();
  878. try self.writeIdxState();
  879. }
  880. fn catchUp(self: *Db, file_size: u64) void {
  881. var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);
  882. defer scratch_state.deinit();
  883. const scratch = scratch_state.allocator();
  884. var scanner = Scanner.init(self, scratch, false);
  885. defer scanner.deinit();
  886. scanner.offset = self.covered_bytes;
  887. while (true) {
  888. const maybe = scanner.next(scratch) catch break;
  889. const rec = maybe orelse break;
  890. self.indexInsertRecord(rec.value, .{ .off = rec.off, .len = rec.len }) catch break;
  891. }
  892. self.covered_bytes = @max(self.covered_bytes, file_size);
  893. }
  894. /// A file that already holds these exact bytes is left alone — its mtime
  895. /// included — so closing a handle that changed nothing writes nothing.
  896. fn fileHolds(alloc: std.mem.Allocator, path: []const u8, bytes: []const u8) bool {
  897. const disk = fio.readAll(alloc, path) orelse return false;
  898. defer alloc.free(disk);
  899. return std.mem.eql(u8, disk, bytes);
  900. }
  901. fn writeIndexFiles(self: *Db) Error!void {
  902. for (self.index_fields.items) |field| {
  903. const fi = self.idx.getPtr(field) orelse continue;
  904. var out: std.ArrayList(u8) = .empty;
  905. defer out.deinit(self.alloc);
  906. for (fi.entries.items) |e| {
  907. var kbuf: [32]u8 = undefined;
  908. const ks = switch (e.key) {
  909. .num => |n| jsNumFmt(&kbuf, n),
  910. .str => |s| s,
  911. };
  912. try self.appendFmt(&out, "{s},{d},{d}\n", .{ ks, e.off, e.len });
  913. }
  914. var pbuf: [512]u8 = undefined;
  915. const p = self.indexPath(&pbuf, field);
  916. if (fileHolds(self.alloc, p, out.items)) continue;
  917. var tbuf: [512]u8 = undefined;
  918. const tmp = std.fmt.bufPrint(&tbuf, "{s}.{d}.tmp", .{ p, linux.getpid() }) catch return Error.IoError;
  919. const fd = fio.openTrunc(self.alloc, tmp) orelse return Error.IoError;
  920. const ok = fio.writeAll(fd, out.items);
  921. fio.close(fd);
  922. if (!ok) return Error.IoError;
  923. if (!fio.rename(self.alloc, tmp, p)) return Error.IoError;
  924. }
  925. }
  926. fn readIdxState(self: *Db) ?u64 {
  927. const content = fio.readAll(self.alloc, self.idxstate_path) orelse return null;
  928. defer self.alloc.free(content);
  929. const parsed = std.json.parseFromSlice(std.json.Value, self.alloc, content, .{}) catch return null;
  930. defer parsed.deinit();
  931. if (parsed.value != .object) return null;
  932. const v = parsed.value.object.get("coveredBytes") orelse return null;
  933. return switch (v) {
  934. .integer => |i| @intCast(@max(i, 0)),
  935. .float => |f| @intFromFloat(@max(f, 0)),
  936. else => null,
  937. };
  938. }
  939. fn writeIdxState(self: *Db) Error!void {
  940. var buf: [128]u8 = undefined;
  941. const json = std.fmt.bufPrint(&buf, "{{\"coveredBytes\":{d}}}", .{self.covered_bytes}) catch return Error.IoError;
  942. if (fileHolds(self.alloc, self.idxstate_path, json)) return;
  943. var tbuf: [512]u8 = undefined;
  944. const tmp = std.fmt.bufPrint(&tbuf, "{s}.{d}.tmp", .{ self.idxstate_path, linux.getpid() }) catch return Error.IoError;
  945. const fd = fio.openTrunc(self.alloc, tmp) orelse return Error.IoError;
  946. const ok = fio.writeAll(fd, json);
  947. fio.close(fd);
  948. if (!ok) return Error.IoError;
  949. if (!fio.rename(self.alloc, tmp, self.idxstate_path)) return Error.IoError;
  950. }
  951. fn persistIndexesLocked(self: *Db) Error!void {
  952. try self.writeIndexFiles();
  953. // coverage can only grow (JS: max of ours and on-disk state)
  954. if (self.readIdxState()) |cb| self.covered_bytes = @max(self.covered_bytes, cb);
  955. try self.writeIdxState();
  956. }
  957. // =====================================================================
  958. // Operations
  959. // =====================================================================
  960. /// A mutable record under construction (op-arena entries).
  961. pub const MutableRecord = struct {
  962. entries: std.ArrayList(msgpack.Entry) = .empty,
  963. pub fn get(self: *const MutableRecord, key: []const u8) ?msgpack.Value {
  964. for (self.entries.items) |e| {
  965. if (e.key == .str and std.mem.eql(u8, e.key.str, key)) return e.value;
  966. }
  967. return null;
  968. }
  969. pub fn set(self: *MutableRecord, alloc: std.mem.Allocator, key: []const u8, v: msgpack.Value) Error!void {
  970. for (self.entries.items) |*e| {
  971. if (e.key == .str and std.mem.eql(u8, e.key.str, key)) {
  972. e.value = v;
  973. return;
  974. }
  975. }
  976. self.entries.append(alloc, .{ .key = .{ .str = key }, .value = v }) catch return Error.OutOfMemory;
  977. }
  978. pub fn toValue(self: *const MutableRecord) msgpack.Value {
  979. return .{ .map = self.entries.items };
  980. }
  981. };
  982. pub const InsertResult = union(enum) {
  983. pk_num: f64,
  984. pk_str: []const u8, // op-arena
  985. record: msgpack.Value,
  986. };
  987. /// insert() — auto primary key, unique checks, append, index.
  988. /// `record` must be a map value; `arena` is the op arena (record memory).
  989. pub fn insert(self: *Db, arena: std.mem.Allocator, record: msgpack.Value, skip_pk: bool) Error!InsertResult {
  990. if (record != .map) return Error.CorruptRecord;
  991. try self.acquireLock();
  992. defer self.releaseLock();
  993. self.refresh();
  994. return self.insertLocked(arena, record, skip_pk);
  995. }
  996. fn insertLocked(self: *Db, arena: std.mem.Allocator, record: msgpack.Value, skip_pk: bool) Error!InsertResult {
  997. // copy into mutable form
  998. var rec = MutableRecord{};
  999. for (record.map) |e| rec.entries.append(arena, e) catch return Error.OutOfMemory;
  1000. // _hasPrimaryKeyValue: present and not null/undefined
  1001. var has_pk_value = false;
  1002. if (self.pk) |p| {
  1003. if (rec.get(p)) |v| {
  1004. has_pk_value = (v != .nil and v != .undef);
  1005. }
  1006. }
  1007. var auto_gen = false;
  1008. if (self.pk != null and !skip_pk and !has_pk_value) {
  1009. if (self.pk_type == .number) {
  1010. rec.set(arena, self.pk.?, .{ .number = self.next_id }) catch return Error.OutOfMemory;
  1011. self.next_id += 1;
  1012. auto_gen = true;
  1013. } else if (self.pk_type == .uuid) {
  1014. // A GENERATED ID IS CHECKED, NOT TRUSTED (ticket #113): the format is a
  1015. // millisecond stamp + 3 base36 chars, so a bulk put collides within a
  1016. // millisecond (two of 7800 in a measured run) and the unique check below
  1017. // skips generated keys — two live rows then shared one pk and the next
  1018. // update() deleted both and re-inserted one. Draw again until it is free.
  1019. var ubuf: [12]u8 = undefined;
  1020. var tries: usize = 0;
  1021. while (true) : (tries += 1) {
  1022. const u = genUuid(&ubuf);
  1023. var taken: std.ArrayList(Loc) = .empty;
  1024. defer taken.deinit(arena);
  1025. try self.findPkLocs(arena, .{ .str = u }, &taken);
  1026. if (taken.items.len == 0 or tries >= 64) break;
  1027. }
  1028. const owned = arena.dupe(u8, ubuf[0..12]) catch return Error.OutOfMemory;
  1029. rec.set(arena, self.pk.?, .{ .str = owned }) catch return Error.OutOfMemory;
  1030. auto_gen = true;
  1031. }
  1032. }
  1033. // unique constraints (skip auto-generated pk)
  1034. if (self.index_fields.items.len > 0 and self.unique_fields.items.len > 0) {
  1035. var del_set = self.deletedSet(arena);
  1036. for (self.unique_fields.items) |field| {
  1037. if (auto_gen and self.pk != null and std.mem.eql(u8, field, self.pk.?)) continue;
  1038. const v = rec.get(field) orelse continue;
  1039. if (v == .undef) continue;
  1040. const key = keyFromValue(v) orelse continue;
  1041. var hits: std.ArrayList(Entry) = .empty;
  1042. defer hits.deinit(arena);
  1043. try self.indexGet(field, key, arena, &hits);
  1044. var live: usize = 0;
  1045. for (hits.items) |h| {
  1046. if (!del_set.contains(h.off)) live += 1;
  1047. }
  1048. if (live > 0) {
  1049. var kbuf: [32]u8 = undefined;
  1050. const ks = switch (key) {
  1051. .num => |n| jsNumFmt(&kbuf, n),
  1052. .str => |s| s,
  1053. };
  1054. self.last_error = std.fmt.bufPrint(&self.last_error_buf, "Duplicate key: {s}={s}", .{ field, ks }) catch "Duplicate key";
  1055. return Error.DuplicateKey;
  1056. }
  1057. }
  1058. }
  1059. // serialize + append
  1060. const framed = msgpack.serialize(arena, rec.toValue()) catch return Error.OutOfMemory;
  1061. const offset = fio.fileSize(self.alloc, self.data_path) orelse 0;
  1062. const fd = fio.openAppend(self.alloc, self.data_path) orelse return Error.IoError;
  1063. const ok = fio.writeAll(fd, framed);
  1064. fio.close(fd);
  1065. if (!ok) return Error.IoError;
  1066. const loc = Loc{ .off = offset, .len = framed.len };
  1067. if (self.index_fields.items.len > 0) {
  1068. try self.indexInsertRecord(rec.toValue(), loc);
  1069. }
  1070. if (self.pk != null and self.pk_type == .number and !skip_pk and !has_pk_value) {
  1071. try self.persistMetaLocked();
  1072. }
  1073. if (self.pk) |p| {
  1074. const v = rec.get(p) orelse return Error.CorruptRecord;
  1075. return switch (v) {
  1076. .number => |n| InsertResult{ .pk_num = n },
  1077. .str => |s| InsertResult{ .pk_str = s },
  1078. else => InsertResult{ .record = rec.toValue() },
  1079. };
  1080. }
  1081. return InsertResult{ .record = rec.toValue() };
  1082. }
  1083. /// Locate live records by primary key (index-backed when available).
  1084. pub fn findPkLocs(self: *Db, arena: std.mem.Allocator, key: Key, out: *std.ArrayList(Loc)) Error!void {
  1085. if (self.pk == null) return Error.NoPrimaryKey;
  1086. self.refresh();
  1087. var del_set = self.deletedSet(arena);
  1088. if (self.index_fields.items.len > 0) {
  1089. var hits: std.ArrayList(Entry) = .empty;
  1090. defer hits.deinit(arena);
  1091. try self.indexGet(self.pk.?, key, arena, &hits);
  1092. for (hits.items) |h| {
  1093. if (del_set.contains(h.off)) continue;
  1094. out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;
  1095. }
  1096. return;
  1097. }
  1098. // no indexes: full scan comparing the pk field
  1099. var scanner = Scanner.init(self, arena, true);
  1100. defer scanner.deinit();
  1101. while (try scanner.next(arena)) |rec| {
  1102. const v = rec.value.get(self.pk.?) orelse continue;
  1103. const k = keyFromValue(v) orelse continue;
  1104. if (cmpKeys(self.parseKeyForField(self.pk.?, k), self.parseKeyForField(self.pk.?, key)) == 0) {
  1105. out.append(arena, .{ .off = rec.off, .len = rec.len }) catch return Error.OutOfMemory;
  1106. }
  1107. }
  1108. }
  1109. /// Locate live records by secondary index equality.
  1110. pub fn findIndexLocs(self: *Db, arena: std.mem.Allocator, field: []const u8, key: Key, out: *std.ArrayList(Loc)) Error!void {
  1111. self.refresh();
  1112. var del_set = self.deletedSet(arena);
  1113. var hits: std.ArrayList(Entry) = .empty;
  1114. defer hits.deinit(arena);
  1115. try self.indexGet(field, key, arena, &hits);
  1116. for (hits.items) |h| {
  1117. if (del_set.contains(h.off)) continue;
  1118. out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;
  1119. }
  1120. }
  1121. /// Locate live records by index range [from, to] ascending.
  1122. pub fn findRangeLocs(self: *Db, arena: std.mem.Allocator, field: []const u8, from: ?Key, to: ?Key, out: *std.ArrayList(Loc)) Error!void {
  1123. self.refresh();
  1124. var del_set = self.deletedSet(arena);
  1125. var hits: std.ArrayList(Entry) = .empty;
  1126. defer hits.deinit(arena);
  1127. try self.indexRange(field, from, to, arena, &hits);
  1128. var seen: std.AutoHashMapUnmanaged(u64, void) = .empty;
  1129. for (hits.items) |h| {
  1130. if (del_set.contains(h.off)) continue;
  1131. if (seen.contains(h.off)) continue;
  1132. seen.put(arena, h.off, {}) catch {};
  1133. out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;
  1134. }
  1135. }
  1136. /// delete by primary key. Returns the deleted records (decoded into arena).
  1137. pub fn deletePk(self: *Db, arena: std.mem.Allocator, key: Key, out_records: *std.ArrayList(msgpack.Value)) Error!usize {
  1138. try self.acquireLock();
  1139. defer self.releaseLock();
  1140. self.refresh();
  1141. return self.deletePkLocked(arena, key, out_records);
  1142. }
  1143. fn deletePkLocked(self: *Db, arena: std.mem.Allocator, key: Key, out_records: *std.ArrayList(msgpack.Value)) Error!usize {
  1144. var locs: std.ArrayList(Loc) = .empty;
  1145. defer locs.deinit(arena);
  1146. try self.findPkLocs(arena, key, &locs);
  1147. for (locs.items) |loc| {
  1148. const rec = try self.readAt(arena, loc);
  1149. out_records.append(arena, rec) catch return Error.OutOfMemory;
  1150. self.deleted.append(self.alloc, loc.off) catch return Error.OutOfMemory;
  1151. self.indexRemoveOffset(loc.off);
  1152. }
  1153. try self.persistMetaLocked();
  1154. return locs.items.len;
  1155. }
  1156. /// update by primary key: delete + insert(skip_pk) — JS update() semantics.
  1157. /// `new_record` must already carry the primary key (the JS callback
  1158. /// contract: the record keeps its pk unless the caller removes it).
  1159. pub fn updatePk(self: *Db, arena: std.mem.Allocator, key: Key, new_record: msgpack.Value) Error!usize {
  1160. if (self.pk) |p| {
  1161. const carried = for (new_record.map) |e| {
  1162. if (e.key == .str and std.mem.eql(u8, e.key.str, p)) break e.value != .nil and e.value != .undef;
  1163. } else false;
  1164. if (!carried) return Error.RecordLacksPrimaryKey;
  1165. }
  1166. try self.acquireLock();
  1167. defer self.releaseLock();
  1168. self.refresh();
  1169. var old: std.ArrayList(msgpack.Value) = .empty;
  1170. defer old.deinit(arena);
  1171. const n = try self.deletePkLocked(arena, key, &old);
  1172. if (n == 0) return 0;
  1173. var i: usize = 0;
  1174. while (i < n) : (i += 1) {
  1175. _ = try self.insertLocked(arena, new_record, true);
  1176. }
  1177. return n;
  1178. }
  1179. /// compact() — rewrite the data file without tombstones, rebuild indexes.
  1180. pub fn compact(self: *Db) Error!void {
  1181. try self.acquireLock();
  1182. defer self.releaseLock();
  1183. self.refresh();
  1184. try self.compactLocked();
  1185. if (self.index_fields.items.len > 0) {
  1186. try self.rebuildIndexes();
  1187. }
  1188. }
  1189. fn compactLocked(self: *Db) Error!void {
  1190. var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);
  1191. defer scratch_state.deinit();
  1192. const scratch = scratch_state.allocator();
  1193. var tbuf: [512]u8 = undefined;
  1194. const tmp = std.fmt.bufPrint(&tbuf, "{s}.tmp", .{self.data_path}) catch return Error.IoError;
  1195. const out_fd = fio.openTrunc(self.alloc, tmp) orelse return Error.IoError;
  1196. var scanner = Scanner.init(self, scratch, true);
  1197. var ok = true;
  1198. while (true) {
  1199. const maybe = scanner.nextRaw(scratch) catch {
  1200. ok = false;
  1201. break;
  1202. };
  1203. const raw = maybe orelse break;
  1204. if (!fio.writeAll(out_fd, raw.bytes)) {
  1205. ok = false;
  1206. break;
  1207. }
  1208. }
  1209. scanner.deinit();
  1210. fio.close(out_fd);
  1211. if (!ok) return Error.IoError;
  1212. if (!fio.rename(self.alloc, tmp, self.data_path)) return Error.IoError;
  1213. self.deleted.clearRetainingCapacity();
  1214. try self.persistMetaLocked();
  1215. }
  1216. pub fn persistNow(self: *Db) Error!void {
  1217. try self.acquireLock();
  1218. defer self.releaseLock();
  1219. self.refresh();
  1220. try self.persistIndexesLocked();
  1221. }
  1222. };
  1223. // =========================================================================
  1224. // Tests
  1225. // =========================================================================
  1226. const testing = std.testing;
  1227. fn tmpBase(buf: []u8, comptime name: []const u8) []const u8 {
  1228. return std.fmt.bufPrint(buf, "/tmp/mpackdb-zigtest-{d}-" ++ name, .{linux.getpid()}) catch unreachable;
  1229. }
  1230. fn cleanup(alloc: std.mem.Allocator, base: []const u8) void {
  1231. var buf: [512]u8 = undefined;
  1232. const suffixes = [_][]const u8{ ".mpack", ".meta.json", ".idxstate.json", ".lock", ".id.txt", ".email.txt", ".age.txt", ".uuid.txt" };
  1233. for (suffixes) |suffix| {
  1234. const p = std.fmt.bufPrint(&buf, "{s}{s}", .{ base, suffix }) catch continue;
  1235. _ = fio.unlink(alloc, p);
  1236. }
  1237. }
  1238. fn strKey(s: []const u8) Key {
  1239. return .{ .str = s };
  1240. }
  1241. fn numKey(n: f64) Key {
  1242. return .{ .num = n };
  1243. }
  1244. fn makeUser(arena: std.mem.Allocator, name: []const u8, email: []const u8, age: f64) !msgpack.Value {
  1245. const entries = try arena.alloc(msgpack.Entry, 3);
  1246. entries[0] = .{ .key = .{ .str = "name" }, .value = .{ .str = name } };
  1247. entries[1] = .{ .key = .{ .str = "email" }, .value = .{ .str = email } };
  1248. entries[2] = .{ .key = .{ .str = "age" }, .value = .{ .number = age } };
  1249. return .{ .map = entries };
  1250. }
  1251. test "engine: insert/find/delete/update with numeric pk + indexes" {
  1252. var base_buf: [128]u8 = undefined;
  1253. const base = tmpBase(&base_buf, "crud");
  1254. cleanup(testing.allocator, base);
  1255. defer cleanup(testing.allocator, base);
  1256. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1257. defer arena_state.deinit();
  1258. const arena = arena_state.allocator();
  1259. const db = try Db.open(testing.allocator, base, .{
  1260. .primary_key = "*id",
  1261. .indexes = &.{ "email", "*age" },
  1262. });
  1263. // insert three — a fresh table counts from 1 (ticket #4)
  1264. const r1 = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 30), false);
  1265. try testing.expectEqual(@as(f64, 1), r1.pk_num);
  1266. const r2 = try db.insert(arena, try makeUser(arena, "Bob", "[email protected]", 25), false);
  1267. try testing.expectEqual(@as(f64, 2), r2.pk_num);
  1268. _ = try db.insert(arena, try makeUser(arena, "Carol", "[email protected]", 35), false);
  1269. // find by pk
  1270. var locs: std.ArrayList(Loc) = .empty;
  1271. try db.findPkLocs(arena, numKey(2), &locs);
  1272. try testing.expectEqual(@as(usize, 1), locs.items.len);
  1273. const bob = try db.readAt(arena, locs.items[0]);
  1274. try testing.expectEqualStrings("Bob", bob.get("name").?.str);
  1275. // find by secondary index
  1276. var locs2: std.ArrayList(Loc) = .empty;
  1277. try db.findIndexLocs(arena, "email", strKey("[email protected]"), &locs2);
  1278. try testing.expectEqual(@as(usize, 1), locs2.items.len);
  1279. // range on numeric index: age 26..40 → Alice(30), Carol(35)
  1280. var locs3: std.ArrayList(Loc) = .empty;
  1281. try db.findRangeLocs(arena, "age", numKey(26), numKey(40), &locs3);
  1282. try testing.expectEqual(@as(usize, 2), locs3.items.len);
  1283. // update Bob's age
  1284. var bob_new = Db.MutableRecord{};
  1285. for (bob.map) |e| try bob_new.entries.append(arena, e);
  1286. try bob_new.set(arena, "age", .{ .number = 26 });
  1287. const updated = try db.updatePk(arena, numKey(2), bob_new.toValue());
  1288. try testing.expectEqual(@as(usize, 1), updated);
  1289. var locs4: std.ArrayList(Loc) = .empty;
  1290. try db.findPkLocs(arena, numKey(2), &locs4);
  1291. try testing.expectEqual(@as(usize, 1), locs4.items.len);
  1292. const bob2 = try db.readAt(arena, locs4.items[0]);
  1293. try testing.expectEqual(@as(f64, 26), bob2.get("age").?.number);
  1294. // delete Alice
  1295. var deleted_recs: std.ArrayList(msgpack.Value) = .empty;
  1296. const dn = try db.deletePk(arena, numKey(1), &deleted_recs);
  1297. try testing.expectEqual(@as(usize, 1), dn);
  1298. try testing.expectEqualStrings("Alice", deleted_recs.items[0].get("name").?.str);
  1299. var locs5: std.ArrayList(Loc) = .empty;
  1300. try db.findPkLocs(arena, numKey(1), &locs5);
  1301. try testing.expectEqual(@as(usize, 0), locs5.items.len);
  1302. db.close();
  1303. // reopen (compaction drops the tombstone) and verify persistence
  1304. const db2 = try Db.open(testing.allocator, base, .{
  1305. .primary_key = "*id",
  1306. .indexes = &.{ "email", "*age" },
  1307. });
  1308. defer db2.close();
  1309. var locs6: std.ArrayList(Loc) = .empty;
  1310. try db2.findPkLocs(arena, numKey(2), &locs6);
  1311. try testing.expectEqual(@as(usize, 1), locs6.items.len);
  1312. const bob3 = try db2.readAt(arena, locs6.items[0]);
  1313. try testing.expectEqual(@as(f64, 26), bob3.get("age").?.number);
  1314. // nextId continues after reopen
  1315. const r4 = try db2.insert(arena, try makeUser(arena, "Dan", "[email protected]", 40), false);
  1316. try testing.expectEqual(@as(f64, 4), r4.pk_num);
  1317. }
  1318. test "engine: a stored counter wins over FIRST_ID (ticket #4)" {
  1319. var base_buf: [128]u8 = undefined;
  1320. const base = tmpBase(&base_buf, "counter");
  1321. cleanup(testing.allocator, base);
  1322. defer cleanup(testing.allocator, base);
  1323. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1324. defer arena_state.deinit();
  1325. const arena = arena_state.allocator();
  1326. // A table the JS mpackdb created and never wrote to: its meta says 0, and
  1327. // an existing table keeps its own counter.
  1328. var mbuf: [160]u8 = undefined;
  1329. const meta_path = try std.fmt.bufPrint(&mbuf, "{s}.meta.json", .{base});
  1330. const fd = fio.openTrunc(testing.allocator, meta_path).?;
  1331. try testing.expect(fio.writeAll(fd, "{\"nextId\":0,\"deleted\":[],\"version\":1}"));
  1332. fio.close(fd);
  1333. const db = try Db.open(testing.allocator, base, .{ .primary_key = "*id" });
  1334. defer db.close();
  1335. const r = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 1), false);
  1336. try testing.expectEqual(@as(f64, 0), r.pk_num);
  1337. }
  1338. test "engine: unique index rejects duplicates" {
  1339. var base_buf: [128]u8 = undefined;
  1340. const base = tmpBase(&base_buf, "uniq");
  1341. cleanup(testing.allocator, base);
  1342. defer cleanup(testing.allocator, base);
  1343. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1344. defer arena_state.deinit();
  1345. const arena = arena_state.allocator();
  1346. const db = try Db.open(testing.allocator, base, .{
  1347. .primary_key = "*id",
  1348. .indexes = &.{"!email"},
  1349. });
  1350. defer db.close();
  1351. _ = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 30), false);
  1352. const dup = db.insert(arena, try makeUser(arena, "Evil", "[email protected]", 31), false);
  1353. try testing.expectError(Error.DuplicateKey, dup);
  1354. try testing.expect(std.mem.indexOf(u8, db.last_error, "[email protected]") != null);
  1355. }
  1356. test "engine: uuid primary key" {
  1357. var base_buf: [128]u8 = undefined;
  1358. const base = tmpBase(&base_buf, "uuid");
  1359. cleanup(testing.allocator, base);
  1360. defer cleanup(testing.allocator, base);
  1361. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1362. defer arena_state.deinit();
  1363. const arena = arena_state.allocator();
  1364. const db = try Db.open(testing.allocator, base, .{ .primary_key = "@uuid" });
  1365. defer db.close();
  1366. const r = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 1), false);
  1367. try testing.expectEqual(@as(usize, 12), r.pk_str.len);
  1368. var locs: std.ArrayList(Loc) = .empty;
  1369. try db.findPkLocs(arena, strKey(r.pk_str), &locs);
  1370. try testing.expectEqual(@as(usize, 1), locs.items.len);
  1371. }
  1372. test "jsNumFmt / jsParseInt / uuid shape" {
  1373. var buf: [32]u8 = undefined;
  1374. try testing.expectEqualStrings("42", jsNumFmt(&buf, 42));
  1375. try testing.expectEqualStrings("-7", jsNumFmt(&buf, -7));
  1376. try testing.expectEqualStrings("1.5", jsNumFmt(&buf, 1.5));
  1377. try testing.expectEqual(@as(f64, 1), jsParseInt("1.5").?);
  1378. try testing.expectEqual(@as(f64, -12), jsParseInt("-12abc").?);
  1379. try testing.expect(jsParseInt("abc") == null);
  1380. var ubuf: [12]u8 = undefined;
  1381. const u = genUuid(&ubuf);
  1382. try testing.expectEqual(@as(usize, 12), u.len);
  1383. for (u) |ch| try testing.expect((ch >= '0' and ch <= '9') or (ch >= 'a' and ch <= 'z'));
  1384. }

Branches

Latest commits

  • 2149e902notes mission 002 (4/4): code order — README file map + same-output test, STATUS, LOG, report; tests/letcount.py, tests/realdata-baseline.mjs, tests/realdata-compare.pymre
  • 47db4fc1notes mission 002 (3/4): code order — let only where reassigned (58 dropped; 34 left: 19 reassigned, 15 loop-bound); gates 50/0 + 18/0, live-data run = step 2mre
  • 3f8cb383notes mission 002 (2/4): code order — topics, map, thin wrappers: lib/util.hl, lib/notes(-helpers).hl, lib/users.hl (+userOfLoginCode, tagOf), lib/api(-helpers).hl; project.hl = map; login routes take &req/&sessions (failed-login reason now kept); dead notes#1 functions removed; gates 50/0 + 18/0mre
  • 8cda5412notes mission 002 (1/4): code order — files moved: lib/notes.hl, lib/users.hl, lib/jsoncheck.hl, components/styles.hl (imports only); gates 50/0 + 18/0, live-data run identicalmre
  • 47dad68bnotes: Hybriel master 06617221 (plugin allocators 3a781359 + 413f60e4, mpackdb 2cb7ae5e, http1 773de63e); gates 50/0 + 18/0mre
  • a4a2b2aeantcolony#40: tracker missions moved too — references to them in missions/reports/LOG.md updatedmre
  • 3dea0ef2notes: Hybriel master 190aa11d (fc838894 GC correctness, #126 closure scopes, #127); gate 50/0mre
  • e27c7d71notes: Hybriel master 8efba065 (#126 GC by bytes, #48 lambda params copy; audit: nothing to fix; gate 50/0)mre
  • 124613b6antcolony#40: mission references point to the moved missionsmre
  • 1ca2f34dantcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmre
  • 2bebebdanotes: Hybriel master ff51cf46 (re-vendor round)mre
  • 3eff126dnotes#3: installable app (manifest + own icon/favicon; notes' own sw.js kept)mre
  • 9883c540deploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
  • eee693b8deploy.sh: never send .git or .gitignore to Byrodinmre
  • c8904061State of 2026-09-27, before the move to gitoriamre