forked from
karlseguin.tngl.sh/http.zig
An HTTP/1.1 server for zig
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};