gitoriaLog in with ident

gitoria

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commitfdfb4b1bfdfb4b1bgitoria: Hybriel master 06617221 (plugin allocators 3a781359 + 413f60e4, mpackdb 2cb7ae5e, http1 773de63e); gates 200/0, 46/0, 44/0mrefdfb4b1b/plugins/mpackdb/engine.zig

63.0 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. // THE BUFFERS OF ONE OPERATION ARE NOT THE HANDLE'S (hybriel#126): the meta read
  296. // back and rendered anew, the index files, the paths. On the arena they stayed for
  297. // the handle's life, and the meta grows with every tombstone, so each update leaked
  298. // more than the last (notes: 2000 updates of one note, +100 MB). They are freed now.
  299. tmp: std.mem.Allocator,
  300. // paths
  301. data_path: []const u8,
  302. meta_path: []const u8,
  303. lock_path: []const u8,
  304. idxstate_path: []const u8,
  305. // schema
  306. pk: ?[]const u8 = null,
  307. pk_type: PkType = .string,
  308. index_fields: std.ArrayList([]const u8) = .empty, // pk first (when set), like JS
  309. unique_fields: std.ArrayList([]const u8) = .empty,
  310. idx: std.StringArrayHashMapUnmanaged(FieldIndex) = .empty,
  311. // meta
  312. next_id: f64 = FIRST_ID,
  313. deleted: std.ArrayList(u64) = .empty,
  314. version: u64 = 0,
  315. covered_bytes: u64 = 0,
  316. compact_on_open: bool = true,
  317. stale_lock_timeout_ms: i64 = 30000,
  318. debug: bool = false,
  319. // last DuplicateKey detail for the plugin surface
  320. last_error_buf: [256]u8 = undefined,
  321. last_error: []const u8 = "",
  322. pub fn open(gpa: std.mem.Allocator, db_file: []const u8, opts: OpenOptions) Error!*Db {
  323. if (db_file.len == 0) return Error.NoDbFile;
  324. const self = gpa.create(Db) catch return Error.OutOfMemory;
  325. self.* = .{
  326. .arena_state = std.heap.ArenaAllocator.init(gpa),
  327. .alloc = undefined,
  328. .tmp = gpa,
  329. .data_path = undefined,
  330. .meta_path = undefined,
  331. .lock_path = undefined,
  332. .idxstate_path = undefined,
  333. };
  334. self.alloc = self.arena_state.allocator();
  335. errdefer {
  336. self.arena_state.deinit();
  337. gpa.destroy(self);
  338. }
  339. const a = self.alloc;
  340. self.compact_on_open = opts.compact;
  341. self.stale_lock_timeout_ms = opts.stale_lock_timeout_ms;
  342. self.debug = opts.debug;
  343. // dirname / basename (extension stripped, like JS basename(f, extname(f)))
  344. const dir = std.fs.path.dirname(db_file) orelse ".";
  345. var base = std.fs.path.basename(db_file);
  346. if (std.mem.lastIndexOfScalar(u8, base, '.')) |dot| {
  347. if (dot > 0) base = base[0..dot];
  348. }
  349. fio.mkdirAll(a, dir);
  350. self.data_path = std.fmt.allocPrint(a, "{s}/{s}.mpack", .{ dir, base }) catch return Error.OutOfMemory;
  351. self.meta_path = std.fmt.allocPrint(a, "{s}/{s}.meta.json", .{ dir, base }) catch return Error.OutOfMemory;
  352. self.lock_path = std.fmt.allocPrint(a, "{s}/{s}.lock", .{ dir, base }) catch return Error.OutOfMemory;
  353. self.idxstate_path = std.fmt.allocPrint(a, "{s}/{s}.idxstate.json", .{ dir, base }) catch return Error.OutOfMemory;
  354. // ---- parse primary key ----
  355. if (opts.primary_key) |pk_raw| {
  356. if (pk_raw.len > 0) {
  357. if (pk_raw[0] == '*') {
  358. self.pk = a.dupe(u8, pk_raw[1..]) catch return Error.OutOfMemory;
  359. self.pk_type = .number;
  360. } else if (pk_raw[0] == '@') {
  361. self.pk = a.dupe(u8, pk_raw[1..]) catch return Error.OutOfMemory;
  362. self.pk_type = .uuid;
  363. } else {
  364. self.pk = a.dupe(u8, pk_raw) catch return Error.OutOfMemory;
  365. self.pk_type = .string;
  366. }
  367. }
  368. }
  369. if (self.pk) |p| self.unique_fields.append(a, p) catch return Error.OutOfMemory;
  370. // ---- parse indexes (pk first, then configured; dedupe) ----
  371. if (self.pk) |p| {
  372. self.index_fields.append(a, p) catch return Error.OutOfMemory;
  373. const t: IdxType = if (self.pk_type == .number) .numeric else .lexical;
  374. self.idx.put(a, p, .{ .typ = t }) catch return Error.OutOfMemory;
  375. }
  376. for (opts.indexes) |raw| {
  377. var clean = raw;
  378. var is_unique = false;
  379. if (clean.len > 0 and clean[0] == '!') {
  380. is_unique = true;
  381. clean = clean[1..];
  382. }
  383. var typ: IdxType = .lexical;
  384. if (clean.len > 0 and clean[0] == '*') {
  385. typ = .numeric;
  386. clean = clean[1..];
  387. } else if (clean.len > 0 and clean[0] == '@') {
  388. clean = clean[1..];
  389. }
  390. if (self.idx.contains(clean)) continue;
  391. const owned = a.dupe(u8, clean) catch return Error.OutOfMemory;
  392. self.index_fields.append(a, owned) catch return Error.OutOfMemory;
  393. self.idx.put(a, owned, .{ .typ = typ }) catch return Error.OutOfMemory;
  394. if (is_unique) self.unique_fields.append(a, owned) catch return Error.OutOfMemory;
  395. }
  396. // ---- load meta ----
  397. self.loadMeta();
  398. // ---- compact on init (JS default; also creates an empty data file) ----
  399. // AN OPEN WITH NOTHING TO CHANGE WRITES NOTHING (ticket #21): a store
  400. // with no tombstones has nothing to compact, and rewriting it anyway
  401. // made opening a backup change it.
  402. var did_compact = false;
  403. const nothing_to_compact = self.deleted.items.len == 0 and
  404. fio.fileSize(self.tmp, self.data_path) != null;
  405. if (self.compact_on_open and !nothing_to_compact) {
  406. try self.acquireLock();
  407. const cr = self.compactLocked();
  408. self.releaseLock();
  409. try cr;
  410. did_compact = true;
  411. }
  412. // ---- persist schema into meta (JS: init-schema under lock).
  413. // compactLocked already wrote meta (with schema); skip the extra bump.
  414. if (!did_compact and (self.pk != null or self.index_fields.items.len > 0) and
  415. !self.metaOnDiskIsCurrent())
  416. {
  417. try self.acquireLock();
  418. const mr = self.persistMetaLocked();
  419. self.releaseLock();
  420. try mr;
  421. }
  422. // ---- indexes (JS: (re)build under the lock; force after compaction) ----
  423. if (self.index_fields.items.len > 0) {
  424. try self.acquireLock();
  425. const ir = self.initIndexes(did_compact);
  426. self.releaseLock();
  427. try ir;
  428. }
  429. return self;
  430. }
  431. pub fn close(self: *Db) void {
  432. // persist indexes + idxstate (JS: IndexManager.close → persist under lock)
  433. if (self.index_fields.items.len > 0) {
  434. if (self.acquireLock()) {
  435. self.persistIndexesLocked() catch {};
  436. self.releaseLock();
  437. } else |_| {}
  438. }
  439. const gpa = self.arena_state.child_allocator;
  440. self.arena_state.deinit();
  441. gpa.destroy(self);
  442. }
  443. fn dbg(self: *Db, comptime fmt: []const u8, args: anytype) void {
  444. if (self.debug) std.debug.print("[mpackdb] " ++ fmt ++ "\n", args);
  445. }
  446. // =====================================================================
  447. // Locking (JS _acquireFileLock protocol)
  448. // =====================================================================
  449. fn acquireLock(self: *Db) Error!void {
  450. var retries: u32 = 0;
  451. while (true) {
  452. if (fio.openExcl(self.tmp, self.lock_path)) |fd| {
  453. var pid_buf: [16]u8 = undefined;
  454. const pid_s = std.fmt.bufPrint(&pid_buf, "{d}", .{linux.getpid()}) catch "0";
  455. _ = fio.writeAll(fd, pid_s);
  456. fio.close(fd);
  457. return;
  458. }
  459. // stale lock takeover
  460. if (self.stale_lock_timeout_ms > 0) {
  461. if (fio.mtimeMs(self.tmp, self.lock_path)) |mt| {
  462. if (fio.nowMs() - mt > self.stale_lock_timeout_ms) {
  463. _ = fio.unlink(self.tmp, self.lock_path);
  464. continue;
  465. }
  466. }
  467. }
  468. if (retries > 480) return Error.LockTimeout; // ~12s at 25ms
  469. fio.sleepMs(25);
  470. retries += 1;
  471. }
  472. }
  473. fn releaseLock(self: *Db) void {
  474. _ = fio.unlink(self.tmp, self.lock_path);
  475. }
  476. // =====================================================================
  477. // Meta (.meta.json)
  478. // =====================================================================
  479. fn appendFmt(self: *Db, list: *std.ArrayList(u8), comptime fmt: []const u8, args: anytype) Error!void {
  480. const s = std.fmt.allocPrint(self.tmp, fmt, args) catch return Error.OutOfMemory;
  481. defer self.tmp.free(s);
  482. list.appendSlice(self.tmp, s) catch return Error.OutOfMemory;
  483. }
  484. fn loadMeta(self: *Db) void {
  485. const content = fio.readAll(self.tmp, self.meta_path) orelse return;
  486. defer self.tmp.free(content);
  487. self.applyMetaJson(content);
  488. }
  489. fn applyMetaJson(self: *Db, content: []const u8) void {
  490. const parsed = std.json.parseFromSlice(std.json.Value, self.tmp, content, .{}) catch return;
  491. defer parsed.deinit();
  492. if (parsed.value != .object) return;
  493. const obj = parsed.value.object;
  494. self.next_id = FIRST_ID;
  495. self.deleted.clearRetainingCapacity();
  496. self.version = 0;
  497. if (obj.get("nextId")) |v| {
  498. self.next_id = switch (v) {
  499. .integer => |i| @floatFromInt(i),
  500. .float => |f| f,
  501. else => FIRST_ID,
  502. };
  503. }
  504. if (obj.get("version")) |v| {
  505. if (v == .integer) self.version = @intCast(@max(v.integer, 0));
  506. }
  507. if (obj.get("deleted")) |v| {
  508. if (v == .array) {
  509. for (v.array.items) |it| {
  510. const off: u64 = switch (it) {
  511. .integer => |i| @intCast(@max(i, 0)),
  512. .float => |f| @intFromFloat(@max(f, 0)),
  513. else => continue,
  514. };
  515. self.deleted.append(self.alloc, off) catch {};
  516. }
  517. }
  518. }
  519. }
  520. /// Pick up other processes' persisted state (JS refresh()): adopt disk meta
  521. /// when its version is newer; index appended tail records.
  522. pub fn refresh(self: *Db) void {
  523. if (fio.readAll(self.tmp, self.meta_path)) |content| {
  524. defer self.tmp.free(content);
  525. const parsed = std.json.parseFromSlice(std.json.Value, self.tmp, content, .{}) catch return;
  526. defer parsed.deinit();
  527. if (parsed.value == .object) {
  528. var disk_version: u64 = 0;
  529. if (parsed.value.object.get("version")) |v| {
  530. if (v == .integer) disk_version = @intCast(@max(v.integer, 0));
  531. }
  532. if (disk_version > self.version) self.applyMetaJson(content);
  533. }
  534. }
  535. // catchUp: index records another process appended
  536. if (self.index_fields.items.len > 0) {
  537. const size = fio.fileSize(self.tmp, self.data_path) orelse 0;
  538. if (size > self.covered_bytes) self.catchUp(size);
  539. }
  540. }
  541. fn persistMetaLocked(self: *Db) Error!void {
  542. self.version += 1;
  543. var out: std.ArrayList(u8) = .empty;
  544. defer out.deinit(self.tmp);
  545. try self.renderMeta(&out);
  546. // atomic tmp + rename (JS: `${metaPath}.${pid}.tmp`)
  547. var tmp_buf: [512]u8 = undefined;
  548. const tmp_path = std.fmt.bufPrint(&tmp_buf, "{s}.{d}.tmp", .{ self.meta_path, linux.getpid() }) catch return Error.IoError;
  549. const fd = fio.openTrunc(self.tmp, tmp_path) orelse return Error.IoError;
  550. const ok = fio.writeAll(fd, out.items);
  551. fio.close(fd);
  552. if (!ok) return Error.IoError;
  553. if (!fio.rename(self.tmp, tmp_path, self.meta_path)) return Error.IoError;
  554. }
  555. /// Does .meta.json already say, byte for byte, what this handle would write
  556. /// at its current version? Then an open has nothing to persist (ticket #21).
  557. fn metaOnDiskIsCurrent(self: *Db) bool {
  558. const disk = fio.readAll(self.tmp, self.meta_path) orelse return false;
  559. defer self.tmp.free(disk);
  560. var out: std.ArrayList(u8) = .empty;
  561. defer out.deinit(self.tmp);
  562. self.renderMeta(&out) catch return false;
  563. return std.mem.eql(u8, disk, out.items);
  564. }
  565. fn renderMeta(self: *Db, out: *std.ArrayList(u8)) Error!void {
  566. try self.appendFmt(out, "{{", .{});
  567. if (self.pk != null and self.pk_type == .number) {
  568. var nbuf: [32]u8 = undefined;
  569. try self.appendFmt(out, "\"nextId\":{s},", .{jsNumFmt(&nbuf, self.next_id)});
  570. }
  571. try self.appendFmt(out, "\"deleted\":[", .{});
  572. for (self.deleted.items, 0..) |off, i| {
  573. if (i > 0) try self.appendFmt(out, ",", .{});
  574. try self.appendFmt(out, "{d}", .{off});
  575. }
  576. try self.appendFmt(out, "],\"version\":{d}", .{self.version});
  577. if (self.pk != null or self.index_fields.items.len > 0) {
  578. try self.appendFmt(out, ",\"schema\":{{", .{});
  579. var first = true;
  580. if (self.pk) |p| {
  581. const prefix: []const u8 = switch (self.pk_type) {
  582. .number => "*",
  583. .uuid => "@",
  584. .string => "",
  585. };
  586. try self.appendFmt(out, "\"primaryKey\":\"{s}{s}\"", .{ prefix, p });
  587. first = false;
  588. }
  589. if (!first) try self.appendFmt(out, ",", .{});
  590. try self.appendFmt(out, "\"indexes\":[", .{});
  591. var n: usize = 0;
  592. for (self.index_fields.items) |field| {
  593. if (self.pk != null and std.mem.eql(u8, field, self.pk.?)) continue;
  594. if (n > 0) try self.appendFmt(out, ",", .{});
  595. const uni = for (self.unique_fields.items) |u| {
  596. if (std.mem.eql(u8, u, field)) break true;
  597. } else false;
  598. const numeric = (self.idx.get(field) orelse FieldIndex{ .typ = .lexical }).typ == .numeric;
  599. try self.appendFmt(out, "\"{s}{s}{s}\"", .{
  600. if (uni) "!" else "",
  601. if (numeric) "*" else "",
  602. field,
  603. });
  604. n += 1;
  605. }
  606. try self.appendFmt(out, "]}}", .{});
  607. }
  608. try self.appendFmt(out, "}}", .{});
  609. }
  610. // =====================================================================
  611. // Data file scanning
  612. // =====================================================================
  613. fn deletedSet(self: *Db, alloc: std.mem.Allocator) std.AutoHashMapUnmanaged(u64, void) {
  614. var set: std.AutoHashMapUnmanaged(u64, void) = .empty;
  615. for (self.deleted.items) |off| set.put(alloc, off, {}) catch {};
  616. return set;
  617. }
  618. /// Sequential scan yielding all live records. Caller supplies an arena for
  619. /// decoded values. Used for rebuilds, compaction and unindexed finds.
  620. pub const Scanner = struct {
  621. db: *Db,
  622. fd: i32 = -1,
  623. offset: u64 = 0,
  624. size: u64 = 0,
  625. skip_deleted: bool = true,
  626. deleted_set: std.AutoHashMapUnmanaged(u64, void) = .empty,
  627. scratch: std.mem.Allocator,
  628. pub fn init(db: *Db, scratch: std.mem.Allocator, skip_deleted: bool) Scanner {
  629. var s = Scanner{ .db = db, .scratch = scratch, .skip_deleted = skip_deleted };
  630. s.size = fio.fileSize(scratch, db.data_path) orelse 0;
  631. if (s.size > 0) {
  632. s.fd = fio.openRead(scratch, db.data_path) orelse -1;
  633. }
  634. if (skip_deleted) s.deleted_set = db.deletedSet(scratch);
  635. return s;
  636. }
  637. pub fn deinit(self: *Scanner) void {
  638. if (self.fd >= 0) fio.close(self.fd);
  639. self.fd = -1;
  640. }
  641. /// Returns the next record (decoded into `arena`) or null at EOF.
  642. pub fn next(self: *Scanner, arena: std.mem.Allocator) Error!?Record {
  643. while (true) {
  644. if (self.fd < 0 or self.offset + 4 > self.size) return null;
  645. var hdr: [4]u8 = undefined;
  646. const got = fio.pread(self.fd, &hdr, self.offset) orelse return Error.IoError;
  647. if (got < 4) return null;
  648. const rec_size = std.mem.readInt(u32, &hdr, .little);
  649. if (rec_size <= 4) return Error.CorruptRecord;
  650. if (self.offset + rec_size > self.size) return null; // incomplete tail
  651. const off = self.offset;
  652. self.offset += rec_size;
  653. if (self.skip_deleted and self.deleted_set.contains(off)) continue;
  654. const buf = arena.alloc(u8, rec_size - 4) catch return Error.OutOfMemory;
  655. const got2 = fio.pread(self.fd, buf, off + 4) orelse return Error.IoError;
  656. if (got2 < rec_size - 4) return Error.Truncated;
  657. const d = msgpack.decode(arena, buf) catch return Error.CorruptRecord;
  658. return Record{ .value = d.value, .off = off, .len = rec_size };
  659. }
  660. }
  661. /// Like next() but without decoding — yields the raw framed bytes.
  662. pub fn nextRaw(self: *Scanner, arena: std.mem.Allocator) Error!?struct { bytes: []u8, off: u64 } {
  663. while (true) {
  664. if (self.fd < 0 or self.offset + 4 > self.size) return null;
  665. var hdr: [4]u8 = undefined;
  666. const got = fio.pread(self.fd, &hdr, self.offset) orelse return Error.IoError;
  667. if (got < 4) return null;
  668. const rec_size = std.mem.readInt(u32, &hdr, .little);
  669. if (rec_size <= 4) return Error.CorruptRecord;
  670. if (self.offset + rec_size > self.size) return null;
  671. const off = self.offset;
  672. self.offset += rec_size;
  673. if (self.skip_deleted and self.deleted_set.contains(off)) continue;
  674. const buf = arena.alloc(u8, rec_size) catch return Error.OutOfMemory;
  675. const got2 = fio.pread(self.fd, buf, off) orelse return Error.IoError;
  676. if (got2 < rec_size) return Error.Truncated;
  677. return .{ .bytes = buf, .off = off };
  678. }
  679. }
  680. };
  681. /// Read + decode one record by location.
  682. pub fn readAt(self: *Db, arena: std.mem.Allocator, loc: Loc) Error!msgpack.Value {
  683. const fd = fio.openRead(self.tmp, self.data_path) orelse return Error.IoError;
  684. defer fio.close(fd);
  685. if (loc.len <= 4) return Error.CorruptRecord;
  686. const buf = arena.alloc(u8, loc.len - 4) catch return Error.OutOfMemory;
  687. const got = fio.pread(fd, buf, loc.off + 4) orelse return Error.IoError;
  688. if (got < buf.len) return Error.Truncated;
  689. const d = msgpack.decode(arena, buf) catch return Error.CorruptRecord;
  690. return d.value;
  691. }
  692. // =====================================================================
  693. // Index management
  694. // =====================================================================
  695. fn keyFromValue(v: msgpack.Value) ?Key {
  696. return switch (v) {
  697. .number => |n| Key{ .num = n },
  698. .str => |s| Key{ .str = s },
  699. else => null, // null/bool/objects are not indexed (see header note)
  700. };
  701. }
  702. /// Duplicate a key's string into the db arena so it outlives the op arena.
  703. fn ownKey(self: *Db, k: Key) Error!Key {
  704. return switch (k) {
  705. .num => k,
  706. .str => |s| Key{ .str = self.alloc.dupe(u8, s) catch return Error.OutOfMemory },
  707. };
  708. }
  709. fn lowerBound(entries: []const Entry, key: Key) usize {
  710. var lo: usize = 0;
  711. var hi: usize = entries.len;
  712. while (lo < hi) {
  713. const mid = lo + (hi - lo) / 2;
  714. if (cmpKeys(entries[mid].key, key) < 0) {
  715. lo = mid + 1;
  716. } else {
  717. hi = mid;
  718. }
  719. }
  720. return lo;
  721. }
  722. /// Binary search: all entries whose key equals `key` (appended to `out`).
  723. pub fn indexGet(self: *Db, field: []const u8, key: Key, alloc: std.mem.Allocator, out: *std.ArrayList(Entry)) Error!void {
  724. const fi = self.idx.getPtr(field) orelse return Error.NoSuchIndex;
  725. const parsed = self.parseKeyForField(field, key);
  726. const entries = fi.entries.items;
  727. var i = lowerBound(entries, parsed);
  728. while (i < entries.len and cmpKeys(entries[i].key, parsed) == 0) : (i += 1) {
  729. out.append(alloc, entries[i]) catch return Error.OutOfMemory;
  730. }
  731. }
  732. /// Range scan [from, to] (either side optional), ascending.
  733. pub fn indexRange(self: *Db, field: []const u8, from: ?Key, to: ?Key, alloc: std.mem.Allocator, out: *std.ArrayList(Entry)) Error!void {
  734. const fi = self.idx.getPtr(field) orelse return Error.NoSuchIndex;
  735. const entries = fi.entries.items;
  736. var i: usize = if (from) |f| lowerBound(entries, self.parseKeyForField(field, f)) else 0;
  737. const to_key: ?Key = if (to) |t| self.parseKeyForField(field, t) else null;
  738. while (i < entries.len) : (i += 1) {
  739. if (to_key) |t| {
  740. if (cmpKeys(entries[i].key, t) > 0) break;
  741. }
  742. out.append(alloc, entries[i]) catch return Error.OutOfMemory;
  743. }
  744. }
  745. /// _parseKey semantics: numeric index fields parse string keys via
  746. /// parseInt; on NaN the key stays a string.
  747. fn parseKeyForField(self: *Db, field: []const u8, key: Key) Key {
  748. const fi = self.idx.get(field) orelse return key;
  749. if (fi.typ == .numeric and key == .str) {
  750. if (jsParseInt(key.str)) |n| return Key{ .num = n };
  751. }
  752. return key;
  753. }
  754. fn indexInsertEntry(self: *Db, field: []const u8, key: Key, loc: Loc) Error!void {
  755. const fi = self.idx.getPtr(field) orelse return;
  756. const owned = try self.ownKey(key);
  757. const e = Entry{ .key = owned, .off = loc.off, .len = loc.len };
  758. // insert at sorted position
  759. const pos = blk: {
  760. var lo: usize = 0;
  761. var hi: usize = fi.entries.items.len;
  762. while (lo < hi) {
  763. const mid = lo + (hi - lo) / 2;
  764. if (entryLess({}, fi.entries.items[mid], e)) {
  765. lo = mid + 1;
  766. } else {
  767. hi = mid;
  768. }
  769. }
  770. break :blk lo;
  771. };
  772. fi.entries.insert(self.alloc, pos, e) catch return Error.OutOfMemory;
  773. }
  774. /// Index a record's fields at loc.
  775. fn indexInsertRecord(self: *Db, rec: msgpack.Value, loc: Loc) Error!void {
  776. for (self.index_fields.items) |field| {
  777. const v = rec.get(field) orelse continue;
  778. const key = keyFromValue(v) orelse continue;
  779. try self.indexInsertEntry(field, key, loc);
  780. }
  781. self.covered_bytes = @max(self.covered_bytes, loc.off + loc.len);
  782. }
  783. /// Remove all entries pointing at offset `off` (all fields).
  784. fn indexRemoveOffset(self: *Db, off: u64) void {
  785. var it = self.idx.iterator();
  786. while (it.next()) |kv| {
  787. const list = &kv.value_ptr.entries;
  788. var i: usize = 0;
  789. while (i < list.items.len) {
  790. if (list.items[i].off == off) {
  791. _ = list.orderedRemove(i);
  792. } else {
  793. i += 1;
  794. }
  795. }
  796. }
  797. }
  798. fn indexPath(self: *Db, buf: []u8, field: []const u8) []const u8 {
  799. // <dir>/<base>.<field>.txt — derive from idxstate path (…/base.idxstate.json)
  800. const prefix = self.idxstate_path[0 .. self.idxstate_path.len - "idxstate.json".len];
  801. return std.fmt.bufPrint(buf, "{s}{s}.txt", .{ prefix, field }) catch buf[0..0];
  802. }
  803. fn initIndexes(self: *Db, force_rebuild: bool) Error!void {
  804. var need_rebuild = force_rebuild;
  805. const data_size = fio.fileSize(self.tmp, self.data_path) orelse 0;
  806. if (!need_rebuild) {
  807. for (self.index_fields.items) |field| {
  808. var pbuf: [512]u8 = undefined;
  809. const p = self.indexPath(&pbuf, field);
  810. const isize_ = fio.fileSize(self.tmp, p);
  811. if (isize_ == null or isize_.? == 0) {
  812. if (data_size > 0) {
  813. need_rebuild = true;
  814. break;
  815. }
  816. }
  817. }
  818. }
  819. if (!need_rebuild) {
  820. // coveredBytes from idxstate — no state file means rebuild (JS)
  821. if (self.readIdxState()) |cb| {
  822. self.covered_bytes = cb;
  823. } else {
  824. need_rebuild = true;
  825. }
  826. }
  827. if (need_rebuild) {
  828. try self.rebuildIndexes();
  829. return;
  830. }
  831. // Load index files into memory
  832. for (self.index_fields.items) |field| {
  833. var pbuf: [512]u8 = undefined;
  834. const p = self.indexPath(&pbuf, field);
  835. const content = fio.readAll(self.tmp, p) orelse continue;
  836. defer self.tmp.free(content);
  837. self.loadIndexLines(field, content);
  838. }
  839. // sort (files are sorted by JS localeCompare; re-sort under our order)
  840. var it = self.idx.iterator();
  841. while (it.next()) |kv| {
  842. std.sort.pdq(Entry, kv.value_ptr.entries.items, {}, entryLess);
  843. }
  844. if (data_size > self.covered_bytes) self.catchUp(data_size);
  845. }
  846. fn loadIndexLines(self: *Db, field: []const u8, content: []const u8) void {
  847. const fi = self.idx.getPtr(field) orelse return;
  848. var lines = std.mem.splitScalar(u8, content, '\n');
  849. while (lines.next()) |line| {
  850. if (line.len == 0) continue;
  851. // key = up to FIRST comma (JS split(',')[0]); then offset, length
  852. const c1 = std.mem.indexOfScalar(u8, line, ',') orelse continue;
  853. const rest = line[c1 + 1 ..];
  854. const c2 = std.mem.indexOfScalar(u8, rest, ',') orelse continue;
  855. const off_s = rest[0..c2];
  856. const len_s = rest[c2 + 1 ..];
  857. const off = std.fmt.parseInt(u64, off_s, 10) catch continue;
  858. const len = std.fmt.parseInt(u64, len_s, 10) catch continue;
  859. if (len == 0) continue;
  860. const key_s = line[0..c1];
  861. var key: Key = undefined;
  862. if (fi.typ == .numeric) {
  863. key = if (jsParseInt(key_s)) |n| Key{ .num = n } else Key{ .str = self.alloc.dupe(u8, key_s) catch return };
  864. } else {
  865. key = Key{ .str = self.alloc.dupe(u8, key_s) catch return };
  866. }
  867. fi.entries.append(self.alloc, .{ .key = key, .off = off, .len = len }) catch return;
  868. }
  869. }
  870. fn rebuildIndexes(self: *Db) Error!void {
  871. var it0 = self.idx.iterator();
  872. while (it0.next()) |kv| kv.value_ptr.entries.clearRetainingCapacity();
  873. var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);
  874. defer scratch_state.deinit();
  875. const scratch = scratch_state.allocator();
  876. // JS rebuild scans ALL records (no tombstone filter — deleted offsets
  877. // simply get filtered at read time). Mirror that.
  878. var scanner = Scanner.init(self, scratch, false);
  879. defer scanner.deinit();
  880. while (try scanner.next(scratch)) |rec| {
  881. try self.indexInsertRecord(rec.value, .{ .off = rec.off, .len = rec.len });
  882. }
  883. self.covered_bytes = fio.fileSize(self.tmp, self.data_path) orelse 0;
  884. try self.writeIndexFiles();
  885. try self.writeIdxState();
  886. }
  887. fn catchUp(self: *Db, file_size: u64) void {
  888. var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);
  889. defer scratch_state.deinit();
  890. const scratch = scratch_state.allocator();
  891. var scanner = Scanner.init(self, scratch, false);
  892. defer scanner.deinit();
  893. scanner.offset = self.covered_bytes;
  894. while (true) {
  895. const maybe = scanner.next(scratch) catch break;
  896. const rec = maybe orelse break;
  897. self.indexInsertRecord(rec.value, .{ .off = rec.off, .len = rec.len }) catch break;
  898. }
  899. self.covered_bytes = @max(self.covered_bytes, file_size);
  900. }
  901. /// A file that already holds these exact bytes is left alone — its mtime
  902. /// included — so closing a handle that changed nothing writes nothing.
  903. fn fileHolds(alloc: std.mem.Allocator, path: []const u8, bytes: []const u8) bool {
  904. const disk = fio.readAll(alloc, path) orelse return false;
  905. defer alloc.free(disk);
  906. return std.mem.eql(u8, disk, bytes);
  907. }
  908. fn writeIndexFiles(self: *Db) Error!void {
  909. for (self.index_fields.items) |field| {
  910. const fi = self.idx.getPtr(field) orelse continue;
  911. var out: std.ArrayList(u8) = .empty;
  912. defer out.deinit(self.tmp);
  913. for (fi.entries.items) |e| {
  914. var kbuf: [32]u8 = undefined;
  915. const ks = switch (e.key) {
  916. .num => |n| jsNumFmt(&kbuf, n),
  917. .str => |s| s,
  918. };
  919. try self.appendFmt(&out, "{s},{d},{d}\n", .{ ks, e.off, e.len });
  920. }
  921. var pbuf: [512]u8 = undefined;
  922. const p = self.indexPath(&pbuf, field);
  923. if (fileHolds(self.tmp, p, out.items)) continue;
  924. var tbuf: [512]u8 = undefined;
  925. const tmp = std.fmt.bufPrint(&tbuf, "{s}.{d}.tmp", .{ p, linux.getpid() }) catch return Error.IoError;
  926. const fd = fio.openTrunc(self.tmp, tmp) orelse return Error.IoError;
  927. const ok = fio.writeAll(fd, out.items);
  928. fio.close(fd);
  929. if (!ok) return Error.IoError;
  930. if (!fio.rename(self.tmp, tmp, p)) return Error.IoError;
  931. }
  932. }
  933. fn readIdxState(self: *Db) ?u64 {
  934. const content = fio.readAll(self.tmp, self.idxstate_path) orelse return null;
  935. defer self.tmp.free(content);
  936. const parsed = std.json.parseFromSlice(std.json.Value, self.tmp, content, .{}) catch return null;
  937. defer parsed.deinit();
  938. if (parsed.value != .object) return null;
  939. const v = parsed.value.object.get("coveredBytes") orelse return null;
  940. return switch (v) {
  941. .integer => |i| @intCast(@max(i, 0)),
  942. .float => |f| @intFromFloat(@max(f, 0)),
  943. else => null,
  944. };
  945. }
  946. fn writeIdxState(self: *Db) Error!void {
  947. var buf: [128]u8 = undefined;
  948. const json = std.fmt.bufPrint(&buf, "{{\"coveredBytes\":{d}}}", .{self.covered_bytes}) catch return Error.IoError;
  949. if (fileHolds(self.tmp, self.idxstate_path, json)) return;
  950. var tbuf: [512]u8 = undefined;
  951. const tmp = std.fmt.bufPrint(&tbuf, "{s}.{d}.tmp", .{ self.idxstate_path, linux.getpid() }) catch return Error.IoError;
  952. const fd = fio.openTrunc(self.tmp, tmp) orelse return Error.IoError;
  953. const ok = fio.writeAll(fd, json);
  954. fio.close(fd);
  955. if (!ok) return Error.IoError;
  956. if (!fio.rename(self.tmp, tmp, self.idxstate_path)) return Error.IoError;
  957. }
  958. fn persistIndexesLocked(self: *Db) Error!void {
  959. try self.writeIndexFiles();
  960. // coverage can only grow (JS: max of ours and on-disk state)
  961. if (self.readIdxState()) |cb| self.covered_bytes = @max(self.covered_bytes, cb);
  962. try self.writeIdxState();
  963. }
  964. // =====================================================================
  965. // Operations
  966. // =====================================================================
  967. /// A mutable record under construction (op-arena entries).
  968. pub const MutableRecord = struct {
  969. entries: std.ArrayList(msgpack.Entry) = .empty,
  970. pub fn get(self: *const MutableRecord, key: []const u8) ?msgpack.Value {
  971. for (self.entries.items) |e| {
  972. if (e.key == .str and std.mem.eql(u8, e.key.str, key)) return e.value;
  973. }
  974. return null;
  975. }
  976. pub fn set(self: *MutableRecord, alloc: std.mem.Allocator, key: []const u8, v: msgpack.Value) Error!void {
  977. for (self.entries.items) |*e| {
  978. if (e.key == .str and std.mem.eql(u8, e.key.str, key)) {
  979. e.value = v;
  980. return;
  981. }
  982. }
  983. self.entries.append(alloc, .{ .key = .{ .str = key }, .value = v }) catch return Error.OutOfMemory;
  984. }
  985. pub fn toValue(self: *const MutableRecord) msgpack.Value {
  986. return .{ .map = self.entries.items };
  987. }
  988. };
  989. pub const InsertResult = union(enum) {
  990. pk_num: f64,
  991. pk_str: []const u8, // op-arena
  992. record: msgpack.Value,
  993. };
  994. /// insert() — auto primary key, unique checks, append, index.
  995. /// `record` must be a map value; `arena` is the op arena (record memory).
  996. pub fn insert(self: *Db, arena: std.mem.Allocator, record: msgpack.Value, skip_pk: bool) Error!InsertResult {
  997. if (record != .map) return Error.CorruptRecord;
  998. try self.acquireLock();
  999. defer self.releaseLock();
  1000. self.refresh();
  1001. return self.insertLocked(arena, record, skip_pk);
  1002. }
  1003. fn insertLocked(self: *Db, arena: std.mem.Allocator, record: msgpack.Value, skip_pk: bool) Error!InsertResult {
  1004. // copy into mutable form
  1005. var rec = MutableRecord{};
  1006. for (record.map) |e| rec.entries.append(arena, e) catch return Error.OutOfMemory;
  1007. // _hasPrimaryKeyValue: present and not null/undefined
  1008. var has_pk_value = false;
  1009. if (self.pk) |p| {
  1010. if (rec.get(p)) |v| {
  1011. has_pk_value = (v != .nil and v != .undef);
  1012. }
  1013. }
  1014. var auto_gen = false;
  1015. if (self.pk != null and !skip_pk and !has_pk_value) {
  1016. if (self.pk_type == .number) {
  1017. rec.set(arena, self.pk.?, .{ .number = self.next_id }) catch return Error.OutOfMemory;
  1018. self.next_id += 1;
  1019. auto_gen = true;
  1020. } else if (self.pk_type == .uuid) {
  1021. // A GENERATED ID IS CHECKED, NOT TRUSTED (ticket #113): the format is a
  1022. // millisecond stamp + 3 base36 chars, so a bulk put collides within a
  1023. // millisecond (two of 7800 in a measured run) and the unique check below
  1024. // skips generated keys — two live rows then shared one pk and the next
  1025. // update() deleted both and re-inserted one. Draw again until it is free.
  1026. var ubuf: [12]u8 = undefined;
  1027. var tries: usize = 0;
  1028. while (true) : (tries += 1) {
  1029. const u = genUuid(&ubuf);
  1030. var taken: std.ArrayList(Loc) = .empty;
  1031. defer taken.deinit(arena);
  1032. try self.findPkLocs(arena, .{ .str = u }, &taken);
  1033. if (taken.items.len == 0 or tries >= 64) break;
  1034. }
  1035. const owned = arena.dupe(u8, ubuf[0..12]) catch return Error.OutOfMemory;
  1036. rec.set(arena, self.pk.?, .{ .str = owned }) catch return Error.OutOfMemory;
  1037. auto_gen = true;
  1038. }
  1039. }
  1040. // unique constraints (skip auto-generated pk)
  1041. if (self.index_fields.items.len > 0 and self.unique_fields.items.len > 0) {
  1042. var del_set = self.deletedSet(arena);
  1043. for (self.unique_fields.items) |field| {
  1044. if (auto_gen and self.pk != null and std.mem.eql(u8, field, self.pk.?)) continue;
  1045. const v = rec.get(field) orelse continue;
  1046. if (v == .undef) continue;
  1047. const key = keyFromValue(v) orelse continue;
  1048. var hits: std.ArrayList(Entry) = .empty;
  1049. defer hits.deinit(arena);
  1050. try self.indexGet(field, key, arena, &hits);
  1051. var live: usize = 0;
  1052. for (hits.items) |h| {
  1053. if (!del_set.contains(h.off)) live += 1;
  1054. }
  1055. if (live > 0) {
  1056. var kbuf: [32]u8 = undefined;
  1057. const ks = switch (key) {
  1058. .num => |n| jsNumFmt(&kbuf, n),
  1059. .str => |s| s,
  1060. };
  1061. self.last_error = std.fmt.bufPrint(&self.last_error_buf, "Duplicate key: {s}={s}", .{ field, ks }) catch "Duplicate key";
  1062. return Error.DuplicateKey;
  1063. }
  1064. }
  1065. }
  1066. // serialize + append
  1067. const framed = msgpack.serialize(arena, rec.toValue()) catch return Error.OutOfMemory;
  1068. const offset = fio.fileSize(self.tmp, self.data_path) orelse 0;
  1069. const fd = fio.openAppend(self.tmp, self.data_path) orelse return Error.IoError;
  1070. const ok = fio.writeAll(fd, framed);
  1071. fio.close(fd);
  1072. if (!ok) return Error.IoError;
  1073. const loc = Loc{ .off = offset, .len = framed.len };
  1074. if (self.index_fields.items.len > 0) {
  1075. try self.indexInsertRecord(rec.toValue(), loc);
  1076. }
  1077. if (self.pk != null and self.pk_type == .number and !skip_pk and !has_pk_value) {
  1078. try self.persistMetaLocked();
  1079. }
  1080. if (self.pk) |p| {
  1081. const v = rec.get(p) orelse return Error.CorruptRecord;
  1082. return switch (v) {
  1083. .number => |n| InsertResult{ .pk_num = n },
  1084. .str => |s| InsertResult{ .pk_str = s },
  1085. else => InsertResult{ .record = rec.toValue() },
  1086. };
  1087. }
  1088. return InsertResult{ .record = rec.toValue() };
  1089. }
  1090. /// Locate live records by primary key (index-backed when available).
  1091. pub fn findPkLocs(self: *Db, arena: std.mem.Allocator, key: Key, out: *std.ArrayList(Loc)) Error!void {
  1092. if (self.pk == null) return Error.NoPrimaryKey;
  1093. self.refresh();
  1094. var del_set = self.deletedSet(arena);
  1095. if (self.index_fields.items.len > 0) {
  1096. var hits: std.ArrayList(Entry) = .empty;
  1097. defer hits.deinit(arena);
  1098. try self.indexGet(self.pk.?, key, arena, &hits);
  1099. for (hits.items) |h| {
  1100. if (del_set.contains(h.off)) continue;
  1101. out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;
  1102. }
  1103. return;
  1104. }
  1105. // no indexes: full scan comparing the pk field
  1106. var scanner = Scanner.init(self, arena, true);
  1107. defer scanner.deinit();
  1108. while (try scanner.next(arena)) |rec| {
  1109. const v = rec.value.get(self.pk.?) orelse continue;
  1110. const k = keyFromValue(v) orelse continue;
  1111. if (cmpKeys(self.parseKeyForField(self.pk.?, k), self.parseKeyForField(self.pk.?, key)) == 0) {
  1112. out.append(arena, .{ .off = rec.off, .len = rec.len }) catch return Error.OutOfMemory;
  1113. }
  1114. }
  1115. }
  1116. /// Locate live records by secondary index equality.
  1117. pub fn findIndexLocs(self: *Db, arena: std.mem.Allocator, field: []const u8, key: Key, out: *std.ArrayList(Loc)) Error!void {
  1118. self.refresh();
  1119. var del_set = self.deletedSet(arena);
  1120. var hits: std.ArrayList(Entry) = .empty;
  1121. defer hits.deinit(arena);
  1122. try self.indexGet(field, key, arena, &hits);
  1123. for (hits.items) |h| {
  1124. if (del_set.contains(h.off)) continue;
  1125. out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;
  1126. }
  1127. }
  1128. /// Locate live records by index range [from, to] ascending.
  1129. pub fn findRangeLocs(self: *Db, arena: std.mem.Allocator, field: []const u8, from: ?Key, to: ?Key, out: *std.ArrayList(Loc)) Error!void {
  1130. self.refresh();
  1131. var del_set = self.deletedSet(arena);
  1132. var hits: std.ArrayList(Entry) = .empty;
  1133. defer hits.deinit(arena);
  1134. try self.indexRange(field, from, to, arena, &hits);
  1135. var seen: std.AutoHashMapUnmanaged(u64, void) = .empty;
  1136. for (hits.items) |h| {
  1137. if (del_set.contains(h.off)) continue;
  1138. if (seen.contains(h.off)) continue;
  1139. seen.put(arena, h.off, {}) catch {};
  1140. out.append(arena, .{ .off = h.off, .len = h.len }) catch return Error.OutOfMemory;
  1141. }
  1142. }
  1143. /// delete by primary key. Returns the deleted records (decoded into arena).
  1144. pub fn deletePk(self: *Db, arena: std.mem.Allocator, key: Key, out_records: *std.ArrayList(msgpack.Value)) Error!usize {
  1145. try self.acquireLock();
  1146. defer self.releaseLock();
  1147. self.refresh();
  1148. return self.deletePkLocked(arena, key, out_records);
  1149. }
  1150. fn deletePkLocked(self: *Db, arena: std.mem.Allocator, key: Key, out_records: *std.ArrayList(msgpack.Value)) Error!usize {
  1151. var locs: std.ArrayList(Loc) = .empty;
  1152. defer locs.deinit(arena);
  1153. try self.findPkLocs(arena, key, &locs);
  1154. for (locs.items) |loc| {
  1155. const rec = try self.readAt(arena, loc);
  1156. out_records.append(arena, rec) catch return Error.OutOfMemory;
  1157. self.deleted.append(self.alloc, loc.off) catch return Error.OutOfMemory;
  1158. self.indexRemoveOffset(loc.off);
  1159. }
  1160. try self.persistMetaLocked();
  1161. return locs.items.len;
  1162. }
  1163. /// update by primary key: delete + insert(skip_pk) — JS update() semantics.
  1164. /// `new_record` must already carry the primary key (the JS callback
  1165. /// contract: the record keeps its pk unless the caller removes it).
  1166. pub fn updatePk(self: *Db, arena: std.mem.Allocator, key: Key, new_record: msgpack.Value) Error!usize {
  1167. if (self.pk) |p| {
  1168. const carried = for (new_record.map) |e| {
  1169. if (e.key == .str and std.mem.eql(u8, e.key.str, p)) break e.value != .nil and e.value != .undef;
  1170. } else false;
  1171. if (!carried) return Error.RecordLacksPrimaryKey;
  1172. }
  1173. try self.acquireLock();
  1174. defer self.releaseLock();
  1175. self.refresh();
  1176. var old: std.ArrayList(msgpack.Value) = .empty;
  1177. defer old.deinit(arena);
  1178. const n = try self.deletePkLocked(arena, key, &old);
  1179. if (n == 0) return 0;
  1180. var i: usize = 0;
  1181. while (i < n) : (i += 1) {
  1182. _ = try self.insertLocked(arena, new_record, true);
  1183. }
  1184. return n;
  1185. }
  1186. /// compact() — rewrite the data file without tombstones, rebuild indexes.
  1187. pub fn compact(self: *Db) Error!void {
  1188. try self.acquireLock();
  1189. defer self.releaseLock();
  1190. self.refresh();
  1191. try self.compactLocked();
  1192. if (self.index_fields.items.len > 0) {
  1193. try self.rebuildIndexes();
  1194. }
  1195. }
  1196. fn compactLocked(self: *Db) Error!void {
  1197. var scratch_state = std.heap.ArenaAllocator.init(self.arena_state.child_allocator);
  1198. defer scratch_state.deinit();
  1199. const scratch = scratch_state.allocator();
  1200. var tbuf: [512]u8 = undefined;
  1201. const tmp = std.fmt.bufPrint(&tbuf, "{s}.tmp", .{self.data_path}) catch return Error.IoError;
  1202. const out_fd = fio.openTrunc(self.tmp, tmp) orelse return Error.IoError;
  1203. var scanner = Scanner.init(self, scratch, true);
  1204. var ok = true;
  1205. while (true) {
  1206. const maybe = scanner.nextRaw(scratch) catch {
  1207. ok = false;
  1208. break;
  1209. };
  1210. const raw = maybe orelse break;
  1211. if (!fio.writeAll(out_fd, raw.bytes)) {
  1212. ok = false;
  1213. break;
  1214. }
  1215. }
  1216. scanner.deinit();
  1217. fio.close(out_fd);
  1218. if (!ok) return Error.IoError;
  1219. if (!fio.rename(self.tmp, tmp, self.data_path)) return Error.IoError;
  1220. self.deleted.clearRetainingCapacity();
  1221. try self.persistMetaLocked();
  1222. }
  1223. pub fn persistNow(self: *Db) Error!void {
  1224. try self.acquireLock();
  1225. defer self.releaseLock();
  1226. self.refresh();
  1227. try self.persistIndexesLocked();
  1228. }
  1229. };
  1230. // =========================================================================
  1231. // Tests
  1232. // =========================================================================
  1233. const testing = std.testing;
  1234. fn tmpBase(buf: []u8, comptime name: []const u8) []const u8 {
  1235. return std.fmt.bufPrint(buf, "/tmp/mpackdb-zigtest-{d}-" ++ name, .{linux.getpid()}) catch unreachable;
  1236. }
  1237. fn cleanup(alloc: std.mem.Allocator, base: []const u8) void {
  1238. var buf: [512]u8 = undefined;
  1239. const suffixes = [_][]const u8{ ".mpack", ".meta.json", ".idxstate.json", ".lock", ".id.txt", ".email.txt", ".age.txt", ".uuid.txt" };
  1240. for (suffixes) |suffix| {
  1241. const p = std.fmt.bufPrint(&buf, "{s}{s}", .{ base, suffix }) catch continue;
  1242. _ = fio.unlink(alloc, p);
  1243. }
  1244. }
  1245. fn strKey(s: []const u8) Key {
  1246. return .{ .str = s };
  1247. }
  1248. fn numKey(n: f64) Key {
  1249. return .{ .num = n };
  1250. }
  1251. fn makeUser(arena: std.mem.Allocator, name: []const u8, email: []const u8, age: f64) !msgpack.Value {
  1252. const entries = try arena.alloc(msgpack.Entry, 3);
  1253. entries[0] = .{ .key = .{ .str = "name" }, .value = .{ .str = name } };
  1254. entries[1] = .{ .key = .{ .str = "email" }, .value = .{ .str = email } };
  1255. entries[2] = .{ .key = .{ .str = "age" }, .value = .{ .number = age } };
  1256. return .{ .map = entries };
  1257. }
  1258. test "engine: insert/find/delete/update with numeric pk + indexes" {
  1259. var base_buf: [128]u8 = undefined;
  1260. const base = tmpBase(&base_buf, "crud");
  1261. cleanup(testing.allocator, base);
  1262. defer cleanup(testing.allocator, base);
  1263. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1264. defer arena_state.deinit();
  1265. const arena = arena_state.allocator();
  1266. const db = try Db.open(testing.allocator, base, .{
  1267. .primary_key = "*id",
  1268. .indexes = &.{ "email", "*age" },
  1269. });
  1270. // insert three — a fresh table counts from 1 (ticket #4)
  1271. const r1 = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 30), false);
  1272. try testing.expectEqual(@as(f64, 1), r1.pk_num);
  1273. const r2 = try db.insert(arena, try makeUser(arena, "Bob", "[email protected]", 25), false);
  1274. try testing.expectEqual(@as(f64, 2), r2.pk_num);
  1275. _ = try db.insert(arena, try makeUser(arena, "Carol", "[email protected]", 35), false);
  1276. // find by pk
  1277. var locs: std.ArrayList(Loc) = .empty;
  1278. try db.findPkLocs(arena, numKey(2), &locs);
  1279. try testing.expectEqual(@as(usize, 1), locs.items.len);
  1280. const bob = try db.readAt(arena, locs.items[0]);
  1281. try testing.expectEqualStrings("Bob", bob.get("name").?.str);
  1282. // find by secondary index
  1283. var locs2: std.ArrayList(Loc) = .empty;
  1284. try db.findIndexLocs(arena, "email", strKey("[email protected]"), &locs2);
  1285. try testing.expectEqual(@as(usize, 1), locs2.items.len);
  1286. // range on numeric index: age 26..40 → Alice(30), Carol(35)
  1287. var locs3: std.ArrayList(Loc) = .empty;
  1288. try db.findRangeLocs(arena, "age", numKey(26), numKey(40), &locs3);
  1289. try testing.expectEqual(@as(usize, 2), locs3.items.len);
  1290. // update Bob's age
  1291. var bob_new = Db.MutableRecord{};
  1292. for (bob.map) |e| try bob_new.entries.append(arena, e);
  1293. try bob_new.set(arena, "age", .{ .number = 26 });
  1294. const updated = try db.updatePk(arena, numKey(2), bob_new.toValue());
  1295. try testing.expectEqual(@as(usize, 1), updated);
  1296. var locs4: std.ArrayList(Loc) = .empty;
  1297. try db.findPkLocs(arena, numKey(2), &locs4);
  1298. try testing.expectEqual(@as(usize, 1), locs4.items.len);
  1299. const bob2 = try db.readAt(arena, locs4.items[0]);
  1300. try testing.expectEqual(@as(f64, 26), bob2.get("age").?.number);
  1301. // delete Alice
  1302. var deleted_recs: std.ArrayList(msgpack.Value) = .empty;
  1303. const dn = try db.deletePk(arena, numKey(1), &deleted_recs);
  1304. try testing.expectEqual(@as(usize, 1), dn);
  1305. try testing.expectEqualStrings("Alice", deleted_recs.items[0].get("name").?.str);
  1306. var locs5: std.ArrayList(Loc) = .empty;
  1307. try db.findPkLocs(arena, numKey(1), &locs5);
  1308. try testing.expectEqual(@as(usize, 0), locs5.items.len);
  1309. db.close();
  1310. // reopen (compaction drops the tombstone) and verify persistence
  1311. const db2 = try Db.open(testing.allocator, base, .{
  1312. .primary_key = "*id",
  1313. .indexes = &.{ "email", "*age" },
  1314. });
  1315. defer db2.close();
  1316. var locs6: std.ArrayList(Loc) = .empty;
  1317. try db2.findPkLocs(arena, numKey(2), &locs6);
  1318. try testing.expectEqual(@as(usize, 1), locs6.items.len);
  1319. const bob3 = try db2.readAt(arena, locs6.items[0]);
  1320. try testing.expectEqual(@as(f64, 26), bob3.get("age").?.number);
  1321. // nextId continues after reopen
  1322. const r4 = try db2.insert(arena, try makeUser(arena, "Dan", "[email protected]", 40), false);
  1323. try testing.expectEqual(@as(f64, 4), r4.pk_num);
  1324. }
  1325. test "engine: a stored counter wins over FIRST_ID (ticket #4)" {
  1326. var base_buf: [128]u8 = undefined;
  1327. const base = tmpBase(&base_buf, "counter");
  1328. cleanup(testing.allocator, base);
  1329. defer cleanup(testing.allocator, base);
  1330. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1331. defer arena_state.deinit();
  1332. const arena = arena_state.allocator();
  1333. // A table the JS mpackdb created and never wrote to: its meta says 0, and
  1334. // an existing table keeps its own counter.
  1335. var mbuf: [160]u8 = undefined;
  1336. const meta_path = try std.fmt.bufPrint(&mbuf, "{s}.meta.json", .{base});
  1337. const fd = fio.openTrunc(testing.allocator, meta_path).?;
  1338. try testing.expect(fio.writeAll(fd, "{\"nextId\":0,\"deleted\":[],\"version\":1}"));
  1339. fio.close(fd);
  1340. const db = try Db.open(testing.allocator, base, .{ .primary_key = "*id" });
  1341. defer db.close();
  1342. const r = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 1), false);
  1343. try testing.expectEqual(@as(f64, 0), r.pk_num);
  1344. }
  1345. test "engine: unique index rejects duplicates" {
  1346. var base_buf: [128]u8 = undefined;
  1347. const base = tmpBase(&base_buf, "uniq");
  1348. cleanup(testing.allocator, base);
  1349. defer cleanup(testing.allocator, base);
  1350. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1351. defer arena_state.deinit();
  1352. const arena = arena_state.allocator();
  1353. const db = try Db.open(testing.allocator, base, .{
  1354. .primary_key = "*id",
  1355. .indexes = &.{"!email"},
  1356. });
  1357. defer db.close();
  1358. _ = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 30), false);
  1359. const dup = db.insert(arena, try makeUser(arena, "Evil", "[email protected]", 31), false);
  1360. try testing.expectError(Error.DuplicateKey, dup);
  1361. try testing.expect(std.mem.indexOf(u8, db.last_error, "[email protected]") != null);
  1362. }
  1363. test "engine: uuid primary key" {
  1364. var base_buf: [128]u8 = undefined;
  1365. const base = tmpBase(&base_buf, "uuid");
  1366. cleanup(testing.allocator, base);
  1367. defer cleanup(testing.allocator, base);
  1368. var arena_state = std.heap.ArenaAllocator.init(testing.allocator);
  1369. defer arena_state.deinit();
  1370. const arena = arena_state.allocator();
  1371. const db = try Db.open(testing.allocator, base, .{ .primary_key = "@uuid" });
  1372. defer db.close();
  1373. const r = try db.insert(arena, try makeUser(arena, "Alice", "[email protected]", 1), false);
  1374. try testing.expectEqual(@as(usize, 12), r.pk_str.len);
  1375. var locs: std.ArrayList(Loc) = .empty;
  1376. try db.findPkLocs(arena, strKey(r.pk_str), &locs);
  1377. try testing.expectEqual(@as(usize, 1), locs.items.len);
  1378. }
  1379. test "jsNumFmt / jsParseInt / uuid shape" {
  1380. var buf: [32]u8 = undefined;
  1381. try testing.expectEqualStrings("42", jsNumFmt(&buf, 42));
  1382. try testing.expectEqualStrings("-7", jsNumFmt(&buf, -7));
  1383. try testing.expectEqualStrings("1.5", jsNumFmt(&buf, 1.5));
  1384. try testing.expectEqual(@as(f64, 1), jsParseInt("1.5").?);
  1385. try testing.expectEqual(@as(f64, -12), jsParseInt("-12abc").?);
  1386. try testing.expect(jsParseInt("abc") == null);
  1387. var ubuf: [12]u8 = undefined;
  1388. const u = genUuid(&ubuf);
  1389. try testing.expectEqual(@as(usize, 12), u.len);
  1390. for (u) |ch| try testing.expect((ch >= '0' and ch <= '9') or (ch >= 'a' and ch <= 'z'));
  1391. }

Branches

Latest commits

  • fdfb4b1bgitoria: Hybriel master 06617221 (plugin allocators 3a781359 + 413f60e4, mpackdb 2cb7ae5e, http1 773de63e); gates 200/0, 46/0, 44/0mre
  • 5b46ac84antcolony#40: LOG.md — missions 069/072 are antcolony missions (report paths on Byrodin)mre
  • 5602ff41gitoria: Hybriel master 190aa11d (fc838894 GC correctness, #127 mountKids by reference, #126, #48) — tracker README flat; gates 200/0, 46/0, 44/0mre
  • e85eaf01gitoria: 069 round 2 — hybriel 1a096ad3 not adopted (Markdown SSR still grows); browser gate waits for the server-side logout before restartmre
  • 09ce4f3fgitoria: mission 069 re-vendor hybriel 8efba065 stopped (big SSR pages grow + slow down); lambda audit clean; old vendor keptmre
  • 3dc43108antcolony#40: mission references point to the moved missionsmre
  • 8d9450fdantcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmre
  • 205d5fe4gitoria: Hybriel master ff51cf46; ssh keys/tokens no double rows (session sync); gates follow #20mre
  • 9b27cb26gitoria#21: installable app (manifest, service worker, offline start page), own iconmre
  • 68dcb603deploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
  • e2deed6dgitoria#20: "Add code" only on the Code page of an empty repository, no collapsiblemre
  • 8bb97ffddeploy.sh: never send .git or .gitignore to Byrodinmre
  • fd981932State of 2026-09-27; bin/ no longer tracked (Hybriel commit is in README)mre
  • 4a2d7125initial commitmre