persistence.zig 4.6 KB

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