소스 검색

no wal limit read, unix socket

Danilo Fragoso 4 달 전
부모
커밋
96e2ee1b13
3개의 변경된 파일74개의 추가작업 그리고 31개의 파일을 삭제
  1. 12 3
      main.zig
  2. 53 28
      persistence.zig
  3. 9 0
      socket.zig

+ 12 - 3
main.zig

@@ -26,6 +26,7 @@ var should_exit = std.atomic.Value(bool).init(false);
 var active_connections = std.atomic.Value(u32).init(0);
 var redis_mode = false;
 var instant_wal_mode = false;
+var unix_mode = false;
 
 fn handleSignal(sig: c_int) callconv(.c) void {
     _ = sig;
@@ -40,6 +41,8 @@ pub fn main() !void {
     while (args.next()) |arg| {
         if (std.mem.eql(u8, arg, "-redis")) {
             redis_mode = true;
+        } else if (std.mem.eql(u8, arg, "-unix")) {
+            unix_mode = true;
         } else if (std.mem.eql(u8, arg, "-iwal")) {
             instant_wal_mode = true;
         } else if (arg.len > 6 and std.mem.eql(u8, arg[0..6], "-port=")) {
@@ -75,10 +78,16 @@ pub fn main() !void {
         std.debug.print("\nInstant WAL mode enabled\n", .{});
     }
 
-    const listener = try socket.init(PORT);
+    const unix_path = ".pizzakv.sock";
+    const listener = if (unix_mode) try socket.initUnix(unix_path) else try socket.init(PORT);
     defer posix.close(listener);
+    defer if (unix_mode) posix.unlink(unix_path) catch {};
 
-    std.debug.print("\n2025 pizzakv! TCP Listening on port {any}\n<danilo@fragoso.dev>\n---------\n", .{PORT});
+    if (unix_mode) {
+        std.debug.print("\n2025 pizzakv! Unix socket at {s}\n<danilo@fragoso.dev>\n---------\n", .{unix_path});
+    } else {
+        std.debug.print("\n2025 pizzakv! TCP Listening on port {any}\n<danilo@fragoso.dev>\n---------\n", .{PORT});
+    }
     if (redis_mode) {
         std.debug.print("Mode: Redis Protocol (RESP)\nCommands: SET, GET, DEL\n", .{});
     } else {
@@ -121,7 +130,7 @@ pub fn main() !void {
             break;
         }
 
-        posix.setsockopt(conn, posix.IPPROTO.TCP, TCP.NODELAY, &std.mem.toBytes(@as(c_int, 1))) catch {};
+        if (!unix_mode) posix.setsockopt(conn, posix.IPPROTO.TCP, TCP.NODELAY, &std.mem.toBytes(@as(c_int, 1))) catch {};
         socket.setReadTimeout(conn, 300) catch {}; // 5 minutes
         socket.setWriteTimeout(conn, 300) catch {}; // 5 minutes
 

+ 53 - 28
persistence.zig

@@ -1,7 +1,6 @@
 const std = @import("std");
 const storage = @import("storage.zig");
 
-const MAX_PERSISTENCE_SIZE = 10_000_000 * 100;
 const BUFFER_SIZE = 1024 * 1024 * 8;
 const FLUSH_THRESHOLD = (BUFFER_SIZE * 3) / 4;
 
@@ -35,36 +34,11 @@ pub fn init() !void {
         return;
     };
 
-    const storage_data = storage_file.?.readToEndAlloc(c_allocator, MAX_PERSISTENCE_SIZE) catch |err| {
+    var record_count: usize = 0;
+    restoreFromFile(storage_file.?, &record_count) catch |err| {
         std.debug.print("Failed to read storage file: {any}\n", .{err});
         return;
     };
-    defer c_allocator.free(storage_data);
-    var records = std.mem.splitScalar(u8, storage_data, '\r');
-    var record_count: usize = 0;
-    while (records.next()) |record| {
-        if (record.len == 0) {
-            continue;
-        }
-
-        record_count += 1;
-        std.debug.print("Restoring record N:{d}\r", .{record_count});
-
-        const first_pipe = std.mem.indexOfScalar(u8, record, '|') orelse continue;
-        const opcode = record[0..first_pipe];
-
-        const remaining = record[first_pipe + 1 ..];
-        const second_pipe = std.mem.indexOfScalar(u8, remaining, '|') orelse continue;
-        const key = remaining[0..second_pipe];
-        const value = remaining[second_pipe + 1 ..];
-
-        const opcodeEnum = std.meta.stringToEnum(OPCode, opcode) orelse continue;
-
-        switch (opcodeEnum) {
-            .W => _ = storage.restore(key, value),
-            .D => _ = storage.restoreDelete(key),
-        }
-    }
     std.debug.print("Restored {d} records from persistence", .{record_count});
 
     storage_file.?.close();
@@ -77,6 +51,57 @@ pub fn init() !void {
     return;
 }
 
+fn restoreFromFile(file: std.fs.File, record_count: *usize) !void {
+    var read_buffer: [BUFFER_SIZE]u8 = undefined;
+    var record_buffer = std.ArrayListUnmanaged(u8){};
+    defer record_buffer.deinit(c_allocator);
+
+    while (true) {
+        const n = try file.read(&read_buffer);
+        if (n == 0) break;
+
+        var start: usize = 0;
+        while (std.mem.indexOfScalarPos(u8, read_buffer[0..n], start, '\r')) |end| {
+            try record_buffer.appendSlice(c_allocator, read_buffer[start..end]);
+            restoreRecord(record_buffer.items, record_count);
+            record_buffer.clearRetainingCapacity();
+            start = end + 1;
+        }
+
+        if (start < n) {
+            try record_buffer.appendSlice(c_allocator, read_buffer[start..n]);
+        }
+    }
+
+    if (record_buffer.items.len > 0) {
+        restoreRecord(record_buffer.items, record_count);
+    }
+}
+
+fn restoreRecord(record: []const u8, record_count: *usize) void {
+    if (record.len == 0) {
+        return;
+    }
+
+    record_count.* += 1;
+    std.debug.print("Restoring record N:{d}\r", .{record_count.*});
+
+    const first_pipe = std.mem.indexOfScalar(u8, record, '|') orelse return;
+    const opcode = record[0..first_pipe];
+
+    const remaining = record[first_pipe + 1 ..];
+    const second_pipe = std.mem.indexOfScalar(u8, remaining, '|') orelse return;
+    const key = remaining[0..second_pipe];
+    const value = remaining[second_pipe + 1 ..];
+
+    const opcodeEnum = std.meta.stringToEnum(OPCode, opcode) orelse return;
+
+    switch (opcodeEnum) {
+        .W => _ = storage.restore(key, value),
+        .D => _ = storage.restoreDelete(key),
+    }
+}
+
 pub fn setInstantWal(enabled: bool) void {
     instant_wal = enabled;
 }

+ 9 - 0
socket.zig

@@ -35,6 +35,15 @@ pub fn init(port: u16) !posix.socket_t {
     return listener;
 }
 
+pub fn initUnix(path: []const u8) !posix.socket_t {
+    posix.unlink(path) catch {};
+    const address = try net.Address.initUnix(path);
+    const listener = try posix.socket(posix.AF.UNIX, posix.SOCK.STREAM, 0);
+    try posix.bind(listener, &address.any, address.getOsSockLen());
+    try posix.listen(listener, 1024);
+    return listener;
+}
+
 pub fn readUntilCR(conn: posix.socket_t, buf: []u8) !usize {
     var total: usize = 0;