Massive logging, improve connection closing/timeout handling
This commit is contained in:
@@ -9,6 +9,7 @@ const Request = @import("Request.zig");
|
||||
const RequestHandler = @import("RequestHandler.zig");
|
||||
const Worker = @import("Worker.zig");
|
||||
|
||||
const log = std.log.scoped(.Server);
|
||||
const linux = std.os.linux;
|
||||
const errno = linux.E.init;
|
||||
|
||||
@@ -18,6 +19,7 @@ ssl_ctx: ?*openssl.SslContext,
|
||||
workers: []Worker,
|
||||
threads: []std.Thread,
|
||||
request_handler: RequestHandler,
|
||||
read_timeout_us: u64,
|
||||
|
||||
connection_queue: std.DoublyLinkedList,
|
||||
// NOTE Connection pool has no need for being doubly-linked, but the queue has
|
||||
@@ -65,6 +67,10 @@ pub const Options = struct {
|
||||
/// bodies generated with the body writer and not to bodies sent with
|
||||
/// `sendfile`.
|
||||
body_write_buffer_huge_pages: u32 = 1,
|
||||
/// How much time should a worker wait on an idle connection before closing
|
||||
/// it. Specifically, how much time can a `read` syscall block for, before
|
||||
/// the connection is forcefully closed.
|
||||
read_timeout_us: u64 = 1 * std.time.us_per_s,
|
||||
};
|
||||
|
||||
pub fn init(allocator: std.mem.Allocator, options: Options) !Server {
|
||||
@@ -220,6 +226,7 @@ pub fn init(allocator: std.mem.Allocator, options: Options) !Server {
|
||||
.workers = workers,
|
||||
.threads = threads,
|
||||
.request_handler = options.request_handler,
|
||||
.read_timeout_us = options.read_timeout_us,
|
||||
|
||||
.connection_queue = .{},
|
||||
.connection_pool = connection_pool,
|
||||
@@ -232,6 +239,7 @@ pub fn init(allocator: std.mem.Allocator, options: Options) !Server {
|
||||
}
|
||||
|
||||
pub fn deinit(self: *Server, allocator: std.mem.Allocator) void {
|
||||
log.debug("Deinitializing Server.", .{});
|
||||
const worker_count = self.workers.len;
|
||||
|
||||
const single_read_buffers_size = self.workers[0].read_buffer_size;
|
||||
@@ -273,14 +281,18 @@ pub fn listen(self: *Server, running: *const std.atomic.Value(bool)) !void {
|
||||
|
||||
var spawned: usize = 0;
|
||||
defer {
|
||||
log.debug("Storing `false` into worker_running.", .{});
|
||||
worker_running.store(false, .release);
|
||||
log.debug("Broadcasting connection queued condition variable.", .{});
|
||||
self.cond_connection_queued.broadcast();
|
||||
for (self.threads[0..spawned]) |*thread| {
|
||||
for (self.threads[0..spawned], 0..) |*thread, i| {
|
||||
log.debug("Joining the thread of worker #{d}.", .{i});
|
||||
thread.join();
|
||||
}
|
||||
}
|
||||
|
||||
for (self.workers, 0..) |*worker, i| {
|
||||
log.debug("Spawning thread for worker #{d}.", .{i});
|
||||
self.threads[i] = try std.Thread.spawn(.{}, Worker.worker, .{ worker, self, &worker_running });
|
||||
spawned += 1;
|
||||
}
|
||||
@@ -289,34 +301,53 @@ pub fn listen(self: *Server, running: *const std.atomic.Value(bool)) !void {
|
||||
var address: std.net.Address = undefined;
|
||||
var address_size: u32 = @sizeOf(std.net.Address);
|
||||
|
||||
log.debug("Accepting connection.", .{});
|
||||
const fd = self.fd.accept(&address.any, &address_size) catch |e| {
|
||||
std.log.err("Error while accepting connection: {}", .{e});
|
||||
log.err("Error while accepting connection: {}", .{e});
|
||||
continue;
|
||||
};
|
||||
log.debug("Accepted connection from {f}", .{address});
|
||||
|
||||
const timeout: linux.timeval = .{
|
||||
.sec = @intCast(self.read_timeout_us / std.time.us_per_s),
|
||||
.usec = @intCast(self.read_timeout_us % std.time.us_per_s),
|
||||
};
|
||||
try fd.setsockopt(linux.SOL.SOCKET, linux.SO.RCVTIMEO, std.mem.asBytes(&timeout));
|
||||
|
||||
const ssl: ?*openssl.Ssl = self.maybeInitSsl(fd) catch |e| {
|
||||
std.log.err("Error while estabilishing SSL connection: {}", .{e});
|
||||
log.err("Error while estabilishing SSL connection: {}", .{e});
|
||||
fd.close();
|
||||
continue;
|
||||
};
|
||||
|
||||
{
|
||||
log.debug("Acquiring mutex.", .{});
|
||||
self.mutex.lock();
|
||||
defer self.mutex.unlock();
|
||||
log.debug("Acquired mutex.", .{});
|
||||
defer {
|
||||
log.debug("Unlocking mutex.", .{});
|
||||
self.mutex.unlock();
|
||||
}
|
||||
|
||||
while (true) {
|
||||
if (self.connection_pool.pop()) |node| {
|
||||
const connection: *Connection = @fieldParentPtr("node", node);
|
||||
connection.reinit(address, fd, ssl);
|
||||
log.debug("Adding connection to {f} to the connection queue.", .{connection.address});
|
||||
self.connection_queue.prepend(node);
|
||||
break;
|
||||
}
|
||||
|
||||
log.debug("Waiting on connection freed condition variable.", .{});
|
||||
self.cond_connection_freed.wait(&self.mutex);
|
||||
log.debug("Woken up on connection freed condition variable.", .{});
|
||||
}
|
||||
}
|
||||
|
||||
log.debug("Signaling connection queued condition variable.", .{});
|
||||
self.cond_connection_queued.signal();
|
||||
} else {
|
||||
log.debug("Loaded `false` from running, the accept loop exited.", .{});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -324,7 +355,9 @@ fn maybeInitSsl(self: *const Server, fd: FileDescriptor) !?*openssl.Ssl {
|
||||
if (self.ssl_ctx) |ssl_ctx| {
|
||||
const ssl = try openssl.Ssl.new(ssl_ctx);
|
||||
try ssl.setFd(fd);
|
||||
log.debug("Accepting SSL layer.", .{});
|
||||
try ssl.accept();
|
||||
log.debug("Accepted SSL layer.", .{});
|
||||
return ssl;
|
||||
} else {
|
||||
return null;
|
||||
|
||||
Reference in New Issue
Block a user