main.zig 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117
  1. const std = @import("std");
  2. const net = std.net;
  3. const posix = std.posix;
  4. const fmt = std.fmt;
  5. const socket = @import("socket.zig");
  6. const command = @import("command.zig");
  7. const storage = @import("storage.zig");
  8. const persistence = @import("persistence.zig");
  9. const PORT = 8085;
  10. var should_exit = std.atomic.Value(bool).init(false);
  11. fn handleSignal(sig: c_int) callconv(.c) void {
  12. _ = sig;
  13. should_exit.store(true, .seq_cst);
  14. }
  15. pub fn main() !void {
  16. const empty_mask = std.mem.zeroes(posix.sigset_t);
  17. const act = posix.Sigaction{
  18. .handler = .{ .handler = handleSignal },
  19. .mask = empty_mask,
  20. .flags = 0,
  21. };
  22. _ = posix.sigaction(posix.SIG.TERM, &act, null);
  23. _ = posix.sigaction(posix.SIG.INT, &act, null);
  24. const listener = try socket.init(PORT);
  25. defer posix.close(listener);
  26. std.debug.print("2025 pizzakv! TCP Listening on port {any}\n<danilo@fragoso.dev>\n---------\n", .{PORT});
  27. std.debug.print("Commands:\n\nread key\nwrite key|value\ndelete key\nkeys\nreads prefix\nstatus\n", .{});
  28. std.debug.print("---------\n", .{});
  29. storage.init();
  30. try persistence.init();
  31. while (!should_exit.load(.seq_cst)) {
  32. var poll_fds = [_]posix.pollfd{
  33. .{
  34. .fd = listener,
  35. .events = posix.POLL.IN,
  36. .revents = 0,
  37. },
  38. };
  39. const ready = posix.poll(&poll_fds, 1000) catch |err| {
  40. if (should_exit.load(.seq_cst)) break;
  41. std.debug.print("poll error: {any}\n", .{err});
  42. continue;
  43. };
  44. if (ready == 0) {
  45. continue;
  46. }
  47. if (should_exit.load(.seq_cst)) break;
  48. var client_address: net.Address = undefined;
  49. var client_address_len: posix.socklen_t = @sizeOf(net.Address);
  50. const conn = posix.accept(listener, &client_address.any, &client_address_len, 0) catch |err| {
  51. if (should_exit.load(.seq_cst)) break;
  52. std.debug.print("error accept: {any}\n", .{err});
  53. continue;
  54. };
  55. if (should_exit.load(.seq_cst)) {
  56. posix.close(conn);
  57. break;
  58. }
  59. posix.setsockopt(conn, posix.IPPROTO.TCP, posix.TCP.NODELAY, &std.mem.toBytes(@as(c_int, 1))) catch {};
  60. const thread = try std.Thread.spawn(.{}, handleConnection, .{conn});
  61. thread.detach();
  62. }
  63. std.debug.print("\nShutdown signal received...\n", .{});
  64. persistence.flush() catch |err| {
  65. std.debug.print("Failed to flush persistence: {any}\n", .{err});
  66. };
  67. }
  68. pub fn handleConnection(conn: posix.socket_t) !void {
  69. defer posix.close(conn);
  70. var requestBuffer: [1024 * 1024]u8 = undefined;
  71. while (true) {
  72. const n = socket.readUntilCR(conn, &requestBuffer) catch |err| {
  73. if (err == error.ConnectionClosed) break;
  74. return err;
  75. };
  76. if (n == 0) {
  77. break;
  78. }
  79. const cmdResponse = command.parse(requestBuffer[0..n]) orelse {
  80. socket.write(conn, "error\r") catch |err| {
  81. std.debug.print("error writing: {any}", .{err});
  82. };
  83. continue;
  84. };
  85. const terminator = "\r";
  86. const iovecs = [_]posix.iovec_const{
  87. .{ .base = cmdResponse.ptr, .len = cmdResponse.len },
  88. .{ .base = terminator.ptr, .len = 1 },
  89. };
  90. socket.writev(conn, &iovecs) catch |err| {
  91. std.debug.print("error writing: {any}", .{err});
  92. };
  93. }
  94. }