2
0

engine.zig 76 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597
  1. const std = @import("std");
  2. const pkvdb = @import("pkvdb.zig");
  3. const keydir = @import("keydir.zig");
  4. const ordered_index = @import("ordered_index.zig");
  5. pub const RecordRef = keydir.RecordRef;
  6. const KeyReader = keydir.Reader;
  7. pub const pkvdb_max_frame = pkvdb.max_transaction_size;
  8. pub const Operation = struct {
  9. opcode: pkvdb.Opcode,
  10. key: []const u8,
  11. value: []const u8 = "",
  12. };
  13. pub const Value = struct {
  14. bytes: []u8,
  15. lsn: u64,
  16. };
  17. pub const CompareCheck = struct {
  18. key: []const u8,
  19. expected_lsn: u64,
  20. };
  21. pub const CompareBatchResult = struct {
  22. lsn: u64,
  23. committed: bool,
  24. };
  25. pub const ScanEntry = struct {
  26. key: []u8,
  27. value: ?[]u8,
  28. lsn: u64,
  29. };
  30. pub const ScanBatch = struct {
  31. entries: []ScanEntry,
  32. next_cursor: []u8,
  33. done: bool,
  34. pub fn deinit(self: *ScanBatch, allocator: std.mem.Allocator) void {
  35. for (self.entries) |entry| {
  36. allocator.free(entry.key);
  37. if (entry.value) |value| allocator.free(value);
  38. }
  39. allocator.free(self.entries);
  40. allocator.free(self.next_cursor);
  41. self.* = undefined;
  42. }
  43. };
  44. pub const Status = struct {
  45. uuid: [16]u8,
  46. file_bytes: u64,
  47. latest_lsn: u64,
  48. oldest_lsn: u64,
  49. checkpoint_lsn: u64,
  50. journal_bytes_since_checkpoint: u64,
  51. live_keys: u64,
  52. keydir_bytes: u64,
  53. ordered_index_bytes: u64,
  54. bytes_written: u64,
  55. checksum_failures: u64,
  56. partial_tails: u64,
  57. recovery_ns: u64,
  58. checkpoint_ns: u64,
  59. connection_bytes: u64,
  60. active_requests: u64,
  61. commit_groups: u64,
  62. committed_transactions: u64,
  63. largest_commit_group: u64,
  64. };
  65. const Root = struct {
  66. block: pkvdb.Superblock,
  67. manifest: ?pkvdb.Manifest,
  68. };
  69. const PendingWrite = struct {
  70. operations: []const Operation,
  71. metadata: []const u8,
  72. checks: []const CompareCheck = &.{},
  73. next: ?*PendingWrite = null,
  74. completion: *WriteCompletion,
  75. lsn: u64 = 0,
  76. committed: bool = true,
  77. frame_position: usize = 0,
  78. changed: bool = false,
  79. prepared_position: usize = 0,
  80. prepared_count: usize = 0,
  81. bytes: usize,
  82. };
  83. const WriteCompletion = struct {
  84. condition: std.Thread.Condition = .{},
  85. remaining: usize,
  86. failure: ?anyerror = null,
  87. };
  88. const max_group_transactions = 4096;
  89. const max_group_bytes = 64 * 1024 * 1024;
  90. const max_queued_bytes = 128 * 1024 * 1024;
  91. const group_wait_ns = 250 * std.time.ns_per_us;
  92. const map_interval = 64 * 1024 * 1024;
  93. const Mapping = struct {
  94. bytes: []align(std.heap.page_size_min) u8,
  95. file_start: u64,
  96. logical_start: u64,
  97. logical_end: u64,
  98. };
  99. pub const Engine = struct {
  100. allocator: std.mem.Allocator,
  101. file: std.fs.File,
  102. directory: keydir.KeyDir,
  103. ordered: ordered_index.OrderedIndex,
  104. ordered_ready: bool = false,
  105. mappings: std.ArrayListUnmanaged(Mapping) = .{},
  106. mapping_lock: std.Thread.RwLock = .{},
  107. mapped_length: u64 = 0,
  108. lock: std.Thread.RwLock = .{},
  109. io_mutex: std.Thread.Mutex = .{},
  110. checkpoint_mutex: std.Thread.Mutex = .{},
  111. ordered_gate: std.Thread.Mutex = .{},
  112. queue_mutex: std.Thread.Mutex = .{},
  113. queue_condition: std.Thread.Condition = .{},
  114. queue_head: ?*PendingWrite = null,
  115. queue_tail: ?*PendingWrite = null,
  116. queued_bytes: usize = 0,
  117. writer_thread: ?std.Thread = null,
  118. writer_stopping: bool = false,
  119. writer_failed: bool = false,
  120. commit_groups: u64 = 0,
  121. committed_transactions: u64 = 0,
  122. largest_commit_group: u64 = 0,
  123. file_length: u64,
  124. latest_lsn: u64 = 0,
  125. oldest_lsn: u64 = 0,
  126. checkpoint_lsn: u64 = 0,
  127. journal_bytes_since_checkpoint: u64 = 0,
  128. bytes_written: u64 = 0,
  129. checksum_failures: u64 = 0,
  130. partial_tails: u64 = 0,
  131. recovery_ns: u64 = 0,
  132. checkpoint_ns: u64 = 0,
  133. connection_bytes: std.atomic.Value(u64) = std.atomic.Value(u64).init(0),
  134. active_requests: std.atomic.Value(u64) = std.atomic.Value(u64).init(0),
  135. uuid: [16]u8,
  136. generation: u64,
  137. active_superblock: u1,
  138. created_ns: i64,
  139. pub fn open(allocator: std.mem.Allocator, path: []const u8) !Engine {
  140. const started = std.time.nanoTimestamp();
  141. const file = std.fs.cwd().openFile(path, .{ .mode = .read_write }) catch |err| switch (err) {
  142. error.FileNotFound => try std.fs.cwd().createFile(path, .{ .read = true, .truncate = false }),
  143. else => return err,
  144. };
  145. errdefer file.close();
  146. var directory = try keydir.KeyDir.init(allocator);
  147. errdefer directory.deinit();
  148. var engine = Engine{
  149. .allocator = allocator,
  150. .file = file,
  151. .directory = directory,
  152. .ordered = ordered_index.OrderedIndex.init(allocator),
  153. .file_length = try file.getEndPos(),
  154. .uuid = undefined,
  155. .generation = 0,
  156. .active_superblock = 0,
  157. .created_ns = now(),
  158. };
  159. errdefer engine.ordered.deinit();
  160. if (engine.file_length == 0) {
  161. try engine.initialize();
  162. } else {
  163. try engine.recover();
  164. }
  165. const elapsed = std.time.nanoTimestamp() - started;
  166. engine.recovery_ns = if (elapsed > 0) @intCast(elapsed) else 0;
  167. engine.mapTail(true);
  168. return engine;
  169. }
  170. pub fn close(self: *Engine) void {
  171. self.queue_mutex.lock();
  172. self.writer_stopping = true;
  173. self.queue_condition.broadcast();
  174. self.queue_mutex.unlock();
  175. if (self.writer_thread) |thread| thread.join();
  176. for (self.mappings.items) |mapping| std.posix.munmap(mapping.bytes);
  177. self.mappings.deinit(self.allocator);
  178. self.ordered.deinit();
  179. self.directory.deinit();
  180. self.file.close();
  181. self.* = undefined;
  182. }
  183. fn now() i64 {
  184. const value = std.time.nanoTimestamp();
  185. return std.math.cast(i64, value) orelse if (value < 0) std.math.minInt(i64) else std.math.maxInt(i64);
  186. }
  187. fn mapTail(self: *Engine, force: bool) void {
  188. if (self.file_length <= self.mapped_length or !force and self.file_length - self.mapped_length < map_interval) return;
  189. const page_size = std.heap.pageSize();
  190. const file_start = self.mapped_length - self.mapped_length % page_size;
  191. const length = std.math.cast(usize, self.file_length - file_start) orelse return;
  192. const bytes = std.posix.mmap(null, length, std.posix.PROT.READ, .{ .TYPE = .SHARED }, self.file.handle, file_start) catch return;
  193. if (!force) self.mapping_lock.lock();
  194. defer if (!force) self.mapping_lock.unlock();
  195. self.mappings.append(self.allocator, .{ .bytes = bytes, .file_start = file_start, .logical_start = self.mapped_length, .logical_end = self.file_length }) catch {
  196. std.posix.munmap(bytes);
  197. return;
  198. };
  199. self.mapped_length = self.file_length;
  200. }
  201. fn readBytes(self: *Engine, destination: []u8, file_offset: u64) !usize {
  202. const read_end = try std.math.add(u64, file_offset, destination.len);
  203. self.mapping_lock.lockShared();
  204. var index = self.mappings.items.len;
  205. while (index != 0) {
  206. index -= 1;
  207. const mapping = self.mappings.items[index];
  208. if (file_offset >= mapping.logical_start and read_end <= mapping.logical_end) {
  209. const start: usize = @intCast(file_offset - mapping.file_start);
  210. @memcpy(destination, mapping.bytes[start .. start + destination.len]);
  211. self.mapping_lock.unlockShared();
  212. return destination.len;
  213. }
  214. }
  215. self.mapping_lock.unlockShared();
  216. return self.file.preadAll(destination, file_offset);
  217. }
  218. fn readKeyBytes(context: *const anyopaque, destination: []u8, file_offset: u64) anyerror!usize {
  219. const self: *Engine = @ptrCast(@alignCast(@constCast(context)));
  220. return self.readBytes(destination, file_offset);
  221. }
  222. fn keyReader(self: *Engine) KeyReader {
  223. return .{ .context = self, .readFn = readKeyBytes };
  224. }
  225. fn initialize(self: *Engine) !void {
  226. std.crypto.random.bytes(&self.uuid);
  227. self.generation = 1;
  228. self.created_ns = now();
  229. var first: [4096]u8 = undefined;
  230. var second: [4096]u8 = undefined;
  231. pkvdb.encodeSuperblock(.{
  232. .generation = 1,
  233. .uuid = self.uuid,
  234. .manifest_offset = 0,
  235. .checkpoint_lsn = 0,
  236. .known_lsn = 0,
  237. .known_file_length = pkvdb.data_offset,
  238. .created_ns = self.created_ns,
  239. .updated_ns = self.created_ns,
  240. }, &first);
  241. pkvdb.encodeSuperblock(.{
  242. .generation = 0,
  243. .uuid = self.uuid,
  244. .manifest_offset = 0,
  245. .checkpoint_lsn = 0,
  246. .known_lsn = 0,
  247. .known_file_length = pkvdb.data_offset,
  248. .created_ns = self.created_ns,
  249. .updated_ns = self.created_ns,
  250. }, &second);
  251. try self.file.pwriteAll(&first, 0);
  252. try self.file.pwriteAll(&second, pkvdb.superblock_size);
  253. try self.file.sync();
  254. self.file_length = pkvdb.data_offset;
  255. try self.file.seekTo(self.file_length);
  256. self.bytes_written = pkvdb.data_offset;
  257. }
  258. fn readExtentHeader(self: *Engine, offset: u64) !pkvdb.ExtentHeader {
  259. var bytes: [64]u8 = undefined;
  260. if (try self.file.preadAll(&bytes, offset) != bytes.len) return error.Truncated;
  261. return pkvdb.decodeExtentHeader(&bytes) catch |err| {
  262. if (err == error.ChecksumMismatch) self.checksum_failures += 1;
  263. return err;
  264. };
  265. }
  266. fn validatePayload(self: *Engine, offset: u64, length: u64, expected: u32) !void {
  267. var crc: u32 = 0xffffffff;
  268. var buffer: [64 * 1024]u8 = undefined;
  269. var done: u64 = 0;
  270. while (done < length) {
  271. const amount: usize = @intCast(@min(buffer.len, length - done));
  272. if (try self.file.preadAll(buffer[0..amount], try std.math.add(u64, offset, done)) != amount) return error.Truncated;
  273. crc = pkvdb.crc32cUpdate(crc, buffer[0..amount]);
  274. done += amount;
  275. }
  276. if (~crc != expected) {
  277. self.checksum_failures += 1;
  278. return error.ChecksumMismatch;
  279. }
  280. }
  281. fn readExtentPayload(self: *Engine, offset: u64, header: pkvdb.ExtentHeader, max: u64) ![]u8 {
  282. if (header.payload_length > max or header.payload_length > std.math.maxInt(usize)) return error.InvalidLength;
  283. const payload_offset = try std.math.add(u64, offset, pkvdb.extent_header_size);
  284. const payload = try self.allocator.alloc(u8, @intCast(header.payload_length));
  285. errdefer self.allocator.free(payload);
  286. if (try self.file.preadAll(payload, payload_offset) != payload.len) return error.Truncated;
  287. if (pkvdb.crc32c(payload) != header.payload_crc) {
  288. self.checksum_failures += 1;
  289. return error.ChecksumMismatch;
  290. }
  291. return payload;
  292. }
  293. fn extentEnd(offset: u64, payload_length: u64) !u64 {
  294. return pkvdb.align8(try std.math.add(u64, try std.math.add(u64, offset, pkvdb.extent_header_size), payload_length));
  295. }
  296. fn loadManifest(self: *Engine, sb: pkvdb.Superblock) !?pkvdb.Manifest {
  297. if (sb.manifest_offset == 0) return null;
  298. const header = try self.readExtentHeader(sb.manifest_offset);
  299. if (header.extent_type != .manifest or header.version != 1 or header.payload_length != pkvdb.manifest_size) return error.InvalidManifest;
  300. if (try extentEnd(sb.manifest_offset, header.payload_length) > sb.known_file_length) return error.InvalidManifest;
  301. const payload = try self.readExtentPayload(sb.manifest_offset, header, pkvdb.manifest_size);
  302. defer self.allocator.free(payload);
  303. const manifest = try pkvdb.decodeManifest(payload);
  304. if (!std.mem.eql(u8, &manifest.uuid, &sb.uuid) or manifest.generation != sb.generation or manifest.checkpoint_lsn != sb.checkpoint_lsn) return error.InvalidManifest;
  305. if (manifest.known_tail != sb.known_file_length or manifest.known_lsn != sb.known_lsn) return error.InvalidManifest;
  306. if (manifest.entries_offset == 0 or manifest.entries_offset >= sb.manifest_offset or manifest.replay_offset < pkvdb.data_offset or manifest.replay_offset > manifest.entries_offset) return error.InvalidManifest;
  307. try self.validateCheckpointExtent(manifest.entries_offset, manifest.checkpoint_lsn, sb.known_file_length);
  308. return manifest;
  309. }
  310. fn validateCheckpointExtent(self: *Engine, offset: u64, lsn: u64, file_length: u64) !void {
  311. const header = try self.readExtentHeader(offset);
  312. if (header.extent_type != .checkpoint_entries or header.version != 1 or header.first_lsn != lsn or header.last_lsn != lsn) return error.InvalidCheckpoint;
  313. const end = try extentEnd(offset, header.payload_length);
  314. if (end > file_length) return error.InvalidCheckpoint;
  315. try self.validatePayload(offset + pkvdb.extent_header_size, header.payload_length, header.payload_crc);
  316. var bytes: [56]u8 = undefined;
  317. if (header.payload_length < bytes.len or try self.file.preadAll(&bytes, offset + pkvdb.extent_header_size) != bytes.len) return error.InvalidCheckpoint;
  318. const checkpoint_header = try pkvdb.decodeCheckpointHeaderOnly(&bytes, header.payload_length);
  319. if (checkpoint_header.lsn != lsn) return error.InvalidCheckpoint;
  320. }
  321. fn recover(self: *Engine) !void {
  322. if (self.file_length < pkvdb.data_offset) return error.Truncated;
  323. var blocks: [2][4096]u8 = undefined;
  324. _ = try self.file.preadAll(&blocks[0], 0);
  325. _ = try self.file.preadAll(&blocks[1], pkvdb.superblock_size);
  326. var roots: [2]?Root = .{ null, null };
  327. for (0..2) |index| {
  328. const sb = pkvdb.decodeSuperblock(&blocks[index], self.file_length) catch continue;
  329. const manifest = self.loadManifest(sb) catch continue;
  330. roots[index] = .{ .block = sb, .manifest = manifest };
  331. }
  332. var selected_index: usize = 0;
  333. const selected = if (roots[0] != null and roots[1] != null) blk: {
  334. selected_index = if (roots[1].?.block.generation > roots[0].?.block.generation) 1 else 0;
  335. break :blk roots[selected_index].?;
  336. } else if (roots[0]) |root| root else if (roots[1]) |root| blk: {
  337. selected_index = 1;
  338. break :blk root;
  339. } else return error.NoUsableSuperblock;
  340. self.uuid = selected.block.uuid;
  341. self.generation = selected.block.generation;
  342. self.active_superblock = @intCast(selected_index);
  343. self.created_ns = selected.block.created_ns;
  344. self.latest_lsn = selected.block.checkpoint_lsn;
  345. self.checkpoint_lsn = selected.block.checkpoint_lsn;
  346. var replay_offset = pkvdb.data_offset;
  347. if (selected.manifest) |manifest| {
  348. try self.loadCheckpoint(manifest.entries_offset, selected.block.known_file_length);
  349. replay_offset = manifest.replay_offset;
  350. self.oldest_lsn = manifest.history_start_lsn;
  351. }
  352. try self.scanExtents(replay_offset);
  353. try self.file.setEndPos(self.file_length);
  354. try self.file.seekTo(self.file_length);
  355. self.bytes_written = self.file_length;
  356. }
  357. fn loadCheckpoint(self: *Engine, offset: u64, known_length: u64) !void {
  358. const extent = try self.readExtentHeader(offset);
  359. var header_bytes: [56]u8 = undefined;
  360. const payload_offset = offset + pkvdb.extent_header_size;
  361. if (try self.file.preadAll(&header_bytes, payload_offset) != header_bytes.len) return error.Truncated;
  362. const header = try pkvdb.decodeCheckpointHeaderOnly(&header_bytes, extent.payload_length);
  363. if (header.entry_count > std.math.maxInt(usize)) return error.InvalidLength;
  364. try self.directory.ensureAdditional(self.file, @intCast(header.entry_count));
  365. var index: u64 = 0;
  366. while (index < header.entry_count) : (index += 1) {
  367. var bytes: [48]u8 = undefined;
  368. const entry_offset = try std.math.add(u64, payload_offset + pkvdb.checkpoint_header_size, try std.math.mul(u64, index, pkvdb.checkpoint_entry_size));
  369. if (try self.file.preadAll(&bytes, entry_offset) != bytes.len) return error.Truncated;
  370. const entry = try pkvdb.decodeCheckpointEntry(&bytes, known_length);
  371. const record = RecordRef{ .hash = entry.hash, .lsn = entry.lsn, .key_offset = entry.key_offset, .value_offset = entry.value_offset, .key_len = entry.key_len, .value_len = entry.value_len, .flags = entry.flags };
  372. try self.directory.put(self.file, record);
  373. }
  374. }
  375. fn scanExtents(self: *Engine, start: u64) !void {
  376. var offset = start;
  377. var last_journal_lsn: u64 = 0;
  378. while (offset < self.file_length) {
  379. if (self.file_length - offset < pkvdb.extent_header_size) {
  380. self.partial_tails += 1;
  381. self.file_length = offset;
  382. break;
  383. }
  384. const header = self.readExtentHeader(offset) catch |failure| {
  385. if (self.file_length - offset == pkvdb.extent_header_size) {
  386. self.partial_tails += 1;
  387. self.file_length = offset;
  388. break;
  389. }
  390. return failure;
  391. };
  392. const end = extentEnd(offset, header.payload_length) catch return error.InvalidLength;
  393. if (end > self.file_length) {
  394. self.partial_tails += 1;
  395. self.file_length = offset;
  396. break;
  397. }
  398. try self.validatePayload(offset + pkvdb.extent_header_size, header.payload_length, header.payload_crc);
  399. if (header.extent_type == .journal and header.version == 1) {
  400. if (header.payload_length > pkvdb.max_transaction_size + pkvdb.group_header_size) return error.InvalidLength;
  401. const payload = try self.readExtentPayload(offset, header, pkvdb.max_transaction_size + pkvdb.group_header_size);
  402. defer self.allocator.free(payload);
  403. last_journal_lsn = try self.replayJournal(offset, payload, last_journal_lsn);
  404. self.journal_bytes_since_checkpoint += end - offset;
  405. } else if (header.extent_type == .store_metadata and header.version == 1) {
  406. if (header.payload_length > pkvdb.max_transaction_size) return error.InvalidLength;
  407. const payload = try self.readExtentPayload(offset, header, pkvdb.max_transaction_size);
  408. defer self.allocator.free(payload);
  409. try self.replayBaseline(offset, payload);
  410. }
  411. offset = end;
  412. }
  413. }
  414. fn replayBaseline(self: *Engine, extent_offset: u64, payload: []const u8) !void {
  415. if (payload.len < 24 or readInt(u16, payload, 0) != 1 or readInt(u64, payload, 8) != 1) return error.InvalidBaseline;
  416. const count = readInt(u32, payload, 4);
  417. if (count > pkvdb.max_operations) return error.InvalidBaseline;
  418. try self.directory.ensureAdditional(self.file, count);
  419. var position: usize = 24;
  420. for (0..count) |_| {
  421. if (position > payload.len or payload.len - position < 8) return error.InvalidBaseline;
  422. const key_length = readInt(u32, payload, position);
  423. const value_length = readInt(u32, payload, position + 4);
  424. if (key_length > pkvdb.max_key_size or value_length > pkvdb.max_value_size) return error.InvalidBaseline;
  425. const key_start = position + 8;
  426. const key_end = try std.math.add(usize, key_start, key_length);
  427. const value_end = try std.math.add(usize, key_end, value_length);
  428. if (value_end > payload.len) return error.InvalidBaseline;
  429. const key = payload[key_start..key_end];
  430. const key_offset = extent_offset + pkvdb.extent_header_size + key_start;
  431. try self.replaceRecord(key, .{ .hash = keydir.KeyDir.hash(key), .lsn = 1, .key_offset = key_offset, .value_offset = key_offset + key_length, .key_len = key_length, .value_len = value_length });
  432. position = @intCast(try pkvdb.align8(value_end));
  433. }
  434. if (position != payload.len) return error.InvalidBaseline;
  435. self.latest_lsn = 1;
  436. self.oldest_lsn = 1;
  437. }
  438. fn replayJournal(self: *Engine, extent_offset: u64, payload: []const u8, previous_lsn: u64) !u64 {
  439. if (payload.len < pkvdb.group_header_size or readInt(u16, payload, 0) != 1) return error.InvalidJournal;
  440. const count = readInt(u32, payload, 4);
  441. if (count == 0 or count > pkvdb.max_operations) return error.InvalidJournal;
  442. var position: usize = pkvdb.group_header_size;
  443. var last = previous_lsn;
  444. for (0..count) |_| {
  445. if (position > payload.len or payload.len - position < pkvdb.transaction_header_size) return error.InvalidJournal;
  446. const tx = try pkvdb.decodeTransactionHeader(payload[position..]);
  447. if (tx.total_length > payload.len - position) return error.InvalidJournal;
  448. const tx_end = position + tx.total_length;
  449. const body = payload[position + pkvdb.transaction_header_size .. tx_end];
  450. if (pkvdb.crc32c(body) != tx.payload_crc or tx.metadata_length > body.len) return error.ChecksumMismatch;
  451. if (last != 0 and tx.lsn <= last) return error.InvalidLsn;
  452. last = tx.lsn;
  453. if (tx.lsn > self.latest_lsn) try self.applyTransaction(extent_offset + pkvdb.extent_header_size, position, tx, payload[position..tx_end]);
  454. position = tx_end;
  455. }
  456. if (position != payload.len) return error.InvalidJournal;
  457. return last;
  458. }
  459. fn applyTransaction(self: *Engine, extent_payload_offset: u64, tx_position: usize, tx: pkvdb.TransactionHeader, frame: []const u8) !void {
  460. var position: usize = pkvdb.transaction_header_size + tx.metadata_length;
  461. try self.directory.ensureAdditional(self.file, tx.operation_count);
  462. for (0..tx.operation_count) |_| {
  463. if (position > frame.len or frame.len - position < pkvdb.operation_header_size) return error.InvalidJournal;
  464. const operation = try pkvdb.decodeOperationHeader(frame[position..]);
  465. const data_start = try std.math.add(usize, position, pkvdb.operation_header_size);
  466. const key_end = try std.math.add(usize, data_start, operation.key_length);
  467. const value_end = try std.math.add(usize, key_end, operation.value_length);
  468. const extension_end = try std.math.add(usize, value_end, operation.extension_length);
  469. if (extension_end > frame.len) return error.InvalidJournal;
  470. const key = frame[data_start..key_end];
  471. switch (operation.opcode) {
  472. .put => {
  473. const key_offset = try std.math.add(u64, extent_payload_offset, tx_position + data_start);
  474. const record = RecordRef{ .hash = keydir.KeyDir.hash(key), .lsn = tx.lsn, .key_offset = key_offset, .value_offset = key_offset + operation.key_length, .key_len = operation.key_length, .value_len = operation.value_length };
  475. try self.replaceRecord(key, record);
  476. },
  477. .delete => _ = try self.removeRecord(key),
  478. }
  479. position = @intCast(try pkvdb.align8(extension_end));
  480. }
  481. if (position != frame.len) return error.InvalidJournal;
  482. self.latest_lsn = tx.lsn;
  483. if (self.oldest_lsn == 0) self.oldest_lsn = tx.lsn;
  484. }
  485. fn replaceRecord(self: *Engine, key: []const u8, record: RecordRef) !void {
  486. if (self.ordered_ready) try self.ordered.put(key, record);
  487. try self.directory.putWithKey(self.file, key, record);
  488. }
  489. fn removeRecord(self: *Engine, key: []const u8) !bool {
  490. const removed = try self.directory.remove(self.file, key);
  491. if (removed and self.ordered_ready and !self.ordered.remove(key)) return error.IndexInconsistent;
  492. return removed;
  493. }
  494. fn appendExtent(self: *Engine, extent_type: pkvdb.ExtentType, first_lsn: u64, last_lsn: u64, payload: []const u8) !u64 {
  495. const offset = try pkvdb.align8(self.file_length);
  496. if (offset != self.file_length) {
  497. const padding = [_]u8{0} ** 8;
  498. try self.file.pwriteAll(padding[0 .. offset - self.file_length], self.file_length);
  499. }
  500. var header: [64]u8 = undefined;
  501. pkvdb.encodeExtentHeader(.{ .extent_type = extent_type, .payload_length = payload.len, .first_lsn = first_lsn, .last_lsn = last_lsn, .payload_crc = pkvdb.crc32c(payload) }, &header);
  502. const old_length = self.file_length;
  503. errdefer self.file.setEndPos(old_length) catch {};
  504. try self.file.pwriteAll(&header, offset);
  505. try self.file.pwriteAll(payload, offset + header.len);
  506. const end = try extentEnd(offset, payload.len);
  507. if (end > offset + header.len + payload.len) {
  508. const padding = [_]u8{0} ** 8;
  509. try self.file.pwriteAll(padding[0 .. end - (offset + header.len + payload.len)], offset + header.len + payload.len);
  510. }
  511. self.file_length = end;
  512. self.bytes_written += end - offset;
  513. return offset;
  514. }
  515. fn transactionLength(operations: []const Operation, metadata: []const u8) !usize {
  516. var frame_length: u64 = pkvdb.transaction_header_size + metadata.len;
  517. for (operations) |operation| {
  518. const raw = try std.math.add(u64, pkvdb.operation_header_size + operation.key.len + operation.value.len, 0);
  519. frame_length = try pkvdb.align8(try std.math.add(u64, frame_length, raw));
  520. }
  521. if (frame_length > pkvdb.max_transaction_size) return error.TransactionTooLarge;
  522. return @intCast(frame_length);
  523. }
  524. fn encodeTransaction(frame: []u8, operations: []const Operation, metadata: []const u8, lsn: u64, transaction_id: u64, timestamp_ns: i64) !void {
  525. @memset(frame, 0);
  526. var position: usize = pkvdb.transaction_header_size;
  527. @memcpy(frame[position .. position + metadata.len], metadata);
  528. position += metadata.len;
  529. for (operations) |operation| {
  530. var header: [16]u8 = undefined;
  531. pkvdb.encodeOperationHeader(.{ .opcode = operation.opcode, .key_length = @intCast(operation.key.len), .value_length = @intCast(operation.value.len) }, &header);
  532. @memcpy(frame[position .. position + header.len], &header);
  533. position += header.len;
  534. @memcpy(frame[position .. position + operation.key.len], operation.key);
  535. position += operation.key.len;
  536. @memcpy(frame[position .. position + operation.value.len], operation.value);
  537. position += operation.value.len;
  538. position = @intCast(try pkvdb.align8(position));
  539. }
  540. if (position != frame.len) return error.InvalidLength;
  541. const body = frame[pkvdb.transaction_header_size..];
  542. var tx_header: [56]u8 = undefined;
  543. pkvdb.encodeTransactionHeader(.{ .total_length = @intCast(frame.len), .lsn = lsn, .transaction_id = transaction_id, .timestamp_ns = timestamp_ns, .operation_count = @intCast(operations.len), .metadata_length = @intCast(metadata.len), .payload_crc = pkvdb.crc32c(body) }, &tx_header);
  544. @memcpy(frame[0..tx_header.len], &tx_header);
  545. }
  546. pub fn batchWrite(self: *Engine, operations: []const Operation, metadata: []const u8) !u64 {
  547. if (operations.len == 0 or operations.len > pkvdb.max_operations or metadata.len > pkvdb.max_transaction_size) return error.InvalidLength;
  548. for (operations) |operation| {
  549. if (operation.key.len > pkvdb.max_key_size or operation.value.len > pkvdb.max_value_size) return error.InvalidLength;
  550. if (operation.opcode == .delete and operation.value.len != 0) return error.InvalidLength;
  551. }
  552. const frame_length = try transactionLength(operations, metadata);
  553. var completion = WriteCompletion{ .remaining = 1 };
  554. var pending = PendingWrite{ .operations = operations, .metadata = metadata, .bytes = frame_length, .completion = &completion };
  555. try self.enqueueAndWait(&.{&pending}, &completion);
  556. if (completion.failure) |failure| return failure;
  557. return pending.lsn;
  558. }
  559. pub fn compareBatchWrite(self: *Engine, checks: []const CompareCheck, operations: []const Operation, metadata: []const u8) !CompareBatchResult {
  560. if (operations.len == 0 or operations.len > pkvdb.max_operations or metadata.len > pkvdb.max_transaction_size) return error.InvalidLength;
  561. if (checks.len > pkvdb.max_operations) return error.InvalidLength;
  562. for (operations) |operation| {
  563. if (operation.key.len > pkvdb.max_key_size or operation.value.len > pkvdb.max_value_size) return error.InvalidLength;
  564. if (operation.opcode == .delete and operation.value.len != 0) return error.InvalidLength;
  565. }
  566. for (checks) |check| {
  567. if (check.key.len > pkvdb.max_key_size) return error.InvalidLength;
  568. }
  569. const frame_length = try transactionLength(operations, metadata);
  570. var completion = WriteCompletion{ .remaining = 1 };
  571. var pending = PendingWrite{ .operations = operations, .metadata = metadata, .checks = checks, .bytes = frame_length, .completion = &completion };
  572. try self.enqueueAndWait(&.{&pending}, &completion);
  573. if (completion.failure) |failure| return failure;
  574. return .{ .lsn = pending.lsn, .committed = pending.committed };
  575. }
  576. pub fn putMany(self: *Engine, operations: []const Operation, lsns: []u64) !void {
  577. if (operations.len == 0 or operations.len != lsns.len or operations.len > max_group_transactions) return error.InvalidLength;
  578. const pending = try self.allocator.alloc(PendingWrite, operations.len);
  579. defer self.allocator.free(pending);
  580. const pointers = try self.allocator.alloc(*PendingWrite, operations.len);
  581. defer self.allocator.free(pointers);
  582. var completion = WriteCompletion{ .remaining = operations.len };
  583. for (operations, 0..) |operation, index| {
  584. if (operation.opcode != .put or operation.key.len > pkvdb.max_key_size or operation.value.len > pkvdb.max_value_size) return error.InvalidLength;
  585. const operation_slice = operations[index .. index + 1];
  586. pending[index] = .{ .operations = operation_slice, .metadata = "", .bytes = try transactionLength(operation_slice, ""), .completion = &completion };
  587. pointers[index] = &pending[index];
  588. }
  589. try self.enqueueAndWait(pointers, &completion);
  590. for (pending, 0..) |result, index| {
  591. if (completion.failure) |failure| return failure;
  592. lsns[index] = result.lsn;
  593. }
  594. }
  595. fn enqueueAndWait(self: *Engine, pending: []const *PendingWrite, completion: *WriteCompletion) !void {
  596. self.queue_mutex.lock();
  597. defer self.queue_mutex.unlock();
  598. if (self.writer_failed or self.writer_stopping) return error.StorageUnavailable;
  599. if (self.writer_thread == null) self.writer_thread = try std.Thread.spawn(.{}, writerMain, .{self});
  600. for (pending) |write| {
  601. while (self.queued_bytes > max_queued_bytes - write.bytes and !self.writer_failed and !self.writer_stopping) self.queue_condition.wait(&self.queue_mutex);
  602. if (self.writer_failed or self.writer_stopping) return error.StorageUnavailable;
  603. if (self.queue_tail) |tail| tail.next = write else self.queue_head = write;
  604. self.queue_tail = write;
  605. self.queued_bytes += write.bytes;
  606. }
  607. self.queue_condition.signal();
  608. while (completion.remaining != 0) completion.condition.wait(&self.queue_mutex);
  609. }
  610. fn writerMain(self: *Engine) void {
  611. var group: [max_group_transactions]*PendingWrite = undefined;
  612. while (true) {
  613. self.queue_mutex.lock();
  614. while (self.queue_head == null and !self.writer_stopping) self.queue_condition.wait(&self.queue_mutex);
  615. if (self.queue_head == null and self.writer_stopping) {
  616. self.queue_mutex.unlock();
  617. return;
  618. }
  619. _ = self.queue_condition.timedWait(&self.queue_mutex, group_wait_ns) catch {};
  620. var count: usize = 0;
  621. var bytes: usize = pkvdb.group_header_size;
  622. var conditional: bool = false;
  623. while (self.queue_head) |pending| {
  624. if (count != 0 and (count == group.len or bytes + pending.bytes > max_group_bytes or conditional or pending.checks.len != 0)) break;
  625. self.queue_head = pending.next;
  626. if (self.queue_head == null) self.queue_tail = null;
  627. pending.next = null;
  628. group[count] = pending;
  629. count += 1;
  630. bytes += pending.bytes;
  631. self.queued_bytes -= pending.bytes;
  632. conditional = pending.checks.len != 0;
  633. }
  634. self.queue_condition.broadcast();
  635. self.queue_mutex.unlock();
  636. self.processGroup(group[0..count], bytes) catch |failure| {
  637. self.queue_mutex.lock();
  638. self.writer_failed = true;
  639. for (group[0..count]) |pending| {
  640. pending.completion.failure = failure;
  641. pending.completion.remaining -= 1;
  642. if (pending.completion.remaining == 0) pending.completion.condition.signal();
  643. }
  644. while (self.queue_head) |pending| {
  645. self.queue_head = pending.next;
  646. pending.completion.failure = error.StorageUnavailable;
  647. pending.completion.remaining -= 1;
  648. if (pending.completion.remaining == 0) pending.completion.condition.signal();
  649. }
  650. self.queue_tail = null;
  651. self.queued_bytes = 0;
  652. self.queue_condition.broadcast();
  653. self.queue_mutex.unlock();
  654. continue;
  655. };
  656. self.queue_mutex.lock();
  657. for (group[0..count]) |pending| {
  658. pending.completion.remaining -= 1;
  659. if (pending.completion.remaining == 0) pending.completion.condition.signal();
  660. }
  661. self.queue_mutex.unlock();
  662. }
  663. }
  664. fn checkConditions(self: *Engine, checks: []const CompareCheck) !bool {
  665. for (checks) |check| {
  666. const record = try self.directory.get(self.file, check.key);
  667. if (check.expected_lsn == 0) {
  668. if (record != null) return false;
  669. } else if (record == null or record.?.lsn != check.expected_lsn) {
  670. return false;
  671. }
  672. }
  673. return true;
  674. }
  675. fn processGroup(self: *Engine, group: []*PendingWrite, payload_length: usize) !void {
  676. var operation_count: usize = 0;
  677. for (group) |pending| operation_count = try std.math.add(usize, operation_count, pending.operations.len);
  678. const payload = try self.allocator.alloc(u8, payload_length);
  679. defer self.allocator.free(payload);
  680. self.ordered_gate.lock();
  681. defer self.ordered_gate.unlock();
  682. self.lock.lock();
  683. if (group[0].checks.len != 0) {
  684. if (group.len != 1) {
  685. self.lock.unlock();
  686. return error.InvalidConditionalGroup;
  687. }
  688. const conditions_match = self.checkConditions(group[0].checks) catch |failure| {
  689. self.lock.unlock();
  690. return failure;
  691. };
  692. if (!conditions_match) {
  693. group[0].committed = false;
  694. group[0].lsn = 0;
  695. self.lock.unlock();
  696. return;
  697. }
  698. }
  699. self.directory.ensureAdditional(self.file, operation_count) catch |failure| {
  700. self.lock.unlock();
  701. return failure;
  702. };
  703. const first_lsn = std.math.add(u64, self.latest_lsn, 1) catch |failure| {
  704. self.lock.unlock();
  705. return failure;
  706. };
  707. const prepare_ordered = self.ordered_ready;
  708. self.lock.unlock();
  709. var prepared_storage: []ordered_index.Prepared = &.{};
  710. if (prepare_ordered) prepared_storage = try self.allocator.alloc(ordered_index.Prepared, operation_count);
  711. var prepared_count: usize = 0;
  712. defer {
  713. for (prepared_storage[0..prepared_count]) |*entry| self.ordered.discard(entry);
  714. if (prepared_storage.len != 0) self.allocator.free(prepared_storage);
  715. }
  716. if (prepare_ordered) for (group) |pending| {
  717. pending.prepared_position = prepared_count;
  718. for (pending.operations) |operation| if (operation.opcode == .put) {
  719. prepared_storage[prepared_count] = try self.ordered.prepare(operation.key);
  720. prepared_count += 1;
  721. };
  722. pending.prepared_count = prepared_count - pending.prepared_position;
  723. };
  724. @memset(payload, 0);
  725. const group_timestamp = now();
  726. writeInt(u16, payload, 0, 1);
  727. writeInt(u32, payload, 4, @intCast(group.len));
  728. writeInt(i64, payload, 8, group_timestamp);
  729. var position: usize = pkvdb.group_header_size;
  730. var lsn = first_lsn;
  731. for (group, 0..) |pending, index| {
  732. pending.frame_position = position;
  733. const timestamp = std.math.add(i64, group_timestamp, @intCast(index)) catch std.math.maxInt(i64);
  734. try encodeTransaction(payload[position .. position + pending.bytes], pending.operations, pending.metadata, lsn, lsn, timestamp);
  735. pending.lsn = lsn;
  736. position += pending.bytes;
  737. lsn = try std.math.add(u64, lsn, 1);
  738. }
  739. const last_lsn = lsn - 1;
  740. self.io_mutex.lock();
  741. defer self.io_mutex.unlock();
  742. const offset = try self.appendExtent(.journal, first_lsn, last_lsn, payload);
  743. self.file.sync() catch |failure| {
  744. self.file.setEndPos(offset) catch {};
  745. self.file_length = offset;
  746. return failure;
  747. };
  748. self.mapTail(false);
  749. self.lock.lock();
  750. defer self.lock.unlock();
  751. for (group) |pending| {
  752. try self.publishPending(offset + pkvdb.extent_header_size, pending, prepared_storage[0..prepared_count]);
  753. }
  754. self.journal_bytes_since_checkpoint += self.file_length - offset;
  755. self.commit_groups += 1;
  756. self.committed_transactions += group.len;
  757. self.largest_commit_group = @max(self.largest_commit_group, group.len);
  758. }
  759. fn publishPending(self: *Engine, extent_payload_offset: u64, pending: *PendingWrite, prepared: []ordered_index.Prepared) !void {
  760. var position: usize = pkvdb.transaction_header_size + pending.metadata.len;
  761. var prepared_position = pending.prepared_position;
  762. for (pending.operations) |operation| {
  763. const key_offset = try std.math.add(u64, extent_payload_offset, pending.frame_position + position + pkvdb.operation_header_size);
  764. switch (operation.opcode) {
  765. .put => {
  766. const record = RecordRef{ .hash = keydir.KeyDir.hash(operation.key), .lsn = pending.lsn, .key_offset = key_offset, .value_offset = key_offset + operation.key.len, .key_len = @intCast(operation.key.len), .value_len = @intCast(operation.value.len) };
  767. if (self.ordered_ready) {
  768. self.ordered.putPrepared(&prepared[prepared_position], record);
  769. prepared_position += 1;
  770. }
  771. try self.directory.putWithKey(self.file, operation.key, record);
  772. },
  773. .delete => pending.changed = (try self.removeRecord(operation.key)) or pending.changed,
  774. }
  775. position = @intCast(try pkvdb.align8(position + pkvdb.operation_header_size + operation.key.len + operation.value.len));
  776. }
  777. if (prepared_position != pending.prepared_position + pending.prepared_count) return error.IndexInconsistent;
  778. if (position != pending.bytes) return error.InvalidLength;
  779. self.latest_lsn = pending.lsn;
  780. if (self.oldest_lsn == 0) self.oldest_lsn = pending.lsn;
  781. }
  782. pub fn importBaseline(self: *Engine, operations: []const Operation) !void {
  783. if (operations.len > pkvdb.max_operations) return error.InvalidLength;
  784. self.io_mutex.lock();
  785. defer self.io_mutex.unlock();
  786. self.lock.lock();
  787. defer self.lock.unlock();
  788. if (self.checkpoint_lsn != 0 or self.latest_lsn > 1 or self.directory.count != 0 and self.latest_lsn == 0) return error.InvalidState;
  789. var length: u64 = 24;
  790. for (operations) |operation| {
  791. if (operation.opcode != .put or operation.key.len > pkvdb.max_key_size or operation.value.len > pkvdb.max_value_size) return error.InvalidLength;
  792. length = try pkvdb.align8(try std.math.add(u64, length, 8 + operation.key.len + operation.value.len));
  793. }
  794. if (length > pkvdb.max_transaction_size) return error.TransactionTooLarge;
  795. const payload = try self.allocator.alloc(u8, @intCast(length));
  796. defer self.allocator.free(payload);
  797. @memset(payload, 0);
  798. writeInt(u16, payload, 0, 1);
  799. writeInt(u32, payload, 4, @intCast(operations.len));
  800. writeInt(u64, payload, 8, 1);
  801. writeInt(i64, payload, 16, now());
  802. var position: usize = 24;
  803. for (operations) |operation| {
  804. writeInt(u32, payload, position, @intCast(operation.key.len));
  805. writeInt(u32, payload, position + 4, @intCast(operation.value.len));
  806. position += 8;
  807. @memcpy(payload[position .. position + operation.key.len], operation.key);
  808. position += operation.key.len;
  809. @memcpy(payload[position .. position + operation.value.len], operation.value);
  810. position += operation.value.len;
  811. const aligned: usize = @intCast(try pkvdb.align8(position));
  812. @memset(payload[position..aligned], 0);
  813. position = aligned;
  814. }
  815. try self.directory.ensureAdditional(self.file, operations.len);
  816. const offset = try self.appendExtent(.store_metadata, 1, 1, payload);
  817. try self.file.sync();
  818. try self.replayBaseline(offset, payload);
  819. }
  820. pub fn put(self: *Engine, key: []const u8, value: []const u8) !u64 {
  821. return self.batchWrite(&.{.{ .opcode = .put, .key = key, .value = value }}, "");
  822. }
  823. pub fn delete(self: *Engine, key: []const u8) !bool {
  824. if (key.len > pkvdb.max_key_size) return error.InvalidLength;
  825. const operations = [_]Operation{.{ .opcode = .delete, .key = key }};
  826. var completion = WriteCompletion{ .remaining = 1 };
  827. var pending = PendingWrite{ .operations = &operations, .metadata = "", .bytes = try transactionLength(&operations, ""), .completion = &completion };
  828. try self.enqueueAndWait(&.{&pending}, &completion);
  829. if (completion.failure) |failure| return failure;
  830. return pending.changed;
  831. }
  832. pub fn get(self: *Engine, allocator: std.mem.Allocator, key: []const u8) !?Value {
  833. const record = try self.getRef(key) orelse return null;
  834. const bytes = try allocator.alloc(u8, record.value_len);
  835. errdefer allocator.free(bytes);
  836. _ = try self.readValue(record, bytes, 0);
  837. return .{ .bytes = bytes, .lsn = record.lsn };
  838. }
  839. pub fn getRef(self: *Engine, key: []const u8) !?RecordRef {
  840. if (key.len > pkvdb.max_key_size) return error.InvalidLength;
  841. self.lock.lockShared();
  842. defer self.lock.unlockShared();
  843. return try self.directory.getWithReader(self.keyReader(), key);
  844. }
  845. pub fn readValue(self: *Engine, record: RecordRef, destination: []u8, value_position: u32) !usize {
  846. if (value_position > record.value_len) return error.InvalidOffset;
  847. const amount = @min(destination.len, record.value_len - value_position);
  848. const file_offset = record.value_offset + value_position;
  849. const got = try self.readBytes(destination[0..amount], file_offset);
  850. if (got != amount) return error.Truncated;
  851. return amount;
  852. }
  853. pub fn exists(self: *Engine, key: []const u8) !bool {
  854. self.lock.lockShared();
  855. defer self.lock.unlockShared();
  856. return (try self.directory.getWithReader(self.keyReader(), key)) != null;
  857. }
  858. pub fn multiGet(self: *Engine, allocator: std.mem.Allocator, keys: []const []const u8) ![]?Value {
  859. if (keys.len > pkvdb.max_operations) return error.InvalidLength;
  860. const refs = try allocator.alloc(?RecordRef, keys.len);
  861. defer allocator.free(refs);
  862. @memset(refs, null);
  863. self.lock.lockShared();
  864. for (keys, 0..) |key, index| {
  865. if (key.len > pkvdb.max_key_size) {
  866. self.lock.unlockShared();
  867. return error.InvalidLength;
  868. }
  869. refs[index] = self.directory.getWithReader(self.keyReader(), key) catch |failure| {
  870. self.lock.unlockShared();
  871. return failure;
  872. };
  873. }
  874. self.lock.unlockShared();
  875. const values = try allocator.alloc(?Value, keys.len);
  876. errdefer allocator.free(values);
  877. @memset(values, null);
  878. errdefer for (values) |value| if (value) |present| allocator.free(present.bytes);
  879. for (refs, 0..) |entry, index| {
  880. const record = entry orelse continue;
  881. const bytes = try allocator.alloc(u8, record.value_len);
  882. errdefer allocator.free(bytes);
  883. _ = try self.readValue(record, bytes, 0);
  884. values[index] = .{ .bytes = bytes, .lsn = record.lsn };
  885. }
  886. return values;
  887. }
  888. pub fn scan(self: *Engine, allocator: std.mem.Allocator, prefix: []const u8, cursor: []const u8, limit: u32, include_values: bool, max_bytes: u32) !ScanBatch {
  889. if (prefix.len > pkvdb.max_key_size or cursor.len > pkvdb.max_key_size or limit == 0 or limit > 4096 or max_bytes == 0 or max_bytes > pkvdb.max_key_size + pkvdb.max_value_size + 1024) return error.InvalidLength;
  890. try self.ensureOrdered();
  891. self.lock.lockShared();
  892. defer self.lock.unlockShared();
  893. var node = self.ordered.lowerBound(if (cursor.len == 0) prefix else cursor);
  894. if (cursor.len != 0 and node != null and std.mem.eql(u8, node.?.key, cursor)) node = ordered_index.OrderedIndex.next(node.?);
  895. var entries = std.ArrayListUnmanaged(ScanEntry){};
  896. errdefer {
  897. for (entries.items) |entry| {
  898. allocator.free(entry.key);
  899. if (entry.value) |value| allocator.free(value);
  900. }
  901. entries.deinit(allocator);
  902. }
  903. var bytes_used: usize = 0;
  904. while (node) |current| {
  905. if (entries.items.len >= limit or !std.mem.startsWith(u8, current.key, prefix)) break;
  906. const record = current.record;
  907. const next_size = 16 + current.key.len + if (include_values) record.value_len else 0;
  908. if (next_size > max_bytes) return error.ScanEntryTooLarge;
  909. if (bytes_used + next_size > max_bytes) break;
  910. const key = try allocator.dupe(u8, current.key);
  911. errdefer allocator.free(key);
  912. var value: ?[]u8 = null;
  913. if (include_values) {
  914. value = try allocator.alloc(u8, record.value_len);
  915. errdefer allocator.free(value.?);
  916. _ = try self.readValue(record, value.?, 0);
  917. }
  918. try entries.append(allocator, .{ .key = key, .value = value, .lsn = record.lsn });
  919. bytes_used += next_size;
  920. node = ordered_index.OrderedIndex.next(current);
  921. }
  922. const next_cursor = if (entries.items.len == 0) try allocator.alloc(u8, 0) else try allocator.dupe(u8, entries.items[entries.items.len - 1].key);
  923. errdefer allocator.free(next_cursor);
  924. const done = node == null or !std.mem.startsWith(u8, node.?.key, prefix);
  925. return .{ .entries = try entries.toOwnedSlice(allocator), .next_cursor = next_cursor, .done = done };
  926. }
  927. fn ensureOrdered(self: *Engine) !void {
  928. self.lock.lockShared();
  929. const ready = self.ordered_ready;
  930. self.lock.unlockShared();
  931. if (ready) return;
  932. self.ordered_gate.lock();
  933. defer self.ordered_gate.unlock();
  934. self.lock.lock();
  935. defer self.lock.unlock();
  936. if (self.ordered_ready) return;
  937. const records = try self.directory.records(self.allocator);
  938. defer self.allocator.free(records);
  939. for (records) |record| {
  940. const key = try self.allocator.alloc(u8, record.key_len);
  941. defer self.allocator.free(key);
  942. if (try self.file.preadAll(key, record.key_offset) != key.len) return error.Truncated;
  943. try self.ordered.put(key, record);
  944. }
  945. self.ordered_ready = true;
  946. }
  947. pub fn checkpoint(self: *Engine) !void {
  948. const started = std.time.nanoTimestamp();
  949. self.checkpoint_mutex.lock();
  950. defer self.checkpoint_mutex.unlock();
  951. self.io_mutex.lock();
  952. self.lock.lockShared();
  953. const records = self.directory.records(self.allocator) catch |err| {
  954. self.lock.unlockShared();
  955. self.io_mutex.unlock();
  956. return err;
  957. };
  958. const lsn = self.latest_lsn;
  959. const replay_offset = self.file_length;
  960. self.lock.unlockShared();
  961. self.io_mutex.unlock();
  962. defer self.allocator.free(records);
  963. const payload_length = try std.math.add(usize, pkvdb.checkpoint_header_size, try std.math.mul(usize, records.len, pkvdb.checkpoint_entry_size));
  964. const payload = try self.allocator.alloc(u8, payload_length);
  965. defer self.allocator.free(payload);
  966. var header: [56]u8 = undefined;
  967. pkvdb.encodeCheckpointHeader(.{ .lsn = lsn, .timestamp_ns = now(), .entry_count = records.len, .source_start = pkvdb.data_offset, .source_end = replay_offset }, &header);
  968. @memcpy(payload[0..header.len], &header);
  969. for (records, 0..) |record, index| {
  970. if (record.flags > std.math.maxInt(u16)) return error.InvalidFlags;
  971. var entry: [48]u8 = undefined;
  972. pkvdb.encodeCheckpointEntry(.{ .hash = record.hash, .lsn = record.lsn, .key_offset = record.key_offset, .value_offset = record.value_offset, .key_len = record.key_len, .value_len = record.value_len, .flags = @intCast(record.flags) }, &entry);
  973. @memcpy(payload[pkvdb.checkpoint_header_size + index * pkvdb.checkpoint_entry_size ..][0..pkvdb.checkpoint_entry_size], &entry);
  974. }
  975. self.io_mutex.lock();
  976. defer self.io_mutex.unlock();
  977. self.lock.lockShared();
  978. const known_lsn = self.latest_lsn;
  979. const next_generation = std.math.add(u64, self.generation, 1) catch |failure| {
  980. self.lock.unlockShared();
  981. return failure;
  982. };
  983. const history_start_lsn = self.oldest_lsn;
  984. self.lock.unlockShared();
  985. const entries_offset = try self.appendExtent(.checkpoint_entries, lsn, lsn, payload);
  986. const manifest_offset = try pkvdb.align8(self.file_length);
  987. const known_tail = try extentEnd(manifest_offset, pkvdb.manifest_size);
  988. var manifest_bytes: [104]u8 = undefined;
  989. pkvdb.encodeManifest(.{ .uuid = self.uuid, .generation = next_generation, .checkpoint_lsn = lsn, .entries_offset = entries_offset, .ordered_offset = 0, .replay_offset = replay_offset, .known_tail = known_tail, .known_lsn = known_lsn, .history_start_lsn = history_start_lsn }, &manifest_bytes);
  990. const actual_manifest = try self.appendExtent(.manifest, lsn, known_lsn, &manifest_bytes);
  991. if (actual_manifest != manifest_offset or self.file_length != known_tail) return error.InvalidManifest;
  992. try self.file.sync();
  993. const inactive: u1 = self.active_superblock ^ 1;
  994. var block: [4096]u8 = undefined;
  995. pkvdb.encodeSuperblock(.{ .generation = next_generation, .uuid = self.uuid, .manifest_offset = manifest_offset, .checkpoint_lsn = lsn, .known_lsn = known_lsn, .known_file_length = known_tail, .created_ns = self.created_ns, .updated_ns = now() }, &block);
  996. try self.file.pwriteAll(&block, @as(u64, inactive) * pkvdb.superblock_size);
  997. try self.file.sync();
  998. self.lock.lock();
  999. self.active_superblock = inactive;
  1000. self.generation = next_generation;
  1001. self.checkpoint_lsn = lsn;
  1002. self.journal_bytes_since_checkpoint = self.file_length - replay_offset;
  1003. const elapsed = std.time.nanoTimestamp() - started;
  1004. self.checkpoint_ns = if (elapsed > 0) @intCast(elapsed) else 0;
  1005. self.lock.unlock();
  1006. }
  1007. pub fn status(self: *Engine) Status {
  1008. self.io_mutex.lock();
  1009. defer self.io_mutex.unlock();
  1010. self.lock.lockShared();
  1011. defer self.lock.unlockShared();
  1012. return .{
  1013. .uuid = self.uuid,
  1014. .file_bytes = self.file_length,
  1015. .latest_lsn = self.latest_lsn,
  1016. .oldest_lsn = self.oldest_lsn,
  1017. .checkpoint_lsn = self.checkpoint_lsn,
  1018. .journal_bytes_since_checkpoint = self.journal_bytes_since_checkpoint,
  1019. .live_keys = self.directory.count,
  1020. .keydir_bytes = self.directory.bytes(),
  1021. .ordered_index_bytes = self.ordered.allocated_bytes,
  1022. .bytes_written = self.bytes_written,
  1023. .checksum_failures = self.checksum_failures,
  1024. .partial_tails = self.partial_tails,
  1025. .recovery_ns = self.recovery_ns,
  1026. .checkpoint_ns = self.checkpoint_ns,
  1027. .connection_bytes = self.connection_bytes.load(.monotonic),
  1028. .active_requests = self.active_requests.load(.monotonic),
  1029. .commit_groups = self.commit_groups,
  1030. .committed_transactions = self.committed_transactions,
  1031. .largest_commit_group = self.largest_commit_group,
  1032. };
  1033. }
  1034. pub fn addConnectionBytes(self: *Engine, amount: u64) void {
  1035. _ = self.connection_bytes.fetchAdd(amount, .monotonic);
  1036. }
  1037. pub fn removeConnectionBytes(self: *Engine, amount: u64) void {
  1038. _ = self.connection_bytes.fetchSub(amount, .monotonic);
  1039. }
  1040. pub fn beginRequest(self: *Engine) void {
  1041. _ = self.active_requests.fetchAdd(1, .monotonic);
  1042. }
  1043. pub fn endRequest(self: *Engine) void {
  1044. _ = self.active_requests.fetchSub(1, .monotonic);
  1045. }
  1046. };
  1047. fn readInt(comptime T: type, bytes: []const u8, offset: usize) T {
  1048. return std.mem.readInt(T, bytes[offset..][0..@sizeOf(T)], .little);
  1049. }
  1050. fn writeInt(comptime T: type, bytes: []u8, offset: usize, value: T) void {
  1051. std.mem.writeInt(T, bytes[offset..][0..@sizeOf(T)], value, .little);
  1052. }
  1053. test "insert overwrite delete recreate and recovery" {
  1054. var tmp = std.testing.tmpDir(.{});
  1055. defer tmp.cleanup();
  1056. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1057. const directory = try tmp.dir.realpath(".", &path_buffer);
  1058. const path = try std.fmt.allocPrint(std.testing.allocator, "{s}/test.pkvdb", .{directory});
  1059. defer std.testing.allocator.free(path);
  1060. var engine = try Engine.open(std.testing.allocator, path);
  1061. _ = try engine.put("key", "one");
  1062. _ = try engine.put("key", "two");
  1063. try std.testing.expect(try engine.delete("key"));
  1064. _ = try engine.put("key", "three");
  1065. var value = (try engine.get(std.testing.allocator, "key")).?;
  1066. try std.testing.expectEqualStrings("three", value.bytes);
  1067. std.testing.allocator.free(value.bytes);
  1068. engine.close();
  1069. engine = try Engine.open(std.testing.allocator, path);
  1070. defer engine.close();
  1071. value = (try engine.get(std.testing.allocator, "key")).?;
  1072. defer std.testing.allocator.free(value.bytes);
  1073. try std.testing.expectEqualStrings("three", value.bytes);
  1074. try std.testing.expectEqual(@as(u64, 4), value.lsn);
  1075. }
  1076. test "atomic batch checkpoint tail and ordered scan" {
  1077. var tmp = std.testing.tmpDir(.{});
  1078. defer tmp.cleanup();
  1079. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1080. const directory = try tmp.dir.realpath(".", &path_buffer);
  1081. const path = try std.fmt.allocPrint(std.testing.allocator, "{s}/test.pkvdb", .{directory});
  1082. defer std.testing.allocator.free(path);
  1083. var engine = try Engine.open(std.testing.allocator, path);
  1084. _ = try engine.batchWrite(&.{ .{ .opcode = .put, .key = "p/2", .value = "b" }, .{ .opcode = .put, .key = "p/1", .value = "a" } }, "meta");
  1085. try engine.checkpoint();
  1086. _ = try engine.put("p/3", "c");
  1087. engine.close();
  1088. engine = try Engine.open(std.testing.allocator, path);
  1089. defer engine.close();
  1090. var batch = try engine.scan(std.testing.allocator, "p/", "", 2, true, 1024);
  1091. try std.testing.expectEqual(@as(usize, 2), batch.entries.len);
  1092. try std.testing.expectEqualStrings("p/1", batch.entries[0].key);
  1093. const cursor = try std.testing.allocator.dupe(u8, batch.next_cursor);
  1094. batch.deinit(std.testing.allocator);
  1095. defer std.testing.allocator.free(cursor);
  1096. batch = try engine.scan(std.testing.allocator, "p/", cursor, 2, true, 1024);
  1097. defer batch.deinit(std.testing.allocator);
  1098. try std.testing.expectEqual(@as(usize, 1), batch.entries.len);
  1099. try std.testing.expectEqualStrings("p/3", batch.entries[0].key);
  1100. }
  1101. test "scan supports full-size PKBFI pages and accounts for entry framing" {
  1102. var tmp = std.testing.tmpDir(.{});
  1103. defer tmp.cleanup();
  1104. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1105. const path = try testPath(&tmp, "large-scan.pkvdb", &path_buffer);
  1106. defer std.testing.allocator.free(path);
  1107. var engine = try Engine.open(std.testing.allocator, path);
  1108. defer engine.close();
  1109. const value = try std.testing.allocator.alloc(u8, 2 * 1024 * 1024);
  1110. defer std.testing.allocator.free(value);
  1111. @memset(value, 'x');
  1112. _ = try engine.put("large/key", value);
  1113. const entry_size = 16 + "large/key".len + value.len;
  1114. var batch = try engine.scan(std.testing.allocator, "large/", "", 1, true, @intCast(entry_size));
  1115. defer batch.deinit(std.testing.allocator);
  1116. try std.testing.expectEqual(@as(usize, 1), batch.entries.len);
  1117. try std.testing.expectEqual(value.len, batch.entries[0].value.?.len);
  1118. try std.testing.expectError(
  1119. error.ScanEntryTooLarge,
  1120. engine.scan(std.testing.allocator, "large/", "", 1, true, @intCast(entry_size - 1)),
  1121. );
  1122. }
  1123. fn testPath(tmp: *std.testing.TmpDir, name: []const u8, buffer: *[std.fs.max_path_bytes]u8) ![]const u8 {
  1124. const directory = try tmp.dir.realpath(".", buffer);
  1125. return std.fs.path.join(std.testing.allocator, &.{ directory, name });
  1126. }
  1127. test "partial final extent is ignored and removed" {
  1128. var tmp = std.testing.tmpDir(.{});
  1129. defer tmp.cleanup();
  1130. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1131. const path = try testPath(&tmp, "partial.pkvdb", &path_buffer);
  1132. defer std.testing.allocator.free(path);
  1133. var engine = try Engine.open(std.testing.allocator, path);
  1134. _ = try engine.put("durable", "value");
  1135. const valid_length = engine.status().file_bytes;
  1136. engine.close();
  1137. const file = try std.fs.cwd().openFile(path, .{ .mode = .read_write });
  1138. var header: [64]u8 = undefined;
  1139. pkvdb.encodeExtentHeader(.{ .extent_type = .journal, .payload_length = 100, .first_lsn = 2, .last_lsn = 2, .payload_crc = 0 }, &header);
  1140. try file.pwriteAll(&header, valid_length);
  1141. try file.pwriteAll("partial transaction", valid_length + header.len);
  1142. file.close();
  1143. engine = try Engine.open(std.testing.allocator, path);
  1144. defer engine.close();
  1145. const value = (try engine.get(std.testing.allocator, "durable")).?;
  1146. defer std.testing.allocator.free(value.bytes);
  1147. try std.testing.expectEqualStrings("value", value.bytes);
  1148. try std.testing.expectEqual(valid_length, engine.status().file_bytes);
  1149. try std.testing.expectEqual(@as(u64, 1), engine.status().partial_tails);
  1150. }
  1151. test "corruption in committed journal is explicit" {
  1152. var tmp = std.testing.tmpDir(.{});
  1153. defer tmp.cleanup();
  1154. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1155. const path = try testPath(&tmp, "corrupt.pkvdb", &path_buffer);
  1156. defer std.testing.allocator.free(path);
  1157. var engine = try Engine.open(std.testing.allocator, path);
  1158. _ = try engine.put("key", "value");
  1159. engine.close();
  1160. const file = try std.fs.cwd().openFile(path, .{ .mode = .read_write });
  1161. var byte: [1]u8 = undefined;
  1162. const offset = pkvdb.data_offset + pkvdb.extent_header_size + pkvdb.group_header_size + pkvdb.transaction_header_size + pkvdb.operation_header_size + 3;
  1163. _ = try file.preadAll(&byte, offset);
  1164. byte[0] ^= 1;
  1165. try file.pwriteAll(&byte, offset);
  1166. file.close();
  1167. try std.testing.expectError(error.ChecksumMismatch, Engine.open(std.testing.allocator, path));
  1168. }
  1169. test "one corrupted superblock falls back and adopts tail" {
  1170. var tmp = std.testing.tmpDir(.{});
  1171. defer tmp.cleanup();
  1172. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1173. const path = try testPath(&tmp, "root.pkvdb", &path_buffer);
  1174. defer std.testing.allocator.free(path);
  1175. var engine = try Engine.open(std.testing.allocator, path);
  1176. _ = try engine.put("before", "one");
  1177. try engine.checkpoint();
  1178. _ = try engine.put("after", "two");
  1179. engine.close();
  1180. const file = try std.fs.cwd().openFile(path, .{ .mode = .read_write });
  1181. var byte: [1]u8 = undefined;
  1182. _ = try file.preadAll(&byte, pkvdb.superblock_size + 24);
  1183. byte[0] ^= 1;
  1184. try file.pwriteAll(&byte, pkvdb.superblock_size + 24);
  1185. file.close();
  1186. engine = try Engine.open(std.testing.allocator, path);
  1187. defer engine.close();
  1188. const value = (try engine.get(std.testing.allocator, "after")).?;
  1189. defer std.testing.allocator.free(value.bytes);
  1190. try std.testing.expectEqualStrings("two", value.bytes);
  1191. try std.testing.expectEqual(@as(u64, 2), engine.status().latest_lsn);
  1192. }
  1193. test "unrooted checkpoint and manifest do not hide journal" {
  1194. var tmp = std.testing.tmpDir(.{});
  1195. defer tmp.cleanup();
  1196. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1197. const path = try testPath(&tmp, "interrupted-checkpoint.pkvdb", &path_buffer);
  1198. defer std.testing.allocator.free(path);
  1199. var engine = try Engine.open(std.testing.allocator, path);
  1200. _ = try engine.put("key", "value");
  1201. var roots: [8192]u8 = undefined;
  1202. _ = try engine.file.preadAll(&roots, 0);
  1203. try engine.checkpoint();
  1204. engine.close();
  1205. const file = try std.fs.cwd().openFile(path, .{ .mode = .read_write });
  1206. try file.pwriteAll(&roots, 0);
  1207. file.close();
  1208. engine = try Engine.open(std.testing.allocator, path);
  1209. defer engine.close();
  1210. const value = (try engine.get(std.testing.allocator, "key")).?;
  1211. defer std.testing.allocator.free(value.bytes);
  1212. try std.testing.expectEqualStrings("value", value.bytes);
  1213. try std.testing.expectEqual(@as(u64, 1), engine.status().latest_lsn);
  1214. }
  1215. test "unknown compatible extent is skipped by length" {
  1216. var tmp = std.testing.tmpDir(.{});
  1217. defer tmp.cleanup();
  1218. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1219. const path = try testPath(&tmp, "unknown.pkvdb", &path_buffer);
  1220. defer std.testing.allocator.free(path);
  1221. var engine = try Engine.open(std.testing.allocator, path);
  1222. _ = try engine.put("key", "value");
  1223. const offset = engine.status().file_bytes;
  1224. engine.close();
  1225. const file = try std.fs.cwd().openFile(path, .{ .mode = .read_write });
  1226. var header: [64]u8 = undefined;
  1227. pkvdb.encodeExtentHeader(.{ .extent_type = @enumFromInt(99), .version = 9, .payload_length = 3, .first_lsn = 0, .last_lsn = 0, .payload_crc = pkvdb.crc32c("new") }, &header);
  1228. try file.pwriteAll(&header, offset);
  1229. try file.pwriteAll("new", offset + header.len);
  1230. try file.pwriteAll(&([_]u8{0} ** 5), offset + header.len + 3);
  1231. file.close();
  1232. engine = try Engine.open(std.testing.allocator, path);
  1233. defer engine.close();
  1234. try std.testing.expect(try engine.exists("key"));
  1235. }
  1236. test "repeated overwrite and delete keep directory memory bounded" {
  1237. var tmp = std.testing.tmpDir(.{});
  1238. defer tmp.cleanup();
  1239. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1240. const path = try testPath(&tmp, "memory.pkvdb", &path_buffer);
  1241. defer std.testing.allocator.free(path);
  1242. var engine = try Engine.open(std.testing.allocator, path);
  1243. defer engine.close();
  1244. _ = try engine.put("same", "first");
  1245. const initial = engine.status().keydir_bytes;
  1246. for (0..100) |index| {
  1247. var value: [16]u8 = undefined;
  1248. const encoded = try std.fmt.bufPrint(&value, "{d}", .{index});
  1249. _ = try engine.put("same", encoded);
  1250. }
  1251. try std.testing.expect(try engine.delete("same"));
  1252. _ = try engine.put("same", "last");
  1253. try std.testing.expectEqual(initial, engine.status().keydir_bytes);
  1254. try std.testing.expectEqual(@as(u64, 1), engine.status().live_keys);
  1255. }
  1256. test "concurrent reads and writes preserve complete values" {
  1257. var tmp = std.testing.tmpDir(.{});
  1258. defer tmp.cleanup();
  1259. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1260. const path = try testPath(&tmp, "concurrent.pkvdb", &path_buffer);
  1261. defer std.testing.allocator.free(path);
  1262. var engine = try Engine.open(std.testing.allocator, path);
  1263. defer engine.close();
  1264. _ = try engine.put("key", "00000000");
  1265. var failed = std.atomic.Value(bool).init(false);
  1266. const Writer = struct {
  1267. fn run(target: *Engine, failure: *std.atomic.Value(bool)) void {
  1268. for (0..50) |index| {
  1269. var value: [8]u8 = undefined;
  1270. _ = std.fmt.bufPrint(&value, "{d:0>8}", .{index}) catch {
  1271. failure.store(true, .release);
  1272. return;
  1273. };
  1274. _ = target.put("key", &value) catch {
  1275. failure.store(true, .release);
  1276. return;
  1277. };
  1278. }
  1279. }
  1280. };
  1281. const thread = try std.Thread.spawn(.{}, Writer.run, .{ &engine, &failed });
  1282. for (0..50) |_| {
  1283. const value = try engine.get(std.testing.allocator, "key") orelse return error.TestUnexpectedResult;
  1284. try std.testing.expectEqual(@as(usize, 8), value.bytes.len);
  1285. std.testing.allocator.free(value.bytes);
  1286. }
  1287. thread.join();
  1288. try std.testing.expect(!failed.load(.acquire));
  1289. }
  1290. fn copyPrefix(source_path: []const u8, destination_path: []const u8) !void {
  1291. const source = try std.fs.cwd().openFile(source_path, .{ .mode = .read_only });
  1292. defer source.close();
  1293. const length = try source.getEndPos();
  1294. const destination = try std.fs.cwd().createFile(destination_path, .{ .read = true, .truncate = true });
  1295. defer destination.close();
  1296. var buffer: [4096]u8 = undefined;
  1297. var offset: u64 = 0;
  1298. while (offset < length) {
  1299. const amount: usize = @intCast(@min(buffer.len, length - offset));
  1300. const got = try source.preadAll(buffer[0..amount], offset);
  1301. if (got == 0) break;
  1302. try destination.writeAll(buffer[0..got]);
  1303. offset += got;
  1304. }
  1305. }
  1306. test "every live copy recovers while writes continue" {
  1307. var tmp = std.testing.tmpDir(.{});
  1308. defer tmp.cleanup();
  1309. var source_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1310. const source = try testPath(&tmp, "live.pkvdb", &source_buffer);
  1311. defer std.testing.allocator.free(source);
  1312. var engine = try Engine.open(std.testing.allocator, source);
  1313. defer engine.close();
  1314. _ = try engine.put("seed", "value");
  1315. var failed = std.atomic.Value(bool).init(false);
  1316. const Writer = struct {
  1317. fn run(target: *Engine, failure: *std.atomic.Value(bool)) void {
  1318. for (0..30) |index| {
  1319. var key: [16]u8 = undefined;
  1320. const encoded = std.fmt.bufPrint(&key, "key-{d}", .{index}) catch {
  1321. failure.store(true, .release);
  1322. return;
  1323. };
  1324. _ = target.put(encoded, "value") catch {
  1325. failure.store(true, .release);
  1326. return;
  1327. };
  1328. if (index == 15) target.checkpoint() catch {
  1329. failure.store(true, .release);
  1330. return;
  1331. };
  1332. }
  1333. }
  1334. };
  1335. const thread = try std.Thread.spawn(.{}, Writer.run, .{ &engine, &failed });
  1336. for (0..8) |index| {
  1337. var name: [32]u8 = undefined;
  1338. const filename = try std.fmt.bufPrint(&name, "copy-{d}.pkvdb", .{index});
  1339. var copy_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1340. const destination = try testPath(&tmp, filename, &copy_buffer);
  1341. defer std.testing.allocator.free(destination);
  1342. try copyPrefix(source, destination);
  1343. var copy = try Engine.open(std.testing.allocator, destination);
  1344. try std.testing.expect(try copy.exists("seed"));
  1345. copy.close();
  1346. }
  1347. thread.join();
  1348. try std.testing.expect(!failed.load(.acquire));
  1349. }
  1350. test "group commit preserves independent transaction LSNs" {
  1351. var tmp = std.testing.tmpDir(.{});
  1352. defer tmp.cleanup();
  1353. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1354. const path = try testPath(&tmp, "groups.pkvdb", &path_buffer);
  1355. defer std.testing.allocator.free(path);
  1356. var engine = try Engine.open(std.testing.allocator, path);
  1357. var operations: [256]Operation = undefined;
  1358. var keys: [256][8]u8 = undefined;
  1359. var lsns: [256]u64 = undefined;
  1360. for (&operations, 0..) |*operation, index| {
  1361. const key = try std.fmt.bufPrint(&keys[index], "k{d}", .{index});
  1362. operation.* = .{ .opcode = .put, .key = key, .value = "value" };
  1363. }
  1364. try engine.putMany(&operations, &lsns);
  1365. const status = engine.status();
  1366. try std.testing.expectEqual(@as(u64, 1), status.commit_groups);
  1367. try std.testing.expectEqual(@as(u64, 256), status.committed_transactions);
  1368. for (lsns, 0..) |lsn, index| try std.testing.expectEqual(index + 1, lsn);
  1369. engine.close();
  1370. engine = try Engine.open(std.testing.allocator, path);
  1371. defer engine.close();
  1372. try std.testing.expectEqual(@as(u64, 256), engine.status().latest_lsn);
  1373. for (operations) |operation| try std.testing.expect(try engine.exists(operation.key));
  1374. }
  1375. test "ordered index stays current after lazy construction" {
  1376. var tmp = std.testing.tmpDir(.{});
  1377. defer tmp.cleanup();
  1378. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1379. const path = try testPath(&tmp, "lazy-order.pkvdb", &path_buffer);
  1380. defer std.testing.allocator.free(path);
  1381. var engine = try Engine.open(std.testing.allocator, path);
  1382. defer engine.close();
  1383. _ = try engine.put("p/a", "one");
  1384. var batch = try engine.scan(std.testing.allocator, "p/", "", 10, false, 1024);
  1385. batch.deinit(std.testing.allocator);
  1386. _ = try engine.put("p/b", "two");
  1387. _ = try engine.put("p/a", "updated");
  1388. try std.testing.expect(try engine.delete("p/b"));
  1389. _ = try engine.put("p/c", "three");
  1390. batch = try engine.scan(std.testing.allocator, "p/", "", 10, true, 1024);
  1391. defer batch.deinit(std.testing.allocator);
  1392. try std.testing.expectEqual(@as(usize, 2), batch.entries.len);
  1393. try std.testing.expectEqualStrings("p/a", batch.entries[0].key);
  1394. try std.testing.expectEqualStrings("updated", batch.entries[0].value.?);
  1395. try std.testing.expectEqualStrings("p/c", batch.entries[1].key);
  1396. }
  1397. test "concurrent delete reports one removal" {
  1398. var tmp = std.testing.tmpDir(.{});
  1399. defer tmp.cleanup();
  1400. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1401. const path = try testPath(&tmp, "delete-race.pkvdb", &path_buffer);
  1402. defer std.testing.allocator.free(path);
  1403. var engine = try Engine.open(std.testing.allocator, path);
  1404. defer engine.close();
  1405. _ = try engine.put("key", "value");
  1406. var results: [2]bool = undefined;
  1407. var failed = std.atomic.Value(bool).init(false);
  1408. const Deleter = struct {
  1409. fn run(target: *Engine, result: *bool, failure: *std.atomic.Value(bool)) void {
  1410. result.* = target.delete("key") catch {
  1411. failure.store(true, .release);
  1412. return;
  1413. };
  1414. }
  1415. };
  1416. const first = try std.Thread.spawn(.{}, Deleter.run, .{ &engine, &results[0], &failed });
  1417. const second = try std.Thread.spawn(.{}, Deleter.run, .{ &engine, &results[1], &failed });
  1418. first.join();
  1419. second.join();
  1420. try std.testing.expect(!failed.load(.acquire));
  1421. try std.testing.expect(results[0] != results[1]);
  1422. }
  1423. test "compare batch write absent match and stale conflict" {
  1424. var tmp = std.testing.tmpDir(.{});
  1425. defer tmp.cleanup();
  1426. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1427. const path = try testPath(&tmp, "compare.pkvdb", &path_buffer);
  1428. defer std.testing.allocator.free(path);
  1429. var engine = try Engine.open(std.testing.allocator, path);
  1430. const absent_checks = [_]CompareCheck{.{ .key = "k", .expected_lsn = 0 }};
  1431. var result = try engine.compareBatchWrite(&absent_checks, &.{.{ .opcode = .put, .key = "k", .value = "one" }}, "meta");
  1432. try std.testing.expect(result.committed);
  1433. try std.testing.expectEqual(@as(u64, 1), result.lsn);
  1434. const match_checks = [_]CompareCheck{.{ .key = "k", .expected_lsn = result.lsn }};
  1435. result = try engine.compareBatchWrite(&match_checks, &.{.{ .opcode = .put, .key = "k", .value = "two" }}, "");
  1436. try std.testing.expect(result.committed);
  1437. try std.testing.expectEqual(@as(u64, 2), result.lsn);
  1438. const stale_checks = [_]CompareCheck{.{ .key = "k", .expected_lsn = 1 }};
  1439. result = try engine.compareBatchWrite(&stale_checks, &.{.{ .opcode = .put, .key = "k", .value = "three" }}, "");
  1440. try std.testing.expect(!result.committed);
  1441. try std.testing.expectEqual(@as(u64, 0), result.lsn);
  1442. var value = (try engine.get(std.testing.allocator, "k")).?;
  1443. try std.testing.expectEqualStrings("two", value.bytes);
  1444. try std.testing.expectEqual(@as(u64, 2), value.lsn);
  1445. std.testing.allocator.free(value.bytes);
  1446. engine.close();
  1447. engine = try Engine.open(std.testing.allocator, path);
  1448. defer engine.close();
  1449. value = (try engine.get(std.testing.allocator, "k")).?;
  1450. defer std.testing.allocator.free(value.bytes);
  1451. try std.testing.expectEqualStrings("two", value.bytes);
  1452. try std.testing.expectEqual(@as(u64, 2), value.lsn);
  1453. }
  1454. test "compare absent check rejects existing key" {
  1455. var tmp = std.testing.tmpDir(.{});
  1456. defer tmp.cleanup();
  1457. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1458. const path = try testPath(&tmp, "compare-absent.pkvdb", &path_buffer);
  1459. defer std.testing.allocator.free(path);
  1460. var engine = try Engine.open(std.testing.allocator, path);
  1461. defer engine.close();
  1462. _ = try engine.put("k", "seed");
  1463. const checks = [_]CompareCheck{.{ .key = "k", .expected_lsn = 0 }};
  1464. const result = try engine.compareBatchWrite(&checks, &.{.{ .opcode = .put, .key = "k", .value = "x" }}, "");
  1465. try std.testing.expect(!result.committed);
  1466. try std.testing.expectEqual(@as(u64, 0), result.lsn);
  1467. const value = (try engine.get(std.testing.allocator, "k")).?;
  1468. defer std.testing.allocator.free(value.bytes);
  1469. try std.testing.expectEqualStrings("seed", value.bytes);
  1470. }
  1471. test "concurrent compare batch writes expecting same lsn commit exactly one" {
  1472. var tmp = std.testing.tmpDir(.{});
  1473. defer tmp.cleanup();
  1474. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  1475. const path = try testPath(&tmp, "compare-race.pkvdb", &path_buffer);
  1476. defer std.testing.allocator.free(path);
  1477. var engine = try Engine.open(std.testing.allocator, path);
  1478. defer engine.close();
  1479. _ = try engine.put("key", "seed");
  1480. var results: [2]CompareBatchResult = undefined;
  1481. var failed = std.atomic.Value(bool).init(false);
  1482. const Worker = struct {
  1483. fn run(target: *Engine, result: *CompareBatchResult, failure: *std.atomic.Value(bool)) void {
  1484. const checks = [_]CompareCheck{.{ .key = "key", .expected_lsn = 1 }};
  1485. const operations = [_]Operation{.{ .opcode = .put, .key = "key", .value = "winner" }};
  1486. result.* = target.compareBatchWrite(&checks, &operations, "") catch {
  1487. failure.store(true, .release);
  1488. return;
  1489. };
  1490. }
  1491. };
  1492. const first = try std.Thread.spawn(.{}, Worker.run, .{ &engine, &results[0], &failed });
  1493. const second = try std.Thread.spawn(.{}, Worker.run, .{ &engine, &results[1], &failed });
  1494. first.join();
  1495. second.join();
  1496. try std.testing.expect(!failed.load(.acquire));
  1497. const committed = @as(u64, @intFromBool(results[0].committed)) + @as(u64, @intFromBool(results[1].committed));
  1498. try std.testing.expectEqual(@as(u64, 1), committed);
  1499. const winner = if (results[0].committed) results[0] else results[1];
  1500. try std.testing.expectEqual(@as(u64, 2), winner.lsn);
  1501. const value = (try engine.get(std.testing.allocator, "key")).?;
  1502. defer std.testing.allocator.free(value.bytes);
  1503. try std.testing.expectEqualStrings("winner", value.bytes);
  1504. }