gitoriaLog in with ident

tickets

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commite01c2b1de01c2b1dtickets#24: installable app (manifest, service worker, offline list), own icon; gate waits for the hello's pongmree01c2b1d/plugins/mpackdb/engine.zig

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

Branches

Latest commits

  • 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