An HTTP/1.1 server for zig
0

Configure Feed

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

http.zig / src / httpz.zig
123 kB 3184 lines
1const std = @import("std"); 2const builtin = @import("builtin"); 3 4pub const testing = @import("testing.zig"); 5pub const websocket = @import("websocket"); 6 7const posix = @import("posix.zig"); 8pub const routing = @import("router.zig"); 9pub const request = @import("request.zig"); 10pub const response = @import("response.zig"); 11pub const key_value = @import("key_value.zig"); 12pub const middleware = @import("middleware/middleware.zig"); 13 14pub const Router = routing.Router; 15pub const Request = request.Request; 16pub const Response = response.Response; 17pub const Url = @import("url.zig").Url; 18pub const Config = @import("config.zig").Config; 19 20const Thread = std.Thread; 21const Io = std.Io; 22const net = std.net; 23const Allocator = std.mem.Allocator; 24const FixedBufferAllocator = std.heap.FixedBufferAllocator; 25 26const log = std.log.scoped(.httpz); 27 28const worker = @import("worker.zig"); 29const HTTPConn = worker.HTTPConn; 30 31const build = @import("build"); 32const force_blocking: bool = if (@hasDecl(build, "httpz_blocking")) build.httpz_blocking else false; 33 34const MAX_REQUEST_COUNT = std.math.maxInt(usize); 35 36pub fn writeMetrics(writer: *std.Io.Writer) !void { 37 return @import("metrics.zig").write(writer); 38} 39 40pub const Protocol = enum { 41 HTTP10, 42 HTTP11, 43}; 44 45pub const Method = enum { 46 GET, 47 HEAD, 48 POST, 49 PUT, 50 PATCH, 51 DELETE, 52 OPTIONS, 53 CONNECT, 54 OTHER, 55}; 56 57pub const ContentType = enum { 58 BINARY, 59 CSS, 60 CSV, 61 EOT, 62 EVENTS, 63 GIF, 64 GZ, 65 HTML, 66 ICO, 67 JPG, 68 JS, 69 JSON, 70 OTF, 71 PDF, 72 PNG, 73 SVG, 74 TAR, 75 TEXT, 76 TTF, 77 WASM, 78 WEBP, 79 WOFF, 80 WOFF2, 81 XML, 82 UNKNOWN, 83 84 const asUint = @import("url.zig").asUint; 85 86 pub fn forExtension(ext: []const u8) ContentType { 87 if (ext.len == 0) return .UNKNOWN; 88 const temp = if (ext[0] == '.') ext[1..] else ext; 89 if (temp.len > 5) return .UNKNOWN; 90 91 var normalized: [5]u8 = undefined; 92 for (temp, 0..) |c, i| { 93 normalized[i] = std.ascii.toLower(c); 94 } 95 96 switch (temp.len) { 97 2 => { 98 switch (@as(u16, @bitCast(normalized[0..2].*))) { 99 asUint("js") => return .JS, 100 asUint("gz") => return .GZ, 101 else => return .UNKNOWN, 102 } 103 }, 104 3 => { 105 switch (@as(u24, @bitCast(normalized[0..3].*))) { 106 asUint("css") => return .CSS, 107 asUint("csv") => return .CSV, 108 asUint("eot") => return .EOT, 109 asUint("gif") => return .GIF, 110 asUint("htm") => return .HTML, 111 asUint("ico") => return .ICO, 112 asUint("jpg") => return .JPG, 113 asUint("otf") => return .OTF, 114 asUint("pdf") => return .PDF, 115 asUint("png") => return .PNG, 116 asUint("svg") => return .SVG, 117 asUint("tar") => return .TAR, 118 asUint("ttf") => return .TTF, 119 asUint("xml") => return .XML, 120 else => return .UNKNOWN, 121 } 122 }, 123 4 => { 124 switch (@as(u32, @bitCast(normalized[0..4].*))) { 125 asUint("jpeg") => return .JPG, 126 asUint("json") => return .JSON, 127 asUint("html") => return .HTML, 128 asUint("text") => return .TEXT, 129 asUint("wasm") => return .WASM, 130 asUint("woff") => return .WOFF, 131 asUint("webp") => return .WEBP, 132 else => return .UNKNOWN, 133 } 134 }, 135 5 => { 136 switch (@as(u40, @bitCast(normalized[0..5].*))) { 137 asUint("woff2") => return .WOFF2, 138 else => return .UNKNOWN, 139 } 140 }, 141 else => return .UNKNOWN, 142 } 143 return .UNKNOWN; 144 } 145 146 pub fn forFile(file_name: []const u8) ContentType { 147 return forExtension(std.fs.path.extension(file_name)); 148 } 149}; 150 151// When we initialize our Server(handler: type) with a non-void handler, 152// the ActionContext will either be defined by the handler or it'll be the 153// handler itself. So, for this type, "ActionContext" can be either 154// the Handler or ActionContext from the Server. 155pub fn Action(comptime ActionContext: type) type { 156 if (ActionContext == void) { 157 return *const fn (*Request, *Response) anyerror!void; 158 } 159 return *const fn (ActionContext, *Request, *Response) anyerror!void; 160} 161 162pub fn Dispatcher(comptime Handler: type, comptime ActionArg: type) type { 163 if (Handler == void) { 164 return *const fn (Action(void), *Request, *Response) anyerror!void; 165 } 166 return *const fn (Handler, ActionArg, *Request, *Response) anyerror!void; 167} 168 169pub fn DispatchableAction(comptime Handler: type, comptime ActionArg: type) type { 170 return struct { 171 data: ?*const anyopaque, 172 handler: Handler, 173 action: ActionArg, 174 dispatcher: Dispatcher(Handler, ActionArg), 175 middlewares: []const Middleware(Handler) = &.{}, 176 }; 177} 178 179pub fn Middleware(comptime H: type) type { 180 return struct { 181 ptr: *anyopaque, 182 deinitFn: *const fn (ptr: *anyopaque) void, 183 executeFn: *const fn (ptr: *anyopaque, req: *Request, res: *Response, executor: *Server(H).Executor) anyerror!void, 184 185 const Self = @This(); 186 187 pub fn init(ptr: anytype) Self { 188 const T = @TypeOf(ptr); 189 const ptr_info = @typeInfo(T); 190 191 const gen = struct { 192 pub fn deinit(pointer: *anyopaque) void { 193 const self: T = @ptrCast(@alignCast(pointer)); 194 if (std.meta.hasMethod(T, "deinit")) { 195 return ptr_info.pointer.child.deinit(self); 196 } 197 } 198 199 pub fn execute(pointer: *anyopaque, req: *Request, res: *Response, executor: *Server(H).Executor) anyerror!void { 200 const self: T = @ptrCast(@alignCast(pointer)); 201 return ptr_info.pointer.child.execute(self, req, res, executor); 202 } 203 }; 204 205 return .{ 206 .ptr = ptr, 207 .deinitFn = gen.deinit, 208 .executeFn = gen.execute, 209 }; 210 } 211 212 pub fn deinit(self: Self) void { 213 self.deinitFn(self.ptr); 214 } 215 216 pub fn execute(self: Self, req: *Request, res: *Response, executor: *Server(H).Executor) !void { 217 return self.executeFn(self.ptr, req, res, executor); 218 } 219 }; 220} 221 222// When no WebsocketHandler is specified, we give it a dummy handler just to get 223// the code to compile. 224pub const DummyWebsocketHandler = struct { 225 pub fn clientMessage(_: DummyWebsocketHandler, _: []const u8) !void {} 226}; 227 228pub const MiddlewareConfig = struct { 229 arena: Allocator, 230 allocator: Allocator, 231}; 232 233pub fn Server(comptime H: type) type { 234 const Handler = switch (@typeInfo(H)) { 235 .@"struct" => H, 236 .pointer => |ptr| ptr.child, 237 .void => void, 238 else => @compileError("Server handler must be a struct, got: " ++ @tagName(@typeInfo(H))), 239 }; 240 241 const ActionArg = if (comptime std.meta.hasFn(Handler, "dispatch")) @typeInfo(@TypeOf(Handler.dispatch)).@"fn".params[1].type.? else Action(H); 242 243 const has_websocket = Handler != void and @hasDecl(Handler, "WebsocketHandler"); 244 const WebsocketHandler = if (has_websocket) Handler.WebsocketHandler else DummyWebsocketHandler; 245 246 const RouterConfig = struct { 247 middlewares: []const Middleware(H) = &.{}, 248 }; 249 250 const MiddlewareItem = struct { 251 middleware: Middleware(H), 252 node: std.SinglyLinkedList.Node = .{}, 253 }; 254 255 return struct { 256 io: Io, 257 handler: H, 258 config: Config, 259 arena: Allocator, 260 allocator: Allocator, 261 _router: Router(H, ActionArg), 262 _mut: Io.Mutex, 263 _workers: []Worker, 264 _cond: Io.Condition, 265 _listener: ?posix.fd_t, 266 _max_request_per_connection: usize, 267 _middlewares: []const Middleware(H), 268 _websocket_state: websocket.server.WorkerState, 269 _middleware_registry: std.SinglyLinkedList, 270 271 const Self = @This(); 272 const Worker = if (blockingMode()) worker.Blocking(*Self, WebsocketHandler) else worker.NonBlocking(*Self, WebsocketHandler); 273 274 pub fn init(io: Io, allocator: Allocator, config: Config, handler: H) !Self { 275 // Be mindful about where we pass this arena. Most things are able to 276 // do dynamic allocation, and need to be able to free when they're 277 // done with their memory. Only use this for stuff that's created on 278 // startup and won't dynamically need to grow/shrink. 279 const arena = try allocator.create(std.heap.ArenaAllocator); 280 errdefer allocator.destroy(arena); 281 arena.* = std.heap.ArenaAllocator.init(allocator); 282 errdefer arena.deinit(); 283 284 const default_dispatcher = if (comptime Handler == void) defaultDispatcher else defaultDispatcherWithHandler; 285 286 // do not pass arena.allocator to WorkerState, it needs to be able to 287 // allocate and free at will. 288 289 const ws_config = config.websocket; 290 var websocket_state = try websocket.server.WorkerState.init(allocator, io, .{ 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(); 311 312 const workers = try arena.allocator().alloc(Worker, config.workerCount()); 313 314 return .{ 315 .io = io, 316 .config = config, 317 .handler = handler, 318 .allocator = allocator, 319 .arena = arena.allocator(), 320 ._mut = .init, 321 ._cond = .init, 322 ._workers = workers, 323 ._listener = null, 324 ._middlewares = &.{}, 325 ._middleware_registry = .{}, 326 ._websocket_state = websocket_state, 327 ._router = try Router(H, ActionArg).init(arena.allocator(), default_dispatcher, handler), 328 ._max_request_per_connection = config.timeout.request_count orelse MAX_REQUEST_COUNT, 329 }; 330 } 331 332 pub fn deinit(self: *Self) void { 333 self._websocket_state.deinit(); 334 335 var node = self._middleware_registry.first; 336 while (node) |n| { 337 const item: *MiddlewareItem = @fieldParentPtr("node", n); 338 item.middleware.deinit(); 339 node = n.next; 340 } 341 342 const arena: *std.heap.ArenaAllocator = @ptrCast(@alignCast(self.arena.ptr)); 343 arena.deinit(); 344 self.allocator.destroy(arena); 345 } 346 347 pub fn listen(self: *Self) !void { 348 // incase "stop" is waiting 349 const io = self.io; 350 351 defer self._cond.signal(io); 352 self._mut.lockUncancelable(io); 353 errdefer self._mut.unlock(io); 354 355 const config = self.config; 356 const address = try config.address.toPosix(io); 357 const is_unix_socket = address.any.family == posix.AF.UNIX; 358 359 const listener = blk: { 360 var sock_flags: u32 = posix.SOCK.STREAM | posix.CLOEXEC; 361 if (blockingMode() == false) sock_flags |= posix.NONBLOCK; 362 363 const proto = if (is_unix_socket) @as(u32, 0) else posix.IPPROTO.TCP; 364 break :blk try posix.socket(address.any.family, sock_flags, proto); 365 }; 366 367 errdefer { 368 posix.close(listener); 369 self._listener = null; 370 } 371 372 if (is_unix_socket) { 373 // TODO: Broken on darwin: 374 // https://github.com/ziglang/zig/issues/17260 375 // if (@hasDecl(os.TCP, "NODELAY")) { 376 // try os.setsockopt(socket.sockfd.?, os.IPPROTO.TCP, os.TCP.NODELAY, &std.mem.toBytes(@as(c_int, 1))); 377 // } 378 try posix.setsockopt(listener, posix.IPPROTO.TCP, 1, &std.mem.toBytes(@as(c_int, 1))); 379 } 380 381 try posix.setsockopt(listener, posix.SOL.SOCKET, posix.SO.REUSEADDR, &std.mem.toBytes(@as(c_int, 1))); 382 383 if (is_unix_socket == false and self._workers.len > 1) { 384 if (@hasDecl(posix.SO, "REUSEPORT_LB")) { 385 try posix.setsockopt(listener, posix.SOL.SOCKET, posix.SO.REUSEPORT_LB, &std.mem.toBytes(@as(c_int, 1))); 386 } else if (@hasDecl(posix.SO, "REUSEPORT")) { 387 try posix.setsockopt(listener, posix.SOL.SOCKET, posix.SO.REUSEPORT, &std.mem.toBytes(@as(c_int, 1))); 388 } 389 } 390 391 { 392 const socklen = address.getOsSockLen(); 393 try posix.bind(listener, &address.any, socklen); 394 try posix.listen(listener, 1024); // kernel backlog 395 } 396 397 self._listener = listener; 398 399 var workers = self._workers; 400 const allocator = self.allocator; 401 402 if (comptime blockingMode()) { 403 workers[0] = try worker.Blocking(*Self, WebsocketHandler).init(io, allocator, self, &config); 404 defer workers[0].deinit(); 405 406 const thrd = try Thread.spawn(.{}, worker.Blocking(*Self, WebsocketHandler).listen, .{ &workers[0], listener }); 407 408 // incase listenInNewThread was used and is waiting for us to start 409 self._cond.signal(io); 410 self._mut.unlock(io); 411 412 // This will unblock when server.stop() is called and the listening 413 // socket is closed. 414 thrd.join(); 415 } else { 416 var started: usize = 0; 417 defer for (0..started) |i| { 418 workers[i].deinit(); 419 }; 420 421 errdefer for (0..started) |i| { 422 workers[i].stop(); 423 }; 424 425 var ready_sem = Io.Semaphore{}; 426 const threads = try self.arena.alloc(Thread, workers.len); 427 for (0..workers.len) |i| { 428 workers[i] = try Worker.init(io, allocator, self, &config); 429 errdefer { 430 workers[i].stop(); 431 workers[i].deinit(); 432 } 433 threads[i] = try Thread.spawn(.{}, Worker.run, .{ &workers[i], listener, &ready_sem }); 434 started += 1; 435 } 436 437 for (0..workers.len) |_| { 438 ready_sem.waitUncancelable(io); 439 } 440 441 // incase listenInNewThread was used and is waiting for us to start 442 self._cond.signal(io); 443 self._mut.unlock(io); 444 445 for (threads) |thrd| { 446 thrd.join(); 447 } 448 } 449 } 450 451 pub fn listenInNewThread(self: *Self) !std.Thread { 452 const io = self.io; 453 self._mut.lockUncancelable(io); 454 defer self._mut.unlock(io); 455 const thrd = try std.Thread.spawn(.{}, listen, .{self}); 456 457 // we don't return until listen() signals us that the server is up 458 self._cond.waitUncancelable(io, &self._mut); 459 460 return thrd; 461 } 462 463 pub fn stop(self: *Self) void { 464 const io = self.io; 465 self._mut.lockUncancelable(io); 466 defer self._mut.unlock(io); 467 468 for (self._workers) |*w| { 469 if (self._listener == null) { 470 log.err("Cannot stop server, .listen() was never called", .{}); 471 break; 472 } 473 474 w.stop(); 475 } 476 477 if (self._listener) |l| { 478 if (comptime blockingMode()) { 479 // necessary to unblock accept on linux 480 // (which might not be that necessary since, on Linux, 481 // NonBlocking should be used) 482 posix.shutdown(l, .recv) catch {}; 483 } 484 posix.close(l); 485 } 486 } 487 488 pub fn router(self: *Self, config: RouterConfig) !*Router(H, ActionArg) { 489 // we store this in self for us when no route is found (these will 490 // still be executed). 491 492 const owned = try self.arena.dupe(Middleware(H), config.middlewares); 493 self._middlewares = owned; 494 495 // we store this in router to append to add/append to created routes 496 self._router.middlewares = owned; 497 498 return &self._router; 499 } 500 501 fn defaultDispatcher(action: ActionArg, req: *Request, res: *Response) !void { 502 return action(req, res); 503 } 504 505 fn defaultDispatcherWithHandler(handler: H, action: ActionArg, req: *Request, res: *Response) !void { 506 if (comptime std.meta.hasFn(Handler, "dispatch")) { 507 return handler.dispatch(action, req, res); 508 } 509 return action(handler, req, res); 510 } 511 512 // This is always called from within a threadpool thread. For nonblocking, 513 // notifyingHandler (above) was the threadpool's main entry and it called this. 514 // For blocking, the threadpool was directed to the worker's handleConnection 515 // which eventually called this. 516 // thread_buf is a thread-specific configurable-sized buffer that we're 517 // free to use as we want. This is, by far, the most efficient memory 518 // we can use because it's allocated on server start and re-used on 519 // each request (which is safe, because, in blocking or nonblocking, once 520 // a request reaches this point, processing is blocking from the point 521 // of view of the server). 522 // We'll use thread_buf as part of a FallBackAllocator with the conn 523 // arena for our request ONLY. We cannot use thread_buf for the response 524 // because the response data must outlive the execution of this function 525 // (and thus, in nonblocking, outlives this threadpool's execution unit). 526 pub fn handleRequest(self: *Self, conn: *HTTPConn, thread_buf: []u8) void { 527 const aa = conn.req_arena.allocator(); 528 529 var fba = FixedBufferAllocator.init(thread_buf); 530 var fb = FallbackAllocator{ 531 .fba = &fba, 532 .fallback = aa, 533 .fixed = fba.allocator(), 534 }; 535 536 const allocator = fb.allocator(); 537 var req = Request.init(allocator, conn); 538 var res = Response.init(allocator, conn); 539 540 defer std.debug.assert(res.written == true); 541 542 if (comptime std.meta.hasFn(Handler, "handle")) { 543 if (comptime @typeInfo(@TypeOf(Handler.handle)).@"fn".return_type != void) { 544 @compileError(@typeName(Handler) ++ ".handle must return 'void'"); 545 } 546 self.handler.handle(&req, &res); 547 } else { 548 const dispatchable_action = self._router.route(req.method, req.method_string, req.url.path, req.params); 549 550 var executor = Executor{ 551 .index = 0, 552 .req = &req, 553 .res = &res, 554 .handler = self.handler, 555 .middlewares = undefined, 556 .dispatchable_action = dispatchable_action, 557 }; 558 559 if (dispatchable_action) |da| { 560 req.route_data = da.data; 561 executor.middlewares = da.middlewares; 562 } else { 563 req.route_data = null; 564 executor.middlewares = self._middlewares; 565 } 566 567 executor.next() catch |err| { 568 if (comptime std.meta.hasFn(Handler, "uncaughtError")) { 569 self.handler.uncaughtError(&req, &res, err); 570 } else { 571 res.status = 500; 572 res.body = "Internal Server Error"; 573 std.log.warn("httpz: unhandled exception for request: {s}\nErr: {}", .{ req.url.raw, err }); 574 } 575 }; 576 } 577 578 if (conn.handover == .unknown) { 579 // close is the default 580 conn.handover = if (req.canKeepAlive() and conn.request_count < self._max_request_per_connection) .keepalive else .close; 581 } 582 583 res.write() catch { 584 conn.handover = .close; 585 }; 586 587 if (req.unread_body > 0 and conn.handover == .keepalive) { 588 drain(&req) catch { 589 conn.handover = .close; 590 }; 591 } 592 } 593 594 pub fn middleware(self: *Self, comptime M: type, config: M.Config) !Middleware(H) { 595 const arena = self.arena; 596 597 const node = try arena.create(MiddlewareItem); 598 errdefer arena.destroy(node); 599 600 const m = try arena.create(M); 601 errdefer arena.destroy(m); 602 switch (comptime @typeInfo(@TypeOf(M.init)).@"fn".params.len) { 603 1 => m.* = try M.init(config), 604 2 => m.* = try M.init(config, MiddlewareConfig{ 605 .arena = arena, 606 .allocator = self.allocator, 607 }), 608 else => @compileError(@typeName(M) ++ ".init should accept 1 or 2 parameters"), 609 } 610 611 const iface = Middleware(H).init(m); 612 node.*.middleware = iface; 613 self._middleware_registry.prepend(&node.node); 614 615 return iface; 616 } 617 618 pub const Executor = struct { 619 index: usize, 620 req: *Request, 621 res: *Response, 622 handler: H, 623 // pull this out of da since we'll access it a lot (not really, but w/e) 624 middlewares: []const Middleware(H), 625 dispatchable_action: ?*const DispatchableAction(H, ActionArg), 626 627 pub fn next(self: *Executor) !void { 628 const index = self.index; 629 const middlewares = self.middlewares; 630 631 if (index < middlewares.len) { 632 self.index = index + 1; 633 return middlewares[index].execute(self.req, self.res, self); 634 } 635 636 // done executing our middlewares, now we either execute the 637 // dispatcher or not found. 638 if (self.dispatchable_action) |da| { 639 if (comptime H == void) { 640 return da.dispatcher(da.action, self.req, self.res); 641 } 642 return da.dispatcher(da.handler, da.action, self.req, self.res); 643 } 644 645 if (comptime std.meta.hasFn(Handler, "notFound")) { 646 return self.handler.notFound(self.req, self.res); 647 } 648 self.res.status = 404; 649 self.res.body = "Not Found"; 650 return; 651 } 652 }; 653 }; 654} 655 656pub fn blockingMode() bool { 657 if (force_blocking) { 658 return true; 659 } 660 return switch (builtin.os.tag) { 661 .linux, .macos, .ios, .tvos, .watchos, .freebsd, .netbsd, .dragonfly, .openbsd => false, 662 else => true, 663 }; 664} 665 666pub 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 } 671 672 const version = req.header("sec-websocket-version") orelse return false; 673 if (std.ascii.eqlIgnoreCase(version, "13") == false) { 674 return false; 675 } 676 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 } 682 683 const key = req.header("sec-websocket-key") orelse return false; 684 685 const http_conn = res.conn; 686 const ws_worker: *websocket.server.Worker(H) = @ptrCast(@alignCast(http_conn.ws_worker)); 687 688 // websocket's Address wraps sockaddr.storage (128B) since v0.1.7; 689 // widen our bare sockaddr into a zeroed storage 690 const ws_compat = posix.Address.fromIOAddress(http_conn.address); 691 var ws_addr_storage = std.mem.zeroes(std.posix.sockaddr.storage); 692 @memcpy( 693 std.mem.asBytes(&ws_addr_storage)[0..@sizeOf(@TypeOf(ws_compat.any))], 694 std.mem.asBytes(&ws_compat.any), 695 ); 696 var hc = try ws_worker.createConn( 697 http_conn.stream.socket.handle, 698 .{ .any = ws_addr_storage }, 699 worker.timestamp(http_conn.io), 700 ); 701 // The HTTP connection still owns the socket until handover is committed. 702 // If handler initialization or handshake setup fails, discard only the 703 // speculative websocket state so the caller can send an HTTP error response. 704 errdefer discardWebsocketConn(ws_worker, hc); 705 706 hc.handler = try H.init(&hc.conn, ctx); 707 708 var compression = false; 709 if (ws_worker.canCompress()) { 710 if (req.header("sec-websocket-extensions")) |ext| { 711 compression = try websocket.Handshake.parseExtension(ext) != null; 712 } 713 } 714 715 // Optional subprotocol negotiation. If the caller's ctx has a 716 // `.subprotocol` field set to a non-null value, echo it in the 101 717 // response header so clients (e.g. python's `websockets`) accept the 718 // upgrade. Per RFC 6455, the server MUST include a subprotocol it 719 // chose from the client's offer; without echo, those clients close. 720 var subprotocol_keys: [1][]const u8 = undefined; 721 var subprotocol_values: [1][]const u8 = undefined; 722 var subprotocol_kv: websocket.Handshake.KeyValue = undefined; 723 const headers_for_reply: ?*websocket.Handshake.KeyValue = blk: { 724 if (comptime @typeInfo(@TypeOf(ctx)) == .pointer) { 725 const CtxT = @TypeOf(ctx.*); 726 if (comptime @hasField(CtxT, "subprotocol")) { 727 if (ctx.subprotocol) |sp| { 728 subprotocol_keys[0] = "Sec-WebSocket-Protocol"; 729 subprotocol_values[0] = sp; 730 subprotocol_kv = .{ 731 .len = 1, 732 .keys = subprotocol_keys[0..], 733 .values = subprotocol_values[0..], 734 }; 735 break :blk &subprotocol_kv; 736 } 737 } 738 } 739 break :blk null; 740 }; 741 742 var reply_buf: [512]u8 = undefined; 743 const reply = try websocket.Handshake.createReply(key, headers_for_reply, compression, &reply_buf); 744 var writer = http_conn.stream.writer(http_conn.io, &.{}); 745 const w = &writer.interface; 746 try w.writeAll(reply); 747 try w.flush(); 748 749 if (comptime std.meta.hasFn(H, "afterInit")) { 750 const params = @typeInfo(@TypeOf(H.afterInit)).@"fn".params; 751 try if (comptime params.len == 1) hc.handler.?.afterInit() else hc.handler.?.afterInit(ctx); 752 } 753 try ws_worker.setupConnection(hc); 754 res.written = true; 755 http_conn.handover = .{ .websocket = hc }; 756 return true; 757} 758 759fn discardWebsocketConn(ws_worker: anytype, hc: anytype) void { 760 const WorkerType = @TypeOf(ws_worker.*); 761 if (comptime @hasDecl(WorkerType, "discardConn")) { 762 ws_worker.discardConn(hc); 763 } else { 764 // websocket.zig <=0.1.8 exposes the underlying public manager cleanup 765 // but not the ownership-naming facade. No reader or compression state 766 // exists before setupConnection, so this is the complete no-close path. 767 ws_worker.worker.conn_manager.cleanup(hc); 768 } 769} 770 771// std.heap.StackFallbackAllocator is very specific. It's really _stack_ as it 772// requires a comptime size. Also, it uses non-public calls from the FixedBufferAllocator. 773// There should be a more generic FallbackAllocator that just takes 2 allocators... 774// which is what this is. 775const FallbackAllocator = struct { 776 fixed: Allocator, 777 fallback: Allocator, 778 fba: *FixedBufferAllocator, 779 780 pub fn allocator(self: *FallbackAllocator) Allocator { 781 return .{ 782 .ptr = self, 783 .vtable = &.{ 784 .alloc = alloc, 785 .resize = resize, 786 .free = free, 787 .remap = remap, 788 }, 789 }; 790 } 791 792 fn alloc(ctx: *anyopaque, len: usize, alignment: std.mem.Alignment, ra: usize) ?[*]u8 { 793 const self: *FallbackAllocator = @ptrCast(@alignCast(ctx)); 794 return self.fixed.rawAlloc(len, alignment, ra) orelse self.fallback.rawAlloc(len, alignment, ra); 795 } 796 797 fn resize(ctx: *anyopaque, buf: []u8, alignment: std.mem.Alignment, new_len: usize, ra: usize) bool { 798 const self: *FallbackAllocator = @ptrCast(@alignCast(ctx)); 799 if (self.fba.ownsPtr(buf.ptr)) { 800 // A failed resize must leave the allocation valid. Callers like 801 // ArenaAllocator keep using the buffer after a false return, so 802 // we cannot reclaim it here (that caused memory to be handed out 803 // twice and corrupted live arena nodes). 804 return self.fixed.rawResize(buf, alignment, new_len, ra); 805 } 806 return self.fallback.rawResize(buf, alignment, new_len, ra); 807 } 808 809 fn free(ctx: *anyopaque, buf: []u8, alignment: std.mem.Alignment, ra: usize) void { 810 const self: *FallbackAllocator = @ptrCast(@alignCast(ctx)); 811 if (self.fba.ownsPtr(buf.ptr)) { 812 // reclaims thread_buf space when buf is the fba's last allocation 813 self.fixed.rawFree(buf, alignment, ra); 814 } 815 // else: fallback is an arena in our specific usage, freeing is a noop 816 } 817 818 fn remap(ctx: *anyopaque, memory: []u8, alignment: std.mem.Alignment, new_len: usize, ret_addr: usize) ?[*]u8 { 819 if (resize(ctx, memory, alignment, new_len, ret_addr)) { 820 return memory.ptr; 821 } 822 return null; 823 } 824}; 825 826// Called when we have unread bytes on the request and want to keepalive the 827// connection. Only happens when lazy_read_size is configured and the client 828// didn't read the [whole] body 829// There should already be a receive timeout on the socket since the only 830// way for this to be 831fn drain(req: *Request) !void { 832 var r = try req.reader(2000); 833 var buf: [4096]u8 = undefined; 834 while (true) { 835 if (try r.read(&buf) == 0) { 836 return; 837 } 838 } 839} 840 841const t = @import("t.zig"); 842var global_test_allocator = std.heap.DebugAllocator(.{}){}; 843 844var test_handler_dispatch = TestHandlerDispatch{ .state = 10 }; 845var test_handler_disaptch_context = TestHandlerDispatchContext{ .state = 20 }; 846var test_handler_default_dispatch1 = TestHandlerDefaultDispatch{ .state = 3 }; 847var test_handler_default_dispatch2 = TestHandlerDefaultDispatch{ .state = 99 }; 848var test_handler_default_dispatch3 = TestHandlerDefaultDispatch{ .state = 20 }; 849 850var default_server: Server(void) = undefined; 851var dispatch_default_server: Server(*TestHandlerDefaultDispatch) = undefined; 852var dispatch_server: Server(*TestHandlerDispatch) = undefined; 853var dispatch_action_context_server: Server(*TestHandlerDispatchContext) = undefined; 854var reuse_server: Server(void) = undefined; 855var handle_server: Server(TestHandlerHandle) = undefined; 856var websocket_server: Server(TestWebsocketHandler) = undefined; 857var cors_wildcard_server: Server(void) = undefined; 858var cors_single_server: Server(void) = undefined; 859var cors_multiple_server: Server(void) = undefined; 860var chunked_server: Server(void) = undefined; 861var snapshot_server: Server(TestWebsocketHandler) = undefined; 862var ws_slot_server: Server(TestWebsocketHandler) = undefined; 863 864var test_server_threads: [13]Thread = undefined; 865 866test "tests:beforeAll" { 867 // this will leak since the server will run until the process exits. If we use 868 // our testing allocator, it'll report the leak. 869 const ga = global_test_allocator.allocator(); 870 871 { 872 default_server = try Server(void).init(t.io, ga, .{ 873 .address = .localhost(5992), 874 .request = .{ 875 .lazy_read_size = 4_096, 876 .max_body_size = 1_048_576, 877 }, 878 }, {}); 879 880 // only need to do this because we're using listenInNewThread instead 881 // of blocking here. So the array to hold the middleware needs to outlive 882 // this function. 883 var cors = try default_server.arena.alloc(Middleware(void), 1); 884 cors[0] = try default_server.middleware(middleware.Cors, .{ 885 .max_age = "300", 886 .methods = "GET,POST", 887 .origin = "httpz.local", 888 .headers = "content-type", 889 }); 890 891 var middlewares = try default_server.arena.alloc(Middleware(void), 2); 892 middlewares[0] = try default_server.middleware(TestMiddleware, .{ .id = 100 }); 893 middlewares[1] = cors[0]; 894 895 var router = try default_server.router(.{}); 896 // router.get("/test/ws", testWS); 897 router.get("/fail", TestDummyHandler.fail, .{}); 898 router.get("/test/json", TestDummyHandler.jsonRes, .{}); 899 router.get("/test/method", TestDummyHandler.method, .{}); 900 router.put("/test/method", TestDummyHandler.method, .{}); 901 router.method("TEA", "/test/method", TestDummyHandler.method, .{}); 902 router.method("PING", "/test/method", TestDummyHandler.method, .{}); 903 router.get("/test/query", TestDummyHandler.reqQuery, .{}); 904 router.get("/test/stream", TestDummyHandler.eventStream, .{}); 905 router.get("/test/streamsync", TestDummyHandler.eventStreamSync, .{}); 906 router.get("/test/req_reader", TestDummyHandler.reqReader, .{}); 907 router.get("/test/chunked", TestDummyHandler.chunked, .{}); 908 router.get("/test/route_data", TestDummyHandler.routeData, .{ .data = &TestDummyHandler.RouteData{ .power = 12345 } }); 909 router.all("/test/cors", TestDummyHandler.jsonRes, .{ .middlewares = cors }); 910 router.all("/test/middlewares", TestDummyHandler.middlewares, .{ .middlewares = middlewares }); 911 router.all("/test/dispatcher", TestDummyHandler.dispatchedAction, .{ .dispatcher = TestDummyHandler.routeSpecificDispacthcer }); 912 test_server_threads[0] = try default_server.listenInNewThread(); 913 } 914 915 { 916 dispatch_default_server = try Server(*TestHandlerDefaultDispatch).init(t.io, ga, .{ .address = .localhost(5993) }, &test_handler_default_dispatch1); 917 var router = try dispatch_default_server.router(.{}); 918 router.get("/", TestHandlerDefaultDispatch.echo, .{}); 919 router.get("/write/*", TestHandlerDefaultDispatch.echoWrite, .{}); 920 router.get("/fail", TestHandlerDefaultDispatch.fail, .{}); 921 router.post("/login", TestHandlerDefaultDispatch.echo, .{}); 922 router.get("/test/body/cl", TestHandlerDefaultDispatch.clBody, .{}); 923 router.get("/test/headers", TestHandlerDefaultDispatch.headers, .{}); 924 router.all("/api/:version/users/:UserId", TestHandlerDefaultDispatch.params, .{}); 925 926 var admin_routes = router.group("/admin/", .{ .dispatcher = TestHandlerDefaultDispatch.dispatch2, .handler = &test_handler_default_dispatch2 }); 927 admin_routes.get("/users", TestHandlerDefaultDispatch.echo, .{}); 928 admin_routes.put("/users/:id", TestHandlerDefaultDispatch.echo, .{}); 929 930 var debug_routes = router.group("/debug", .{ .dispatcher = TestHandlerDefaultDispatch.dispatch3, .handler = &test_handler_default_dispatch3 }); 931 debug_routes.head("/ping", TestHandlerDefaultDispatch.echo, .{}); 932 debug_routes.options("/stats", TestHandlerDefaultDispatch.echo, .{}); 933 934 test_server_threads[1] = try dispatch_default_server.listenInNewThread(); 935 } 936 937 { 938 dispatch_server = try Server(*TestHandlerDispatch).init(t.io, ga, .{ .address = .localhost(5994) }, &test_handler_dispatch); 939 var router = try dispatch_server.router(.{}); 940 router.get("/", TestHandlerDispatch.root, .{}); 941 test_server_threads[2] = try dispatch_server.listenInNewThread(); 942 } 943 944 { 945 dispatch_action_context_server = try Server(*TestHandlerDispatchContext).init(t.io, ga, .{ .address = .localhost(5995) }, &test_handler_disaptch_context); 946 var router = try dispatch_action_context_server.router(.{}); 947 router.get("/", TestHandlerDispatchContext.root, .{}); 948 test_server_threads[3] = try dispatch_action_context_server.listenInNewThread(); 949 } 950 951 { 952 // with only 1 worker, and a min/max conn of 1, each request should 953 // hit our reset path. 954 reuse_server = try Server(void).init(t.io, ga, .{ .address = .localhost(5996), .workers = .{ .count = 1, .min_conn = 1, .max_conn = 1 } }, {}); 955 var router = try reuse_server.router(.{}); 956 router.get("/test/writer", TestDummyHandler.reuseWriter, .{}); 957 test_server_threads[4] = try reuse_server.listenInNewThread(); 958 } 959 960 { 961 handle_server = try Server(TestHandlerHandle).init(t.io, ga, .{ .address = .localhost(5997) }, TestHandlerHandle{}); 962 test_server_threads[5] = try handle_server.listenInNewThread(); 963 } 964 965 { 966 websocket_server = try Server(TestWebsocketHandler).init(t.io, ga, .{ .address = .localhost(5998) }, TestWebsocketHandler{}); 967 var router = try websocket_server.router(.{}); 968 router.get("/ws", TestWebsocketHandler.upgrade, .{}); 969 router.get("/ws-reject", TestWebsocketHandler.rejectUpgrade, .{}); 970 test_server_threads[6] = try websocket_server.listenInNewThread(); 971 } 972 973 { 974 cors_wildcard_server = try Server(void).init(t.io, ga, .{ .address = .localhost(5999) }, {}); 975 var cors_wildcard = try cors_wildcard_server.arena.alloc(Middleware(void), 1); 976 cors_wildcard[0] = try cors_wildcard_server.middleware(middleware.Cors, .{ 977 .origin = "*", 978 .max_age = "600", 979 .methods = "GET,POST,PUT", 980 .headers = "authorization,content-type", 981 }); 982 var router = try cors_wildcard_server.router(.{}); 983 router.all("/test/cors", TestDummyHandler.jsonRes, .{ .middlewares = cors_wildcard }); 984 test_server_threads[7] = try cors_wildcard_server.listenInNewThread(); 985 } 986 987 { 988 cors_single_server = try Server(void).init(t.io, ga, .{ .address = .localhost(6000) }, {}); 989 var cors_single = try cors_single_server.arena.alloc(Middleware(void), 1); 990 cors_single[0] = try cors_single_server.middleware(middleware.Cors, .{ 991 .origin = "https://example.com", 992 .credentials = "true", 993 .max_age = "3600", 994 }); 995 var router = try cors_single_server.router(.{}); 996 router.all("/test/cors", TestDummyHandler.jsonRes, .{ .middlewares = cors_single }); 997 test_server_threads[8] = try cors_single_server.listenInNewThread(); 998 } 999 1000 { 1001 cors_multiple_server = try Server(void).init(t.io, ga, .{ .address = .localhost(6001) }, {}); 1002 var cors_multiple = try cors_multiple_server.arena.alloc(Middleware(void), 1); 1003 cors_multiple[0] = try cors_multiple_server.middleware(middleware.Cors, .{ 1004 .origin = "https://example.com, https://api.example.com, https://test.local", 1005 .credentials = "true", 1006 .methods = "GET,POST,DELETE", 1007 .headers = "x-custom-header", 1008 }); 1009 var router = try cors_multiple_server.router(.{}); 1010 router.all("/test/cors", TestDummyHandler.jsonRes, .{ .middlewares = cors_multiple }); 1011 test_server_threads[9] = try cors_multiple_server.listenInNewThread(); 1012 } 1013 1014 { 1015 // Dedicated server for chunked-encoding integration tests. Small 1016 // max_body_size lets us cover BodyTooBig, and a small pool buffer 1017 // forces growth on bodies larger than 1KB. 1018 chunked_server = try Server(void).init(t.io, ga, .{ 1019 .address = .localhost(6002), 1020 .request = .{ .max_body_size = 32_768 }, 1021 .workers = .{ .large_buffer_size = 1024 }, 1022 }, {}); 1023 var router = try chunked_server.router(.{}); 1024 router.post("/echo_body", TestDummyHandler.echoBody, .{}); 1025 test_server_threads[10] = try chunked_server.listenInNewThread(); 1026 } 1027 1028 { 1029 // Single-worker server mixing plain http and websocket upgrades, for 1030 // the handover-snapshot regression test. One worker means all 1031 // handovers funnel through one reactor's handover_list. 1032 snapshot_server = try Server(TestWebsocketHandler).init(t.io, ga, .{ 1033 .address = .localhost(6003), 1034 .workers = .{ .count = 1 }, 1035 }, TestWebsocketHandler{}); 1036 var router = try snapshot_server.router(.{}); 1037 router.get("/echo", TestWebsocketHandler.echoTag, .{}); 1038 router.get("/ws", TestWebsocketHandler.upgrade, .{}); 1039 test_server_threads[11] = try snapshot_server.listenInNewThread(); 1040 } 1041 1042 { 1043 // Tiny max_conn so leaked conn slots wedge the worker within a few 1044 // websocket connections, for the ws slot-release regression test. 1045 ws_slot_server = try Server(TestWebsocketHandler).init(t.io, ga, .{ 1046 .address = .localhost(6004), 1047 .workers = .{ .count = 1, .max_conn = 4, .min_conn = 4 }, 1048 }, TestWebsocketHandler{}); 1049 var router = try ws_slot_server.router(.{}); 1050 router.get("/echo", TestWebsocketHandler.echoTag, .{}); 1051 router.get("/ws", TestWebsocketHandler.upgrade, .{}); 1052 test_server_threads[12] = try ws_slot_server.listenInNewThread(); 1053 } 1054 1055 std.testing.refAllDecls(@This()); 1056} 1057 1058test "tests:afterAll" { 1059 default_server.stop(); 1060 dispatch_default_server.stop(); 1061 dispatch_server.stop(); 1062 dispatch_action_context_server.stop(); 1063 reuse_server.stop(); 1064 handle_server.stop(); 1065 websocket_server.stop(); 1066 cors_wildcard_server.stop(); 1067 cors_single_server.stop(); 1068 cors_multiple_server.stop(); 1069 chunked_server.stop(); 1070 snapshot_server.stop(); 1071 ws_slot_server.stop(); 1072 1073 for (test_server_threads) |thread| { 1074 thread.join(); 1075 } 1076 1077 default_server.deinit(); 1078 dispatch_default_server.deinit(); 1079 dispatch_server.deinit(); 1080 dispatch_action_context_server.deinit(); 1081 reuse_server.deinit(); 1082 handle_server.deinit(); 1083 websocket_server.deinit(); 1084 cors_wildcard_server.deinit(); 1085 cors_single_server.deinit(); 1086 cors_multiple_server.deinit(); 1087 chunked_server.deinit(); 1088 snapshot_server.deinit(); 1089 ws_slot_server.deinit(); 1090 1091 try t.expectEqual(0, global_test_allocator.detectLeaks()); 1092} 1093 1094test "httpz: quick shutdown" { 1095 var server = try Server(void).init(t.io, t.allocator, .{ .address = .localhost(6992) }, {}); 1096 const thrd = try server.listenInNewThread(); 1097 server.stop(); 1098 thrd.join(); 1099 server.deinit(); 1100} 1101 1102test "httpz: bind failure releases mutex" { 1103 // Start a server to occupy the port 1104 var server1 = try Server(void).init(t.io, t.allocator, .{ .address = .localhost(6993) }, {}); 1105 const thrd1 = try server1.listenInNewThread(); 1106 defer { 1107 server1.stop(); 1108 thrd1.join(); 1109 server1.deinit(); 1110 } 1111 1112 // Try to start another server on the same port - will fail to bind 1113 var server2 = try Server(void).init(t.io, t.allocator, .{ .address = .localhost(6993) }, {}); 1114 defer server2.deinit(); 1115 1116 // First call fails with AddressInUse 1117 try t.expectError(error.AddressInUse, server2.listen()); 1118 1119 // Before fix: second call would deadlock because mutex left locked 1120 // After fix: errdefer releases mutex, so second call can proceed 1121 try t.expectError(error.AddressInUse, server2.listen()); 1122} 1123 1124test "httpz: shutdown without listen" { 1125 // Should not throw a .BADF (unreachable) error 1126 var server = try Server(void).init(t.io, t.allocator, .{ .address = .localhost(6992) }, {}); 1127 server.stop(); 1128 server.deinit(); 1129} 1130 1131test "httpz: invalid request" { 1132 const stream = testStream(5992); 1133 defer stream.close(t.io); 1134 var writer = stream.writer(t.io, &.{}); 1135 const w = &writer.interface; 1136 try w.writeAll("TEA HTTP/1.1\r\n\r\n"); 1137 try w.flush(); 1138 1139 var buf: [100]u8 = undefined; 1140 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1141} 1142 1143test "httpz: invalid request path" { 1144 const stream = testStream(5992); 1145 defer stream.close(t.io); 1146 var writer = stream.writer(t.io, &.{}); 1147 const w = &writer.interface; 1148 try w.writeAll("TEA /hello\rn\nWorld:test HTTP/1.1\r\n\r\n"); 1149 try w.flush(); 1150 1151 var buf: [100]u8 = undefined; 1152 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1153} 1154 1155test "httpz: invalid header name" { 1156 const stream = testStream(5992); 1157 defer stream.close(t.io); 1158 var writer = stream.writer(t.io, &.{}); 1159 const w = &writer.interface; 1160 try w.writeAll("GET / HTTP/1.1\r\nOver: 9000\r\nHel\tlo:World\r\n\r\n"); 1161 try w.flush(); 1162 1163 var buf: [100]u8 = undefined; 1164 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1165} 1166 1167test "httpz: invalid content length value (1)" { 1168 const stream = testStream(5992); 1169 defer stream.close(t.io); 1170 var writer = stream.writer(t.io, &.{}); 1171 const w = &writer.interface; 1172 try w.writeAll("GET / HTTP/1.1\r\nContent-Length: HaHA\r\n\r\n"); 1173 try w.flush(); 1174 1175 var buf: [100]u8 = undefined; 1176 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1177} 1178 1179test "httpz: invalid content length value (2)" { 1180 const stream = testStream(5992); 1181 defer stream.close(t.io); 1182 var writer = stream.writer(t.io, &.{}); 1183 const w = &writer.interface; 1184 try w.writeAll("GET / HTTP/1.1\r\nContent-Length: 1.0\r\n\r\n"); 1185 try w.flush(); 1186 1187 var buf: [100]u8 = undefined; 1188 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1189} 1190 1191test "httpz: body too big" { 1192 const stream = testStream(5993); 1193 defer stream.close(t.io); 1194 var writer = stream.writer(t.io, &.{}); 1195 const w = &writer.interface; 1196 try w.writeAll("POST / HTTP/1.1\r\nContent-Length: 999999999999999999\r\n\r\n"); 1197 try w.flush(); 1198 1199 var buf: [100]u8 = undefined; 1200 try t.expectString("HTTP/1.1 413 \r\nConnection: Close\r\nContent-Length: 23\r\n\r\nRequest body is too big", testReadAll(stream, &buf)); 1201} 1202 1203test "httpz: overflow content length" { 1204 const stream = testStream(5992); 1205 defer stream.close(t.io); 1206 var writer = stream.writer(t.io, &.{}); 1207 const w = &writer.interface; 1208 try w.writeAll("GET / HTTP/1.1\r\nContent-Length: 999999999999999999999999999\r\n\r\n"); 1209 try w.flush(); 1210 1211 var buf: [100]u8 = undefined; 1212 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1213} 1214 1215test "httpz: chunked body" { 1216 // small body — fits entirely in the initial pooled buffer (1024 bytes 1217 // on this test server). Split into two chunks plus a chunk extension, 1218 // which must be ignored. 1219 const stream = testStream(6002); 1220 defer stream.close(t.io); 1221 var writer = stream.writer(t.io, &.{}); 1222 const w = &writer.interface; 1223 try w.writeAll("POST /echo_body HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n5;ext=1\r\nHello\r\n6\r\n World\r\n0\r\n\r\n"); 1224 try w.flush(); 1225 1226 var buf: [200]u8 = undefined; 1227 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 11\r\n\r\nHello World", testReadAll(stream, &buf)); 1228} 1229 1230test "httpz: chunked body grows past pool buffer" { 1231 // pool buffer_size is 1024 on this server; send a 5000-byte body so the 1232 // grow path runs at least twice (1024 -> 2048 -> 4096 -> needed). 1233 const total: usize = 5000; 1234 const aa = t.arena.allocator(); 1235 defer t.reset(); 1236 1237 const body = aa.alloc(u8, total) catch unreachable; 1238 for (body, 0..) |*b, i| { 1239 b.* = @intCast('A' + (i % 26)); 1240 } 1241 1242 // Build the whole request up-front so it goes out in one shot — avoids 1243 // any partial-write/timeout interaction with the small SNDTIMEO on the 1244 // test socket. 1245 var aw: std.Io.Writer.Allocating = .init(aa); 1246 try aw.writer.writeAll("POST /echo_body HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n"); 1247 var sent: usize = 0; 1248 while (sent < total) { 1249 const take = @min(@as(usize, 1250), total - sent); 1250 try aw.writer.print("{x}\r\n", .{take}); 1251 try aw.writer.writeAll(body[sent .. sent + take]); 1252 try aw.writer.writeAll("\r\n"); 1253 sent += take; 1254 } 1255 try aw.writer.writeAll("0\r\n\r\n"); 1256 1257 const stream = testStream(6002); 1258 defer stream.close(t.io); 1259 var writer = stream.writer(t.io, &.{}); 1260 const w = &writer.interface; 1261 try w.writeAll(aw.written()); 1262 try w.flush(); 1263 1264 const buf = aa.alloc(u8, total + 1024) catch unreachable; 1265 const expected = std.fmt.allocPrint(aa, "HTTP/1.1 200 \r\nContent-Length: {d}\r\n\r\n{s}", .{ total, body }) catch unreachable; 1266 try t.expectString(expected, testReadAll(stream, buf)); 1267} 1268 1269test "httpz: chunked body too big" { 1270 // max_body_size is 32768; send 40000 bytes split into chunks. Server 1271 // should fail with 413 once the cap is exceeded. 1272 const total: usize = 40_000; 1273 const aa = t.arena.allocator(); 1274 defer t.reset(); 1275 1276 const body = aa.alloc(u8, total) catch unreachable; 1277 @memset(body, 'x'); 1278 1279 const stream = testStream(6002); 1280 defer stream.close(t.io); 1281 var writer = stream.writer(t.io, &.{}); 1282 const w = &writer.interface; 1283 try w.writeAll("POST /echo_body HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\n"); 1284 var sent: usize = 0; 1285 while (sent < total) { 1286 const take = @min(@as(usize, 4096), total - sent); 1287 try w.print("{x}\r\n", .{take}); 1288 try w.writeAll(body[sent .. sent + take]); 1289 try w.writeAll("\r\n"); 1290 sent += take; 1291 } 1292 try w.writeAll("0\r\n\r\n"); 1293 // The server may close the connection mid-write once it errors; ignore 1294 // a broken pipe. 1295 w.flush() catch {}; 1296 1297 var buf: [100]u8 = undefined; 1298 try t.expectString("HTTP/1.1 413 \r\nConnection: Close\r\nContent-Length: 23\r\n\r\nRequest body is too big", testReadAll(stream, &buf)); 1299} 1300 1301test "httpz: invalid transfer-encoding" { 1302 const stream = testStream(6002); 1303 defer stream.close(t.io); 1304 var writer = stream.writer(t.io, &.{}); 1305 const w = &writer.interface; 1306 try w.writeAll("POST /echo_body HTTP/1.1\r\nTransfer-Encoding: gzip\r\n\r\n"); 1307 try w.flush(); 1308 1309 var buf: [100]u8 = undefined; 1310 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1311} 1312 1313test "httpz: chunked with content-length is rejected" { 1314 const stream = testStream(6002); 1315 defer stream.close(t.io); 1316 var writer = stream.writer(t.io, &.{}); 1317 const w = &writer.interface; 1318 try w.writeAll("POST /echo_body HTTP/1.1\r\nContent-Length: 5\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nHello\r\n0\r\n\r\n"); 1319 try w.flush(); 1320 1321 var buf: [100]u8 = undefined; 1322 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1323} 1324 1325test "httpz: chunked with malformed chunk size" { 1326 const stream = testStream(6002); 1327 defer stream.close(t.io); 1328 var writer = stream.writer(t.io, &.{}); 1329 const w = &writer.interface; 1330 try w.writeAll("POST /echo_body HTTP/1.1\r\nTransfer-Encoding: chunked\r\n\r\nZZ\r\nnope\r\n0\r\n\r\n"); 1331 try w.flush(); 1332 1333 var buf: [100]u8 = undefined; 1334 try t.expectString("HTTP/1.1 400 \r\nConnection: Close\r\nContent-Length: 15\r\n\r\nInvalid Request", testReadAll(stream, &buf)); 1335} 1336 1337test "httpz: no route" { 1338 const stream = testStream(5992); 1339 defer stream.close(t.io); 1340 1341 var writer = stream.writer(t.io, &.{}); 1342 const w = &writer.interface; 1343 try w.writeAll("GET / HTTP/1.1\r\n\r\n"); 1344 try w.flush(); 1345 1346 var buf: [100]u8 = undefined; 1347 try t.expectString("HTTP/1.1 404 \r\nContent-Length: 9\r\n\r\nNot Found", testReadAll(stream, &buf)); 1348} 1349 1350test "httpz: no route with custom notFound handler" { 1351 const stream = testStream(5993); 1352 defer stream.close(t.io); 1353 var writer = stream.writer(t.io, &.{}); 1354 const w = &writer.interface; 1355 try w.writeAll("GET /not_found HTTP/1.1\r\n\r\n"); 1356 try w.flush(); 1357 1358 var buf: [100]u8 = undefined; 1359 try t.expectString("HTTP/1.1 404 \r\nstate: 3\r\nContent-Length: 10\r\n\r\nwhere lah?", testReadAll(stream, &buf)); 1360} 1361 1362test "httpz: unhandled exception" { 1363 std.testing.log_level = .err; 1364 defer std.testing.log_level = .warn; 1365 1366 const stream = testStream(5992); 1367 defer stream.close(t.io); 1368 var writer = stream.writer(t.io, &.{}); 1369 const w = &writer.interface; 1370 try w.writeAll("GET /fail HTTP/1.1\r\n\r\n"); 1371 try w.flush(); 1372 1373 var buf: [150]u8 = undefined; 1374 try t.expectString("HTTP/1.1 500 \r\nContent-Length: 21\r\n\r\nInternal Server Error", testReadAll(stream, &buf)); 1375} 1376 1377test "httpz: unhandled exception with custom error handler" { 1378 std.testing.log_level = .err; 1379 defer std.testing.log_level = .warn; 1380 1381 const stream = testStream(5993); 1382 defer stream.close(t.io); 1383 var writer = stream.writer(t.io, &.{}); 1384 const w = &writer.interface; 1385 try w.writeAll("GET /fail HTTP/1.1\r\n\r\n"); 1386 try w.flush(); 1387 1388 var buf: [150]u8 = undefined; 1389 try t.expectString("HTTP/1.1 500 \r\nstate: 3\r\nerr: TestUnhandledError\r\nContent-Length: 29\r\n\r\n#/why/arent/tags/hierarchical", testReadAll(stream, &buf)); 1390} 1391 1392test "httpz: custom methods" { 1393 const stream = testStream(5992); 1394 defer stream.close(t.io); 1395 1396 { 1397 var writer = stream.writer(t.io, &.{}); 1398 const w = &writer.interface; 1399 try w.writeAll("GET /test/method HTTP/1.1\r\n\r\n"); 1400 try w.flush(); 1401 var res = testReadParsed(stream); 1402 defer res.deinit(); 1403 try res.expectJson(.{ .method = "GET", .string = "" }); 1404 } 1405 1406 { 1407 var writer = stream.writer(t.io, &.{}); 1408 const w = &writer.interface; 1409 try w.writeAll("PUT /test/method HTTP/1.1\r\n\r\n"); 1410 try w.flush(); 1411 var res = testReadParsed(stream); 1412 defer res.deinit(); 1413 try res.expectJson(.{ .method = "PUT", .string = "" }); 1414 } 1415 1416 { 1417 var writer = stream.writer(t.io, &.{}); 1418 const w = &writer.interface; 1419 try w.writeAll("TEA /test/method HTTP/1.1\r\n\r\n"); 1420 try w.flush(); 1421 var res = testReadParsed(stream); 1422 defer res.deinit(); 1423 try res.expectJson(.{ .method = "OTHER", .string = "TEA" }); 1424 } 1425 1426 { 1427 var writer = stream.writer(t.io, &.{}); 1428 const w = &writer.interface; 1429 try w.writeAll("PING /test/method HTTP/1.1\r\n\r\n"); 1430 try w.flush(); 1431 var res = testReadParsed(stream); 1432 defer res.deinit(); 1433 try res.expectJson(.{ .method = "OTHER", .string = "PING" }); 1434 } 1435 1436 { 1437 var writer = stream.writer(t.io, &.{}); 1438 const w = &writer.interface; 1439 try w.writeAll("TEA /test/other HTTP/1.1\r\n\r\n"); 1440 try w.flush(); 1441 var buf: [100]u8 = undefined; 1442 try t.expectString("HTTP/1.1 404 \r\nContent-Length: 9\r\n\r\nNot Found", testReadAll(stream, &buf)); 1443 } 1444} 1445 1446test "httpz: route params" { 1447 const stream = testStream(5993); 1448 defer stream.close(t.io); 1449 var writer = stream.writer(t.io, &.{}); 1450 const w = &writer.interface; 1451 try w.writeAll("GET /api/v2/users/9001 HTTP/1.1\r\n\r\n"); 1452 try w.flush(); 1453 1454 var buf: [100]u8 = undefined; 1455 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 20\r\n\r\nversion=v2,user=9001", testReadAll(stream, &buf)); 1456} 1457 1458test "httpz: request and response headers" { 1459 const stream = testStream(5993); 1460 defer stream.close(t.io); 1461 var writer = stream.writer(t.io, &.{}); 1462 const w = &writer.interface; 1463 try w.writeAll("GET /test/headers HTTP/1.1\r\nHeader-Name: Header-Value\r\n\r\n"); 1464 try w.flush(); 1465 1466 var buf: [100]u8 = undefined; 1467 try t.expectString("HTTP/1.1 200 \r\nstate: 3\r\nEcho: Header-Value\r\nother: test-value\r\nContent-Length: 0\r\n\r\n", testReadAll(stream, &buf)); 1468} 1469 1470test "httpz: content-length body" { 1471 const stream = testStream(5993); 1472 defer stream.close(t.io); 1473 var writer = stream.writer(t.io, &.{}); 1474 const w = &writer.interface; 1475 try w.writeAll("GET /test/body/cl HTTP/1.1\r\nHeader-Name: Header-Value\r\nContent-Length: 4\r\n\r\nabcz"); 1476 try w.flush(); 1477 1478 var buf: [100]u8 = undefined; 1479 try t.expectString("HTTP/1.1 200 \r\nEcho-Body: abcz\r\nContent-Length: 0\r\n\r\n", testReadAll(stream, &buf)); 1480} 1481 1482test "httpz: json response" { 1483 const stream = testStream(5992); 1484 defer stream.close(t.io); 1485 var writer = stream.writer(t.io, &.{}); 1486 const w = &writer.interface; 1487 try w.writeAll("GET /test/json HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1488 try w.flush(); 1489 1490 var buf: [200]u8 = undefined; 1491 try t.expectString("HTTP/1.1 201 \r\nContent-Type: application/json; charset=UTF-8\r\nContent-Length: 26\r\n\r\n{\"over\":9000,\"teg\":\"soup\"}", testReadAll(stream, &buf)); 1492} 1493 1494test "httpz: query" { 1495 const stream = testStream(5992); 1496 defer stream.close(t.io); 1497 var writer = stream.writer(t.io, &.{}); 1498 const w = &writer.interface; 1499 try w.writeAll("GET /test/query?fav=keemun%20te%61%21 HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1500 try w.flush(); 1501 1502 var buf: [200]u8 = undefined; 1503 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 11\r\n\r\nkeemun tea!", testReadAll(stream, &buf)); 1504} 1505 1506test "httpz: chunked" { 1507 const stream = testStream(5992); 1508 defer stream.close(t.io); 1509 var writer = stream.writer(t.io, &.{}); 1510 const w = &writer.interface; 1511 try w.writeAll("GET /test/chunked HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1512 try w.flush(); 1513 1514 var buf: [1000]u8 = undefined; 1515 try t.expectString("HTTP/1.1 200 \r\nOver: 9000!\r\nTransfer-Encoding: chunked\r\n\r\n7\r\nChunk 1\r\n11\r\nand another chunk\r\n0\r\n\r\n", testReadAll(stream, &buf)); 1516} 1517 1518test "httpz: route-specific dispatcher" { 1519 const stream = testStream(5992); 1520 defer stream.close(t.io); 1521 var writer = stream.writer(t.io, &.{}); 1522 const w = &writer.interface; 1523 try w.writeAll("HEAD /test/dispatcher HTTP/1.1\r\n\r\n"); 1524 try w.flush(); 1525 1526 var buf: [200]u8 = undefined; 1527 try t.expectString("HTTP/1.1 200 \r\ndispatcher: test-dispatcher-1\r\nContent-Length: 6\r\n\r\naction", testReadAll(stream, &buf)); 1528} 1529 1530test "httpz: middlewares" { 1531 const stream = testStream(5992); 1532 defer stream.close(t.io); 1533 var writer = stream.writer(t.io, &.{}); 1534 const w = &writer.interface; 1535 1536 { 1537 try w.writeAll("GET /test/middlewares HTTP/1.1\r\nOrigin: httpz.local\r\n\r\n"); 1538 try w.flush(); 1539 var res = testReadParsed(stream); 1540 defer res.deinit(); 1541 1542 try res.expectJson(.{ .v1 = "tm1-100", .v2 = "tm2-100" }); 1543 try t.expectString("httpz.local", res.headers.get("Access-Control-Allow-Origin").?); 1544 } 1545} 1546 1547test "httpz: CORS" { 1548 const stream = testStream(5992); 1549 defer stream.close(t.io); 1550 1551 { 1552 var writer = stream.writer(t.io, &.{}); 1553 const w = &writer.interface; 1554 try w.writeAll("GET /echo HTTP/1.1\r\n\r\n"); 1555 try w.flush(); 1556 var res = testReadParsed(stream); 1557 defer res.deinit(); 1558 try t.expectEqual(null, res.headers.get("Access-Control-Max-Age")); 1559 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Methods")); 1560 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Headers")); 1561 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Origin")); 1562 } 1563 1564 { 1565 // cors endpoint but not cors options 1566 var writer = stream.writer(t.io, &.{}); 1567 const w = &writer.interface; 1568 try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nOrigin: httpz.local\r\nSec-Fetch-Mode: navigate\r\n\r\n"); 1569 try w.flush(); 1570 var res = testReadParsed(stream); 1571 defer res.deinit(); 1572 1573 try t.expectEqual(null, res.headers.get("Access-Control-Max-Age")); 1574 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Methods")); 1575 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Headers")); 1576 try t.expectString("httpz.local", res.headers.get("Access-Control-Allow-Origin").?); 1577 } 1578 1579 { 1580 // cors request 1581 var writer = stream.writer(t.io, &.{}); 1582 const w = &writer.interface; 1583 try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nOrigin: httpz.local\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1584 try w.flush(); 1585 var res = testReadParsed(stream); 1586 defer res.deinit(); 1587 1588 try t.expectString("300", res.headers.get("Access-Control-Max-Age").?); 1589 try t.expectString("GET,POST", res.headers.get("Access-Control-Allow-Methods").?); 1590 try t.expectString("content-type", res.headers.get("Access-Control-Allow-Headers").?); 1591 try t.expectString("httpz.local", res.headers.get("Access-Control-Allow-Origin").?); 1592 } 1593 1594 { 1595 // cors request, non-options 1596 var writer = stream.writer(t.io, &.{}); 1597 const w = &writer.interface; 1598 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: httpz.local\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1599 try w.flush(); 1600 var res = testReadParsed(stream); 1601 defer res.deinit(); 1602 1603 try t.expectEqual(null, res.headers.get("Access-Control-Max-Age")); 1604 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Methods")); 1605 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Headers")); 1606 try t.expectString("httpz.local", res.headers.get("Access-Control-Allow-Origin").?); 1607 } 1608} 1609 1610test "httpz: CORS wildcard origin" { 1611 const stream = testStream(5999); 1612 defer stream.close(t.io); 1613 1614 { 1615 var writer = stream.writer(t.io, &.{}); 1616 const w = &writer.interface; 1617 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://example.com\r\n\r\n"); 1618 try w.flush(); 1619 var res = testReadParsed(stream); 1620 defer res.deinit(); 1621 1622 try t.expectString("*", res.headers.get("Access-Control-Allow-Origin").?); 1623 } 1624 1625 { 1626 var writer = stream.writer(t.io, &.{}); 1627 const w = &writer.interface; 1628 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://any-domain.com\r\n\r\n"); 1629 try w.flush(); 1630 var res = testReadParsed(stream); 1631 defer res.deinit(); 1632 1633 try t.expectString("*", res.headers.get("Access-Control-Allow-Origin").?); 1634 } 1635 1636 { 1637 var writer = stream.writer(t.io, &.{}); 1638 const w = &writer.interface; 1639 try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nOrigin: https://test.com\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1640 try w.flush(); 1641 var res = testReadParsed(stream); 1642 defer res.deinit(); 1643 1644 try t.expectString("*", res.headers.get("Access-Control-Allow-Origin").?); 1645 try t.expectString("600", res.headers.get("Access-Control-Max-Age").?); 1646 try t.expectString("GET,POST,PUT", res.headers.get("Access-Control-Allow-Methods").?); 1647 try t.expectString("authorization,content-type", res.headers.get("Access-Control-Allow-Headers").?); 1648 try t.expectEqual(@as(u16, 204), res.status); 1649 } 1650 1651 { 1652 var writer = stream.writer(t.io, &.{}); 1653 const w = &writer.interface; 1654 try w.writeAll("GET /test/cors HTTP/1.1\r\n\r\n"); 1655 try w.flush(); 1656 var res = testReadParsed(stream); 1657 defer res.deinit(); 1658 1659 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Origin")); 1660 } 1661} 1662 1663test "httpz: CORS single origin" { 1664 const stream = testStream(6000); 1665 defer stream.close(t.io); 1666 1667 { 1668 var writer = stream.writer(t.io, &.{}); 1669 const w = &writer.interface; 1670 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://example.com\r\n\r\n"); 1671 try w.flush(); 1672 var res = testReadParsed(stream); 1673 defer res.deinit(); 1674 1675 try t.expectString("https://example.com", res.headers.get("Access-Control-Allow-Origin").?); 1676 try t.expectString("true", res.headers.get("Access-Control-Allow-Credentials").?); 1677 } 1678 1679 { 1680 var writer = stream.writer(t.io, &.{}); 1681 const w = &writer.interface; 1682 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://attacker.com\r\n\r\n"); 1683 try w.flush(); 1684 var res = testReadParsed(stream); 1685 defer res.deinit(); 1686 1687 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Origin")); 1688 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Credentials")); 1689 } 1690 1691 { 1692 var writer = stream.writer(t.io, &.{}); 1693 const w = &writer.interface; 1694 try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nOrigin: https://example.com\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1695 try w.flush(); 1696 var res = testReadParsed(stream); 1697 defer res.deinit(); 1698 1699 try t.expectString("https://example.com", res.headers.get("Access-Control-Allow-Origin").?); 1700 try t.expectString("true", res.headers.get("Access-Control-Allow-Credentials").?); 1701 try t.expectString("3600", res.headers.get("Access-Control-Max-Age").?); 1702 try t.expectEqual(@as(u16, 204), res.status); 1703 } 1704 1705 { 1706 var writer = stream.writer(t.io, &.{}); 1707 const w = &writer.interface; 1708 try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nOrigin: https://wrong.com\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1709 try w.flush(); 1710 var res = testReadParsed(stream); 1711 defer res.deinit(); 1712 1713 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Origin")); 1714 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Credentials")); 1715 } 1716} 1717 1718test "httpz: CORS multiple origins" { 1719 const stream = testStream(6001); 1720 defer stream.close(t.io); 1721 1722 { 1723 var writer = stream.writer(t.io, &.{}); 1724 const w = &writer.interface; 1725 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://example.com\r\n\r\n"); 1726 try w.flush(); 1727 var res = testReadParsed(stream); 1728 defer res.deinit(); 1729 1730 try t.expectString("https://example.com", res.headers.get("Access-Control-Allow-Origin").?); 1731 try t.expectString("true", res.headers.get("Access-Control-Allow-Credentials").?); 1732 } 1733 1734 { 1735 var writer = stream.writer(t.io, &.{}); 1736 const w = &writer.interface; 1737 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://api.example.com\r\n\r\n"); 1738 try w.flush(); 1739 var res = testReadParsed(stream); 1740 defer res.deinit(); 1741 1742 try t.expectString("https://api.example.com", res.headers.get("Access-Control-Allow-Origin").?); 1743 try t.expectString("true", res.headers.get("Access-Control-Allow-Credentials").?); 1744 } 1745 1746 { 1747 var writer = stream.writer(t.io, &.{}); 1748 const w = &writer.interface; 1749 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://test.local\r\n\r\n"); 1750 try w.flush(); 1751 var res = testReadParsed(stream); 1752 defer res.deinit(); 1753 1754 try t.expectString("https://test.local", res.headers.get("Access-Control-Allow-Origin").?); 1755 try t.expectString("true", res.headers.get("Access-Control-Allow-Credentials").?); 1756 } 1757 1758 { 1759 var writer = stream.writer(t.io, &.{}); 1760 const w = &writer.interface; 1761 try w.writeAll("GET /test/cors HTTP/1.1\r\nOrigin: https://attacker.com\r\n\r\n"); 1762 try w.flush(); 1763 var res = testReadParsed(stream); 1764 defer res.deinit(); 1765 1766 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Origin")); 1767 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Credentials")); 1768 } 1769 1770 { 1771 var writer = stream.writer(t.io, &.{}); 1772 const w = &writer.interface; 1773 try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nOrigin: https://api.example.com\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1774 try w.flush(); 1775 var res = testReadParsed(stream); 1776 defer res.deinit(); 1777 1778 try t.expectString("https://api.example.com", res.headers.get("Access-Control-Allow-Origin").?); 1779 try t.expectString("true", res.headers.get("Access-Control-Allow-Credentials").?); 1780 try t.expectString("GET,POST,DELETE", res.headers.get("Access-Control-Allow-Methods").?); 1781 try t.expectString("x-custom-header", res.headers.get("Access-Control-Allow-Headers").?); 1782 try t.expectEqual(@as(u16, 204), res.status); 1783 } 1784 1785 { 1786 var writer = stream.writer(t.io, &.{}); 1787 const w = &writer.interface; 1788 try w.writeAll("OPTIONS /test/cors HTTP/1.1\r\nOrigin: https://not-in-list.com\r\nSec-Fetch-Mode: cors\r\n\r\n"); 1789 try w.flush(); 1790 var res = testReadParsed(stream); 1791 defer res.deinit(); 1792 1793 try t.expectEqual(null, res.headers.get("Access-Control-Allow-Origin")); 1794 } 1795} 1796 1797test "httpz: router groups" { 1798 const stream = testStream(5993); 1799 defer stream.close(t.io); 1800 1801 { 1802 var writer = stream.writer(t.io, &.{}); 1803 const w = &writer.interface; 1804 try w.writeAll("GET / HTTP/1.1\r\n\r\n"); 1805 try w.flush(); 1806 var res = testReadParsed(stream); 1807 defer res.deinit(); 1808 1809 try res.expectJson(.{ .state = 3, .method = "GET", .path = "/" }); 1810 try t.expectEqual(true, res.headers.get("dispatcher") == null); 1811 } 1812 1813 { 1814 var writer = stream.writer(t.io, &.{}); 1815 const w = &writer.interface; 1816 try w.writeAll("GET /admin/users HTTP/1.1\r\n\r\n"); 1817 try w.flush(); 1818 var res = testReadParsed(stream); 1819 defer res.deinit(); 1820 1821 try res.expectJson(.{ .state = 99, .method = "GET", .path = "/admin/users" }); 1822 try t.expectString("test-dispatcher-2", res.headers.get("dispatcher").?); 1823 } 1824 1825 { 1826 var writer = stream.writer(t.io, &.{}); 1827 const w = &writer.interface; 1828 try w.writeAll("PUT /admin/users/:id HTTP/1.1\r\n\r\n"); 1829 try w.flush(); 1830 var res = testReadParsed(stream); 1831 defer res.deinit(); 1832 1833 try res.expectJson(.{ .state = 99, .method = "PUT", .path = "/admin/users/:id" }); 1834 try t.expectString("test-dispatcher-2", res.headers.get("dispatcher").?); 1835 } 1836 1837 { 1838 var writer = stream.writer(t.io, &.{}); 1839 const w = &writer.interface; 1840 try w.writeAll("HEAD /debug/ping HTTP/1.1\r\n\r\n"); 1841 try w.flush(); 1842 var res = testReadParsed(stream); 1843 defer res.deinit(); 1844 1845 try res.expectJson(.{ .state = 20, .method = "HEAD", .path = "/debug/ping" }); 1846 try t.expectString("test-dispatcher-3", res.headers.get("dispatcher").?); 1847 } 1848 1849 { 1850 var writer = stream.writer(t.io, &.{}); 1851 const w = &writer.interface; 1852 try w.writeAll("OPTIONS /debug/stats HTTP/1.1\r\n\r\n"); 1853 try w.flush(); 1854 var res = testReadParsed(stream); 1855 defer res.deinit(); 1856 1857 try res.expectJson(.{ .state = 20, .method = "OPTIONS", .path = "/debug/stats" }); 1858 try t.expectString("test-dispatcher-3", res.headers.get("dispatcher").?); 1859 } 1860 1861 { 1862 var writer = stream.writer(t.io, &.{}); 1863 const w = &writer.interface; 1864 try w.writeAll("POST /login HTTP/1.1\r\n\r\n"); 1865 try w.flush(); 1866 var res = testReadParsed(stream); 1867 defer res.deinit(); 1868 1869 try res.expectJson(.{ .state = 3, .method = "POST", .path = "/login" }); 1870 try t.expectEqual(true, res.headers.get("dispatcher") == null); 1871 } 1872} 1873 1874test "httpz: event stream" { 1875 const stream = testStream(5992); 1876 defer stream.close(t.io); 1877 var writer = stream.writer(t.io, &.{}); 1878 const w = &writer.interface; 1879 try w.writeAll("GET /test/stream HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1880 try w.flush(); 1881 1882 var res = testReadParsed(stream); 1883 defer res.deinit(); 1884 1885 try t.expectEqual(818, res.status); 1886 try t.expectEqual(true, res.headers.get("Content-Length") == null); 1887 try t.expectString("text/event-stream; charset=UTF-8", res.headers.get("Content-Type").?); 1888 try t.expectString("no-cache", res.headers.get("Cache-Control").?); 1889 try t.expectString("keep-alive", res.headers.get("Connection").?); 1890 try t.expectString("helloa message", res.body); 1891} 1892 1893test "httpz: event stream sync" { 1894 const stream = testStream(5992); 1895 defer stream.close(t.io); 1896 var writer = stream.writer(t.io, &.{}); 1897 const w = &writer.interface; 1898 try w.writeAll("GET /test/streamsync HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1899 try w.flush(); 1900 1901 var res = testReadParsed(stream); 1902 defer res.deinit(); 1903 1904 try t.expectEqual(818, res.status); 1905 try t.expectEqual(true, res.headers.get("Content-Length") == null); 1906 try t.expectString("text/event-stream; charset=UTF-8", res.headers.get("Content-Type").?); 1907 try t.expectString("no-cache", res.headers.get("Cache-Control").?); 1908 try t.expectString("keep-alive", res.headers.get("Connection").?); 1909 try t.expectString("helloa sync message", res.body); 1910} 1911 1912test "httpz: keepalive" { 1913 const stream = testStream(5993); 1914 defer stream.close(t.io); 1915 var writer = stream.writer(t.io, &.{}); 1916 const w = &writer.interface; 1917 try w.writeAll("GET /api/v2/users/9001 HTTP/1.1\r\n\r\n"); 1918 try w.flush(); 1919 1920 var buf: [100]u8 = undefined; 1921 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 20\r\n\r\nversion=v2,user=9001", testReadAll(stream, &buf)); 1922 1923 try w.writeAll("GET /api/v2/users/123 HTTP/1.1\r\n\r\n"); 1924 try w.flush(); 1925 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 19\r\n\r\nversion=v2,user=123", testReadAll(stream, &buf)); 1926} 1927 1928test "httpz: route data" { 1929 const stream = testStream(5992); 1930 defer stream.close(t.io); 1931 var writer = stream.writer(t.io, &.{}); 1932 const w = &writer.interface; 1933 try w.writeAll("GET /test/route_data HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 1934 try w.flush(); 1935 1936 var res = testReadParsed(stream); 1937 defer res.deinit(); 1938 try res.expectJson(.{ .power = 12345 }); 1939} 1940 1941test "httpz: keepalive with explicit write" { 1942 const stream = testStream(5993); 1943 defer stream.close(t.io); 1944 var writer = stream.writer(t.io, &.{}); 1945 const w = &writer.interface; 1946 try w.writeAll("GET /write/9001 HTTP/1.1\r\n\r\n"); 1947 try w.flush(); 1948 1949 var buf: [1000]u8 = undefined; 1950 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 47\r\n\r\n{\"state\":3,\"method\":\"GET\",\"path\":\"/write/9001\"}", testReadAll(stream, &buf)); 1951 1952 try w.writeAll("GET /write/123 HTTP/1.1\r\n\r\n"); 1953 try w.flush(); 1954 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 46\r\n\r\n{\"state\":3,\"method\":\"GET\",\"path\":\"/write/123\"}", testReadAll(stream, &buf)); 1955} 1956 1957test "httpz: request in chunks" { 1958 const stream = testStream(5993); 1959 defer stream.close(t.io); 1960 var writer = stream.writer(t.io, &.{}); 1961 const w = &writer.interface; 1962 try w.writeAll("GET /api/v2/use"); 1963 try w.flush(); 1964 try t.io.sleep(.fromMilliseconds(10), .awake); 1965 try w.writeAll("rs/11 HTTP/1.1\r\n\r\n"); 1966 try w.flush(); 1967 1968 var buf: [100]u8 = undefined; 1969 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 18\r\n\r\nversion=v2,user=11", testReadAll(stream, &buf)); 1970} 1971 1972test "httpz: writer re-use" { 1973 defer t.reset(); 1974 1975 const stream = testStream(5996); 1976 defer stream.close(t.io); 1977 var writer = stream.writer(t.io, &.{}); 1978 const w = &writer.interface; 1979 1980 var expected: [10]TestUser = undefined; 1981 1982 var buf: [100]u8 = undefined; 1983 for (0..10) |i| { 1984 expected[i] = .{ 1985 .id = try std.fmt.allocPrint(t.arena.allocator(), "id-{d}", .{i}), 1986 .power = i, 1987 }; 1988 try w.writeAll(try std.fmt.bufPrint(&buf, "GET /test/writer?count={d} HTTP/1.1\r\nContent-Length: 0\r\n\r\n", .{i + 1})); 1989 try w.flush(); 1990 1991 var res = testReadParsed(stream); 1992 defer res.deinit(); 1993 1994 try res.expectJson(.{ .data = expected[0 .. i + 1] }); 1995 } 1996} 1997 1998test "httpz: custom dispatch without action context" { 1999 const stream = testStream(5994); 2000 defer stream.close(t.io); 2001 var writer = stream.writer(t.io, &.{}); 2002 const w = &writer.interface; 2003 try w.writeAll("GET / HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 2004 try w.flush(); 2005 2006 var buf: [200]u8 = undefined; 2007 try t.expectString("HTTP/1.1 200 \r\nContent-Type: application/json; charset=UTF-8\r\ndstate: 10\r\ndispatch: TestHandlerDispatch\r\nContent-Length: 12\r\n\r\n{\"state\":10}", testReadAll(stream, &buf)); 2008} 2009 2010test "httpz: custom dispatch with action context" { 2011 const stream = testStream(5995); 2012 defer stream.close(t.io); 2013 var writer = stream.writer(t.io, &.{}); 2014 const w = &writer.interface; 2015 try w.writeAll("GET /?name=teg HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 2016 try w.flush(); 2017 2018 var buf: [200]u8 = undefined; 2019 try t.expectString("HTTP/1.1 200 \r\nContent-Type: application/json; charset=UTF-8\r\ndstate: 20\r\ndispatch: TestHandlerDispatchContext\r\nContent-Length: 12\r\n\r\n{\"other\":30}", testReadAll(stream, &buf)); 2020} 2021 2022test "httpz: custom handle" { 2023 const stream = testStream(5997); 2024 defer stream.close(t.io); 2025 var writer = stream.writer(t.io, &.{}); 2026 const w = &writer.interface; 2027 try w.writeAll("GET /whatever?name=teg HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 2028 try w.flush(); 2029 2030 var buf: [100]u8 = undefined; 2031 try t.expectString("HTTP/1.1 200 \r\nContent-Length: 9\r\n\r\nhello teg", testReadAll(stream, &buf)); 2032} 2033 2034test "httpz: request body reader" { 2035 { 2036 // no body 2037 const stream = testStream(5992); 2038 defer stream.close(t.io); 2039 2040 var writer = stream.writer(t.io, &.{}); 2041 const w = &writer.interface; 2042 2043 try w.writeAll("GET /test/req_reader HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 2044 try w.flush(); 2045 2046 var res = testReadParsed(stream); 2047 defer res.deinit(); 2048 try res.expectJson(.{ .length = 0 }); 2049 } 2050 2051 { 2052 // small body 2053 const stream = testStream(5992); 2054 defer stream.close(t.io); 2055 2056 var writer = stream.writer(t.io, &.{}); 2057 const w = &writer.interface; 2058 2059 try w.writeAll("GET /test/req_reader HTTP/1.1\r\nContent-Length: 4\r\n\r\n123z"); 2060 try w.flush(); 2061 2062 var res = testReadParsed(stream); 2063 defer res.deinit(); 2064 try res.expectJson(.{ .length = 4 }); 2065 } 2066 2067 comptime var length = 4043; 2068 { 2069 // medium body 2070 const stream = testStream(5992); 2071 defer stream.close(t.io); 2072 2073 var writer = stream.writer(t.io, &.{}); 2074 const w = &writer.interface; 2075 2076 try w.writeAll(std.fmt.comptimePrint("GET /test/req_reader HTTP/1.1\r\nContent-Length: {d}\r\n\r\n" ++ ("a" ** length), .{length})); 2077 try w.flush(); 2078 2079 var res = testReadParsed(stream); 2080 defer res.deinit(); 2081 try res.expectJson(.{ .length = length }); 2082 } 2083 2084 var r = t.getRandom(); 2085 const random = r.random(); 2086 length = 20000; 2087 2088 // a bit of fuzzing 2089 for (0..10) |_| { 2090 const stream = testStream(5992); 2091 defer stream.close(t.io); 2092 2093 var buf: [1024]u8 = undefined; 2094 var writer = stream.writer(t.io, &buf); 2095 const w = &writer.interface; 2096 2097 var req: []const u8 = std.fmt.comptimePrint("GET /test/req_reader HTTP/1.1\r\nContent-Length: {d}\r\n\r\n" ++ ("a" ** length), .{length}); 2098 while (req.len > 0) { 2099 const len = random.uintAtMost(usize, req.len - 1) + 1; 2100 try w.writeAll(req[0..len]); 2101 try t.io.sleep(.fromMilliseconds(2), .awake); 2102 req = req[len..]; 2103 } 2104 2105 try w.flush(); 2106 2107 var res = testReadParsed(stream); 2108 defer res.deinit(); 2109 try res.expectJson(.{ .length = length }); 2110 } 2111} 2112 2113test "websocket: invalid request" { 2114 const stream = testStream(5998); 2115 defer stream.close(t.io); 2116 var writer = stream.writer(t.io, &.{}); 2117 const w = &writer.interface; 2118 try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n\r\n"); 2119 try w.flush(); 2120 2121 var res = testReadParsed(stream); 2122 defer res.deinit(); 2123 try t.expectString("invalid websocket", res.body); 2124} 2125 2126test "websocket: upgrade" { 2127 const stream = testStream(5998); 2128 defer stream.close(t.io); 2129 var writer = stream.writer(t.io, &.{}); 2130 const w = &writer.interface; 2131 try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n"); 2132 try w.writeAll("upgrade: WEBsocket\r\n"); 2133 try w.writeAll("Sec-Websocket-verSIon: 13\r\n"); 2134 try w.writeAll("ConnectioN: abc,upgrade,123\r\n"); 2135 try w.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n"); 2136 try w.flush(); 2137 2138 var res = testReadHeader(stream); 2139 defer res.deinit(); 2140 try t.expectEqual(101, res.status); 2141 try t.expectString("websocket", res.headers.get("Upgrade").?); 2142 try t.expectString("upgrade", res.headers.get("Connection").?); 2143 try t.expectString("55eM2SNGu+68v5XXrr982mhPFkU=", res.headers.get("Sec-Websocket-Accept").?); 2144 2145 try w.writeAll(&websocket.frameText("over 9000!")); 2146 2147 // https://github.com/karlseguin/http.zig/pull/188 2148 try w.flush(); 2149 try t.io.sleep(.fromMilliseconds(5), .awake); 2150 2151 try w.writeAll(&websocket.frameText("close")); 2152 try w.flush(); 2153 2154 var pos: usize = 0; 2155 var buf: [100]u8 = undefined; 2156 var wait_count: usize = 0; 2157 2158 while (pos < 16) { 2159 const n = posix.read(stream.socket.handle, buf[pos..]) catch |err| switch (err) { 2160 error.WouldBlock => { 2161 if (wait_count == 100) { 2162 break; 2163 } 2164 wait_count += 1; 2165 try t.io.sleep(.fromMilliseconds(1), .awake); 2166 continue; 2167 }, 2168 else => return err, 2169 }; 2170 if (n == 0) { 2171 break; 2172 } 2173 pos += n; 2174 } 2175 try t.expectEqual(16, pos); 2176 try t.expectEqual(129, buf[0]); 2177 try t.expectEqual(10, buf[1]); 2178 try t.expectString("over 9000!", buf[2..12]); 2179 try t.expectString(&.{ 136, 2, 3, 232 }, buf[12..16]); 2180} 2181 2182test "websocket: rejected upgrade preserves HTTP socket ownership" { 2183 const stream = testStream(5998); 2184 defer stream.close(t.io); 2185 var writer = stream.writer(t.io, &.{}); 2186 try writer.interface.writeAll("GET /ws-reject HTTP/1.1\r\nUpgrade: websocket\r\nSec-WebSocket-Version: 13\r\nConnection: upgrade\r\nSec-WebSocket-Key: a-secret-key\r\n\r\n"); 2187 try writer.interface.flush(); 2188 2189 var res = testReadParsed(stream); 2190 defer res.deinit(); 2191 try t.expectEqual(400, res.status); 2192 try t.expectString("rejected upgrade", res.body); 2193} 2194 2195// Regression for EPOLLONESHOT/KQUEUE one-shot re-arm ordering. The client 2196// sends the next frame immediately after receiving the previous echo, keeping 2197// reads near the handler's release/re-arm boundary. Re-arming while the Conn's 2198// `processing` flag is still held can consume that read, fail dispatch, and 2199// leave the websocket permanently unarmed. 2200test "websocket: sequential echoes survive every rearm boundary" { 2201 if (force_blocking) return; 2202 2203 const stream = testStream(5998); 2204 defer stream.close(t.io); 2205 var writer = stream.writer(t.io, &.{}); 2206 const w = &writer.interface; 2207 try w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n"); 2208 try w.writeAll("Upgrade: websocket\r\n"); 2209 try w.writeAll("Sec-WebSocket-Version: 13\r\n"); 2210 try w.writeAll("Connection: upgrade\r\n"); 2211 try w.writeAll("Sec-WebSocket-Key: a-secret-key\r\n\r\n"); 2212 try w.flush(); 2213 2214 var res = testReadHeader(stream); 2215 defer res.deinit(); 2216 try t.expectEqual(101, res.status); 2217 2218 for (0..500) |_| { 2219 try w.writeAll(&websocket.frameText("x")); 2220 try w.flush(); 2221 2222 var echo: [3]u8 = undefined; 2223 var read: usize = 0; 2224 while (read < echo.len) { 2225 const n = posix.read(stream.socket.handle, echo[read..]) catch |err| switch (err) { 2226 error.WouldBlock => return error.WebsocketRearmLostRead, 2227 else => return err, 2228 }; 2229 if (n == 0) return error.WebsocketClosedBeforeEcho; 2230 read += n; 2231 } 2232 try t.expectString(&.{ 0x81, 0x01, 'x' }, &echo); 2233 } 2234 2235 try w.writeAll(&websocket.frameText("close")); 2236 try w.flush(); 2237} 2238 2239// Stress test: multiple concurrent websocket clients sending many messages each. 2240// Run repeatedly (e.g. zig build test -Dtest-filter="websocket: stress" or run 50x) 2241// to verify no race in reader.done() / allocator free when using thread pool. 2242test "websocket: stress" { 2243 if (force_blocking) return; // non-blocking mode only (thread pool) 2244 const num_clients = 8; 2245 const messages_per_client = 150; 2246 2247 // When run with -Dtest-filter="websocket: stress", tests:beforeAll may not run, 2248 // so nothing is listening on 5998. Wait for port and start our own server if needed. 2249 var stress_server: ?Server(TestWebsocketHandler) = null; 2250 var stress_listen_thread: ?Thread = null; 2251 testing.waitForPort(5998) catch { 2252 stress_server = try Server(TestWebsocketHandler).init(t.io, t.allocator, .{ .address = .localhost(5998) }, TestWebsocketHandler{}); 2253 var router = try stress_server.?.router(.{}); 2254 router.get("/ws", TestWebsocketHandler.upgrade, .{}); 2255 stress_listen_thread = try stress_server.?.listenInNewThread(); 2256 try testing.waitForPort(5998); 2257 }; 2258 defer if (stress_server) |*srv| { 2259 srv.stop(); 2260 if (stress_listen_thread) |thrd| thrd.join(); 2261 srv.deinit(); 2262 }; 2263 2264 var threads: [num_clients]Thread = undefined; 2265 for (0..num_clients) |i| { 2266 threads[i] = Thread.spawn(.{}, struct { 2267 fn run(_: usize) void { 2268 const stream = testStream(5998); 2269 defer stream.close(t.io); 2270 2271 var writer = stream.writer(t.io, &.{}); 2272 const w = &writer.interface; 2273 w.writeAll("GET /ws HTTP/1.1\r\nContent-Length: 0\r\n") catch return; 2274 w.writeAll("upgrade: WEBsocket\r\n") catch return; 2275 w.writeAll("Sec-Websocket-verSIon: 13\r\n") catch return; 2276 w.writeAll("ConnectioN: upgrade\r\n") catch return; 2277 w.writeAll("SEC-WEBSOCKET-KeY: a-secret-key\r\n\r\n") catch return; 2278 w.flush() catch return; 2279 2280 var buf: [1024]u8 = undefined; 2281 var pos: usize = 0; 2282 2283 while (!std.mem.endsWith(u8, buf[0..pos], "\r\n\r\n")) { 2284 if (pos >= buf.len) { 2285 return; 2286 } 2287 const n = posix.read(stream.socket.handle, buf[pos..]) catch return; 2288 if (n == 0) { 2289 return; 2290 } 2291 pos += n; 2292 } 2293 if (pos < 12 or !std.mem.startsWith(u8, buf[0..12], "HTTP/1.1 101")) return; 2294 2295 for (0..messages_per_client) |_| { 2296 const frame = websocket.frameText("stress"); 2297 w.writeAll(&frame) catch return; 2298 } 2299 w.writeAll(&websocket.frameText("close")) catch return; 2300 w.flush() catch return; 2301 } 2302 }.run, .{i}) catch return; 2303 } 2304 for (&threads) |*th| th.join(); 2305} 2306 2307// Regression: processSignal used to release handover-snapshot nodes via 2308// disown(), whose List.remove ran against the node's stale snapshot links and 2309// grafted the rest of the snapshot back onto the live handover_list. When a 2310// close-handover preceded a websocket upgrade in one signal batch, the next 2311// signal re-processed the (now websocket) node: union confusion in Debug, and 2312// in ReleaseFast a double pool-release that let two live sockets share one 2313// HTTPConn — responses interleaved across unrelated connections. 2314// 2315// Each round races a Connection:close request against a websocket upgrade so 2316// both handovers land in one snapshot, then fires more signals and verifies 2317// on fresh connections that every response still matches its request. Under 2318// the old code this panics within a few rounds ("access of union field 'http' 2319// while field 'websocket' is active"). 2320test "httpz: handover snapshot survives close + websocket upgrade in one batch" { 2321 // Keep the regression independently runnable with -Dtest-filter; the 2322 // shared beforeAll server is absent when Zig filters the suite. 2323 var own_server: ?Server(TestWebsocketHandler) = null; 2324 var own_thread: ?Thread = null; 2325 testing.waitForPort(6003) catch { 2326 own_server = try Server(TestWebsocketHandler).init(t.io, t.allocator, .{ 2327 .address = .localhost(6003), 2328 .workers = .{ .count = 1 }, 2329 }, TestWebsocketHandler{}); 2330 var router = try own_server.?.router(.{}); 2331 router.get("/echo", TestWebsocketHandler.echoTag, .{}); 2332 router.get("/ws", TestWebsocketHandler.upgrade, .{}); 2333 own_thread = try own_server.?.listenInNewThread(); 2334 try testing.waitForPort(6003); 2335 }; 2336 defer if (own_server) |*srv| { 2337 srv.stop(); 2338 if (own_thread) |thread| thread.join(); 2339 srv.deinit(); 2340 }; 2341 var tag_buf: [64]u8 = undefined; 2342 var req_buf: [128]u8 = undefined; 2343 for (0..200) |round| { 2344 // 1) close-handover, sent without reading so it stays in flight 2345 const c_close = testStreamWithTimeout(6003, 250_000); 2346 defer c_close.close(t.io); 2347 var close_writer = c_close.writer(t.io, &.{}); 2348 const cw = &close_writer.interface; 2349 const close_tag = std.fmt.bufPrint(&tag_buf, "close-{d}", .{round}) catch unreachable; 2350 try cw.writeAll(std.fmt.bufPrint(&req_buf, "GET /echo?tag={s} HTTP/1.1\r\nConnection: close\r\n\r\n", .{close_tag}) catch unreachable); 2351 try cw.flush(); 2352 2353 // 2) websocket upgrade racing into the same signal batch 2354 const c_ws = testStreamWithTimeout(6003, 250_000); 2355 var ws_writer = c_ws.writer(t.io, &.{}); 2356 const ww = &ws_writer.interface; 2357 try ww.writeAll("GET /ws HTTP/1.1\r\nupgrade: websocket\r\nSec-Websocket-Version: 13\r\nConnection: upgrade\r\nSec-Websocket-Key: a-secret-key\r\n\r\n"); 2358 try ww.flush(); 2359 2360 var close_res = testReadParsedChecked(c_close) catch |err| { 2361 std.debug.print("snapshot regression timed out at round {d} close response: {}\n", .{ round, err }); 2362 return err; 2363 }; 2364 defer close_res.deinit(); 2365 try t.expectString(close_tag, close_res.body); 2366 2367 var ws_res = testReadHeader(c_ws); 2368 defer ws_res.deinit(); 2369 try t.expectEqual(101, ws_res.status); 2370 c_ws.close(t.io); 2371 2372 // 3) more close-handovers: each signals the reactor again, which under 2373 // the old code re-walked the grafted snapshot remainder 2374 for (0..2) |extra| { 2375 const c3 = testStreamWithTimeout(6003, 250_000); 2376 defer c3.close(t.io); 2377 var w3 = c3.writer(t.io, &.{}); 2378 const tag3 = std.fmt.bufPrint(&tag_buf, "extra-{d}-{d}", .{ round, extra }) catch unreachable; 2379 try w3.interface.writeAll(std.fmt.bufPrint(&req_buf, "GET /echo?tag={s} HTTP/1.1\r\nConnection: close\r\n\r\n", .{tag3}) catch unreachable); 2380 try w3.interface.flush(); 2381 var res3 = testReadParsedChecked(c3) catch |err| { 2382 std.debug.print("snapshot regression timed out at round {d} extra {d}: {}\n", .{ round, extra, err }); 2383 return err; 2384 }; 2385 defer res3.deinit(); 2386 try t.expectString(tag3, res3.body); 2387 } 2388 2389 // 4) keepalive canary: both sequential responses must match their tags 2390 const canary = testStreamWithTimeout(6003, 250_000); 2391 defer canary.close(t.io); 2392 var cwr = canary.writer(t.io, &.{}); 2393 for (0..2) |i| { 2394 const tag = std.fmt.bufPrint(&tag_buf, "canary-{d}-{d}", .{ round, i }) catch unreachable; 2395 try cwr.interface.writeAll(std.fmt.bufPrint(&req_buf, "GET /echo?tag={s} HTTP/1.1\r\n\r\n", .{tag}) catch unreachable); 2396 try cwr.interface.flush(); 2397 var res = testReadParsedChecked(canary) catch |err| { 2398 std.debug.print("snapshot regression timed out at round {d} canary {d}: {}\n", .{ round, i, err }); 2399 return err; 2400 }; 2401 defer res.deinit(); 2402 try t.expectString(tag, res.body); 2403 } 2404 } 2405} 2406 2407// Regression: a websocket connection's Conn node was never handed back to the 2408// reactor when the connection died, so every websocket permanently consumed a 2409// conn slot. Once a worker's len reached max_conn it paused accept forever — 2410// the server wedged. With max_conn=4, twenty ws connect/close cycles wedge the 2411// old code almost immediately; the fix drains dead ws nodes via processSignal. 2412test "httpz: websocket connections release their conn slot" { 2413 var tag_buf: [32]u8 = undefined; 2414 var req_buf: [96]u8 = undefined; 2415 for (0..20) |i| { 2416 const c_ws = testStream(6004); 2417 var ws_writer = c_ws.writer(t.io, &.{}); 2418 const ww = &ws_writer.interface; 2419 try ww.writeAll("GET /ws HTTP/1.1\r\nupgrade: websocket\r\nSec-Websocket-Version: 13\r\nConnection: upgrade\r\nSec-Websocket-Key: a-secret-key\r\n\r\n"); 2420 try ww.flush(); 2421 var ws_res = testReadHeader(c_ws); 2422 defer ws_res.deinit(); 2423 try t.expectEqual(101, ws_res.status); 2424 // server-initiated close via the handler's "close" message, so the 2425 // dead connection goes through the thread pool cleanup path 2426 try ww.writeAll(&websocket.frameText("close")); 2427 try ww.flush(); 2428 var drain_buf: [64]u8 = undefined; 2429 _ = testReadAll(c_ws, &drain_buf); 2430 c_ws.close(t.io); 2431 2432 // give the reactor a beat to drain the graveyard, then prove the 2433 // worker still accepts and serves 2434 try t.io.sleep(.fromMilliseconds(5), .awake); 2435 const canary = testStream(6004); 2436 defer canary.close(t.io); 2437 var cw = canary.writer(t.io, &.{}); 2438 const tag = std.fmt.bufPrint(&tag_buf, "slot-{d}", .{i}) catch unreachable; 2439 try cw.interface.writeAll(std.fmt.bufPrint(&req_buf, "GET /echo?tag={s} HTTP/1.1\r\nConnection: close\r\n\r\n", .{tag}) catch unreachable); 2440 try cw.interface.flush(); 2441 var res = testReadParsed(canary); 2442 defer res.deinit(); 2443 try t.expectString(tag, res.body); 2444 } 2445} 2446 2447// A websocket closes its socket on a pool thread. The reactor used to issue a 2448// second EV_DELETE/EPOLL_CTL_DEL later while draining its graveyard. If accept 2449// reused the numeric fd in between, that stale delete removed the new HTTP 2450// connection's read monitor and left it permanently stuck in request_list. 2451test "websocket: closed fd cleanup cannot delete reused HTTP monitor" { 2452 if (force_blocking) return; 2453 2454 var server = try Server(TestWebsocketHandler).init(t.io, t.allocator, .{ 2455 .address = .localhost(6005), 2456 .workers = .{ .count = 1, .max_conn = 2, .min_conn = 2 }, 2457 }, TestWebsocketHandler{}); 2458 var router = try server.router(.{}); 2459 router.get("/echo", TestWebsocketHandler.echoTag, .{}); 2460 router.get("/ws", TestWebsocketHandler.upgrade, .{}); 2461 const thread = try server.listenInNewThread(); 2462 defer { 2463 server.stop(); 2464 thread.join(); 2465 server.deinit(); 2466 } 2467 try testing.waitForPort(6005); 2468 2469 var tag_buf: [32]u8 = undefined; 2470 var req_buf: [96]u8 = undefined; 2471 for (0..200) |round| { 2472 const ws_stream = testStreamWithTimeout(6005, 250_000); 2473 var ws_writer = ws_stream.writer(t.io, &.{}); 2474 try ws_writer.interface.writeAll("GET /ws HTTP/1.1\r\nUpgrade: websocket\r\nSec-WebSocket-Version: 13\r\nConnection: upgrade\r\nSec-WebSocket-Key: a-secret-key\r\n\r\n"); 2475 try ws_writer.interface.flush(); 2476 var ws_res = testReadHeader(ws_stream); 2477 try t.expectEqual(101, ws_res.status); 2478 ws_res.deinit(); 2479 ws_stream.close(t.io); 2480 2481 const http_stream = testStreamWithTimeout(6005, 250_000); 2482 defer http_stream.close(t.io); 2483 var http_writer = http_stream.writer(t.io, &.{}); 2484 const tag = std.fmt.bufPrint(&tag_buf, "reused-{d}", .{round}) catch unreachable; 2485 try http_writer.interface.writeAll(std.fmt.bufPrint(&req_buf, "GET /echo?tag={s} HTTP/1.1\r\nConnection: close\r\n\r\n", .{tag}) catch unreachable); 2486 try http_writer.interface.flush(); 2487 var res = testReadParsedChecked(http_stream) catch |err| { 2488 std.debug.print("fd-reuse regression timed out at round {d}: {}\n", .{ round, err }); 2489 return err; 2490 }; 2491 defer res.deinit(); 2492 try t.expectString(tag, res.body); 2493 } 2494} 2495 2496// Timed-out HTTP connections used to skip explicit reactor deregistration and 2497// immediately return both Conn and HTTPConn storage to their pools. Under fd 2498// churn, a queued readiness event could then carry the retired pointer after 2499// that storage had been assigned to a different socket. Exercise a whole 2500// expired keepalive cohort followed by aggressive reuse on a tiny one-worker 2501// server; every response must remain paired with its request. 2502test "httpz: keepalive timeout deregisters before connection reuse" { 2503 if (force_blocking) return; 2504 2505 var server = try Server(TestWebsocketHandler).init(t.io, t.allocator, .{ 2506 .address = .localhost(6006), 2507 .workers = .{ .count = 1, .max_conn = 64, .min_conn = 64 }, 2508 .timeout = .{ .keepalive = 0 }, 2509 }, TestWebsocketHandler{}); 2510 var router = try server.router(.{}); 2511 router.get("/echo", TestWebsocketHandler.echoTag, .{}); 2512 const thread = try server.listenInNewThread(); 2513 defer { 2514 server.stop(); 2515 thread.join(); 2516 server.deinit(); 2517 } 2518 try testing.waitForPort(6006); 2519 2520 var idle: [48]Io.net.Stream = undefined; 2521 var opened: usize = 0; 2522 defer for (idle[0..opened]) |stream| stream.close(t.io); 2523 2524 var req_buf: [128]u8 = undefined; 2525 var tag_buf: [48]u8 = undefined; 2526 for (&idle, 0..) |*stream, i| { 2527 stream.* = testStreamWithTimeout(6006, 500_000); 2528 opened += 1; 2529 var writer = stream.writer(t.io, &.{}); 2530 const tag = std.fmt.bufPrint(&tag_buf, "idle-{d}", .{i}) catch unreachable; 2531 try writer.interface.writeAll(std.fmt.bufPrint(&req_buf, "GET /echo?tag={s} HTTP/1.1\r\n\r\n", .{tag}) catch unreachable); 2532 try writer.interface.flush(); 2533 var res = try testReadParsedChecked(stream.*); 2534 defer res.deinit(); 2535 try t.expectString(tag, res.body); 2536 } 2537 2538 // Timeout collection runs on the worker's one-second maintenance tick. 2539 try t.io.sleep(.fromMilliseconds(2200), .awake); 2540 2541 for (0..1000) |round| { 2542 const stream = testStreamWithTimeout(6006, 500_000); 2543 defer stream.close(t.io); 2544 var writer = stream.writer(t.io, &.{}); 2545 const tag = std.fmt.bufPrint(&tag_buf, "reuse-{d}", .{round}) catch unreachable; 2546 try writer.interface.writeAll(std.fmt.bufPrint(&req_buf, "GET /echo?tag={s} HTTP/1.1\r\nConnection: close\r\n\r\n", .{tag}) catch unreachable); 2547 try writer.interface.flush(); 2548 var res = testReadParsedChecked(stream) catch |err| { 2549 std.debug.print("timeout-reuse regression failed at round {d}: {}\n", .{ round, err }); 2550 return err; 2551 }; 2552 defer res.deinit(); 2553 try t.expectString(tag, res.body); 2554 } 2555} 2556 2557// A request timeout can reclaim an HTTPConn while its parser is waiting for a 2558// body. Returning that object to the pool without resetting Request.State makes 2559// the next socket's request bytes become the old request's body. 2560test "httpz: request timeout resets parser before HTTPConn reuse" { 2561 if (force_blocking) return; 2562 2563 var server = try Server(TestWebsocketHandler).init(t.io, t.allocator, .{ 2564 .address = .localhost(6007), 2565 .workers = .{ .count = 1, .max_conn = 1, .min_conn = 1 }, 2566 .timeout = .{ .request = 0 }, 2567 }, TestWebsocketHandler{}); 2568 var router = try server.router(.{}); 2569 router.get("/echo", TestWebsocketHandler.echoTag, .{}); 2570 router.post("/echo", TestWebsocketHandler.echoTag, .{}); 2571 const thread = try server.listenInNewThread(); 2572 defer { 2573 server.stop(); 2574 thread.join(); 2575 server.deinit(); 2576 } 2577 try testing.waitForPort(6007); 2578 2579 const fresh_request = "GET /echo?tag=fresh HTTP/1.1\r\nConnection: close\r\n\r\n"; 2580 const stale = testStreamWithTimeout(6007, 500_000); 2581 var stale_writer = stale.writer(t.io, &.{}); 2582 var header_buf: [128]u8 = undefined; 2583 try stale_writer.interface.writeAll(try std.fmt.bufPrint( 2584 &header_buf, 2585 "POST /echo?tag=stale HTTP/1.1\r\nContent-Length: {d}\r\n\r\n", 2586 .{fresh_request.len}, 2587 )); 2588 try stale_writer.interface.flush(); 2589 2590 // Timeout collection runs on the worker's one-second maintenance tick. 2591 try t.io.sleep(.fromMilliseconds(2200), .awake); 2592 stale.close(t.io); 2593 2594 const fresh = testStreamWithTimeout(6007, 500_000); 2595 defer fresh.close(t.io); 2596 var fresh_writer = fresh.writer(t.io, &.{}); 2597 try fresh_writer.interface.writeAll(fresh_request); 2598 try fresh_writer.interface.flush(); 2599 var res = try testReadParsedChecked(fresh); 2600 defer res.deinit(); 2601 try t.expectString("fresh", res.body); 2602} 2603 2604// Once bytes for a new keepalive request arrive, the idle keepalive deadline 2605// must be replaced by the request deadline. Otherwise a header/body split near 2606// the idle deadline is killed while the body is legitimately still arriving. 2607test "httpz: partial keepalive request receives request timeout" { 2608 if (force_blocking) return; 2609 2610 var server = try Server(TestWebsocketHandler).init(t.io, t.allocator, .{ 2611 .address = .localhost(6008), 2612 .workers = .{ .count = 1, .max_conn = 1, .min_conn = 1 }, 2613 .timeout = .{ .request = 5, .keepalive = 1 }, 2614 }, TestWebsocketHandler{}); 2615 var router = try server.router(.{}); 2616 router.get("/echo", TestWebsocketHandler.echoTag, .{}); 2617 router.post("/echo", TestWebsocketHandler.echoTag, .{}); 2618 const thread = try server.listenInNewThread(); 2619 defer { 2620 server.stop(); 2621 thread.join(); 2622 server.deinit(); 2623 } 2624 try testing.waitForPort(6008); 2625 2626 const stream = testStreamWithTimeout(6008, 500_000); 2627 defer stream.close(t.io); 2628 var writer = stream.writer(t.io, &.{}); 2629 try writer.interface.writeAll("GET /echo?tag=first HTTP/1.1\r\n\r\n"); 2630 try writer.interface.flush(); 2631 var first = try testReadParsedChecked(stream); 2632 try t.expectString("first", first.body); 2633 first.deinit(); 2634 2635 try writer.interface.writeAll("POST /echo?tag=alive HTTP/1.1\r\nContent-Length: 5\r\n\r\n"); 2636 try writer.interface.flush(); 2637 try t.io.sleep(.fromMilliseconds(2200), .awake); 2638 try writer.interface.writeAll("abcde"); 2639 try writer.interface.flush(); 2640 2641 var second = try testReadParsedChecked(stream); 2642 defer second.deinit(); 2643 try t.expectString("alive", second.body); 2644} 2645 2646test "FallbackAllocator: failed resize leaves allocation valid" { 2647 var thread_buf: [1024]u8 = undefined; 2648 var req_arena = std.heap.ArenaAllocator.init(t.allocator); 2649 defer req_arena.deinit(); 2650 2651 var fba = FixedBufferAllocator.init(&thread_buf); 2652 var fb = FallbackAllocator{ 2653 .fba = &fba, 2654 .fallback = req_arena.allocator(), 2655 .fixed = fba.allocator(), 2656 }; 2657 const allocator = fb.allocator(); 2658 2659 const a = try allocator.alloc(u8, 100); 2660 @memset(a, 1); 2661 2662 // can't grow past thread_buf, must fail AND leave `a` valid 2663 try t.expectEqual(false, allocator.resize(a, 2000)); 2664 2665 // `a` is still live, a new allocation must not alias it 2666 const b = try allocator.alloc(u8, 100); 2667 @memset(b, 2); 2668 try t.expectEqual(true, a.ptr != b.ptr); 2669 try t.expectEqual(1, a[0]); 2670 2671 // an explicit free should reclaim the fba space 2672 allocator.free(b); 2673 allocator.free(a); 2674 const c = try allocator.alloc(u8, 800); 2675 try t.expectEqual(true, fba.ownsPtr(c.ptr)); 2676} 2677 2678test "FallbackAllocator: nested arena survives node resize failure" { 2679 // Simulates a handler that creates its own arena on top of req.arena 2680 // while the response writer grows in the same allocator. The arena 2681 // grows its node via resize; when that fails, the node must remain 2682 // valid or its header gets overwritten by subsequent allocations. 2683 var thread_buf: [2048]u8 = undefined; 2684 var req_arena = std.heap.ArenaAllocator.init(t.allocator); 2685 defer req_arena.deinit(); 2686 2687 var fba = FixedBufferAllocator.init(&thread_buf); 2688 var fb = FallbackAllocator{ 2689 .fba = &fba, 2690 .fallback = req_arena.allocator(), 2691 .fixed = fba.allocator(), 2692 }; 2693 const allocator = fb.allocator(); 2694 2695 var nested = std.heap.ArenaAllocator.init(allocator); 2696 defer nested.deinit(); 2697 2698 // first allocation puts an arena node in thread_buf, the second forces 2699 // a node resize that fails (larger than thread_buf) 2700 _ = try nested.allocator().alloc(u8, 64); 2701 _ = try nested.allocator().alloc(u8, 4096); 2702 2703 // response writer keeps allocating from the same request allocator. If 2704 // the failed resize freed the node, this overwrites the node header and 2705 // nested.deinit() crashes. 2706 var out = std.Io.Writer.Allocating.init(allocator); 2707 defer out.deinit(); 2708 for (0..100) |_| { 2709 try out.writer.writeAll("{\"key\": \"some json value written by the response writer\"}"); 2710 } 2711} 2712 2713test "ContentType: forX" { 2714 inline for (@typeInfo(ContentType).@"enum".fields) |field| { 2715 if (comptime std.mem.eql(u8, "BINARY", field.name)) continue; 2716 if (comptime std.mem.eql(u8, "EVENTS", field.name)) continue; 2717 try t.expectEqual(@field(ContentType, field.name), ContentType.forExtension(field.name)); 2718 try t.expectEqual(@field(ContentType, field.name), ContentType.forExtension("." ++ field.name)); 2719 try t.expectEqual(@field(ContentType, field.name), ContentType.forFile("some_file." ++ field.name)); 2720 } 2721 // variations 2722 try t.expectEqual(ContentType.HTML, ContentType.forExtension(".htm")); 2723 try t.expectEqual(ContentType.JPG, ContentType.forExtension(".jpeg")); 2724 2725 try t.expectEqual(ContentType.UNKNOWN, ContentType.forExtension(".spice")); 2726 try t.expectEqual(ContentType.UNKNOWN, ContentType.forExtension("")); 2727 try t.expectEqual(ContentType.UNKNOWN, ContentType.forExtension(".x")); 2728 try t.expectEqual(ContentType.UNKNOWN, ContentType.forFile("")); 2729 try t.expectEqual(ContentType.UNKNOWN, ContentType.forFile("css")); 2730 try t.expectEqual(ContentType.UNKNOWN, ContentType.forFile("css")); 2731 try t.expectEqual(ContentType.UNKNOWN, ContentType.forFile("must.spice")); 2732} 2733 2734fn testStream(port: u16) Io.net.Stream { 2735 return testStreamWithTimeout(port, 20_000); 2736} 2737 2738fn testStreamWithTimeout(port: u16, timeout_usec: i32) Io.net.Stream { 2739 const timeout = std.mem.toBytes(posix.timeval{ 2740 .sec = 0, 2741 .usec = timeout_usec, 2742 }); 2743 2744 const address = Io.net.IpAddress.parse("127.0.0.1", port) catch unreachable; 2745 const stream = address.connect(t.io, .{ .mode = .stream }) catch unreachable; 2746 posix.setsockopt(stream.socket.handle, posix.SOL.SOCKET, posix.SO.RCVTIMEO, &timeout) catch unreachable; 2747 posix.setsockopt(stream.socket.handle, posix.SOL.SOCKET, posix.SO.SNDTIMEO, &timeout) catch unreachable; 2748 return stream; 2749} 2750 2751fn testReadAll(stream: Io.net.Stream, buf: []u8) []u8 { 2752 var pos: usize = 0; 2753 var blocked = false; 2754 while (true) { 2755 std.debug.assert(pos < buf.len); 2756 const n = posix.read(stream.socket.handle, buf[pos..]) catch |err| switch (err) { 2757 error.WouldBlock => { 2758 if (blocked) { 2759 return buf[0..pos]; 2760 } 2761 blocked = true; 2762 std.Io.sleep(t.io, .fromMilliseconds(1), .awake) catch unreachable; 2763 continue; 2764 }, 2765 error.ConnectionResetByPeer => return buf[0..pos], 2766 else => @panic(@errorName(err)), 2767 }; 2768 2769 if (n == 0) { 2770 return buf[0..pos]; 2771 } 2772 pos += n; 2773 blocked = false; 2774 } 2775 unreachable; 2776} 2777 2778fn testReadParsed(stream: Io.net.Stream) testing.Testing.Response { 2779 var buf: [4096]u8 = undefined; 2780 const data = testReadAll(stream, &buf); 2781 return testing.parse(data) catch unreachable; 2782} 2783 2784fn testReadParsedChecked(stream: Io.net.Stream) !testing.Testing.Response { 2785 var buf: [4096]u8 = undefined; 2786 var pos: usize = 0; 2787 var blocked = false; 2788 var expected_total: ?usize = null; 2789 2790 while (pos < buf.len) { 2791 const n = posix.read(stream.socket.handle, buf[pos..]) catch |err| switch (err) { 2792 error.WouldBlock => { 2793 if (blocked) return error.ResponseTimedOut; 2794 blocked = true; 2795 std.Io.sleep(t.io, .fromMilliseconds(1), .awake) catch unreachable; 2796 continue; 2797 }, 2798 error.ConnectionResetByPeer => return error.ResponseReset, 2799 else => return err, 2800 }; 2801 if (n == 0) return error.ResponseClosedEarly; 2802 pos += n; 2803 blocked = false; 2804 2805 if (expected_total == null) { 2806 if (std.mem.indexOf(u8, buf[0..pos], "\r\n\r\n")) |header_end| { 2807 const marker = "Content-Length: "; 2808 const marker_start = std.mem.indexOf(u8, buf[0..header_end], marker) orelse 2809 return error.ResponseMissingContentLength; 2810 const value_start = marker_start + marker.len; 2811 const value_end = std.mem.indexOfPos(u8, buf[0 .. header_end + 2], value_start, "\r\n") orelse 2812 return error.InvalidContentLength; 2813 const body_len = try std.fmt.parseInt(usize, buf[value_start..value_end], 10); 2814 expected_total = header_end + 4 + body_len; 2815 } 2816 } 2817 2818 if (expected_total) |total| { 2819 if (pos >= total) return try testing.parse(buf[0..total]); 2820 } 2821 } 2822 return error.ResponseTooLarge; 2823} 2824 2825fn testReadHeader(stream: Io.net.Stream) testing.Testing.Response { 2826 var pos: usize = 0; 2827 var blocked = false; 2828 var buf: [1024]u8 = undefined; 2829 var reader = stream.reader(t.io, &.{}); 2830 const r = &reader.interface; 2831 while (true) { 2832 std.debug.assert(pos < buf.len); 2833 var vecs: [1][]u8 = .{buf[pos..]}; 2834 const n = r.readVec(&vecs) catch |err| 2835 switch (err) { 2836 error.ReadFailed => { 2837 if (reader.err) |e| { 2838 switch (e) { 2839 // @ZIG016 2840 // error.WouldBlock => { 2841 // if (blocked) unreachable; 2842 // blocked = true; 2843 // std.Thread.sleep(std.time.ns_per_ms); 2844 // continue; 2845 // }, 2846 else => @panic(@errorName(e)), 2847 } 2848 } 2849 @panic(@errorName(err)); 2850 }, 2851 error.EndOfStream => 0, 2852 }; 2853 2854 if (n == 0) unreachable; 2855 2856 pos += n; 2857 if (std.mem.endsWith(u8, buf[0..pos], "\r\n\r\n")) { 2858 return testing.parse(buf[0..pos]) catch unreachable; 2859 } 2860 blocked = false; 2861 } 2862 unreachable; 2863} 2864 2865const TestUser = struct { 2866 id: []const u8, 2867 power: usize, 2868}; 2869 2870// simulates having a void handler, but keeps the test actions organized within 2871// this namespace. 2872const TestDummyHandler = struct { 2873 const RouteData = struct { 2874 power: usize, 2875 }; 2876 2877 fn fail(_: *Request, _: *Response) !void { 2878 return error.Failure; 2879 } 2880 2881 fn reqQuery(req: *Request, res: *Response) !void { 2882 res.status = 200; 2883 const query = try req.query(); 2884 res.body = query.get("fav").?; 2885 } 2886 2887 fn method(req: *Request, res: *Response) !void { 2888 try res.json(.{ .method = req.method, .string = req.method_string }, .{}); 2889 } 2890 2891 fn chunked(_: *Request, res: *Response) !void { 2892 res.header("Over", "9000!"); 2893 res.status = 200; 2894 try res.chunk("Chunk 1"); 2895 try res.chunk("and another chunk"); 2896 } 2897 2898 fn jsonRes(_: *Request, res: *Response) !void { 2899 res.setStatus(.created); 2900 try res.json(.{ .over = 9000, .teg = "soup" }, .{}); 2901 } 2902 2903 fn echoBody(req: *Request, res: *Response) !void { 2904 res.status = 200; 2905 res.body = req.body() orelse ""; 2906 } 2907 2908 fn routeData(req: *Request, res: *Response) !void { 2909 const rd: *const RouteData = @ptrCast(@alignCast(req.route_data.?)); 2910 try res.json(.{ .power = rd.power }, .{}); 2911 } 2912 2913 fn eventStream(_: *Request, res: *Response) !void { 2914 res.status = 818; 2915 try res.startEventStream(StreamContext{ .data = "hello" }, StreamContext.handle); 2916 } 2917 2918 fn eventStreamSync(_: *Request, res: *Response) !void { 2919 res.status = 818; 2920 const stream = try res.startEventStreamSync(); 2921 var w = stream.writer(res.conn.io, &.{}); 2922 w.interface.writeAll("hello") catch unreachable; 2923 w.interface.writeAll("a sync message") catch unreachable; 2924 } 2925 2926 fn reqReader(req: *Request, res: *Response) !void { 2927 var reader = try req.reader(2000); 2928 2929 var l: usize = 0; 2930 var buf: [1024]u8 = undefined; 2931 while (true) { 2932 const n = try reader.read(&buf); 2933 if (n == 0) { 2934 break; 2935 } 2936 if (req.body_len > 10 and std.mem.indexOfNonePos(u8, buf[0..n], 0, "a") != null) { 2937 return error.InvalidData; 2938 } 2939 l += n; 2940 } 2941 return res.json(.{ .length = l }, .{}); 2942 } 2943 2944 const StreamContext = struct { 2945 data: []const u8, 2946 2947 fn handle(self: StreamContext, stream: Io.net.Stream) void { 2948 var writer = stream.writer(t.io, &.{}); 2949 const w = &writer.interface; 2950 w.writeAll(self.data) catch unreachable; 2951 w.writeAll("a message") catch unreachable; 2952 w.flush() catch unreachable; 2953 } 2954 }; 2955 2956 fn routeSpecificDispacthcer(action: Action(void), req: *Request, res: *Response) !void { 2957 res.header("dispatcher", "test-dispatcher-1"); 2958 return action(req, res); 2959 } 2960 2961 fn dispatchedAction(_: *Request, res: *Response) !void { 2962 const writer = res.writer(); 2963 return writer.writeAll("action"); 2964 } 2965 2966 fn middlewares(req: *Request, res: *Response) !void { 2967 return res.json(.{ 2968 .v1 = TestMiddleware.value1(req), 2969 .v2 = TestMiddleware.value2(req), 2970 }, .{}); 2971 } 2972 2973 // called by the re-use server, but put here because, like the default server 2974 // this is a handler-less server 2975 fn reuseWriter(req: *Request, res: *Response) !void { 2976 res.status = 200; 2977 const query = try req.query(); 2978 const count = try std.fmt.parseInt(u16, query.get("count").?, 10); 2979 2980 var data = try res.arena.alloc(TestUser, count); 2981 for (0..count) |i| { 2982 data[i] = .{ 2983 .id = try std.fmt.allocPrint(res.arena, "id-{d}", .{i}), 2984 .power = i, 2985 }; 2986 } 2987 return res.json(.{ .data = data }, .{}); 2988 } 2989}; 2990 2991const TestHandlerDefaultDispatch = struct { 2992 state: usize, 2993 2994 fn dispatch2(h: *TestHandlerDefaultDispatch, action: Action(*TestHandlerDefaultDispatch), req: *Request, res: *Response) !void { 2995 res.header("dispatcher", "test-dispatcher-2"); 2996 return action(h, req, res); 2997 } 2998 2999 fn dispatch3(h: *TestHandlerDefaultDispatch, action: Action(*TestHandlerDefaultDispatch), req: *Request, res: *Response) !void { 3000 res.header("dispatcher", "test-dispatcher-3"); 3001 return action(h, req, res); 3002 } 3003 3004 fn echo(h: *TestHandlerDefaultDispatch, req: *Request, res: *Response) !void { 3005 return res.json(.{ 3006 .state = h.state, 3007 .method = @tagName(req.method), 3008 .path = req.url.path, 3009 }, .{}); 3010 } 3011 3012 fn echoWrite(h: *TestHandlerDefaultDispatch, req: *Request, res: *Response) !void { 3013 const json_writer = std.json.fmt(.{ 3014 .state = h.state, 3015 .method = @tagName(req.method), 3016 .path = req.url.path, 3017 }, .{}); 3018 3019 var aw: std.Io.Writer.Allocating = .init(res.arena); 3020 try json_writer.format(&aw.writer); 3021 3022 res.body = aw.written(); 3023 return res.write(); 3024 } 3025 3026 fn params(_: *TestHandlerDefaultDispatch, req: *Request, res: *Response) !void { 3027 const args = .{ req.param("version").?, req.param("UserId").? }; 3028 res.body = try std.fmt.allocPrint(req.arena, "version={s},user={s}", args); 3029 } 3030 3031 fn headers(h: *TestHandlerDefaultDispatch, req: *Request, res: *Response) !void { 3032 res.header("state", try std.fmt.allocPrint(res.arena, "{d}", .{h.state})); 3033 res.header("Echo", req.header("header-name").?); 3034 res.header("other", "test-value"); 3035 } 3036 3037 fn clBody(_: *TestHandlerDefaultDispatch, req: *Request, res: *Response) !void { 3038 res.header("Echo-Body", req.body().?); 3039 } 3040 3041 fn fail(_: *TestHandlerDefaultDispatch, _: *Request, _: *Response) !void { 3042 return error.TestUnhandledError; 3043 } 3044 3045 pub fn notFound(h: *TestHandlerDefaultDispatch, _: *Request, res: *Response) !void { 3046 res.status = 404; 3047 res.header("state", try std.fmt.allocPrint(res.arena, "{d}", .{h.state})); 3048 res.body = "where lah?"; 3049 } 3050 3051 pub fn uncaughtError(h: *TestHandlerDefaultDispatch, _: *Request, res: *Response, err: anyerror) void { 3052 res.status = 500; 3053 res.header("state", std.fmt.allocPrint(res.arena, "{d}", .{h.state}) catch unreachable); 3054 res.header("err", @errorName(err)); 3055 res.body = "#/why/arent/tags/hierarchical"; 3056 } 3057}; 3058 3059const TestHandlerDispatch = struct { 3060 state: usize, 3061 3062 pub fn dispatch(self: *TestHandlerDispatch, action: Action(*TestHandlerDispatch), req: *Request, res: *Response) !void { 3063 res.header("dstate", try std.fmt.allocPrint(res.arena, "{d}", .{self.state})); 3064 res.header("dispatch", "TestHandlerDispatch"); 3065 return action(self, req, res); 3066 } 3067 3068 fn root(h: *TestHandlerDispatch, _: *Request, res: *Response) !void { 3069 return res.json(.{ .state = h.state }, .{}); 3070 } 3071}; 3072 3073const TestHandlerDispatchContext = struct { 3074 state: usize, 3075 3076 const ActionContext = struct { 3077 other: usize, 3078 }; 3079 3080 pub fn dispatch(self: *TestHandlerDispatchContext, action: Action(*ActionContext), req: *Request, res: *Response) !void { 3081 res.header("dstate", try std.fmt.allocPrint(res.arena, "{d}", .{self.state})); 3082 res.header("dispatch", "TestHandlerDispatchContext"); 3083 var action_context = ActionContext{ .other = self.state + 10 }; 3084 return action(&action_context, req, res); 3085 } 3086 3087 pub fn root(a: *const ActionContext, _: *Request, res: *Response) !void { 3088 return res.json(.{ .other = a.other }, .{}); 3089 } 3090}; 3091 3092const TestHandlerHandle = struct { 3093 pub fn handle(_: TestHandlerHandle, req: *Request, res: *Response) void { 3094 const query = req.query() catch unreachable; 3095 const writer = res.writer(); 3096 writer.print("hello {s}", .{query.get("name") orelse "world"}) catch unreachable; 3097 } 3098}; 3099 3100const TestWebsocketHandler = struct { 3101 pub const WebsocketHandler = struct { 3102 ctx: u32, 3103 conn: *websocket.Conn, 3104 3105 pub fn init(conn: *websocket.Conn, ctx: u32) !WebsocketHandler { 3106 if (ctx == 9999) return error.RejectedUpgrade; 3107 return .{ 3108 .ctx = ctx, 3109 .conn = conn, 3110 }; 3111 } 3112 3113 pub fn afterInit(self: *WebsocketHandler, ctx: u32) !void { 3114 try t.expectEqual(self.ctx, ctx); 3115 } 3116 3117 pub fn clientMessage(self: *WebsocketHandler, data: []const u8) !void { 3118 if (std.mem.eql(u8, data, "close")) { 3119 self.conn.close(.{}) catch {}; 3120 return; 3121 } 3122 try self.conn.write(data); 3123 } 3124 }; 3125 3126 pub fn echoTag(_: TestWebsocketHandler, req: *Request, res: *Response) !void { 3127 const query = try req.query(); 3128 res.body = query.get("tag") orelse "?"; 3129 } 3130 3131 pub fn upgrade(_: TestWebsocketHandler, req: *Request, res: *Response) !void { 3132 if (try upgradeWebsocket(WebsocketHandler, req, res, 9001) == false) { 3133 res.status = 500; 3134 res.body = "invalid websocket"; 3135 } 3136 } 3137 3138 pub fn rejectUpgrade(_: TestWebsocketHandler, req: *Request, res: *Response) !void { 3139 _ = upgradeWebsocket(WebsocketHandler, req, res, 9999) catch { 3140 res.status = 400; 3141 res.body = "rejected upgrade"; 3142 return; 3143 }; 3144 return error.ExpectedRejectedUpgrade; 3145 } 3146}; 3147 3148const TestMiddleware = struct { 3149 const Config = struct { 3150 id: i32, 3151 }; 3152 3153 allocator: Allocator, 3154 v1: []const u8, 3155 v2: []const u8, 3156 3157 fn init(config: TestMiddleware.Config, mc: MiddlewareConfig) !TestMiddleware { 3158 return .{ 3159 .allocator = mc.allocator, 3160 .v1 = try std.fmt.allocPrint(mc.arena, "tm1-{d}", .{config.id}), 3161 .v2 = try std.fmt.allocPrint(mc.allocator, "tm2-{d}", .{config.id}), 3162 }; 3163 } 3164 3165 pub fn deinit(self: *const TestMiddleware) void { 3166 self.allocator.free(self.v2); 3167 } 3168 3169 fn value1(req: *const Request) []const u8 { 3170 const v: [*]u8 = @ptrCast(req.middlewares.get("text_middleware_1").?); 3171 return v[0..7]; 3172 } 3173 3174 fn value2(req: *const Request) []const u8 { 3175 const v: [*]u8 = @ptrCast(req.middlewares.get("text_middleware_2").?); 3176 return v[0..7]; 3177 } 3178 3179 fn execute(self: *const TestMiddleware, req: *Request, _: *Response, executor: anytype) !void { 3180 try req.middlewares.put("text_middleware_1", (try req.arena.dupe(u8, self.v1)).ptr); 3181 try req.middlewares.put("text_middleware_2", (try req.arena.dupe(u8, self.v2)).ptr); 3182 return executor.next(); 3183 } 3184};