engine.zig 68 KB

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