migration.zig 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146
  1. const std = @import("std");
  2. const engine_mod = @import("engine.zig");
  3. pub const Result = struct {
  4. keys: u64,
  5. records: u64,
  6. checksum: u64,
  7. };
  8. pub fn migrate(allocator: std.mem.Allocator, source_path: []const u8, destination_path: []const u8) !Result {
  9. std.fs.cwd().access(destination_path, .{}) catch |err| switch (err) {
  10. error.FileNotFound => {},
  11. else => return err,
  12. };
  13. if (std.fs.cwd().openFile(destination_path, .{})) |file| {
  14. file.close();
  15. return error.DestinationExists;
  16. } else |_| {}
  17. const source = try std.fs.cwd().openFile(source_path, .{ .mode = .read_only });
  18. defer source.close();
  19. var state = std.StringHashMap([]u8).init(allocator);
  20. defer {
  21. var iterator = state.iterator();
  22. while (iterator.next()) |entry| {
  23. allocator.free(entry.key_ptr.*);
  24. allocator.free(entry.value_ptr.*);
  25. }
  26. state.deinit();
  27. }
  28. var record = std.ArrayListUnmanaged(u8){};
  29. defer record.deinit(allocator);
  30. var buffer: [64 * 1024]u8 = undefined;
  31. var records: u64 = 0;
  32. while (true) {
  33. const amount = try source.read(&buffer);
  34. if (amount == 0) break;
  35. var start: usize = 0;
  36. while (std.mem.indexOfScalarPos(u8, buffer[0..amount], start, '\r')) |end| {
  37. try record.appendSlice(allocator, buffer[start..end]);
  38. if (record.items.len != 0) {
  39. try applyRecord(allocator, &state, record.items);
  40. records += 1;
  41. }
  42. record.clearRetainingCapacity();
  43. start = end + 1;
  44. }
  45. if (start < amount) {
  46. if (record.items.len + amount - start > 65 * 1024 * 1024) return error.RecordTooLarge;
  47. try record.appendSlice(allocator, buffer[start..amount]);
  48. }
  49. }
  50. if (record.items.len != 0) {
  51. try applyRecord(allocator, &state, record.items);
  52. records += 1;
  53. }
  54. var engine = try engine_mod.Engine.open(allocator, destination_path);
  55. defer engine.close();
  56. var operations = std.ArrayListUnmanaged(engine_mod.Operation){};
  57. defer operations.deinit(allocator);
  58. var bytes: usize = 24;
  59. var iterator = state.iterator();
  60. while (iterator.next()) |entry| {
  61. const needed = 8 + entry.key_ptr.*.len + entry.value_ptr.*.len + 7;
  62. if (operations.items.len == 65535 or bytes + needed > 60 * 1024 * 1024) {
  63. try engine.importBaseline(operations.items);
  64. operations.clearRetainingCapacity();
  65. bytes = 24;
  66. }
  67. try operations.append(allocator, .{ .opcode = .put, .key = entry.key_ptr.*, .value = entry.value_ptr.* });
  68. bytes += needed;
  69. }
  70. if (operations.items.len != 0) {
  71. try engine.importBaseline(operations.items);
  72. } else if (state.count() == 0) {
  73. try engine.importBaseline(&.{});
  74. }
  75. try engine.checkpoint();
  76. const status = engine.status();
  77. if (status.live_keys != state.count() or status.latest_lsn != 1 or status.checkpoint_lsn != 1) return error.VerificationFailed;
  78. var checksum: u64 = 14695981039346656037;
  79. iterator = state.iterator();
  80. while (iterator.next()) |entry| {
  81. const value = try engine.get(allocator, entry.key_ptr.*) orelse return error.VerificationFailed;
  82. defer allocator.free(value.bytes);
  83. if (!std.mem.eql(u8, value.bytes, entry.value_ptr.*)) return error.VerificationFailed;
  84. for (entry.key_ptr.*) |byte| {
  85. checksum ^= byte;
  86. checksum *%= 1099511628211;
  87. }
  88. for (value.bytes) |byte| {
  89. checksum ^= byte;
  90. checksum *%= 1099511628211;
  91. }
  92. }
  93. return .{ .keys = status.live_keys, .records = records, .checksum = checksum };
  94. }
  95. fn applyRecord(allocator: std.mem.Allocator, state: *std.StringHashMap([]u8), record: []const u8) !void {
  96. const first = std.mem.indexOfScalar(u8, record, '|') orelse return error.InvalidLegacyRecord;
  97. const second_relative = std.mem.indexOfScalar(u8, record[first + 1 ..], '|') orelse return error.InvalidLegacyRecord;
  98. const second = first + 1 + second_relative;
  99. const opcode = record[0..first];
  100. const key = record[first + 1 .. second];
  101. const value = record[second + 1 ..];
  102. if (std.mem.eql(u8, opcode, "W")) {
  103. if (state.getPtr(key)) |current| {
  104. const replacement = try allocator.dupe(u8, value);
  105. allocator.free(current.*);
  106. current.* = replacement;
  107. } else {
  108. const owned_key = try allocator.dupe(u8, key);
  109. errdefer allocator.free(owned_key);
  110. const owned_value = try allocator.dupe(u8, value);
  111. errdefer allocator.free(owned_value);
  112. try state.put(owned_key, owned_value);
  113. }
  114. } else if (std.mem.eql(u8, opcode, "D")) {
  115. if (state.fetchRemove(key)) |removed| {
  116. allocator.free(removed.key);
  117. allocator.free(removed.value);
  118. }
  119. } else return error.InvalidLegacyRecord;
  120. }
  121. test "legacy migration creates LSN one checkpoint and leaves source" {
  122. var tmp = std.testing.tmpDir(.{});
  123. defer tmp.cleanup();
  124. try tmp.dir.writeFile(.{ .sub_path = "old.db", .data = "W|a|one\rW|b|two\rW|a|three\rD|b|\r" });
  125. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  126. const directory = try tmp.dir.realpath(".", &path_buffer);
  127. const source = try std.fmt.allocPrint(std.testing.allocator, "{s}/old.db", .{directory});
  128. defer std.testing.allocator.free(source);
  129. const destination = try std.fmt.allocPrint(std.testing.allocator, "{s}/new.pkvdb", .{directory});
  130. defer std.testing.allocator.free(destination);
  131. const result = try migrate(std.testing.allocator, source, destination);
  132. try std.testing.expectEqual(@as(u64, 1), result.keys);
  133. const original = try tmp.dir.readFileAlloc(std.testing.allocator, "old.db", 1024);
  134. defer std.testing.allocator.free(original);
  135. try std.testing.expectEqualStrings("W|a|one\rW|b|two\rW|a|three\rD|b|\r", original);
  136. var engine = try engine_mod.Engine.open(std.testing.allocator, destination);
  137. defer engine.close();
  138. const value = (try engine.get(std.testing.allocator, "a")).?;
  139. defer std.testing.allocator.free(value.bytes);
  140. try std.testing.expectEqualStrings("three", value.bytes);
  141. try std.testing.expectEqual(@as(u64, 1), engine.status().checkpoint_lsn);
  142. }