An HTTP/1.1 server for zig
0

Configure Feed

Select the types of activity you want to include in your feed.

integrate websocket 0.16

+352 -401
+3 -3
build.zig
··· 6 6 7 7 const dep_opts = .{ .target = target, .optimize = optimize }; 8 8 const metrics_module = b.dependency("metrics", dep_opts).module("metrics"); 9 - // const websocket_module = b.dependency("websocket", dep_opts).module("websocket"); 9 + const websocket_module = b.dependency("websocket", dep_opts).module("websocket"); 10 10 11 11 const enable_tsan = b.option(bool, "tsan", "Enable ThreadSanitizer"); 12 12 ··· 18 18 .sanitize_thread = enable_tsan, 19 19 .imports = &.{ 20 20 .{ .name = "metrics", .module = metrics_module }, 21 - // .{ .name = "websocket", .module = websocket_module }, 21 + .{ .name = "websocket", .module = websocket_module }, 22 22 }, 23 23 }); 24 24 { ··· 67 67 .{ .file = "examples/05_request_takeover.zig", .name = "example_5" }, 68 68 .{ .file = "examples/06_middleware.zig", .name = "example_6" }, 69 69 .{ .file = "examples/07_advanced_routing.zig", .name = "example_7" }, 70 + .{ .file = "examples/08_websocket.zig", .name = "example_8" }, 70 71 // @ZIG016 71 - // .{ .file = "examples/08_websocket.zig", .name = "example_8" }, 72 72 // .{ .file = "examples/09_shutdown.zig", .name = "example_9", .libc = true }, 73 73 .{ .file = "examples/10_file_upload.zig", .name = "example_10" }, 74 74 .{ .file = "examples/11_html_streaming.zig", .name = "example_11" },
+4 -4
build.zig.zon
··· 8 8 .url = "https://github.com/karlseguin/metrics.zig/archive/6de29b83a750a06c438d268543e0e3c3c1b309da.tar.gz", 9 9 .hash = "metrics-0.0.0-W7G4eIegAQD4XxA9Co7Atbw59u_2zvxYf406AZuoAHPM", 10 10 }, 11 - // .websocket = .{ 12 - // .url = "https://github.com/karlseguin/websocket.zig/archive/4deaaef2b4475a63f19c5e2f43e38fd55464b118.tar.gz", 13 - // .hash = "websocket-0.1.0-ZPISdZJxAwAt6Ys_JpoHQQV3NpWCof_N9Jg-Ul2g7OoV", 14 - // }, 11 + .websocket = .{ 12 + .url = "https://github.com/karlseguin/websocket.zig/archive/3be6210f53297fb4b458d88562047ff4d69629a3.tar.gz", 13 + .hash = "websocket-0.1.0-ZPISdUU6BAAPe0iZ_JHMVAXaBlz327xZRBrRY06-Vw5h", 14 + }, 15 15 // .websocket = .{ .path = "../websocket.zig" }, 16 16 }, 17 17 }
+228 -248
src/httpz.zig
··· 2 2 const builtin = @import("builtin"); 3 3 4 4 pub const testing = @import("testing.zig"); 5 - // @ZIG016 6 - // pub const websocket = @import("websocket"); 5 + pub const websocket = @import("websocket"); 7 6 8 7 const posix = @import("posix.zig"); 9 8 pub const routing = @import("router.zig"); ··· 266 265 _listener: ?posix.fd_t, 267 266 _max_request_per_connection: usize, 268 267 _middlewares: []const Middleware(H), 269 - // ZIG016 270 - // _websocket_state: websocket.server.WorkerState, 268 + _websocket_state: websocket.server.WorkerState, 271 269 _middleware_registry: std.SinglyLinkedList, 272 270 273 271 const Self = @This(); ··· 288 286 // do not pass arena.allocator to WorkerState, it needs to be able to 289 287 // allocate and free at will. 290 288 291 - // @ZIG016 292 - // const ws_config = config.websocket; 293 - // var websocket_state = try websocket.server.WorkerState.init(allocator, .{ 294 - // .max_message_size = ws_config.max_message_size, 295 - // .buffers = .{ 296 - // .small_size = if (has_websocket) ws_config.small_buffer_size else 0, 297 - // .small_pool = if (has_websocket) ws_config.small_buffer_pool else 0, 298 - // .large_size = if (has_websocket) ws_config.large_buffer_size else 0, 299 - // .large_pool = if (has_websocket) ws_config.large_buffer_pool else 0, 300 - // }, 301 - // // disable handshake memory allocation since httpz is handling 302 - // // the handshake request directly 303 - // .handshake = .{ 304 - // .count = 0, 305 - // .max_size = 0, 306 - // .max_headers = 0, 307 - // }, 308 - // .compression = if (ws_config.compression) .{ 309 - // .write_threshold = ws_config.compression_write_treshold, 310 - // .retain_write_buffer = ws_config.compression_retain_writer, 311 - // } else null, 312 - // }); 313 - // errdefer websocket_state.deinit(); 289 + const ws_config = config.websocket; 290 + var websocket_state = try websocket.server.WorkerState.init(io, allocator, .{ 291 + .max_message_size = ws_config.max_message_size, 292 + .buffers = .{ 293 + .small_size = if (has_websocket) ws_config.small_buffer_size else 0, 294 + .small_pool = if (has_websocket) ws_config.small_buffer_pool else 0, 295 + .large_size = if (has_websocket) ws_config.large_buffer_size else 0, 296 + .large_pool = if (has_websocket) ws_config.large_buffer_pool else 0, 297 + }, 298 + // disable handshake memory allocation since httpz is handling 299 + // the handshake request directly 300 + .handshake = .{ 301 + .count = 0, 302 + .max_size = 0, 303 + .max_headers = 0, 304 + }, 305 + .compression = if (ws_config.compression) .{ 306 + .write_threshold = ws_config.compression_write_treshold, 307 + .retain_write_buffer = ws_config.compression_retain_writer, 308 + } else null, 309 + }); 310 + errdefer websocket_state.deinit(); 314 311 315 312 const workers = try arena.allocator().alloc(Worker, config.workerCount()); 316 313 ··· 326 323 ._listener = null, 327 324 ._middlewares = &.{}, 328 325 ._middleware_registry = .{}, 329 - // ._websocket_state = websocket_state, 326 + ._websocket_state = websocket_state, 330 327 ._router = try Router(H, ActionArg).init(arena.allocator(), default_dispatcher, handler), 331 328 ._max_request_per_connection = config.timeout.request_count orelse MAX_REQUEST_COUNT, 332 329 }; 333 330 } 334 331 335 332 pub fn deinit(self: *Self) void { 336 - // @ZIG016 337 - // self._websocket_state.deinit(); 333 + self._websocket_state.deinit(); 338 334 339 335 var node = self._middleware_registry.first; 340 336 while (node) |n| { ··· 662 658 }; 663 659 } 664 660 665 - // @ZIG016 666 - // pub fn upgradeWebsocket(comptime H: type, req: *Request, res: *Response, ctx: anytype) !bool { 667 - // const upgrade = req.header("upgrade") orelse return false; 668 - // if (std.ascii.eqlIgnoreCase(upgrade, "websocket") == false) { 669 - // return false; 670 - // } 661 + pub fn upgradeWebsocket(comptime H: type, req: *Request, res: *Response, ctx: anytype) !bool { 662 + const upgrade = req.header("upgrade") orelse return false; 663 + if (std.ascii.eqlIgnoreCase(upgrade, "websocket") == false) { 664 + return false; 665 + } 671 666 672 - // const version = req.header("sec-websocket-version") orelse return false; 673 - // if (std.ascii.eqlIgnoreCase(version, "13") == false) { 674 - // return false; 675 - // } 667 + const version = req.header("sec-websocket-version") orelse return false; 668 + if (std.ascii.eqlIgnoreCase(version, "13") == false) { 669 + return false; 670 + } 676 671 677 - // // firefox will send multiple values for this header 678 - // const connection = req.header("connection") orelse return false; 679 - // if (std.ascii.indexOfIgnoreCase(connection, "upgrade") == null) { 680 - // return false; 681 - // } 672 + // firefox will send multiple values for this header 673 + const connection = req.header("connection") orelse return false; 674 + if (std.ascii.indexOfIgnoreCase(connection, "upgrade") == null) { 675 + return false; 676 + } 682 677 683 - // const key = req.header("sec-websocket-key") orelse return false; 678 + const key = req.header("sec-websocket-key") orelse return false; 684 679 685 - // const http_conn = res.conn; 686 - // const ws_worker: *websocket.server.Worker(H) = @ptrCast(@alignCast(http_conn.ws_worker)); 680 + const http_conn = res.conn; 681 + const ws_worker: *websocket.server.Worker(H) = @ptrCast(@alignCast(http_conn.ws_worker)); 687 682 688 - // var hc = try ws_worker.createConn(http_conn.stream.handle, http_conn.address, worker.timestamp(0)); 689 - // errdefer ws_worker.cleanupConn(hc); 683 + var hc = try ws_worker.createConn(http_conn.stream.socket.handle, http_conn.address, worker.timestamp(http_conn.io)); 684 + errdefer ws_worker.cleanupConn(hc); 690 685 691 - // hc.handler = try H.init(&hc.conn, ctx); 686 + hc.handler = try H.init(&hc.conn, ctx); 692 687 693 - // var compression = false; 694 - // if (ws_worker.canCompress()) { 695 - // if (req.header("sec-websocket-extensions")) |ext| { 696 - // compression = try websocket.Handshake.parseExtension(ext) != null; 697 - // } 698 - // } 688 + var compression = false; 689 + if (ws_worker.canCompress()) { 690 + if (req.header("sec-websocket-extensions")) |ext| { 691 + compression = try websocket.Handshake.parseExtension(ext) != null; 692 + } 693 + } 699 694 700 - // var reply_buf: [512]u8 = undefined; 701 - // const reply = try websocket.Handshake.createReply(key, null, compression, &reply_buf); 702 - // var writer = http_conn.stream.writer(&.{}); 703 - // const w = &writer.interface; 704 - // try w.writeAll(reply); 705 - // try w.flush(); 695 + var reply_buf: [512]u8 = undefined; 696 + const reply = try websocket.Handshake.createReply(key, null, compression, &reply_buf); 697 + var writer = http_conn.stream.writer(http_conn.io, &.{}); 698 + const w = &writer.interface; 699 + try w.writeAll(reply); 700 + try w.flush(); 706 701 707 - // if (comptime std.meta.hasFn(H, "afterInit")) { 708 - // const params = @typeInfo(@TypeOf(H.afterInit)).@"fn".params; 709 - // try if (comptime params.len == 1) hc.handler.?.afterInit() else hc.handler.?.afterInit(ctx); 710 - // } 711 - // try ws_worker.setupConnection(hc); 712 - // res.written = true; 713 - // http_conn.handover = .{ .websocket = hc }; 714 - // return true; 715 - // } 702 + if (comptime std.meta.hasFn(H, "afterInit")) { 703 + const params = @typeInfo(@TypeOf(H.afterInit)).@"fn".params; 704 + try if (comptime params.len == 1) hc.handler.?.afterInit() else hc.handler.?.afterInit(ctx); 705 + } 706 + try ws_worker.setupConnection(hc); 707 + res.written = true; 708 + http_conn.handover = .{ .websocket = hc }; 709 + return true; 710 + } 716 711 717 712 // std.heap.StackFallbackAllocator is very specific. It's really _stack_ as it 718 713 // requires a comptime size. Also, it uses non-public calls from the FixedBufferAllocator. ··· 797 792 var dispatch_action_context_server: Server(*TestHandlerDispatchContext) = undefined; 798 793 var reuse_server: Server(void) = undefined; 799 794 var handle_server: Server(TestHandlerHandle) = undefined; 800 - // @ZIG016 801 - // var websocket_server: Server(TestWebsocketHandler) = undefined; 795 + var websocket_server: Server(TestWebsocketHandler) = undefined; 802 796 var cors_wildcard_server: Server(void) = undefined; 803 797 var cors_single_server: Server(void) = undefined; 804 798 var cors_multiple_server: Server(void) = undefined; ··· 904 898 test_server_threads[5] = try handle_server.listenInNewThread(); 905 899 } 906 900 907 - // @ZIG016 908 - // { 909 - // websocket_server = try Server(TestWebsocketHandler).init(ga, .{ .address = .localhost(5998) }, TestWebsocketHandler{}); 910 - // var router = try websocket_server.router(.{}); 911 - // router.get("/ws", TestWebsocketHandler.upgrade, .{}); 912 - // test_server_threads[6] = try websocket_server.listenInNewThread(); 913 - // } 914 - test_server_threads[6] = try Thread.spawn(.{}, struct { 915 - fn dummy() void {} 916 - }.dummy, .{}); 901 + { 902 + websocket_server = try Server(TestWebsocketHandler).init(t.io, ga, .{ .address = .localhost(5998) }, TestWebsocketHandler{}); 903 + var router = try websocket_server.router(.{}); 904 + router.get("/ws", TestWebsocketHandler.upgrade, .{}); 905 + test_server_threads[6] = try websocket_server.listenInNewThread(); 906 + } 917 907 918 908 { 919 909 cors_wildcard_server = try Server(void).init(t.io, ga, .{ .address = .localhost(5999) }, {}); ··· 966 956 dispatch_action_context_server.stop(); 967 957 reuse_server.stop(); 968 958 handle_server.stop(); 969 - // @ZIG016 970 - // websocket_server.stop(); 959 + websocket_server.stop(); 971 960 cors_wildcard_server.stop(); 972 961 cors_single_server.stop(); 973 962 cors_multiple_server.stop(); ··· 982 971 dispatch_action_context_server.deinit(); 983 972 reuse_server.deinit(); 984 973 handle_server.deinit(); 985 - // @ZIG016 986 - // websocket_server.deinit(); 974 + websocket_server.deinit(); 987 975 cors_wildcard_server.deinit(); 988 976 cors_single_server.deinit(); 989 977 cors_multiple_server.deinit(); ··· 1888 1876 } 1889 1877 } 1890 1878 1891 - // @ZIG016 1892 - // test "websocket: invalid request" { 1893 - // const stream = testStream(5998); 1894 - // defer stream.close(t.io); 1895 - // var writer = stream.writer(t.io, &.{}); 1896 - // const w = &writer.interface; 1897 - // try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1898 - // try w.flush(); 1879 + test "websocket: invalid request" { 1880 + const stream = testStream(5998); 1881 + defer stream.close(t.io); 1882 + var writer = stream.writer(t.io, &.{}); 1883 + const w = &writer.interface; 1884 + try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1885 + try w.flush(); 1899 1886 1900 - // var res = testReadParsed(stream); 1901 - // defer res.deinit(); 1902 - // try t.expectString("invalid websocket", res.body); 1903 - // } 1887 + var res = testReadParsed(stream); 1888 + defer res.deinit(); 1889 + try t.expectString("invalid websocket", res.body); 1890 + } 1904 1891 1905 - // test "websocket: upgrade" { 1906 - // const stream = testStream(5998); 1907 - // defer stream.close(t.io); 1908 - // var writer = stream.writer(t.io, &.{}); 1909 - // const w = &writer.interface; 1910 - // try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n"); 1911 - // try w.writeAll("upgrade: WEBsocket\r\n"); 1912 - // try w.writeAll("Sec-Websocket-verSIon: 13\r\n"); 1913 - // try w.writeAll("ConnectioN: abc,upgrade,123\r\n"); 1914 - // try w.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n"); 1915 - // try w.flush(); 1892 + test "websocket: upgrade" { 1893 + const stream = testStream(5998); 1894 + defer stream.close(t.io); 1895 + var writer = stream.writer(t.io, &.{}); 1896 + const w = &writer.interface; 1897 + try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n"); 1898 + try w.writeAll("upgrade: WEBsocket\r\n"); 1899 + try w.writeAll("Sec-Websocket-verSIon: 13\r\n"); 1900 + try w.writeAll("ConnectioN: abc,upgrade,123\r\n"); 1901 + try w.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n"); 1902 + try w.flush(); 1916 1903 1917 - // var res = testReadHeader(stream); 1918 - // defer res.deinit(); 1919 - // try t.expectEqual(101, res.status); 1920 - // try t.expectString("websocket", res.headers.get("Upgrade").?); 1921 - // try t.expectString("upgrade", res.headers.get("Connection").?); 1922 - // try t.expectString("55eM2SNGu+68v5XXrr982mhPFkU=", res.headers.get("Sec-Websocket-Accept").?); 1904 + var res = testReadHeader(stream); 1905 + defer res.deinit(); 1906 + try t.expectEqual(101, res.status); 1907 + try t.expectString("websocket", res.headers.get("Upgrade").?); 1908 + try t.expectString("upgrade", res.headers.get("Connection").?); 1909 + try t.expectString("55eM2SNGu+68v5XXrr982mhPFkU=", res.headers.get("Sec-Websocket-Accept").?); 1923 1910 1924 - // try w.writeAll(&websocket.frameText("over 9000!")); 1911 + try w.writeAll(&websocket.frameText("over 9000!")); 1925 1912 1926 - // // https://github.com/karlseguin/http.zig/pull/188 1927 - // try w.flush(); 1928 - // std.Thread.sleep(std.time.ns_per_ms * 5); 1913 + // https://github.com/karlseguin/http.zig/pull/188 1914 + try w.flush(); 1915 + try t.io.sleep(.fromMilliseconds(5), .awake); 1929 1916 1930 - // try w.writeAll(&websocket.frameText("close")); 1931 - // try w.flush(); 1917 + try w.writeAll(&websocket.frameText("close")); 1918 + try w.flush(); 1932 1919 1933 - // var pos: usize = 0; 1934 - // var buf: [100]u8 = undefined; 1935 - // var wait_count: usize = 0; 1936 - // var reader = stream.reader(&.{}); 1937 - // const r = reader.interface(); 1938 - // while (pos < 16) { 1939 - // const n = r.readSliceShort(buf[pos..]) catch |err| 1940 - // switch (err) { 1941 - // error.ReadFailed => { 1942 - // if (reader.getError()) |e| { 1943 - // switch (e) { 1944 - // error.WouldBlock => { 1945 - // if (wait_count == 100) { 1946 - // break; 1947 - // } 1948 - // wait_count += 1; 1949 - // std.Thread.sleep(std.time.ns_per_ms); 1950 - // continue; 1951 - // }, 1952 - // else => {}, 1953 - // } 1954 - // } 1955 - // return err; 1956 - // }, 1957 - // }; 1920 + var pos: usize = 0; 1921 + var buf: [100]u8 = undefined; 1922 + var wait_count: usize = 0; 1923 + 1924 + while (pos < 16) { 1925 + const n = posix.read(stream.socket.handle, buf[pos..]) catch |err| { 1926 + switch (err) { 1927 + error.WouldBlock => { 1928 + if (wait_count == 100) { 1929 + break; 1930 + } 1931 + wait_count += 1; 1932 + try t.io.sleep(.fromMilliseconds(1), .awake); 1933 + continue; 1934 + }, 1935 + else => return err, 1936 + } 1937 + }; 1938 + if (n == 0) { 1939 + break; 1940 + } 1941 + pos += n; 1942 + } 1943 + try t.expectEqual(16, pos); 1944 + try t.expectEqual(129, buf[0]); 1945 + try t.expectEqual(10, buf[1]); 1946 + try t.expectString("over 9000!", buf[2..12]); 1947 + try t.expectString(&.{ 136, 2, 3, 232 }, buf[12..16]); 1948 + } 1958 1949 1959 - // if (n == 0) { 1960 - // break; 1961 - // } 1962 - // pos += n; 1963 - // } 1964 - // try t.expectEqual(16, pos); 1965 - // try t.expectEqual(129, buf[0]); 1966 - // try t.expectEqual(10, buf[1]); 1967 - // try t.expectString("over 9000!", buf[2..12]); 1968 - // try t.expectString(&.{ 136, 2, 3, 232 }, buf[12..16]); 1969 - // } 1950 + // Stress test: multiple concurrent websocket clients sending many messages each. 1951 + // Run repeatedly (e.g. zig build test -Dtest-filter="websocket: stress" or run 50x) 1952 + // to verify no race in reader.done() / allocator free when using thread pool. 1953 + test "websocket: stress" { 1954 + if (force_blocking) return; // non-blocking mode only (thread pool) 1955 + const num_clients = 8; 1956 + const messages_per_client = 150; 1970 1957 1971 - // // Stress test: multiple concurrent websocket clients sending many messages each. 1972 - // // Run repeatedly (e.g. zig build test -Dtest-filter="websocket: stress" or run 50x) 1973 - // // to verify no race in reader.done() / allocator free when using thread pool. 1974 - // test "websocket: stress" { 1975 - // if (force_blocking) return; // non-blocking mode only (thread pool) 1976 - // const num_clients = 8; 1977 - // const messages_per_client = 150; 1958 + // When run with -Dtest-filter="websocket: stress", tests:beforeAll may not run, 1959 + // so nothing is listening on 5998. Wait for port and start our own server if needed. 1960 + var stress_server: ?Server(TestWebsocketHandler) = null; 1961 + var stress_listen_thread: ?Thread = null; 1962 + testing.waitForPort(5998) catch { 1963 + stress_server = try Server(TestWebsocketHandler).init(t.io, t.allocator, .{ .address = .localhost(5998) }, TestWebsocketHandler{}); 1964 + var router = try stress_server.?.router(.{}); 1965 + router.get("/ws", TestWebsocketHandler.upgrade, .{}); 1966 + stress_listen_thread = try stress_server.?.listenInNewThread(); 1967 + try testing.waitForPort(5998); 1968 + }; 1969 + defer if (stress_server) |*srv| { 1970 + srv.stop(); 1971 + if (stress_listen_thread) |thrd| thrd.join(); 1972 + srv.deinit(); 1973 + }; 1978 1974 1979 - // // When run with -Dtest-filter="websocket: stress", tests:beforeAll may not run, 1980 - // // so nothing is listening on 5998. Wait for port and start our own server if needed. 1981 - // var stress_server: ?Server(TestWebsocketHandler) = null; 1982 - // var stress_listen_thread: ?Thread = null; 1983 - // testing.waitForPort(5998) catch { 1984 - // stress_server = try Server(TestWebsocketHandler).init(t.allocator, .{ .address = .localhost(5998) }, TestWebsocketHandler{}); 1985 - // var router = try stress_server.?.router(.{}); 1986 - // router.get("/ws", TestWebsocketHandler.upgrade, .{}); 1987 - // stress_listen_thread = try stress_server.?.listenInNewThread(); 1988 - // try testing.waitForPort(5998); 1989 - // }; 1990 - // defer if (stress_server) |*srv| { 1991 - // srv.stop(); 1992 - // if (stress_listen_thread) |thrd| thrd.join(); 1993 - // srv.deinit(); 1994 - // }; 1975 + var threads: [num_clients]Thread = undefined; 1976 + for (0..num_clients) |i| { 1977 + threads[i] = Thread.spawn(.{}, struct { 1978 + fn run(_: usize) void { 1979 + const stream = testStream(5998); 1980 + defer stream.close(t.io); 1995 1981 1996 - // var threads: [num_clients]Thread = undefined; 1997 - // for (0..num_clients) |i| { 1998 - // threads[i] = Thread.spawn(.{}, struct { 1999 - // fn run(_: usize) void { 2000 - // const stream = testStream(5998); 2001 - // defer stream.close(); 1982 + var writer = stream.writer(t.io, &.{}); 1983 + const w = &writer.interface; 1984 + w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n") catch return; 1985 + w.writeAll("upgrade: WEBsocket\r\n") catch return; 1986 + w.writeAll("Sec-Websocket-verSIon: 13\r\n") catch return; 1987 + w.writeAll("ConnectioN: upgrade\r\n") catch return; 1988 + w.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n") catch return; 1989 + w.flush() catch return; 2002 1990 2003 - // var writer = stream.writer(t.io, &.{}); 2004 - // const w = &writer.interface; 2005 - // w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n") catch return; 2006 - // w.writeAll("upgrade: WEBsocket\r\n") catch return; 2007 - // w.writeAll("Sec-Websocket-verSIon: 13\r\n") catch return; 2008 - // w.writeAll("ConnectioN: upgrade\r\n") catch return; 2009 - // w.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n") catch return; 2010 - // w.flush() catch return; 1991 + var buf: [1024]u8 = undefined; 1992 + var pos: usize = 0; 2011 1993 2012 - // var buf: [1024]u8 = undefined; 2013 - // var pos: usize = 0; 2014 - // var reader = stream.reader(&.{}); 2015 - // const r = reader.interface(); 2016 - // while (!std.mem.endsWith(u8, buf[0..pos], "\r\n\r\n")) { 2017 - // if (pos >= buf.len) return; 2018 - // var vecs: [1][]u8 = .{buf[pos..]}; 2019 - // const n = r.readVec(&vecs) catch return; 2020 - // if (n == 0) return; 2021 - // pos += n; 2022 - // } 2023 - // if (pos < 12 or !std.mem.startsWith(u8, buf[0..12], "HTTP/1.1 101")) return; 1994 + while (!std.mem.endsWith(u8, buf[0..pos], "\r\n\r\n")) { 1995 + if (pos >= buf.len) { 1996 + return; 1997 + } 1998 + const n = posix.read(stream.socket.handle, buf[pos..]) catch return; 1999 + if (n == 0) { 2000 + return; 2001 + } 2002 + pos += n; 2003 + } 2004 + if (pos < 12 or !std.mem.startsWith(u8, buf[0..12], "HTTP/1.1 101")) return; 2024 2005 2025 - // for (0..messages_per_client) |_| { 2026 - // const frame = websocket.frameText("stress"); 2027 - // w.writeAll(&frame) catch return; 2028 - // } 2029 - // w.writeAll(&websocket.frameText("close")) catch return; 2030 - // w.flush() catch return; 2031 - // } 2032 - // }.run, .{i}) catch return; 2033 - // } 2034 - // for (&threads) |*th| th.join(); 2035 - // } 2006 + for (0..messages_per_client) |_| { 2007 + const frame = websocket.frameText("stress"); 2008 + w.writeAll(&frame) catch return; 2009 + } 2010 + w.writeAll(&websocket.frameText("close")) catch return; 2011 + w.flush() catch return; 2012 + } 2013 + }.run, .{i}) catch return; 2014 + } 2015 + for (&threads) |*th| th.join(); 2016 + } 2036 2017 2037 2018 test "ContentType: forX" { 2038 2019 inline for (@typeInfo(ContentType).@"enum".fields) |field| { ··· 2371 2352 } 2372 2353 }; 2373 2354 2374 - // @ZIG016 2375 - // const TestWebsocketHandler = struct { 2376 - // pub const WebsocketHandler = struct { 2377 - // ctx: u32, 2378 - // conn: *websocket.Conn, 2355 + const TestWebsocketHandler = struct { 2356 + pub const WebsocketHandler = struct { 2357 + ctx: u32, 2358 + conn: *websocket.Conn, 2379 2359 2380 - // pub fn init(conn: *websocket.Conn, ctx: u32) !WebsocketHandler { 2381 - // return .{ 2382 - // .ctx = ctx, 2383 - // .conn = conn, 2384 - // }; 2385 - // } 2360 + pub fn init(conn: *websocket.Conn, ctx: u32) !WebsocketHandler { 2361 + return .{ 2362 + .ctx = ctx, 2363 + .conn = conn, 2364 + }; 2365 + } 2386 2366 2387 - // pub fn afterInit(self: *WebsocketHandler, ctx: u32) !void { 2388 - // try t.expectEqual(self.ctx, ctx); 2389 - // } 2367 + pub fn afterInit(self: *WebsocketHandler, ctx: u32) !void { 2368 + try t.expectEqual(self.ctx, ctx); 2369 + } 2390 2370 2391 - // pub fn clientMessage(self: *WebsocketHandler, data: []const u8) !void { 2392 - // if (std.mem.eql(u8, data, "close")) { 2393 - // self.conn.close(.{}) catch {}; 2394 - // return; 2395 - // } 2396 - // try self.conn.write(data); 2397 - // } 2398 - // }; 2371 + pub fn clientMessage(self: *WebsocketHandler, data: []const u8) !void { 2372 + if (std.mem.eql(u8, data, "close")) { 2373 + self.conn.close(.{}) catch {}; 2374 + return; 2375 + } 2376 + try self.conn.write(data); 2377 + } 2378 + }; 2399 2379 2400 - // pub fn upgrade(_: TestWebsocketHandler, req: *Request, res: *Response) !void { 2401 - // if (try upgradeWebsocket(WebsocketHandler, req, res, 9001) == false) { 2402 - // res.status = 500; 2403 - // res.body = "invalid websocket"; 2404 - // } 2405 - // } 2406 - // }; 2380 + pub fn upgrade(_: TestWebsocketHandler, req: *Request, res: *Response) !void { 2381 + if (try upgradeWebsocket(WebsocketHandler, req, res, 9001) == false) { 2382 + res.status = 500; 2383 + res.body = "invalid websocket"; 2384 + } 2385 + } 2386 + }; 2407 2387 2408 2388 const TestMiddleware = struct { 2409 2389 const Config = struct {
+1 -1
src/posix.zig
··· 611 611 .MFILE => return error.ProcessFdQuotaExceeded, 612 612 .NFILE => return error.SystemFdQuotaExceeded, 613 613 .NOMEM => return error.SystemResources, 614 - else => return error.Unexpected 614 + else => return error.Unexpected, 615 615 } 616 616 } 617 617
+4 -4
src/testing.zig
··· 290 290 /// Waits until a TCP port is accepting connections (e.g. after starting a server in another thread). 291 291 /// Tries up to 100 times with 20ms sleep between attempts. 292 292 pub fn waitForPort(port: u16) !void { 293 - const address = std.net.Address.parseIp("127.0.0.1", port) catch unreachable; 293 + const address = std.Io.net.IpAddress.parse("127.0.0.1", port) catch unreachable; 294 294 for (0..100) |_| { 295 - if (std.net.tcpConnectToAddress(address)) |stream| { 296 - stream.close(); 295 + if (address.connect(t.io, .{ .mode = .stream })) |stream| { 296 + stream.close(t.io); 297 297 return; 298 298 } else |err| { 299 299 if (err != error.ConnectionRefused) return err; 300 - std.Thread.sleep(20 * std.time.ns_per_ms); 300 + try t.io.sleep(.fromMilliseconds(20), .awake); 301 301 } 302 302 } 303 303 return error.ConnectionRefused;
+108 -141
src/worker.zig
··· 4 4 const posix = @import("posix.zig"); 5 5 const httpz = @import("httpz.zig"); 6 6 const metrics = @import("metrics.zig"); 7 - // @ZIG016 8 - // const ws = @import("websocket").server; 7 + const ws = @import("websocket").server; 9 8 10 9 const Config = httpz.Config; 11 10 const Request = httpz.Request; ··· 26 25 // This is our Blocking worker. It's very different than NonBlocking and much 27 26 // simpler. (WSH is our websocket handler, and can be void) 28 27 pub fn Blocking(comptime S: type, comptime WSH: type) type { 29 - // @ZIG016 30 - _ = WSH; 31 28 return struct { 32 29 io: Io, 33 30 server: S, ··· 36 33 allocator: Allocator, 37 34 buffer_pool: *BufferPool, 38 35 http_conn_pool: HTTPConnPool, 39 - // @ZIG016 40 - // websocket: *ws.Worker(WSH), 36 + websocket: *ws.Worker(WSH), 41 37 timeout_request: ?Timeout, 42 38 timeout_keepalive: ?Timeout, 43 39 timeout_write_error: Timeout, ··· 91 87 timeout_keepalive = Timeout.init(0); 92 88 } 93 89 94 - // @ZIG016 95 - // const websocket = try allocator.create(ws.Worker(WSH)); 96 - // errdefer allocator.destroy(websocket); 97 - // websocket.* = try ws.Worker(WSH).init(allocator, &server._websocket_state); 98 - // errdefer websocket.deinit(); 90 + const websocket = try allocator.create(ws.Worker(WSH)); 91 + errdefer allocator.destroy(websocket); 92 + websocket.* = try ws.Worker(WSH).init(io, allocator, &server._websocket_state); 93 + errdefer websocket.deinit(); 99 94 100 - // @ZIG016 undefined 101 - var http_conn_pool = try HTTPConnPool.init(io, allocator, buffer_pool, undefined, 0, config); 95 + var http_conn_pool = try HTTPConnPool.init(io, allocator, buffer_pool, websocket, 0, config); 102 96 errdefer http_conn_pool.deinit(); 103 97 104 98 const retain_allocated_bytes_keepalive = config.workers.retain_allocated_bytes orelse 8192; ··· 120 114 .config = config, 121 115 .connections = .{}, 122 116 .allocator = allocator, 123 - // @ZIG016 124 - // .websocket = websocket, 117 + .websocket = websocket, 125 118 .buffer_pool = buffer_pool, 126 119 .thread_pool = thread_pool, 127 120 .http_conn_pool = http_conn_pool, ··· 136 129 pub fn deinit(self: *Self) void { 137 130 const allocator = self.allocator; 138 131 139 - // @ZIG016 140 - // self.websocket.deinit(); 132 + self.websocket.deinit(); 141 133 self.thread_pool.deinit(); 142 - // @ZIG016 143 - // allocator.destroy(self.websocket); 134 + allocator.destroy(self.websocket); 144 135 145 136 self.http_conn_pool.deinit(); 146 137 self.conn_node_pool.deinit(allocator); ··· 157 148 var address_len: posix.socklen_t = @sizeOf(posix.Address); 158 149 const socket = posix.accept(listener, &address.any, &address_len, posix.SOCK.CLOEXEC) catch |err| { 159 150 if (err == error.ConnectionAborted or err == error.SocketNotListening) { 160 - // @ZIG016 161 - // self.websocket.shutdown(); 151 + self.websocket.shutdown(); 162 152 break; 163 153 } 164 154 log.err("Failed to accept socket: {}", .{err}); ··· 182 172 } 183 173 184 174 pub fn stop(self: *const Self) void { 185 - _ = self; 186 175 // The HTTP server will stop when the http.Server shutdown the listening socket. 187 - // @ZIG016 188 - // self.websocket.shutdown(); 176 + self.websocket.shutdown(); 189 177 } 190 178 191 179 // Called in a worker thread. `thread_buf` is a thread-specific buffer that ··· 239 227 self.http_conn_pool.release(conn); 240 228 return; 241 229 }, 242 - .websocket => unreachable, 243 - // @ZIG016 244 - // .websocket => |ptr| { 245 - // const hc: *ws.HandlerConn(WSH) = @ptrCast(@alignCast(ptr)); 246 - // // impossible for this to fail in blocking mode 247 - // conn.requestDone(self.retain_allocated_bytes_keepalive, false) catch unreachable; 248 - // self.http_conn_pool.release(conn); 249 - // // blocking read loop 250 - // // will close the connection 251 - // self.handleWebSocket(hc) catch |err| { 252 - // log.err("({f} websocket connection error: {}", .{ address, err }); 253 - // }; 254 - // return; 255 - // }, 230 + .websocket => |ptr| { 231 + const hc: *ws.HandlerConn(WSH) = @ptrCast(@alignCast(ptr)); 232 + // impossible for this to fail in blocking mode 233 + conn.requestDone(self.retain_allocated_bytes_keepalive, false) catch unreachable; 234 + self.http_conn_pool.release(conn); 235 + // blocking read loop 236 + // will close the connection 237 + self.handleWebSocket(hc) catch |err| { 238 + log.err("({f} websocket connection error: {}", .{ address, err }); 239 + }; 240 + return; 241 + }, 256 242 .disown => { 257 243 // impossible for this to fail in blocking mode 258 244 conn.requestDone(self.retain_allocated_bytes_keepalive, false) catch unreachable; ··· 338 324 return conn.handover; 339 325 } 340 326 341 - // @ZIG016 342 - // fn handleWebSocket(self: *const Self, hc: *ws.HandlerConn(WSH)) !void { 343 - // posix.setsockopt(hc.socket, posix.SOL.SOCKET, posix.SO.RCVTIMEO, &std.mem.toBytes(posix.timeval{ .sec = 0, .usec = 0 })) catch |err| { 344 - // self.websocket.cleanupConn(hc); 345 - // return err; 346 - // }; 347 - // // closes the connection before returning 348 - // return self.websocket.worker.readLoop(hc); 349 - // } 327 + fn handleWebSocket(self: *const Self, hc: *ws.HandlerConn(WSH)) !void { 328 + posix.setsockopt(hc.socket, posix.SOL.SOCKET, posix.SO.RCVTIMEO, &std.mem.toBytes(posix.timeval{ .sec = 0, .usec = 0 })) catch |err| { 329 + self.websocket.cleanupConn(hc); 330 + return err; 331 + }; 332 + // closes the connection before returning 333 + return self.websocket.worker.readLoop(hc); 334 + } 350 335 }; 351 336 } 352 337 ··· 388 373 389 374 config: *const Config, 390 375 391 - // @ZIG016 392 - // websocket: *ws.Worker(WSH), 376 + websocket: *ws.Worker(WSH), 393 377 394 378 // how many bytes should we retain in a connection's arena allocator 395 379 retain_allocated_bytes: usize, ··· 465 449 const loop = try Loop.init(); 466 450 errdefer loop.deinit(); 467 451 468 - // @ZIG016 469 - // const websocket = try allocator.create(ws.Worker(WSH)); 470 - // errdefer allocator.destroy(websocket); 471 - // websocket.* = try ws.Worker(WSH).init(allocator, &server._websocket_state); 472 - // errdefer websocket.deinit(); 452 + const websocket = try allocator.create(ws.Worker(WSH)); 453 + errdefer allocator.destroy(websocket); 454 + websocket.* = try ws.Worker(WSH).init(io, allocator, &server._websocket_state); 455 + errdefer websocket.deinit(); 473 456 474 457 var buffer_pool = try initializeBufferPool(io, allocator, config); 475 458 errdefer buffer_pool.deinit(); ··· 477 460 var conn_mem_pool: std.heap.MemoryPool(Conn(WSH)) = .empty; 478 461 errdefer conn_mem_pool.deinit(allocator); 479 462 480 - // @ZIG016 undefined!! 481 - var http_conn_pool = try HTTPConnPool.init(io, allocator, buffer_pool, undefined, loop.fd, config); 463 + var http_conn_pool = try HTTPConnPool.init(io, allocator, buffer_pool, websocket, loop.fd, config); 482 464 errdefer http_conn_pool.deinit(); 483 465 484 466 const thread_pool = try ThreadPool(Self.processData).init(io, allocator, .{ ··· 500 482 .config = config, 501 483 .server = server, 502 484 .allocator = allocator, 503 - // @ZIG016 504 - // .websocket = websocket, 485 + .websocket = websocket, 505 486 .thread_pool = thread_pool, 506 487 .active_list = .{}, 507 488 .request_list = .{}, ··· 520 501 pub fn deinit(self: *Self) void { 521 502 const allocator = self.allocator; 522 503 523 - // @ZIG016 524 - // self.websocket.deinit(); 525 - // allocator.destroy(self.websocket); 504 + self.websocket.deinit(); 505 + allocator.destroy(self.websocket); 526 506 527 507 self.thread_pool.deinit(); 528 508 ··· 561 541 return; 562 542 }; 563 543 ready_sem.post(io); 564 - // @ZIG016 565 - // defer self.websocket.shutdown(); 544 + defer self.websocket.shutdown(); 566 545 567 546 var now = timestamp(io); 568 547 var last_timeout = now; ··· 638 617 self.swapList(conn, .active); 639 618 thread_pool.spawn(.{ self, now, conn }); 640 619 }, 641 - // @ZIG016 642 - // .websocket => { 643 - // if (conn.acquireProcessing() == false) { 644 - // // Connection is already being processed. We need 645 - // // to wait for the current processing to complete. 646 - // // See the processing field in Conn 647 - // continue; 648 - // } 649 - // thread_pool.spawn(.{ self, now, conn }); 650 - // }, 620 + .websocket => { 621 + if (conn.acquireProcessing() == false) { 622 + // Connection is already being processed. We need 623 + // to wait for the current processing to complete. 624 + // See the processing field in Conn 625 + continue; 626 + } 627 + thread_pool.spawn(.{ self, now, conn }); 628 + }, 651 629 }, 652 630 .shutdown => return, 653 631 } ··· 760 738 fn processSignal(self: *Self, closed_bool: *bool) void { 761 739 const io = self.io; 762 740 const loop = &self.loop; 763 - _ = loop; 764 741 var hl = &self.handover_list; 765 742 766 743 // We take the handover list, and then re-initialize it. We do this ··· 798 775 closed_bool.* = true; 799 776 self.disown(conn); 800 777 }, 801 - .websocket => unreachable, 802 - // @ZIG016 803 - // .websocket => |ptr| { 804 - // if (comptime WSH == httpz.DummyWebsocketHandler) { 805 - // std.debug.print("Your httpz handler must have a `WebsocketHandler` declaration. This must be the same type passed to `httpz.upgradeWebsocket`. Closing the connection.\n", .{}); 806 - // closed_bool.* = true; 807 - // conn.close(); 808 - // self.disown(conn); 809 - // continue; 810 - // } 778 + .websocket => |ptr| { 779 + if (comptime WSH == httpz.DummyWebsocketHandler) { 780 + std.debug.print("Your httpz handler must have a `WebsocketHandler` declaration. This must be the same type passed to `httpz.upgradeWebsocket`. Closing the connection.\n", .{}); 781 + closed_bool.* = true; 782 + conn.close(); 783 + self.disown(conn); 784 + continue; 785 + } 811 786 812 - // self.http_conn_pool.release(http_conn); 787 + self.http_conn_pool.release(http_conn); 813 788 814 - // const hc: *ws.HandlerConn(WSH) = @ptrCast(@alignCast(ptr)); 815 - // conn.protocol = .{ .websocket = hc }; 789 + const hc: *ws.HandlerConn(WSH) = @ptrCast(@alignCast(ptr)); 790 + conn.protocol = .{ .websocket = hc }; 816 791 817 - // loop.switchToOneShot(conn) catch { 818 - // metrics.internalError(); 819 - // closed_bool.* = true; 820 - // conn.close(); 821 - // self.disown(conn); 822 - // continue; 823 - // }; 824 - // }, 792 + loop.switchToOneShot(conn) catch { 793 + metrics.internalError(); 794 + closed_bool.* = true; 795 + conn.close(); 796 + self.disown(conn); 797 + continue; 798 + }; 799 + }, 825 800 .keepalive => unreachable, 826 801 } 827 802 } ··· 835 810 pub fn processData(self: *Self, now: u32, conn: *Conn(WSH), thread_buf: []u8) void { 836 811 switch (conn.protocol) { 837 812 .http => |http_conn| self.processHTTPData(now, conn, thread_buf, http_conn), 838 - // @ZIG016 839 - // .websocket => |hc| self.processWebsocketData(conn, thread_buf, hc), 813 + .websocket => |hc| self.processWebsocketData(conn, thread_buf, hc), 840 814 } 841 815 } 842 816 ··· 876 850 self.loop.signal() catch |err| log.err("failed to signal worker: {}", .{err}); 877 851 } 878 852 879 - // @ZIG016 880 - // pub fn processWebsocketData(self: *Self, conn: *Conn(WSH), thread_buf: []u8, hc: *ws.HandlerConn(WSH)) void { 881 - // defer conn.releaseProcessing(); 853 + pub fn processWebsocketData(self: *Self, conn: *Conn(WSH), thread_buf: []u8, hc: *ws.HandlerConn(WSH)) void { 854 + defer conn.releaseProcessing(); 882 855 883 - // var ws_conn = &hc.conn; 884 - // const success = self.websocket.worker.dataAvailable(hc, thread_buf); 885 - // if (success == false) { 886 - // ws_conn.close(.{ .code = 4997, .reason = "wsz" }) catch {}; 887 - // self.websocket.cleanupConn(hc); 888 - // } else if (ws_conn.isClosed()) { 889 - // self.websocket.cleanupConn(hc); 890 - // } else { 891 - // self.loop.rearmRead(conn) catch |err| { 892 - // log.debug("({f}) failed to add read event monitor: {}", .{ ws_conn.address, err }); 893 - // ws_conn.close(.{ .code = 4998, .reason = "wsz" }) catch {}; 894 - // self.websocket.cleanupConn(hc); 895 - // }; 896 - // } 897 - // } 856 + var ws_conn = &hc.conn; 857 + const success = self.websocket.worker.dataAvailable(hc, thread_buf); 858 + if (success == false) { 859 + ws_conn.close(.{ .code = 4997, .reason = "wsz" }) catch {}; 860 + self.websocket.cleanupConn(hc); 861 + } else if (ws_conn.isClosed()) { 862 + self.websocket.cleanupConn(hc); 863 + } else { 864 + self.loop.rearmRead(conn) catch |err| { 865 + log.debug("({f}) failed to add read event monitor: {}", .{ ws_conn.address, err }); 866 + ws_conn.close(.{ .code = 4998, .reason = "wsz" }) catch {}; 867 + self.websocket.cleanupConn(hc); 868 + }; 869 + } 870 + } 898 871 899 872 fn disown(self: *Self, conn: *Conn(WSH)) void { 900 873 const io = self.io; ··· 1511 1484 return struct { 1512 1485 protocol: union(enum) { 1513 1486 http: *HTTPConn, 1514 - // @ZIG016 1515 - // websocket: *ws.HandlerConn(WSH), 1487 + websocket: *ws.HandlerConn(WSH), 1516 1488 }, 1517 1489 1518 1490 // Node in a List(WSH). List is [obviously] intrusive. ··· 1540 1512 fn close(self: *Self) void { 1541 1513 switch (self.protocol) { 1542 1514 .http => |http_conn| posix.close(http_conn.stream.socket.handle), 1543 - // @ZIG016 1544 - // .websocket => |hc| hc.conn.close(.{}) catch {}, 1515 + .websocket => |hc| hc.conn.close(.{}) catch {}, 1545 1516 } 1546 1517 } 1547 1518 1548 1519 pub fn getSocket(self: Self) posix.fd_t { 1549 1520 return switch (self.protocol) { 1550 1521 .http => |hc| hc.stream.socket.handle, 1551 - // @ZIG016 1552 - // .websocket => |hc| hc.socket, 1522 + .websocket => |hc| hc.socket, 1553 1523 }; 1554 1524 } 1555 1525 ··· 1758 1728 } 1759 1729 1760 1730 pub fn writeAll(self: *HTTPConn, data: []const u8) !void { 1761 - var writer = self.stream.writer(self.io, &.{}); 1762 - try writer.interface.writeAll(data); 1763 - // ZIG016 would block 1731 + var i: usize = 0; 1732 + var blocking = false; 1764 1733 1765 - // var i: usize = 0; 1766 - // var blocking = false; 1734 + const socket = self.stream.socket.handle; 1767 1735 1768 - // while (i < data.len) { 1769 - // const remaining = 1770 - // const n = posix.system.write(socket, data[i..]) catch |err| switch (err) { 1771 - // error.WouldBlock => { 1772 - // try self.blockingMode(); 1773 - // blocking = true; 1774 - // continue; 1775 - // }, 1776 - // else => return err, 1777 - // }; 1736 + while (i < data.len) { 1737 + const n = posix.write(socket, data[i..]) catch |err| switch (err) { 1738 + error.WouldBlock => { 1739 + try self.blockingMode(); 1740 + blocking = true; 1741 + continue; 1742 + }, 1743 + else => return err, 1744 + }; 1778 1745 1779 - // // shouldn't be posssible on a correct posix implementation 1780 - // // but let's assert to make sure 1781 - // std.debug.assert(n != 0); 1782 - // i += n; 1783 - // } 1746 + // shouldn't be posssible on a correct posix implementation 1747 + // but let's assert to make sure 1748 + std.debug.assert(n != 0); 1749 + i += n; 1750 + } 1784 1751 } 1785 1752 1786 1753 pub fn writeAllIOVec(self: *HTTPConn, vec: [][]const u8) !void {
+4
test_runner.zig
··· 6 6 7 7 const BORDER = "=" ** 80; 8 8 9 + pub const std_options = std.Options{ .log_scope_levels = &[_]std.log.ScopeLevel{ 10 + .{ .scope = .websocket, .level = .warn }, 11 + } }; 12 + 9 13 // use in custom panic handler 10 14 var current_test: ?[]const u8 = null; 11 15