2
0

persistence.zig 5.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188
  1. const std = @import("std");
  2. const storage = @import("storage.zig");
  3. const BUFFER_SIZE = 1024 * 1024 * 8;
  4. const FLUSH_THRESHOLD = (BUFFER_SIZE * 3) / 4;
  5. var storage_file: ?std.fs.File = null;
  6. const c_allocator = std.heap.c_allocator;
  7. var mutex: std.Thread.Mutex = .{};
  8. var write_buffer: [BUFFER_SIZE]u8 = undefined;
  9. var buffer_position: usize = 0;
  10. var instant_wal: bool = false;
  11. const OPCode = enum {
  12. W,
  13. D,
  14. };
  15. pub fn init() !void {
  16. const cwd = std.fs.cwd();
  17. storage_file = cwd.openFile(".db", .{ .mode = .read_write }) catch |err| {
  18. if (err == std.fs.File.OpenError.FileNotFound) {
  19. std.debug.print("No persisted data found, starting fresh...\n", .{});
  20. storage_file = cwd.createFile(".db", .{ .read = true }) catch |ierr| {
  21. std.debug.print("Failed to create storage file: {any}\n", .{ierr});
  22. return;
  23. };
  24. std.debug.print("Created new storage file .db\n", .{});
  25. }
  26. return;
  27. };
  28. var record_count: usize = 0;
  29. restoreFromFile(storage_file.?, &record_count) catch |err| {
  30. std.debug.print("Failed to read storage file: {any}\n", .{err});
  31. return;
  32. };
  33. std.debug.print("Restored {d} records from persistence", .{record_count});
  34. storage_file.?.close();
  35. storage_file = cwd.openFile(".db", .{ .mode = .write_only }) catch |err| {
  36. std.debug.print("Failed to reopen storage file in append mode: {any}\n", .{err});
  37. return;
  38. };
  39. try storage_file.?.seekFromEnd(0);
  40. return;
  41. }
  42. fn restoreFromFile(file: std.fs.File, record_count: *usize) !void {
  43. var read_buffer: [BUFFER_SIZE]u8 = undefined;
  44. var record_buffer = std.ArrayListUnmanaged(u8){};
  45. defer record_buffer.deinit(c_allocator);
  46. while (true) {
  47. const n = try file.read(&read_buffer);
  48. if (n == 0) break;
  49. var start: usize = 0;
  50. while (std.mem.indexOfScalarPos(u8, read_buffer[0..n], start, '\r')) |end| {
  51. try record_buffer.appendSlice(c_allocator, read_buffer[start..end]);
  52. restoreRecord(record_buffer.items, record_count);
  53. record_buffer.clearRetainingCapacity();
  54. start = end + 1;
  55. }
  56. if (start < n) {
  57. try record_buffer.appendSlice(c_allocator, read_buffer[start..n]);
  58. }
  59. }
  60. if (record_buffer.items.len > 0) {
  61. restoreRecord(record_buffer.items, record_count);
  62. }
  63. }
  64. fn restoreRecord(record: []const u8, record_count: *usize) void {
  65. if (record.len == 0) {
  66. return;
  67. }
  68. record_count.* += 1;
  69. std.debug.print("Restoring record N:{d}\r", .{record_count.*});
  70. const first_pipe = std.mem.indexOfScalar(u8, record, '|') orelse return;
  71. const opcode = record[0..first_pipe];
  72. const remaining = record[first_pipe + 1 ..];
  73. const second_pipe = std.mem.indexOfScalar(u8, remaining, '|') orelse return;
  74. const key = remaining[0..second_pipe];
  75. const value = remaining[second_pipe + 1 ..];
  76. const opcodeEnum = std.meta.stringToEnum(OPCode, opcode) orelse return;
  77. switch (opcodeEnum) {
  78. .W => _ = storage.restore(key, value),
  79. .D => _ = storage.restoreDelete(key),
  80. }
  81. }
  82. pub fn setInstantWal(enabled: bool) void {
  83. instant_wal = enabled;
  84. }
  85. pub fn persist(opcode: u8, key: []const u8, value: []const u8) void {
  86. const record_len = 1 + 1 + key.len + 1 + value.len + 1;
  87. mutex.lock();
  88. defer mutex.unlock();
  89. if (buffer_position + record_len > FLUSH_THRESHOLD) {
  90. flushBuffer() catch |err| {
  91. std.debug.print("Failed to flush buffer: {any}\n", .{err});
  92. return;
  93. };
  94. }
  95. if (record_len > BUFFER_SIZE) {
  96. var temp_buffer: [BUFFER_SIZE]u8 = undefined;
  97. var pos: usize = 0;
  98. temp_buffer[pos] = opcode;
  99. pos += 1;
  100. temp_buffer[pos] = '|';
  101. pos += 1;
  102. @memcpy(temp_buffer[pos .. pos + key.len], key);
  103. pos += key.len;
  104. temp_buffer[pos] = '|';
  105. pos += 1;
  106. @memcpy(temp_buffer[pos .. pos + value.len], value);
  107. pos += value.len;
  108. temp_buffer[pos] = '\r';
  109. pos += 1;
  110. _ = storage_file.?.write(temp_buffer[0..pos]) catch |err| {
  111. std.debug.print("Failed to write large record to storage file: {any}\n", .{err});
  112. return;
  113. };
  114. return;
  115. }
  116. if (buffer_position + record_len > BUFFER_SIZE) {
  117. flushBuffer() catch |err| {
  118. std.debug.print("Failed to flush buffer: {any}\n", .{err});
  119. return;
  120. };
  121. }
  122. var pos = buffer_position;
  123. write_buffer[pos] = opcode;
  124. pos += 1;
  125. write_buffer[pos] = '|';
  126. pos += 1;
  127. @memcpy(write_buffer[pos .. pos + key.len], key);
  128. pos += key.len;
  129. write_buffer[pos] = '|';
  130. pos += 1;
  131. @memcpy(write_buffer[pos .. pos + value.len], value);
  132. pos += value.len;
  133. write_buffer[pos] = '\r';
  134. pos += 1;
  135. buffer_position = pos;
  136. if (instant_wal) {
  137. flushBuffer() catch |err| {
  138. std.debug.print("Failed to flush buffer in instant WAL mode: {any}\n", .{err});
  139. };
  140. }
  141. }
  142. pub fn flush() !void {
  143. mutex.lock();
  144. defer mutex.unlock();
  145. try flushBuffer();
  146. }
  147. fn flushBuffer() !void {
  148. if (buffer_position == 0) {
  149. return;
  150. }
  151. _ = try storage_file.?.write(write_buffer[0..buffer_position]);
  152. buffer_position = 0;
  153. }