web: some stuff
This commit is contained in:
249
packages/web/src/Server.zig
Normal file
249
packages/web/src/Server.zig
Normal file
@@ -0,0 +1,249 @@
|
||||
const std = @import("std");
|
||||
const Server = @This();
|
||||
|
||||
const Connection = @import("Connection.zig");
|
||||
const FileDescriptor = @import("FileDescriptor.zig").FileDescriptor;
|
||||
const http = @import("http.zig");
|
||||
const RequestRouter = @import("RequestRouter.zig");
|
||||
const Worker = @import("Worker.zig");
|
||||
|
||||
const linux = std.os.linux;
|
||||
const errno = linux.E.init;
|
||||
|
||||
fd: FileDescriptor,
|
||||
address: std.net.Address,
|
||||
workers: []Worker,
|
||||
request_router: RequestRouter,
|
||||
|
||||
connection_queue: std.DoublyLinkedList,
|
||||
// NOTE Connection pool has no need for being doubly-linked, but the queue has
|
||||
// (as it's FIFO) and we want a single intrusive Node to be able to participate
|
||||
// in both lists. This is possible because a connection will never belong to
|
||||
// both lists at the same time.
|
||||
connection_pool: std.DoublyLinkedList,
|
||||
connection_buffer: []Connection,
|
||||
|
||||
mutex: std.Thread.Mutex,
|
||||
cond_connection_queued: std.Thread.Condition,
|
||||
cond_connection_freed: std.Thread.Condition,
|
||||
|
||||
/// 2 MiB
|
||||
const huge_page_size = 2 * 1024 * 1024;
|
||||
|
||||
pub const Options = struct {
|
||||
request_router: RequestRouter,
|
||||
address: std.net.Address = .initIp4(.{ 127, 0, 0, 1 }, 80),
|
||||
max_connections: u32 = 128,
|
||||
/// The number of worker threads. If set to `0`, the number of worker
|
||||
/// threads will be equal to the number of logical CPU cores.
|
||||
worker_count: u32 = 0,
|
||||
/// The number of 2 MiB pages reserved for a single read buffer. Each worker
|
||||
/// has its own read buffer. An HTTP request (headers and content combined)
|
||||
/// will be rejected if it is larger than the read buffer.
|
||||
read_buffer_pages: u32 = 1,
|
||||
/// The number of 2 MiB pages reserved for a single write buffer. Each
|
||||
/// worker has its own write buffer. An HTTP response (headers and content
|
||||
/// combined) must be larger than the write buffer.
|
||||
write_buffer_pages: u32 = 1,
|
||||
};
|
||||
|
||||
pub fn init(allocator: std.mem.Allocator, options: Options) !Server {
|
||||
const worker_count = if (options.worker_count > 0) options.worker_count else try std.Thread.getCpuCount();
|
||||
|
||||
// Create socket fd
|
||||
|
||||
const fd: FileDescriptor = try .socket(
|
||||
options.address.any.family,
|
||||
linux.SOCK.STREAM | linux.SOCK.CLOEXEC,
|
||||
if (options.address.any.family == linux.AF.UNIX) 0 else linux.IPPROTO.TCP,
|
||||
);
|
||||
errdefer fd.close();
|
||||
|
||||
const opt = std.mem.toBytes(@as(c_int, 1));
|
||||
try fd.setsockopt(linux.SOL.SOCKET, linux.SO.REUSEADDR, &opt);
|
||||
try fd.setsockopt(linux.SOL.SOCKET, linux.SO.REUSEPORT, &opt);
|
||||
|
||||
var socklen = options.address.getOsSockLen();
|
||||
try fd.bind(&options.address.any, socklen);
|
||||
try fd.listen(options.kernel_backlog);
|
||||
|
||||
var listen_address = options.address;
|
||||
try fd.getsockname(fd, &listen_address, &socklen);
|
||||
|
||||
// Allocate workers
|
||||
|
||||
const workers = try allocator.alloc(Worker, worker_count);
|
||||
errdefer allocator.free(workers);
|
||||
|
||||
// Allocate connection pool
|
||||
|
||||
const connection_buffer = try allocator.alloc(Connection, options.max_connections);
|
||||
errdefer allocator.free(connection_buffer);
|
||||
|
||||
// Allocate and remap read buffers
|
||||
|
||||
const single_read_buffer_size = @as(usize, options.read_buffer_pages) * huge_page_size;
|
||||
const all_read_buffers_size = worker_count * single_read_buffer_size;
|
||||
|
||||
const double_single_read_buffer_size = 2 * single_read_buffer_size;
|
||||
const double_all_read_buffers_size = 2 * all_read_buffers_size;
|
||||
|
||||
const read_buffer_fd: FileDescriptor = try .memfd_create("read_buffer", 0);
|
||||
defer read_buffer_fd.close();
|
||||
|
||||
try read_buffer_fd.ftruncate(all_read_buffers_size);
|
||||
|
||||
const read_buffer_ptr = try errOrPtr(linux.mmap(
|
||||
null,
|
||||
double_all_read_buffers_size,
|
||||
linux.PROT.NONE,
|
||||
linux.MAP{ .TYPE = .PRIVATE, .ANONYMOUS = true },
|
||||
-1,
|
||||
0,
|
||||
));
|
||||
errdefer _ = linux.munmap(read_buffer_ptr, double_all_read_buffers_size);
|
||||
_ = linux.madvise(read_buffer_ptr, double_all_read_buffers_size, linux.MADV.HUGEPAGE);
|
||||
|
||||
for (0..worker_count) |i| {
|
||||
const offset = i * single_read_buffer_size;
|
||||
const double_offset = i * double_single_read_buffer_size;
|
||||
|
||||
try err(linux.mmap(
|
||||
read_buffer_ptr + double_offset,
|
||||
single_read_buffer_size,
|
||||
linux.PROT.READ | linux.PROT.WRITE,
|
||||
linux.MAP{ .TYPE = .SHARED, .FIXED = true },
|
||||
@intFromEnum(read_buffer_fd),
|
||||
offset,
|
||||
));
|
||||
|
||||
try err(linux.mmap(
|
||||
read_buffer_ptr + double_offset + single_read_buffer_size,
|
||||
single_read_buffer_size,
|
||||
linux.PROT.READ | linux.PROT.WRITE,
|
||||
linux.MAP{ .TYPE = .SHARED, .FIXED = true },
|
||||
@intFromEnum(read_buffer_fd),
|
||||
offset,
|
||||
));
|
||||
}
|
||||
|
||||
// Allocate write buffer
|
||||
|
||||
const single_write_buffer_size = @as(usize, options.write_buffer_pages) * huge_page_size;
|
||||
const all_write_buffers_size = worker_count * single_write_buffer_size;
|
||||
|
||||
const write_buffer_ptr = try errOrPtr(linux.mmap(
|
||||
null,
|
||||
all_write_buffers_size,
|
||||
linux.PROT.READ | linux.PROT.WRITE,
|
||||
linux.MAP{ .TYPE = .PRIVATE, .ANONYMOUS = true },
|
||||
-1,
|
||||
0,
|
||||
));
|
||||
errdefer _ = linux.munmap(write_buffer_ptr, all_write_buffers_size);
|
||||
_ = linux.madvise(write_buffer_ptr, all_write_buffers_size, linux.MADV.HUGEPAGE);
|
||||
|
||||
// Initialize workers
|
||||
|
||||
for (workers, 0..) |*worker, i| {
|
||||
const read_offset = i * double_single_read_buffer_size;
|
||||
const write_offset = i * single_write_buffer_size;
|
||||
worker.* = .{
|
||||
.read_buffer_ptr = read_buffer_ptr + read_offset,
|
||||
.read_buffer_size = single_read_buffer_size,
|
||||
.read_head = 0,
|
||||
.read_tail = 0,
|
||||
|
||||
.write_buffer = (write_buffer_ptr + write_offset)[0..single_write_buffer_size],
|
||||
};
|
||||
}
|
||||
|
||||
// Fill connection pool
|
||||
|
||||
var connection_pool: std.DoublyLinkedList = .{};
|
||||
for (connection_buffer) |*c| {
|
||||
connection_pool.prepend(c.node);
|
||||
}
|
||||
|
||||
return .{
|
||||
.fd = fd,
|
||||
.address = listen_address,
|
||||
.workers = workers,
|
||||
.request_router = options.request_router,
|
||||
|
||||
.connection_queue = .{},
|
||||
.connection_pool = connection_pool,
|
||||
.connection_buffer = connection_buffer,
|
||||
|
||||
.mutex = .{},
|
||||
.cond_connection_queued = .{},
|
||||
.cond_connection_freed = .{},
|
||||
};
|
||||
}
|
||||
|
||||
pub fn deinit(self: *Server, allocator: std.mem.Allocator) void {
|
||||
// TODO Deinitialize workers
|
||||
self.fd.close();
|
||||
allocator.free(self.connection_buffer);
|
||||
self.* = undefined;
|
||||
}
|
||||
|
||||
/// This method block until the server is stopped, which is achieved by storing
|
||||
/// `false` into `running`. You should use another thread or interruption
|
||||
/// handler to be able to stop the server.
|
||||
pub fn listen(self: *Server, running: *const std.atomic.Value(bool)) !void {
|
||||
var worker_running: std.atomic.Value(bool) = .init(running.load(.acquire));
|
||||
|
||||
var spawned: usize = 0;
|
||||
defer {
|
||||
worker_running.store(false, .release);
|
||||
self.cond_connection_queued.broadcast();
|
||||
for (self.workers[0..spawned]) |*worker| {
|
||||
worker.thread.join();
|
||||
}
|
||||
}
|
||||
|
||||
for (self.workers) |*worker| {
|
||||
worker.thread = try std.Thread.spawn(.{}, Worker.worker, .{ worker, self, &worker_running });
|
||||
spawned += 1;
|
||||
}
|
||||
|
||||
while (running.load(.acquire)) {
|
||||
var address: std.net.Address = undefined;
|
||||
var address_size: u32 = @sizeOf(std.net.Address);
|
||||
|
||||
const fd = self.fd.accept(&address.any, &address_size) catch |e| {
|
||||
std.log.err("Error while accepting connection: {}", .{e});
|
||||
continue;
|
||||
};
|
||||
|
||||
{
|
||||
self.mutex.lock();
|
||||
defer self.mutex.unlock();
|
||||
|
||||
while (true) {
|
||||
if (self.connection_pool.pop()) |node| {
|
||||
const connection: *Connection = @fieldParentPtr("node", node);
|
||||
connection.fd = fd;
|
||||
connection.address = address;
|
||||
self.connection_queue.prepend(node);
|
||||
break;
|
||||
}
|
||||
|
||||
self.cond_connection_freed.wait(&self.mutex);
|
||||
}
|
||||
}
|
||||
|
||||
self.cond_connection_queued.signal();
|
||||
}
|
||||
}
|
||||
|
||||
fn err(rc: usize) !void {
|
||||
const e = errno(rc);
|
||||
return if (e != .SUCCESS) error.SystemError else rc;
|
||||
}
|
||||
|
||||
fn errOrPtr(rc: usize) ![*]u8 {
|
||||
const e = errno(rc);
|
||||
return if (e != .SUCCESS) error.SystemError else @ptrFromInt(rc);
|
||||
}
|
||||
Reference in New Issue
Block a user