tiny.http.Server
Defined in tiny.http.
API (37)
Actions
Public operations.
connectionCountconnectionInputCapacityconnectionInputStatusdeinitinitlistenrequestCapacityresponseCapacitysetHandlerstop
Types and contracts
Public types and contracts.
Fields and members
Public fields and members.
accept_completionactive_connectionsallocatorcompleted_connection_headcompleted_connection_tailconfigconnection_input_storageconnection_slotsconnections_mutexevent_loopfree_connection_slothandlerhandler_contextidle_completionidle_timerinput_capacity_rejectionslisten_activelistenerlistener_close_completionnext_conn_idpoolrequest_storageresponse_storageshutdownwakewake_completion
Source
Source: lib/http/src/server/runtime.zig:45
zig
pub const Server = struct { allocator: std.mem.Allocator, listener: ?sys.Socket, pool: ThreadPool, event_loop: event.Loop, wake: event.Async, wake_completion: event.Completion, accept_completion: event.Completion, listener_close_completion: event.Completion, idle_timer: event.Timer, idle_completion: event.Completion, connection_slots: []ServerConnectionSlot, connection_input_storage: Connection.InputStorage, request_storage: http.RequestStorage, response_storage: http.ResponseStorage, free_connection_slot: ?usize, completed_connection_head: ?usize, completed_connection_tail: ?usize, active_connections: usize, connections_mutex: std.atomic.Mutex, input_capacity_rejections: u64, shutdown: std.atomic.Value(bool), listen_active: std.atomic.Value(bool), config: Config, next_conn_id: std.atomic.Value(usize), handler: ?*const fn (?*anyopaque, *Connection) void, handler_context: ?*anyopaque, pub const Config = struct { address: []const u8 = "127.0.0.1", port: u16 = 8080, num_workers: usize = 4, max_connections: usize = 1024, read_timeout_ms: u32 = 30_000, write_timeout_ms: u32 = 30_000, boot_clock: time.BootClock = .system(), tcp_no_delay: bool = false, connection_input_bytes_per_connection: usize = http.default_connection_input_bytes_per_connection, request_header_count_per_connection: usize = http.default_request_header_count, request_header_line_bytes: usize = http.default_request_header_line_bytes, request_body_bytes_per_connection: usize = http.default_request_body_bytes, response_header_count_per_connection: usize = http.default_response_header_count, response_head_bytes_per_connection: usize = http.default_response_head_bytes, }; pub fn init(allocator: std.mem.Allocator, config: Config) !*Server { if (config.max_connections > std.math.maxInt(u32) - 4) return error.ConnectionLimitTooLarge; const self = try allocator.create(Server); errdefer allocator.destroy(self); self.pool = try ThreadPool.init(allocator, .{ .workers = config.num_workers, .tasks = config.max_connections, }); errdefer self.pool.deinit(allocator); self.event_loop = try event.Loop.init(.{ .allocator = allocator, .entries = @intCast(config.max_connections + 4), }); errdefer self.event_loop.deinit(); self.wake = try event.Async.init(); errdefer self.wake.deinit(); self.idle_timer = try event.Timer.init(); errdefer self.idle_timer.deinit(); self.connection_slots = try allocator.alloc(ServerConnectionSlot, config.max_connections); errdefer allocator.free(self.connection_slots); self.connection_input_storage = try Connection.InputStorage.init(allocator, .{ .connection_count = config.max_connections, .bytes_per_connection = config.connection_input_bytes_per_connection, }); errdefer self.connection_input_storage.deinit(allocator); self.request_storage = try http.RequestStorage.init(allocator, .{ .request_count = config.max_connections, .header_count_per_request = config.request_header_count_per_connection, .header_line_bytes = config.request_header_line_bytes, .body_bytes_per_request = config.request_body_bytes_per_connection, }); errdefer self.request_storage.deinit(allocator); self.response_storage = try http.ResponseStorage.init(allocator, .{ .response_count = config.max_connections, .header_count_per_response = config.response_header_count_per_connection, .head_bytes_per_response = config.response_head_bytes_per_connection, }); errdefer self.response_storage.deinit(allocator); const listener = try createListenerSocket(config.address, config.port); errdefer sys.close(listener); self.allocator = allocator; self.listener = listener; self.wake_completion = .{}; self.accept_completion = .{}; self.listener_close_completion = .{}; self.idle_completion = .{}; self.free_connection_slot = if (self.connection_slots.len == 0) null else 0; self.completed_connection_head = null; self.completed_connection_tail = null; self.active_connections = 0; self.connections_mutex = .unlocked; self.input_capacity_rejections = 0; self.shutdown = std.atomic.Value(bool).init(false); self.listen_active = std.atomic.Value(bool).init(false); self.config = config; self.next_conn_id = std.atomic.Value(usize).init(1); self.handler = null; self.handler_context = null; self.initializeSlots(); try self.pool.activate(); self.connection_input_storage.activate(); self.request_storage.activate(); self.response_storage.activate(); return self; } fn initializeSlots(self: *Server) void { for (self.connection_slots, 0..) |*slot, index| { slot.* = .{ .server = self, .index = index, .next_free = if (index + 1 < self.connection_slots.len) index + 1 else null, .next_completed = null, .occupied = false, .phase = .free, .connection = undefined, .read_completion = .{}, }; } } fn createListenerSocket(address: []const u8, port: u16) !sys.Socket { const addr = try sys.ip4AddressForHost(address, port); const sock = try sys.tcpStreamSocket(.{}); errdefer sys.close(sock); try sys.setReuseAddress(sock); try sys.bindIp4(sock, addr); try sys.listen(sock, 128); try sys.setNonBlocking(sock); return sock; } pub fn deinit(self: *Server) void { self.stop(); var wait_rounds: usize = 0; while (self.listen_active.load(.acquire) and wait_rounds < 5000) : (wait_rounds += 1) { sleepMillis(1); } std.debug.assert(!self.listen_active.load(.acquire)); if (self.listener) |sock| { sys.close(sock); self.listener = null; } lock(&self.connections_mutex); for (self.connection_slots) |*slot| { if (slot.occupied) slot.connection.interrupt(); } self.connections_mutex.unlock(); _ = self.pool.drain( .system(), .fromNanoseconds(5 * std.time.ns_per_s), ) catch false; self.pool.deinit(self.allocator); lock(&self.connections_mutex); for (self.connection_slots) |*slot| { if (slot.occupied) slot.connection.deinit(); } self.connections_mutex.unlock(); self.idle_timer.deinit(); self.wake.deinit(); self.event_loop.deinit(); self.response_storage.deinit(self.allocator); self.request_storage.deinit(self.allocator); self.connection_input_storage.deinit(self.allocator); self.allocator.free(self.connection_slots); self.allocator.destroy(self); } pub fn setHandler(self: *Server, context: anytype, comptime handler: fn (@TypeOf(context), *Connection) void) void { const Context = @TypeOf(context); comptime std.debug.assert(@typeInfo(Context) == .pointer); const Erased = struct { fn call(erased: ?*anyopaque, conn: *Connection) void { handler(@ptrCast(@alignCast(erased)), conn); } }; self.handler_context = @ptrCast(@constCast(context)); self.handler = Erased.call; } pub fn listen(self: *Server) !void { const listener = self.listener orelse return error.NotListening; if (self.listen_active.swap(true, .acq_rel)) return error.AlreadyListening; defer self.listen_active.store(false, .release); self.wake.wait( &self.event_loop, &self.wake_completion, Server, self, wakeReady, ); event.TCP.initFd(listener).accept( &self.event_loop, &self.accept_completion, Server, self, acceptReady, ); if (self.config.read_timeout_ms != 0) { const interval = time.Duration.fromMilliseconds( idleSweepInterval(self.config.read_timeout_ms), ); try self.idle_timer.run( &self.event_loop, &self.idle_completion, .{ .after = interval, .repeat = .{ .fixed_delay = interval } }, Server, self, idleTimerReady, ); } try self.event_loop.run(.until_done); } fn acquireConnection( self: *Server, socket: sys.Socket, boot_clock: time.BootClock, accepted_at: time.BootInstant, ) ?*ServerConnectionSlot { lock(&self.connections_mutex); defer self.connections_mutex.unlock(); const index = self.free_connection_slot orelse return null; const slot = &self.connection_slots[index]; self.free_connection_slot = slot.next_free; slot.next_free = null; slot.next_completed = null; slot.occupied = true; slot.phase = .waiting; slot.connection = Connection.init( self.next_conn_id.fetchAdd(1, .monotonic), socket, self.connection_input_storage.connection(index) catch unreachable, self.request_storage.request(index) catch unreachable, self.response_storage.response(index) catch unreachable, boot_clock, accepted_at, ); self.active_connections += 1; return slot; } fn releaseConnection(self: *Server, slot: *ServerConnectionSlot) void { lock(&self.connections_mutex); defer self.connections_mutex.unlock(); std.debug.assert(slot.occupied); self.input_capacity_rejections +|= slot.connection.inputStatus().capacity_rejections; slot.occupied = false; slot.phase = .free; slot.next_free = self.free_connection_slot; slot.next_completed = null; self.free_connection_slot = slot.index; self.active_connections -= 1; } fn handleConnection(slot: *ServerConnectionSlot) void { if (slot.server.handler) |handler| { handler(slot.server.handler_context, &slot.connection); } else { slot.connection.markClosing(); } } fn finishConnection(slot: *ServerConnectionSlot) void { const server = slot.server; lock(&server.connections_mutex); std.debug.assert(slot.occupied); std.debug.assert(slot.phase == .running); slot.phase = .completed; slot.next_completed = null; if (server.completed_connection_tail) |tail| { server.connection_slots[tail].next_completed = slot.index; } else { server.completed_connection_head = slot.index; } server.completed_connection_tail = slot.index; server.connections_mutex.unlock(); server.wake.notify() catch {}; } pub fn stop(self: *Server) void { self.shutdown.store(true, .release); self.wake.notify() catch {}; } pub fn connectionCount(self: *Server) usize { lock(&self.connections_mutex); defer self.connections_mutex.unlock(); return self.active_connections; } pub fn connectionInputCapacity(self: *const Server) Connection.InputCapacity { return self.connection_input_storage.capacity; } pub fn requestCapacity(self: *const Server) http.RequestCapacity { return self.request_storage.capacity; } pub fn responseCapacity(self: *const Server) http.ResponseCapacity { return self.response_storage.capacity; } pub fn connectionInputStatus(self: *Server) http.ConnectionInputStatus { lock(&self.connections_mutex); defer self.connections_mutex.unlock(); var status = http.ConnectionInputStatus{ .capacity_rejections = self.input_capacity_rejections, }; for (self.connection_slots) |*slot| { if (!slot.occupied) continue; status.capacity_rejections +|= slot.connection.inputStatus().capacity_rejections; } return status; } fn acceptReady( maybe_server: ?*Server, loop: *event.Loop, _: *event.Completion, result: event.AcceptError!event.TCP, ) event.CallbackAction { const self = maybe_server.?; const client = result catch return if (self.shutdown.load(.acquire)) .disarm else .rearm; if (self.shutdown.load(.acquire)) { sys.close(client.fd); return .disarm; } if (self.config.tcp_no_delay) { sys.setTcpNoDelay(client.fd) catch |err| { log.warn( "failed to disable TCP delay: {s}", .{@errorName(err)}, ); sys.close(client.fd); return .rearm; }; } const accepted_at = self.config.boot_clock.now() catch { sys.close(client.fd); return .rearm; }; const slot = self.acquireConnection(client.fd, self.config.boot_clock, accepted_at) orelse { sys.close(client.fd); return .rearm; }; slot.connection.setWriteTimeout(self.config.write_timeout_ms) catch |err| { log.warn("conn {d}: failed to set write timeout: {s}", .{ slot.connection.id, @errorName(err) }); }; self.armWaiting(loop, slot); return .rearm; } fn readableReady( maybe_slot: ?*ServerConnectionSlot, _: *event.Loop, _: *event.Completion, _: event.File, result: event.PollError!event.PollEvent, ) event.CallbackAction { const slot = maybe_slot.?; _ = result catch { slot.connection.deinit(); slot.server.releaseConnection(slot); return .disarm; }; slot.server.dispatch(slot) catch { slot.connection.deinit(); slot.server.releaseConnection(slot); }; return .disarm; } fn wakeReady( maybe_server: ?*Server, loop: *event.Loop, _: *event.Completion, result: event.Async.WaitError!void, ) event.CallbackAction { const self = maybe_server.?; _ = result catch return .disarm; self.drainCompleted(loop); if (!self.shutdown.load(.acquire)) return .rearm; self.closeListener(loop); self.closeConnections(loop); return .disarm; } fn idleTimerReady( maybe_server: ?*Server, loop: *event.Loop, _: *event.Completion, result: event.Timer.RunError!void, ) event.CallbackAction { const self = maybe_server.?; _ = result catch return .disarm; if (self.shutdown.load(.acquire)) return .disarm; const now = self.config.boot_clock.now() catch return .disarm; for (self.connection_slots) |*slot| { lock(&self.connections_mutex); const expired = slot.occupied and slot.phase == .waiting and idleExpired(now, slot.connection.lastActivity(), self.config.read_timeout_ms); if (expired) slot.phase = .closing; self.connections_mutex.unlock(); if (expired) self.closeWaiting(loop, slot); } return .rearm; } fn listenerClosed( _: ?*Server, _: *event.Loop, _: *event.Completion, _: event.TCP, result: event.CloseError!void, ) event.CallbackAction { _ = result catch {}; return .disarm; } fn connectionClosed( maybe_slot: ?*ServerConnectionSlot, _: *event.Loop, _: *event.Completion, _: event.TCP, result: event.CloseError!void, ) event.CallbackAction { const slot = maybe_slot.?; _ = result catch {}; slot.connection.transportClosed(); slot.connection.deinit(); slot.server.releaseConnection(slot); return .disarm; } fn dispatch(self: *Server, slot: *ServerConnectionSlot) !void { lock(&self.connections_mutex); std.debug.assert(slot.occupied); std.debug.assert(canDispatch(slot.phase)); slot.phase = .running; self.connections_mutex.unlock(); slot.connection.beginTurn(); self.pool.submit(.{ .slot = slot }) catch |err| { lock(&self.connections_mutex); slot.phase = .waiting; self.connections_mutex.unlock(); return err; }; } fn armWaiting(self: *Server, loop: *event.Loop, slot: *ServerConnectionSlot) void { lock(&self.connections_mutex); std.debug.assert(slot.occupied); std.debug.assert(canDispatch(slot.phase)); slot.phase = .waiting; self.connections_mutex.unlock(); event.File.initFd(slot.connection.socket).poll( loop, &slot.read_completion, .read, ServerConnectionSlot, slot, readableReady, ); } fn drainCompleted(self: *Server, loop: *event.Loop) void { while (self.popCompleted()) |slot| { const state = slot.connection.currentState(); if (self.shutdown.load(.acquire) or state == .closing or state == .closed or !slot.connection.shouldWaitForRead()) { lock(&self.connections_mutex); slot.phase = .closing; self.connections_mutex.unlock(); slot.connection.deinit(); self.releaseConnection(slot); } else if (slot.connection.shouldConsumeBufferedInput() and slot.connection.hasBufferedInput()) { self.dispatch(slot) catch { slot.connection.deinit(); self.releaseConnection(slot); }; } else { self.armWaiting(loop, slot); } } } fn popCompleted(self: *Server) ?*ServerConnectionSlot { lock(&self.connections_mutex); defer self.connections_mutex.unlock(); const index = self.completed_connection_head orelse return null; const slot = &self.connection_slots[index]; self.completed_connection_head = slot.next_completed; if (self.completed_connection_head == null) self.completed_connection_tail = null; slot.next_completed = null; return slot; } fn closeListener(self: *Server, loop: *event.Loop) void { const socket = self.listener orelse return; self.listener = null; event.TCP.initFd(socket).close( loop, &self.listener_close_completion, Server, self, listenerClosed, ); } fn closeConnections(self: *Server, loop: *event.Loop) void { for (self.connection_slots) |*slot| { lock(&self.connections_mutex); if (!slot.occupied) { self.connections_mutex.unlock(); continue; } const phase = slot.phase; if (phase == .waiting) slot.phase = .closing; self.connections_mutex.unlock(); switch (phase) { .waiting => self.closeWaiting(loop, slot), .running => slot.connection.interrupt(), .completed => {}, .free, .closing => {}, } } self.drainCompleted(loop); } fn closeWaiting(self: *Server, loop: *event.Loop, slot: *ServerConnectionSlot) void { _ = self; event.TCP.initFd(slot.connection.socket).close( loop, &slot.read_completion, ServerConnectionSlot, slot, connectionClosed, ); }};Source: lib/http/src/root.zig:29
zig
pub const Server = server.Server;Complete caller list for Server.init
11 direct callers.
lib.http.src.server.test.test_Server_accepts_connection[function] — test source atlib/http/src/server/test.zig:114in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_closes_overload_and_reuses_bounded_connection_slots[function] — test source atlib/http/src/server/test.zig:197in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_expires_idle_connections_without_occupying_a_worker[function] — test source atlib/http/src/server/test.zig:376in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_init_and_deinit[function] — test source atlib/http/src/server/test.zig:12in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_init_with_custom_config[function] — test source atlib/http/src/server/test.zig:20in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_HTTP_and_WebSocket_turns_beyond_idle_upgrades[function] — test source atlib/http/src/server/test.zig:401in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_around_a_partial_HTTP_request[function] — test source atlib/http/src/server/test.zig:341in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_fresh_requests_beyond_idle_keepalive_count[function] — test source atlib/http/src/server/test.zig:302in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_stop_unblocks_listen[function] — test source atlib/http/src/server/test.zig:86in nearest public ownerlib.http.src.server.testtools.smg.src.scan.test.test_zig_ast_scan_matches_core_extractor_contract[function] — test; no exact target attools/smg/src/scan/test.zig:3312in nearest public ownertools.smg.src.scan.testtools.smg.src.scan.test.test_zig_tree_scan_extracts_comptime_parameter_functions_and_metrics[function] — test; no exact target attools/smg/src/scan/test.zig:4508in nearest public ownertools.smg.src.scan.test
Complete caller list for Server.setHandler
7 direct callers.
lib.http.src.server.test.test_Server_accepts_connection[function] — test source atlib/http/src/server/test.zig:114in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_closes_overload_and_reuses_bounded_connection_slots[function] — test source atlib/http/src/server/test.zig:197in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_expires_idle_connections_without_occupying_a_worker[function] — test source atlib/http/src/server/test.zig:376in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_HTTP_and_WebSocket_turns_beyond_idle_upgrades[function] — test source atlib/http/src/server/test.zig:401in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_around_a_partial_HTTP_request[function] — test source atlib/http/src/server/test.zig:341in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_fresh_requests_beyond_idle_keepalive_count[function] — test source atlib/http/src/server/test.zig:302in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_stop_unblocks_listen[function] — test source atlib/http/src/server/test.zig:86in nearest public ownerlib.http.src.server.test
Complete caller list for Server.stop
8 direct callers.
tiny.http.Server.deinit[method] atlib/http/src/server/runtime.zig:197lib.http.src.server.test.test_Server_accepts_connection[function] — test source atlib/http/src/server/test.zig:114in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_closes_overload_and_reuses_bounded_connection_slots[function] — test source atlib/http/src/server/test.zig:197in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_expires_idle_connections_without_occupying_a_worker[function] — test source atlib/http/src/server/test.zig:376in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_HTTP_and_WebSocket_turns_beyond_idle_upgrades[function] — test source atlib/http/src/server/test.zig:401in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_around_a_partial_HTTP_request[function] — test source atlib/http/src/server/test.zig:341in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_schedules_fresh_requests_beyond_idle_keepalive_count[function] — test source atlib/http/src/server/test.zig:302in nearest public ownerlib.http.src.server.testlib.http.src.server.test.test_Server_stop_unblocks_listen[function] — test source atlib/http/src/server/test.zig:86in nearest public ownerlib.http.src.server.test
Audit
| Definitions | 12 |
|---|---|
| Public names | 12 |
| Members | 40 |
| Version | 26.7.0 |
| Revision | daab053ee433 |