2
0

command.zig 4.3 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980
  1. const std = @import("std");
  2. const Engine = @import("engine.zig").Engine;
  3. const max_text_response = 1024 * 1024;
  4. pub fn execute(engine: *Engine, allocator: std.mem.Allocator, message: []const u8) ![]u8 {
  5. const clean = std.mem.trim(u8, message, "\r\n ");
  6. const split = std.mem.indexOfScalar(u8, clean, ' ');
  7. const name = if (split) |index| clean[0..index] else clean;
  8. const arguments = if (split) |index| clean[index + 1 ..] else "";
  9. if (std.mem.eql(u8, name, "read")) {
  10. const value = try engine.get(allocator, arguments) orelse return allocator.dupe(u8, "error");
  11. return value.bytes;
  12. }
  13. if (std.mem.eql(u8, name, "write")) {
  14. const separator = std.mem.indexOfScalar(u8, arguments, '|') orelse return allocator.dupe(u8, "error");
  15. _ = engine.put(arguments[0..separator], arguments[separator + 1 ..]) catch return allocator.dupe(u8, "error");
  16. return allocator.dupe(u8, "success");
  17. }
  18. if (std.mem.eql(u8, name, "delete")) {
  19. const deleted = engine.delete(arguments) catch return allocator.dupe(u8, "error");
  20. return allocator.dupe(u8, if (deleted) "success" else "error");
  21. }
  22. if (std.mem.eql(u8, name, "status")) {
  23. const status = engine.status();
  24. return std.fmt.allocPrint(allocator, "well going our operation keys={d} latest_lsn={d} checkpoint_lsn={d} file_bytes={d} groups={d} transactions={d} largest_group={d}", .{ status.live_keys, status.latest_lsn, status.checkpoint_lsn, status.file_bytes, status.commit_groups, status.committed_transactions, status.largest_commit_group });
  25. }
  26. if (std.mem.eql(u8, name, "keys")) return scanText(engine, allocator, "", arguments, false);
  27. if (std.mem.eql(u8, name, "reads")) return scanText(engine, allocator, arguments, "", true);
  28. if (std.mem.eql(u8, name, "scan")) {
  29. var fields = std.mem.splitScalar(u8, arguments, '|');
  30. const prefix = fields.next() orelse "";
  31. const cursor = fields.next() orelse "";
  32. const limit_text = fields.next() orelse "256";
  33. const mode = fields.next() orelse "keys";
  34. const limit = std.fmt.parseInt(u32, limit_text, 10) catch return allocator.dupe(u8, "error");
  35. return scanTextLimit(engine, allocator, prefix, cursor, std.mem.eql(u8, mode, "values"), limit);
  36. }
  37. if (std.mem.eql(u8, name, "checkpoint")) {
  38. engine.checkpoint() catch return allocator.dupe(u8, "error");
  39. return allocator.dupe(u8, "success");
  40. }
  41. return allocator.dupe(u8, "error");
  42. }
  43. fn scanText(engine: *Engine, allocator: std.mem.Allocator, prefix: []const u8, cursor: []const u8, values: bool) ![]u8 {
  44. return scanTextLimit(engine, allocator, prefix, cursor, values, 256);
  45. }
  46. fn scanTextLimit(engine: *Engine, allocator: std.mem.Allocator, prefix: []const u8, cursor: []const u8, values: bool, limit: u32) ![]u8 {
  47. var batch = engine.scan(allocator, prefix, cursor, @min(limit, 4096), values, max_text_response) catch return allocator.dupe(u8, "error");
  48. defer batch.deinit(allocator);
  49. var output = std.ArrayListUnmanaged(u8){};
  50. errdefer output.deinit(allocator);
  51. for (batch.entries, 0..) |entry, index| {
  52. if (index != 0) try output.append(allocator, '\n');
  53. try output.appendSlice(allocator, if (values) entry.value.? else entry.key);
  54. }
  55. return output.toOwnedSlice(allocator);
  56. }
  57. test "Pizzaria point compatibility and bounded scan" {
  58. var tmp = std.testing.tmpDir(.{});
  59. defer tmp.cleanup();
  60. var path_buffer: [std.fs.max_path_bytes]u8 = undefined;
  61. const directory = try tmp.dir.realpath(".", &path_buffer);
  62. const path = try std.fmt.allocPrint(std.testing.allocator, "{s}/protocol.pkvdb", .{directory});
  63. defer std.testing.allocator.free(path);
  64. var engine = try Engine.open(std.testing.allocator, path);
  65. defer engine.close();
  66. var response = try execute(&engine, std.testing.allocator, "write p/1|one\r");
  67. try std.testing.expectEqualStrings("success", response);
  68. std.testing.allocator.free(response);
  69. response = try execute(&engine, std.testing.allocator, "read p/1\r");
  70. try std.testing.expectEqualStrings("one", response);
  71. std.testing.allocator.free(response);
  72. response = try execute(&engine, std.testing.allocator, "scan p/||1|keys\r");
  73. defer std.testing.allocator.free(response);
  74. try std.testing.expectEqualStrings("p/1", response);
  75. }