mirror of
https://github.com/lightpanda-io/browser.git
synced 2026-10-08 20:32:00 -04:00
The http max default was 4K with a 16KB hard limit. The default limit is now 1MB with an initial default of 4K. This is to accommodate larger WebDriver payloads.
1227 lines
48 KiB
Zig
1227 lines
48 KiB
Zig
// Copyright (C) 2023-2026 Lightpanda (Selecy SAS)
|
|
//
|
|
// Francis Bouvier <francis@lightpanda.io>
|
|
// Pierre Tachoire <pierre@lightpanda.io>
|
|
//
|
|
// This program is free software: you can redistribute it and/or modify
|
|
// it under the terms of the GNU Affero General Public License as
|
|
// published by the Free Software Foundation, either version 3 of the
|
|
// License, or (at your option) any later version.
|
|
//
|
|
// This program is distributed in the hope that it will be useful,
|
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
|
// GNU Affero General Public License for more details.
|
|
//
|
|
// You should have received a copy of the GNU Affero General Public License
|
|
// along with this program. If not, see <https://www.gnu.org/licenses/>.
|
|
|
|
const std = @import("std");
|
|
const lp = @import("lightpanda");
|
|
|
|
const App = @import("../App.zig");
|
|
const sys_net = @import("../sys/net.zig");
|
|
const header_parser = @import("../network/header_parser.zig");
|
|
const statusCategory = @import("../network/http.zig").statusCategory;
|
|
|
|
const Server = @import("Server.zig");
|
|
const Driver = @import("Driver.zig");
|
|
const bidi_session = @import("bidi/session.zig");
|
|
const http_command = @import("bidi/http_command.zig");
|
|
const uuidv4 = @import("../id.zig").uuidv4;
|
|
|
|
const log = lp.log;
|
|
const posix = std.posix;
|
|
const Allocator = std.mem.Allocator;
|
|
|
|
// A client connection in its http phase: loop-owned, pooled.
|
|
pub const Connection = struct {
|
|
state: State,
|
|
buffer: Buffer,
|
|
socket: posix.socket_t,
|
|
address: sys_net.IpAddress,
|
|
node: std.DoublyLinkedList.Node,
|
|
|
|
// When a keepalive (or just connected) connection should be closed
|
|
deadline: u64,
|
|
|
|
// Response that couldn't be sent without blocking. Socket will switch to
|
|
// "write-mode" until it's drained.
|
|
pending: ?Writing,
|
|
|
|
pub const Writing = struct {
|
|
pos: usize, // how ,uch of Data we've already written
|
|
data: Data,
|
|
keepalive: bool,
|
|
|
|
pub const Data = union(enum) {
|
|
// copied out of the server's scratch buffer; freed once written
|
|
owned: []const u8,
|
|
|
|
// lives as long as the server; referenced, never freed
|
|
static: []const u8,
|
|
|
|
// built by a worker in a pooled arena; released once written
|
|
pooled: Pooled,
|
|
};
|
|
|
|
pub const Pooled = struct {
|
|
arena: *lp.Arena,
|
|
bytes: []const u8,
|
|
};
|
|
|
|
pub fn remaining(self: *const Writing) []const u8 {
|
|
return switch (self.data) {
|
|
.owned, .static => |d| d[self.pos..],
|
|
.pooled => |p| p.bytes[self.pos..],
|
|
};
|
|
}
|
|
|
|
pub fn deinit(self: *const Writing, allocator: Allocator) void {
|
|
switch (self.data) {
|
|
.static => {},
|
|
.owned => |owned| allocator.free(owned),
|
|
.pooled => |p| p.arena.release(),
|
|
}
|
|
}
|
|
};
|
|
|
|
pub fn deinit(self: *Connection) void {
|
|
self.buffer.deinit();
|
|
}
|
|
|
|
// True if the request is in keepalive state and thus is a candidate to be
|
|
// closed if we need its slot for a new connection.
|
|
pub fn isIdle(self: *const Connection) bool {
|
|
if (self.pending != null or self.buffer.len != 0) {
|
|
// has a pending write, or has extra data to read
|
|
return false;
|
|
}
|
|
return self.state == .header;
|
|
}
|
|
|
|
pub const Request = struct {
|
|
method: Method,
|
|
// origin-form, query string stripped, always starts with '/'
|
|
path: []const u8,
|
|
keepalive: bool,
|
|
body: []const u8,
|
|
|
|
// The raw request head (request line + headers, through the final
|
|
// CRLF CRLF); a slice into the read buffer. Upgrade handlers re-parse
|
|
// it for the WebSocket headers.
|
|
head: []const u8,
|
|
|
|
// Filled in by the router for /session/{id}[/...] routes; points
|
|
// into the read buffer like path does.
|
|
session_id: ?*const [36]u8 = null,
|
|
|
|
// valid for handling a single request up to sending the response
|
|
arena: Allocator,
|
|
};
|
|
|
|
pub const Method = enum {
|
|
GET,
|
|
POST,
|
|
PUT,
|
|
DELETE,
|
|
};
|
|
|
|
pub const State = union(enum) {
|
|
header: void, // still parsing the header
|
|
request: Request,
|
|
|
|
// What the connection still needs before the request can be served.
|
|
const Parsed = union(enum) {
|
|
complete,
|
|
need: usize, // bytes the buffer needs will hold, 0 while parsing the header
|
|
};
|
|
|
|
fn parseHeader(self: *State, arena: Allocator, data: []u8) !Parsed {
|
|
const header_index = std.mem.indexOf(u8, data, "\r\n\r\n") orelse {
|
|
return .{ .need = 0 };
|
|
};
|
|
|
|
// include the last line's \r\n so every line, including the request
|
|
// line of a header-less request, is terminated
|
|
const header = data[0 .. header_index + 2];
|
|
const method, const path, const keepalive, const line_1_end = try parseRequestLine(header);
|
|
|
|
_ = line_1_end;
|
|
const body_start = header_index + 4;
|
|
// large content lenghts will saturate to max(usize) -> 413
|
|
const total = body_start +| try contentLength(header);
|
|
if (data.len < total) {
|
|
// the body is still arriving
|
|
return .{ .need = total };
|
|
}
|
|
// A WebSocket upgrade may be pipelined with its first frames, but every
|
|
// client we care about waits for the 101 first. Anything past the
|
|
// declared body is unsupported (and rejects pipelining).
|
|
if (data.len != total) {
|
|
return error.BodyNotSupported;
|
|
}
|
|
|
|
self.* = .{ .request = .{
|
|
.method = method,
|
|
.path = path,
|
|
.keepalive = keepalive,
|
|
.body = data[body_start..total],
|
|
.head = data[0..body_start],
|
|
.arena = arena,
|
|
} };
|
|
|
|
return .complete;
|
|
}
|
|
|
|
fn contentLength(header: []const u8) !usize {
|
|
const key = "\r\ncontent-length:";
|
|
const at = std.ascii.indexOfIgnoreCase(header, key) orelse return 0;
|
|
const start = at + key.len;
|
|
const end = std.mem.indexOfPos(u8, header, start, "\r\n") orelse return error.InvalidHeader;
|
|
const value = std.mem.trim(u8, header[start..end], " \t");
|
|
return std.fmt.parseInt(usize, value, 10) catch error.InvalidHeader;
|
|
}
|
|
|
|
fn parseRequestLine(header: []const u8) !struct { Method, []const u8, bool, usize } {
|
|
const l1 = std.mem.indexOfScalar(u8, header, '\r') orelse return error.InvalidHeader;
|
|
if (l1 == header.len) {
|
|
return error.InvalidHeader;
|
|
}
|
|
if (header[l1 + 1] != '\n') {
|
|
return error.InvalidHeader;
|
|
}
|
|
|
|
var it = std.mem.tokenizeScalar(u8, header[0..l1], ' ');
|
|
const method = std.meta.stringToEnum(Method, it.next() orelse return error.InvalidHeader) orelse return error.InvalidHTTPMethod;
|
|
|
|
// Only the origin-form request-target is accepted; nothing we serve
|
|
// reads the query string, so it's dropped here.
|
|
const target = it.next() orelse return error.InvalidHeader;
|
|
if (target[0] != '/') {
|
|
return error.InvalidHeader;
|
|
}
|
|
const path = target[0 .. std.mem.indexOfScalar(u8, target, '?') orelse target.len];
|
|
|
|
const protocol = it.next() orelse return error.InvalidHeader;
|
|
const keepalive = std.mem.indexOf(u8, protocol, "1.0") == null;
|
|
|
|
return .{ method, path, keepalive, l1 };
|
|
}
|
|
};
|
|
|
|
const Buffer = struct {
|
|
buf: []u8,
|
|
// position in buf up until where we have valid data
|
|
len: usize,
|
|
max: usize,
|
|
allocator: Allocator,
|
|
|
|
fn init(allocator: Allocator, max: usize) !Buffer {
|
|
const real_max = @max(max, INITIAL_BUFFER_SIZE);
|
|
return .{
|
|
.len = 0,
|
|
.max = real_max,
|
|
.buf = try allocator.alloc(u8, INITIAL_BUFFER_SIZE),
|
|
.allocator = allocator,
|
|
};
|
|
}
|
|
|
|
fn deinit(self: *const Buffer) void {
|
|
self.allocator.free(self.buf);
|
|
}
|
|
|
|
fn reset(self: *Buffer) void {
|
|
self.len = 0;
|
|
if (self.buf.len == INITIAL_BUFFER_SIZE) {
|
|
return;
|
|
}
|
|
// keeping the larger buffer is only wasteful, so failure is fine
|
|
self.buf = self.allocator.realloc(self.buf, INITIAL_BUFFER_SIZE) catch self.buf;
|
|
}
|
|
|
|
fn ensureCapacity(self: *Buffer, needed: usize) !void {
|
|
if (needed <= self.buf.len) {
|
|
return;
|
|
}
|
|
if (needed > self.max) {
|
|
return error.RequestTooLarge;
|
|
}
|
|
self.buf = try self.allocator.realloc(self.buf, needed);
|
|
}
|
|
|
|
pub fn read(self: *Buffer, socket: posix.socket_t) ![]u8 {
|
|
const len = self.len;
|
|
if (len == self.buf.len) {
|
|
if (self.buf.len == self.max) {
|
|
return error.RequestTooLarge;
|
|
}
|
|
// Only the header gets here: its length isn't declared, so we
|
|
// double until it fits. A body is sized from Content-Length.
|
|
try self.ensureCapacity(@min(self.buf.len * 2, self.max));
|
|
}
|
|
|
|
const n = try posix.read(socket, self.buf[len..]);
|
|
if (n == 0) {
|
|
return error.ConnectionClosed;
|
|
}
|
|
const total = len + n;
|
|
self.len = total;
|
|
return self.buf[0..total];
|
|
}
|
|
};
|
|
|
|
pub const Pool = struct {
|
|
allocator: Allocator,
|
|
free: std.DoublyLinkedList,
|
|
live: usize, // acquired and not yet released
|
|
retain: usize, // min # to keep
|
|
free_count: usize, // # of connections available in free
|
|
max_buffer_size: usize, // --cdp-max-http-message-size
|
|
|
|
pub fn init(app: *App) !Pool {
|
|
const retain = app.config.maxConnections();
|
|
var self = Pool{
|
|
.live = 0,
|
|
.free = .{},
|
|
.free_count = 0,
|
|
.retain = retain,
|
|
.allocator = app.allocator,
|
|
.max_buffer_size = app.config.cdpMaxHTTPMessageSize(),
|
|
};
|
|
errdefer self.deinit();
|
|
|
|
for (0..retain) |_| {
|
|
const conn = try self.create();
|
|
self.free.append(&conn.node);
|
|
self.free_count += 1;
|
|
}
|
|
return self;
|
|
}
|
|
|
|
// Every live connection must have been released (the server disconnects
|
|
// them all on deinit).
|
|
pub fn deinit(self: *Pool) void {
|
|
lp.assert(self.live == 0, "Connection.Pool.deinit live", .{ .live = self.live });
|
|
while (self.free.popFirst()) |node| {
|
|
const conn: *Connection = @fieldParentPtr("node", node);
|
|
self.destroy(conn);
|
|
}
|
|
}
|
|
|
|
pub fn acquire(self: *Pool) !*Connection {
|
|
const conn = blk: {
|
|
if (self.free.popFirst()) |node| {
|
|
self.free_count -= 1;
|
|
break :blk @as(*Connection, @fieldParentPtr("node", node));
|
|
}
|
|
break :blk try self.create();
|
|
};
|
|
self.live += 1;
|
|
return conn;
|
|
}
|
|
|
|
pub fn release(self: *Pool, conn: *Connection) void {
|
|
self.live -= 1;
|
|
if (self.free_count == self.retain) {
|
|
return self.destroy(conn);
|
|
}
|
|
|
|
conn.node = .{};
|
|
conn.socket = -1;
|
|
conn.address = .{ .ip4 = .unspecified(0) };
|
|
conn.deadline = 0;
|
|
conn.pending = null;
|
|
conn.buffer.reset();
|
|
conn.state = .header;
|
|
|
|
self.free.prepend(&conn.node);
|
|
self.free_count += 1;
|
|
}
|
|
|
|
fn create(self: *Pool) !*Connection {
|
|
const allocator = self.allocator;
|
|
const conn = try allocator.create(Connection);
|
|
errdefer allocator.destroy(conn);
|
|
conn.* = .{
|
|
.node = .{},
|
|
.socket = -1,
|
|
.address = .{ .ip4 = .unspecified(0) },
|
|
.deadline = 0,
|
|
.pending = null,
|
|
.state = .header,
|
|
.buffer = try .init(allocator, self.max_buffer_size),
|
|
};
|
|
return conn;
|
|
}
|
|
|
|
fn destroy(self: *Pool, conn: *Connection) void {
|
|
conn.deinit();
|
|
self.allocator.destroy(conn);
|
|
}
|
|
};
|
|
};
|
|
|
|
// How long a connection may sit without completing a request before we close it.
|
|
pub const IDLE_TIMEOUT_MS = 10_000;
|
|
|
|
// Default buffer size of a new connection. For CDP connections, this should be
|
|
// enough for the few HTTP requests that it makes. WebDriver can send larger
|
|
// bodies and the buffer will grow up to --cdp-max-http-message-size as needed
|
|
pub const INITIAL_BUFFER_SIZE = 4096;
|
|
|
|
const REQUEST_ARENA_RETAIN = 8192;
|
|
|
|
pub fn processEvent(server: *Server, conn: *Connection, rw: Server.IOEvent.ReadWrite, now: u64) void {
|
|
if (conn.pending != null) {
|
|
// registered for OUT only; a hangup shows up as a write error
|
|
if (rw.writable or rw.hangup) {
|
|
flush(server, conn, now);
|
|
}
|
|
return;
|
|
}
|
|
|
|
if (rw.readable) {
|
|
const keepalive = processHTTP(server, conn, now) catch |err| blk: {
|
|
writeError(conn, err);
|
|
break :blk false;
|
|
};
|
|
if (keepalive == false) {
|
|
disconnect(server, conn);
|
|
}
|
|
// else: the socket is level-triggered and stays registered; the
|
|
// deadline was refreshed by processHTTP when the response went out
|
|
} else if (rw.hangup) {
|
|
disconnect(server, conn);
|
|
}
|
|
}
|
|
|
|
// Continues a write that previously hit WouldBlock.
|
|
fn flush(server: *Server, conn: *Connection, now: u64) void {
|
|
const pending = &conn.pending.?;
|
|
const remaining = pending.remaining();
|
|
const n = write(conn.socket, remaining) catch |err| {
|
|
log.debug(.serve, "flush", .{ .err = err });
|
|
return disconnect(server, conn);
|
|
};
|
|
|
|
if (n < remaining.len) {
|
|
// hit a WouldBlock
|
|
pending.pos += n;
|
|
return;
|
|
}
|
|
|
|
// write is complete
|
|
|
|
const keepalive = pending.keepalive;
|
|
pending.deinit(server.app.allocator);
|
|
conn.pending = null;
|
|
|
|
if (keepalive == false) {
|
|
return disconnect(server, conn);
|
|
}
|
|
server.io_engine.waitReadable(conn) catch |err| {
|
|
log.err(.serve, "wait readable", .{ .err = err });
|
|
return disconnect(server, conn);
|
|
};
|
|
touch(server, conn, now);
|
|
}
|
|
|
|
fn processHTTP(server: *Server, conn: *Connection, now: u64) !bool {
|
|
const http = &conn.state;
|
|
const arena = server.request_arena.allocator();
|
|
while (true) {
|
|
switch (http.*) {
|
|
.header => {
|
|
const data = try conn.buffer.read(conn.socket);
|
|
switch (try http.parseHeader(arena, data)) {
|
|
.need => |needed| {
|
|
try conn.buffer.ensureCapacity(needed);
|
|
return true;
|
|
},
|
|
.complete => {},
|
|
}
|
|
if (comptime lp.IS_DEBUG) {
|
|
// we do have a complete header, the state must have transitioned
|
|
// to .request
|
|
std.debug.assert(http.* == .request);
|
|
}
|
|
},
|
|
.request => |*req| {
|
|
defer _ = server.request_arena.reset(.{ .retain_with_limit = REQUEST_ARENA_RETAIN });
|
|
const served = try serveHTTP(server, conn, req);
|
|
if (served == .upgraded) {
|
|
// The fd moved to a WebSocket (and out of server.http); all
|
|
// that's left of this Connection is to recycle it.
|
|
// upgradeConnection already took it out of http_connections.
|
|
recycle(server, conn);
|
|
return true;
|
|
}
|
|
|
|
// req lives in http.*; read what we need before resetting it
|
|
const keepalive = req.keepalive;
|
|
http.* = .header;
|
|
// safe to free the buffer, a parked command wil have copied
|
|
// what it needed from it.
|
|
conn.buffer.reset();
|
|
|
|
if (served == .parked) {
|
|
// off the loop until its worker answers (resumeParked)
|
|
return true;
|
|
}
|
|
|
|
if (conn.pending != null) {
|
|
// We got a WouldBlock and now have a pending write. The
|
|
// connection stays alive until we flush it. After the write
|
|
// if flushed, we'll apply the keepalive result.
|
|
return true;
|
|
}
|
|
|
|
if (keepalive == false) {
|
|
return false;
|
|
}
|
|
touch(server, conn, now);
|
|
return true;
|
|
},
|
|
}
|
|
}
|
|
}
|
|
|
|
// Error responses use a minimal, uniform shape: no reason phrase, an explicit
|
|
// Connection: Close, and no Content-Type. errorResponse builds it at comptime.
|
|
const invalid_request_response = errorResponse(400, "Invalid request");
|
|
const invalid_protocol_response = errorResponse(400, "Invalid HTTP protocol");
|
|
const missing_header_response = errorResponse(400, "Missing required header");
|
|
const forbidden_origin_response = errorResponse(403, "Origin not allowed");
|
|
const forbidden_host_response = errorResponse(403, "Host not allowed");
|
|
const request_too_large_response = errorResponse(413, "Request too large");
|
|
const not_found_response = errorResponse(404, "Not found");
|
|
const session_connected_response = errorResponse(409, "Session already connected");
|
|
const session_busy_response = errorResponse(429, "Session is releasing its previous connection");
|
|
const method_not_allowed_response = errorResponse(405, "Method not allowed");
|
|
const service_unavailable_response = errorResponse(503, "Too many connections");
|
|
const internal_error_response = errorResponse(500, "Internal server error");
|
|
const empty_json_list_response = staticResponse(.{ .status = "200 OK", .body = "[]", .content_type = "application/json; charset=UTF-8" });
|
|
// WebDriver's discovery endpoint; `ready` is whether a new session can be
|
|
// created, which the bootstrap never refuses.
|
|
const status_response = staticResponse(.{ .status = "200 OK", .body = "{\"value\":{\"ready\":true,\"message\":\"\"}}", .content_type = "application/json; charset=UTF-8" });
|
|
const delete_session_response = staticResponse(.{ .status = "200 OK", .body = "{\"value\":null}", .content_type = "application/json; charset=UTF-8" });
|
|
const protocol_response = staticResponse(.{ .status = "200 OK", .body = @embedFile("../data/protocol.json"), .content_type = "application/json; charset=UTF-8" });
|
|
// A parked command whose worker went away before answering.
|
|
pub const session_ended_response = staticResponse(.{
|
|
.status = "404 Not Found",
|
|
.body = "{\"value\":{\"error\":\"invalid session id\",\"message\":\"session ended\",\"stacktrace\":\"\"}}",
|
|
.content_type = "application/json; charset=UTF-8",
|
|
.close = true,
|
|
});
|
|
|
|
const Served = enum {
|
|
responded,
|
|
upgraded,
|
|
parked,
|
|
};
|
|
|
|
const Route = struct {
|
|
method: Connection.Method,
|
|
// exact match against the normalized path
|
|
path: []const u8,
|
|
gate: Gate = .none,
|
|
handler: *const fn (*Server, *Connection, *Connection.Request) anyerror!Served,
|
|
|
|
// A closed gate makes the route invisible (404), not forbidden.
|
|
const Gate = enum {
|
|
none,
|
|
cdp,
|
|
webdriver,
|
|
metrics,
|
|
};
|
|
};
|
|
|
|
const routes = [_]Route{
|
|
.{ .method = .GET, .path = "/", .gate = .cdp, .handler = upgradeCDP },
|
|
.{ .method = .GET, .path = "/metrics", .gate = .metrics, .handler = serveMetrics },
|
|
.{ .method = .GET, .path = "/json/version", .gate = .cdp, .handler = serveJSONVersion },
|
|
.{ .method = .GET, .path = "/json/list", .gate = .cdp, .handler = serveJSONList },
|
|
.{ .method = .GET, .path = "/json", .gate = .cdp, .handler = serveJSONList },
|
|
.{ .method = .GET, .path = "/json/protocol", .gate = .cdp, .handler = serveJSONProtocol },
|
|
// /session is the path Firefox advertises its BiDi endpoint on
|
|
.{ .method = .GET, .path = "/session", .gate = .webdriver, .handler = upgradeBiDi },
|
|
.{ .method = .POST, .path = "/session", .gate = .webdriver, .handler = newSession },
|
|
.{ .method = .GET, .path = "/status", .gate = .webdriver, .handler = serveStatus },
|
|
};
|
|
|
|
const session_routes = [_]Route{
|
|
.{ .method = .GET, .path = "", .handler = upgradeSession },
|
|
.{ .method = .DELETE, .path = "", .handler = deleteSession },
|
|
};
|
|
|
|
// /session/{id} itself is session_routes; everything under it is a command
|
|
// (http_command.parse).
|
|
const SESSION_PREFIX = "/session/";
|
|
|
|
const SESSION_ID_LEN = 36;
|
|
|
|
fn serveHTTP(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
var path = req.path;
|
|
if (path.len > 1 and path[path.len - 1] == '/') {
|
|
path = path[0 .. path.len - 1];
|
|
}
|
|
|
|
if (std.mem.startsWith(u8, path, SESSION_PREFIX) and path.len >= SESSION_PREFIX.len + SESSION_ID_LEN) {
|
|
if (!server.protocols.webdriver) {
|
|
return serveNotFound(server, conn, req);
|
|
}
|
|
const tail = path[SESSION_PREFIX.len + SESSION_ID_LEN ..];
|
|
if (tail.len != 0 and tail[0] != '/') {
|
|
return serveNotFound(server, conn, req);
|
|
}
|
|
req.session_id = path[SESSION_PREFIX.len..][0..SESSION_ID_LEN];
|
|
if (tail.len == 0) {
|
|
return dispatch(server, &session_routes, conn, req, tail);
|
|
}
|
|
return serveSessionCommand(server, conn, req, tail);
|
|
}
|
|
return dispatch(server, &routes, conn, req, path);
|
|
}
|
|
|
|
fn dispatch(server: *Server, comptime table: []const Route, conn: *Connection, req: *Connection.Request, path: []const u8) !Served {
|
|
var path_matched = false;
|
|
inline for (table) |route| {
|
|
if (std.mem.eql(u8, route.path, path) and gateOpen(server, route.gate)) {
|
|
if (route.method == req.method) {
|
|
return route.handler(server, conn, req);
|
|
}
|
|
path_matched = true;
|
|
}
|
|
}
|
|
if (path_matched) {
|
|
return serveMethodNotAllowed(server, conn, req);
|
|
}
|
|
return serveNotFound(server, conn, req);
|
|
}
|
|
|
|
// Best effort, connection is being closed. A partial write ends up as a partial
|
|
// write: no pending, no retry.
|
|
fn writeError(conn: *Connection, err: anyerror) void {
|
|
const response: []const u8 = switch (err) {
|
|
error.ConnectionClosed, error.ConnectionResetByPeer, error.BrokenPipe => return,
|
|
error.InvalidHeader, error.InvalidHTTPMethod, error.BodyNotSupported => invalid_request_response,
|
|
error.RequestTooLarge => request_too_large_response,
|
|
else => blk: {
|
|
log.warn(.serve, "serve error", .{ .err = err });
|
|
break :blk internal_error_response;
|
|
},
|
|
};
|
|
recordResponse(response);
|
|
_ = write(conn.socket, response) catch {};
|
|
}
|
|
|
|
// Every response starts with the status line our two builders emit, so the
|
|
// category is read straight off the bytes rather than threaded through.
|
|
fn recordResponse(response: []const u8) void {
|
|
const prefix = "HTTP/1.1 ";
|
|
lp.assert(std.mem.startsWith(u8, response, prefix), "Server.recordResponse status line", .{});
|
|
const status = std.fmt.parseInt(u16, response[prefix.len..][0..3], 10) catch 0;
|
|
lp.metrics.serve_http_requests.incr(statusCategory(status));
|
|
}
|
|
|
|
const Response = union(enum) {
|
|
// lives as long as the server; a queued remainder references it
|
|
static: []const u8,
|
|
// lives in server.scratch until the next response; a queued remainder is copied
|
|
dynamic: []const u8,
|
|
};
|
|
|
|
// Can do a partial write
|
|
fn write(socket: posix.socket_t, data: []const u8) !usize {
|
|
var pos: usize = 0;
|
|
while (pos < data.len) {
|
|
const n = sys_net.write(socket, data[pos..]) catch |err| switch (err) {
|
|
error.WouldBlock => break,
|
|
error.Interrupted => continue,
|
|
else => return err,
|
|
};
|
|
pos += n;
|
|
}
|
|
return pos;
|
|
}
|
|
|
|
// Dynamic responses are built in server.scratch with room for the header
|
|
// reserved up front; once the body length is known the header is written
|
|
// right-aligned against it (the same trick as WS.fillHeader).
|
|
const HEADER_RESERVE = 192;
|
|
fn beginBody(server: *Server) !*std.Io.Writer {
|
|
server.scratch.clearRetainingCapacity();
|
|
try server.scratch.writer.splatByteAll(0, HEADER_RESERVE);
|
|
return &server.scratch.writer;
|
|
}
|
|
|
|
fn serveDynamicHTTPResponse(server: *Server, conn: *Connection, req: *const Connection.Request, status: std.http.Status, comptime content_type: []const u8) !Served {
|
|
return serveHTTPResponse(server, conn, req, .{ .dynamic = fillHeader(server.scratch.written(), status, content_type) });
|
|
}
|
|
|
|
const JSON_CONTENT_TYPE = "application/json; charset=UTF-8";
|
|
|
|
// `buf` is HEADER_RESERVE bytes followed by the body; returns the response,
|
|
// its header right-aligned against the body.
|
|
fn fillHeader(buf: []u8, status: std.http.Status, comptime content_type: []const u8) []u8 {
|
|
const header_format = "HTTP/1.1 {d} {s}\r\n" ++
|
|
"Content-Length: {d}\r\n" ++
|
|
"Content-Type: " ++ content_type ++ "\r\n\r\n";
|
|
|
|
// 3 status digits, the longest std.http.Status phrase ("Network
|
|
// Authentication Required"), a usize's 20 digits
|
|
comptime std.debug.assert(header_format.len + 3 + 31 + 20 <= HEADER_RESERVE);
|
|
|
|
var header_buf: [HEADER_RESERVE]u8 = undefined;
|
|
const header = std.fmt.bufPrint(&header_buf, header_format, .{
|
|
@intFromEnum(status),
|
|
status.phrase() orelse "",
|
|
buf.len - HEADER_RESERVE,
|
|
}) catch unreachable;
|
|
const start = HEADER_RESERVE - header.len;
|
|
@memcpy(buf[start..HEADER_RESERVE], header);
|
|
return buf[start..];
|
|
}
|
|
|
|
fn errorResponse(comptime status: u16, comptime body: []const u8) []const u8 {
|
|
return std.fmt.comptimePrint(
|
|
"HTTP/1.1 {d} \r\nConnection: Close\r\nContent-Length: {d}\r\n\r\n{s}",
|
|
.{ status, body.len, body },
|
|
);
|
|
}
|
|
|
|
fn staticResponse(comptime opts: struct {
|
|
status: []const u8,
|
|
body: []const u8,
|
|
content_type: []const u8 = "text/plain",
|
|
close: bool = false,
|
|
}) []const u8 {
|
|
return std.fmt.comptimePrint("HTTP/1.1 " ++ opts.status ++ "\r\n" ++
|
|
"Content-Length: {d}\r\n" ++
|
|
(if (opts.close) "Connection: Close\r\n" else "") ++
|
|
"Content-Type: " ++ opts.content_type ++ "\r\n\r\n", .{opts.body.len}) ++ opts.body;
|
|
}
|
|
|
|
fn gateOpen(server: *const Server, gate: Route.Gate) bool {
|
|
return switch (gate) {
|
|
.none => true,
|
|
.cdp => server.protocols.cdp,
|
|
.webdriver => server.protocols.webdriver,
|
|
.metrics => server.app.config.metricsEndpointEnabled(),
|
|
};
|
|
}
|
|
|
|
// GET / (cdp)
|
|
fn upgradeCDP(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return upgradeSpawn(server, conn, req, .cdp);
|
|
}
|
|
|
|
// GET /json/version (cdp)
|
|
fn serveJSONVersion(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return serveHTTPResponse(server, conn, req, .{ .static = server.json_version_response });
|
|
}
|
|
|
|
// GET /json/list or GET /json (cdp)
|
|
fn serveJSONList(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return serveHTTPResponse(server, conn, req, .{ .static = empty_json_list_response });
|
|
}
|
|
|
|
// GET /json/protocol (cdp)
|
|
fn serveJSONProtocol(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return serveHTTPResponse(server, conn, req, .{ .static = protocol_response });
|
|
}
|
|
|
|
// GET /metrics (internal)
|
|
fn serveMetrics(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
const writer = try beginBody(server);
|
|
lp.metrics.write(writer);
|
|
return serveDynamicHTTPResponse(server, conn, req, .ok, "text/plain; version=0.0.4; charset=utf-8");
|
|
}
|
|
|
|
// GET /status (webdriver)
|
|
fn serveStatus(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return serveHTTPResponse(server, conn, req, .{ .static = status_response });
|
|
}
|
|
|
|
// GET /session (webdriver (direct bidi))
|
|
fn upgradeBiDi(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return upgradeSpawn(server, conn, req, .bidi);
|
|
}
|
|
|
|
// POST /session (webdriver)
|
|
fn newSession(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
const Capability = struct { webSocketUrl: ?bool = null };
|
|
const parsed = std.json.parseFromSliceLeaky(struct {
|
|
capabilities: ?struct {
|
|
alwaysMatch: ?Capability = null,
|
|
firstMatch: ?[]const Capability = null,
|
|
} = null,
|
|
}, req.arena, req.body, .{ .ignore_unknown_fields = true }) catch {
|
|
return serveWebDriverError(server, conn, req, "invalid argument", "invalid JSON body");
|
|
};
|
|
|
|
if (server.worker_pool.isFull()) {
|
|
lp.metrics.serve_connection_limit.incr();
|
|
return serveWebDriverError(server, conn, req, "session not created", "too many sessions");
|
|
}
|
|
|
|
var session_id: [36]u8 = undefined;
|
|
uuidv4(&session_id);
|
|
|
|
const worker = server.spawnWorker(.bidi, .{ .session = session_id }) catch |err| {
|
|
log.err(.serve, "worker spawn", .{ .err = err });
|
|
return serveWebDriverError(server, conn, req, "session not created", "failed to start the session");
|
|
};
|
|
// The client never learns the id if we fail to answer (e.g. it hung up),
|
|
// so nothing would ever DELETE this session.
|
|
errdefer server.quitSession(worker);
|
|
|
|
const is_requesting_websocket_url = blk: {
|
|
const caps = parsed.capabilities orelse break :blk false;
|
|
if (caps.alwaysMatch) |always| {
|
|
if (always.webSocketUrl == true) {
|
|
break :blk true;
|
|
}
|
|
}
|
|
for (caps.firstMatch orelse &.{}) |first| {
|
|
if (first.webSocketUrl == true) {
|
|
break :blk true;
|
|
}
|
|
}
|
|
break :blk false;
|
|
};
|
|
|
|
const url: ?[]const u8 = blk: {
|
|
if (is_requesting_websocket_url) {
|
|
break :blk try std.fmt.allocPrint(req.arena, "{s}{s}", .{ server.bidi_session_url, &session_id });
|
|
}
|
|
break :blk null;
|
|
};
|
|
|
|
return serveWebDriver(server, conn, req, .ok, .{
|
|
.sessionId = &session_id,
|
|
.capabilities = bidi_session.Capabilities{
|
|
.userAgent = server.app.config.http_headers.user_agent,
|
|
.webSocketUrl = url,
|
|
},
|
|
});
|
|
}
|
|
|
|
// GET /session/ID (webdriver (upgrade to bidi))
|
|
fn upgradeSession(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
const worker = server.findSession(req.session_id.?) orelse {
|
|
return serveNotFound(server, conn, req);
|
|
};
|
|
|
|
if (worker.linkDropping()) {
|
|
// The previous connection is gone but the worker hasn't given the
|
|
// link back yet. Dirver can retry.
|
|
return serveHTTPResponse(server, conn, req, .{ .static = session_busy_response });
|
|
}
|
|
|
|
if (worker.link != null) {
|
|
// already joined
|
|
return serveHTTPResponse(server, conn, req, .{ .static = session_connected_response });
|
|
}
|
|
|
|
return upgrade(server, conn, req, .{ .attach = worker });
|
|
}
|
|
|
|
// DELETE /session/ID (webdriver)
|
|
fn deleteSession(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
const worker = server.findSession(req.session_id.?) orelse {
|
|
return serveNoSuchSession(server, conn, req);
|
|
};
|
|
server.quitSession(worker);
|
|
return serveHTTPResponse(server, conn, req, .{ .static = delete_session_response });
|
|
}
|
|
|
|
// /session/ID/... (webdriver): parsed here, answered by the session's worker.
|
|
fn serveSessionCommand(server: *Server, conn: *Connection, req: *Connection.Request, path: []const u8) !Served {
|
|
const worker = server.findSession(req.session_id.?) orelse {
|
|
return serveNoSuchSession(server, conn, req);
|
|
};
|
|
if (worker.http_request != null) {
|
|
return serveCommandInProgress(server, conn, req);
|
|
}
|
|
|
|
const arena = try server.app.arena_pool.acquire(req.body.len, "http command");
|
|
const command = http_command.parse(arena.allocator(), req.method, path, req.body) catch |err| {
|
|
arena.release();
|
|
switch (err) {
|
|
error.OutOfMemory => return err,
|
|
error.UnknownCommand => return serveWebDriverError(server, conn, req, "unknown command", "unknown command"),
|
|
error.UnknownMethod => return serveWebDriverError(server, conn, req, "unknown method", "unknown method"),
|
|
error.InvalidArgument => return serveWebDriverError(server, conn, req, "invalid argument", "invalid body"),
|
|
}
|
|
};
|
|
server.parkRequest(worker, conn, req.keepalive, arena, command);
|
|
return .parked;
|
|
}
|
|
|
|
fn serveCommandInProgress(server: *Server, conn: *Connection, req: *const Connection.Request) !Served {
|
|
return serveWebDriverError(server, conn, req, "unknown error", "a command is already in progress");
|
|
}
|
|
|
|
fn serveNoSuchSession(server: *Server, conn: *Connection, req: *const Connection.Request) !Served {
|
|
return serveWebDriverError(server, conn, req, "invalid session id", "no such session");
|
|
}
|
|
|
|
// CDP or Bidi directly creating a Worker from an websocket upgrade
|
|
fn upgradeSpawn(server: *Server, conn: *Connection, req: *Connection.Request, protocol: Driver.Protocol) !Served {
|
|
if (server.worker_pool.isFull()) {
|
|
lp.metrics.serve_connection_limit.incr();
|
|
return serveHTTPResponse(server, conn, req, .{ .static = service_unavailable_response });
|
|
}
|
|
return upgrade(server, conn, req, .{ .spawn = protocol });
|
|
}
|
|
|
|
// Answers a HTTP WebDriver request with {"value": value}.
|
|
fn serveWebDriver(server: *Server, conn: *Connection, req: *const Connection.Request, status: std.http.Status, value: anytype) !Served {
|
|
const writer = try beginBody(server);
|
|
try std.json.Stringify.value(.{ .value = value }, .{}, writer);
|
|
return serveDynamicHTTPResponse(server, conn, req, status, JSON_CONTENT_TYPE);
|
|
}
|
|
|
|
fn serveWebDriverError(server: *Server, conn: *Connection, req: *const Connection.Request, code: []const u8, message: []const u8) !Served {
|
|
return serveWebDriver(server, conn, req, webDriverErrorStatus(code), .{
|
|
.@"error" = code,
|
|
.message = message,
|
|
.stacktrace = "",
|
|
});
|
|
}
|
|
|
|
fn serveNotFound(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return serveHTTPResponse(server, conn, req, .{ .static = not_found_response });
|
|
}
|
|
|
|
fn serveMethodNotAllowed(server: *Server, conn: *Connection, req: *Connection.Request) !Served {
|
|
return serveHTTPResponse(server, conn, req, .{ .static = method_not_allowed_response });
|
|
}
|
|
|
|
// Writes what the socket will take now. Anything left is queued on the
|
|
// connection, which switches to waiting for writability.
|
|
fn serveHTTPResponse(server: *Server, conn: *Connection, req: *const Connection.Request, response: Response) !Served {
|
|
const data = switch (response) {
|
|
inline else => |d| d,
|
|
};
|
|
recordResponse(data);
|
|
const n = try write(conn.socket, data);
|
|
if (n == data.len) {
|
|
return .responded;
|
|
}
|
|
|
|
lp.assert(conn.pending == null, "Server.send pending", .{});
|
|
conn.pending = .{
|
|
.pos = 0,
|
|
.keepalive = req.keepalive,
|
|
.data = switch (response) {
|
|
.static => .{ .static = data[n..] },
|
|
.dynamic => .{ .owned = try server.app.allocator.dupe(u8, data[n..]) },
|
|
},
|
|
};
|
|
// on failure the caller disconnects, which frees pending
|
|
try server.io_engine.waitWritable(conn);
|
|
return .responded;
|
|
}
|
|
|
|
// A parked connection's response is ready, join the loop so that we can start
|
|
// writing the response.
|
|
pub fn resumeParked(server: *Server, conn: *Connection, req_keepalive: bool, response: Connection.Writing.Data, now: u64) void {
|
|
const allocator = server.app.allocator;
|
|
// beginShutdown closed the http connections, this one was off the loop then
|
|
const keepalive = req_keepalive and server.shutdown_begun == false;
|
|
var writing: Connection.Writing = .{ .pos = 0, .data = response, .keepalive = keepalive };
|
|
|
|
server.io_engine.monitorHTTP(conn) catch |err| {
|
|
log.err(.serve, "resume monitor", .{ .err = err });
|
|
writing.deinit(allocator);
|
|
sys_net.close(conn.socket);
|
|
return recycle(server, conn);
|
|
};
|
|
conn.deadline = now + IDLE_TIMEOUT_MS;
|
|
server.http_connections.append(&conn.node);
|
|
|
|
const data = writing.remaining();
|
|
recordResponse(data);
|
|
writing.pos = write(conn.socket, data) catch |err| {
|
|
log.debug(.serve, "resume write", .{ .err = err });
|
|
writing.deinit(allocator);
|
|
return disconnect(server, conn);
|
|
};
|
|
|
|
if (writing.pos < data.len) {
|
|
// disconnect frees it from here
|
|
conn.pending = writing;
|
|
server.io_engine.waitWritable(conn) catch |err| {
|
|
log.err(.serve, "wait writable", .{ .err = err });
|
|
return disconnect(server, conn);
|
|
};
|
|
return;
|
|
}
|
|
|
|
writing.deinit(allocator);
|
|
if (keepalive == false) {
|
|
disconnect(server, conn);
|
|
}
|
|
}
|
|
|
|
// A complete {"value": value} response, built by a worker for resumeParked,
|
|
// in an arena the loop releases once it's written.
|
|
pub fn webDriverResponse(arena: *lp.Arena, status: std.http.Status, value: anytype) ![]const u8 {
|
|
var aw: std.Io.Writer.Allocating = try .initCapacity(arena.allocator(), 512);
|
|
try aw.writer.splatByteAll(0, HEADER_RESERVE);
|
|
try std.json.Stringify.value(.{ .value = value }, .{}, &aw.writer);
|
|
return fillHeader(aw.written(), status, JSON_CONTENT_TYPE);
|
|
}
|
|
|
|
// W3C WebDriver's error table; every code not listed is a 500.
|
|
pub fn webDriverErrorStatus(code: []const u8) std.http.Status {
|
|
const statuses = std.StaticStringMap(std.http.Status).initComptime(.{
|
|
.{ "detached shadow root", .not_found },
|
|
.{ "element click intercepted", .bad_request },
|
|
.{ "element not interactable", .bad_request },
|
|
.{ "insecure certificate", .bad_request },
|
|
.{ "invalid argument", .bad_request },
|
|
.{ "invalid cookie domain", .bad_request },
|
|
.{ "invalid element state", .bad_request },
|
|
.{ "invalid selector", .bad_request },
|
|
.{ "invalid session id", .not_found },
|
|
.{ "no such alert", .not_found },
|
|
.{ "no such cookie", .not_found },
|
|
.{ "no such element", .not_found },
|
|
.{ "no such frame", .not_found },
|
|
.{ "no such shadow root", .not_found },
|
|
.{ "no such window", .not_found },
|
|
.{ "stale element reference", .not_found },
|
|
.{ "unknown command", .not_found },
|
|
.{ "unknown method", .method_not_allowed },
|
|
});
|
|
return statuses.get(code) orelse .internal_server_error;
|
|
}
|
|
|
|
// HTTP-phase teardown. Websockets tear down via releaseWorker.
|
|
pub fn disconnect(server: *Server, conn: *Connection) void {
|
|
server.io_engine.remove(conn.socket);
|
|
sys_net.close(conn.socket);
|
|
if (conn.pending) |*pending| {
|
|
pending.deinit(server.app.allocator);
|
|
conn.pending = null;
|
|
}
|
|
server.http_connections.remove(&conn.node);
|
|
recycle(server, conn);
|
|
}
|
|
|
|
// Return a connection to the pool; a slot in the fd budget is free.
|
|
fn recycle(server: *Server, conn: *Connection) void {
|
|
server.http_connection_pool.release(conn);
|
|
server.slotFreed();
|
|
}
|
|
|
|
fn touch(server: *Server, conn: *Connection, now: u64) void {
|
|
conn.deadline = now + IDLE_TIMEOUT_MS;
|
|
const node = &conn.node;
|
|
if (server.http_connections.last == node) {
|
|
return;
|
|
}
|
|
|
|
server.http_connections.remove(&conn.node);
|
|
server.http_connections.append(&conn.node);
|
|
}
|
|
|
|
pub fn buildJSONVersionResponse(app: *const App, port: u16) ![]const u8 {
|
|
const host = app.config.advertiseHost();
|
|
if (app.config.bindIsWildcard()) {
|
|
// Serve is bound to INADDR_ANY but no --advertise-host was given;
|
|
// advertiseHost() falls back to 127.0.0.1 so clients can still
|
|
// connect locally. Surface the trade-off so users running
|
|
// outside the same host know they have to opt in.
|
|
log.note(.cdp, "advertising loopback for wildcard bind", .{
|
|
.message = "--host is a wildcard (0.0.0.0 / ::) without --advertise-host; clients on other hosts will need --advertise-host to reach the CDP endpoint",
|
|
});
|
|
}
|
|
const body_format =
|
|
"{{" ++
|
|
"\"Browser\": \"Lightpanda/1.0\", " ++
|
|
"\"Protocol-Version\": \"1.3\", " ++
|
|
"\"User-Agent\": \"Lightpanda/1.0\", " ++
|
|
"\"Lightpanda-Version\": \"" ++ lp.build_config.version ++ "\", " ++
|
|
"\"webSocketDebuggerUrl\": \"ws://{s}:{d}/\"" ++
|
|
"}}";
|
|
const body_len = std.fmt.count(body_format, .{ host, port });
|
|
|
|
const response_format =
|
|
"HTTP/1.1 200 OK\r\n" ++
|
|
"Content-Length: {d}\r\n" ++
|
|
"Content-Type: application/json; charset=UTF-8\r\n\r\n" ++
|
|
body_format;
|
|
return try std.fmt.allocPrint(app.allocator, response_format, .{ body_len, host, port });
|
|
}
|
|
|
|
// Where the upgraded socket goes: a new worker, or an existing session's.
|
|
const Upgrade = union(enum) {
|
|
spawn: Driver.Protocol,
|
|
attach: *Server.Worker,
|
|
};
|
|
|
|
// Shared upgrade path: validate the WebSocket headers, write the 101, and
|
|
// hand the fd to its worker (spawning one for a new connection).
|
|
fn upgrade(server: *Server, conn: *Connection, req: *Connection.Request, target: Upgrade) !Served {
|
|
var accept_buf: [28]u8 = undefined;
|
|
const accept_key = webSocketAccept(req.head, &accept_buf) catch |err| {
|
|
const response: []const u8 = switch (err) {
|
|
error.ForbiddenOrigin => forbidden_origin_response,
|
|
error.ForbiddenHost => forbidden_host_response,
|
|
error.InvalidProtocol => invalid_protocol_response,
|
|
error.MissingHeader => missing_header_response,
|
|
else => invalid_request_response,
|
|
};
|
|
return serveHTTPResponse(server, conn, req, .{ .static = response });
|
|
};
|
|
|
|
// The 101 is ~129 bytes into an empty send buffer, so a single write
|
|
// always completes; a partial write here means the peer is already gone.
|
|
var response_buf: [160]u8 = undefined;
|
|
const response = std.fmt.bufPrint(&response_buf, "HTTP/1.1 101 Switching Protocols\r\n" ++
|
|
"Upgrade: websocket\r\n" ++
|
|
"Connection: upgrade\r\n" ++
|
|
"Sec-Websocket-Accept: {s}\r\n\r\n", .{accept_key}) catch unreachable;
|
|
const n = write(conn.socket, response) catch return error.ConnectionClosed;
|
|
if (n != response.len) {
|
|
return error.ConnectionClosed;
|
|
}
|
|
|
|
switch (target) {
|
|
.spawn => |protocol| server.upgradeConnection(conn, protocol),
|
|
.attach => |worker| server.attachConnection(worker, conn),
|
|
}
|
|
return .upgraded;
|
|
}
|
|
|
|
// Validate an incoming WebSocket upgrade request head and, on success, write
|
|
// the Sec-WebSocket-Accept value into `out`. Mirrors the origin/host defenses
|
|
// from the old Handshake path.
|
|
fn webSocketAccept(head: []const u8, out: *[28]u8) ![]const u8 {
|
|
const FOUND_UPGRADE: u8 = 1 << 0;
|
|
const FOUND_VERSION: u8 = 1 << 1;
|
|
const FOUND_CONNECTION: u8 = 1 << 2;
|
|
const FOUND_KEY: u8 = 1 << 3;
|
|
const FOUND_ALL = FOUND_UPGRADE | FOUND_VERSION | FOUND_CONNECTION | FOUND_KEY;
|
|
|
|
const method, _, const version, var it = header_parser.parseRequest(head) catch return error.InvalidRequest;
|
|
if (method != .get or version != .@"1.1") {
|
|
return error.InvalidProtocol;
|
|
}
|
|
|
|
var found: u8 = 0;
|
|
var key: []const u8 = "";
|
|
while (it.next() catch return error.InvalidRequest) |h| {
|
|
if (std.ascii.eqlIgnoreCase(h.key, "upgrade")) {
|
|
if (!std.ascii.eqlIgnoreCase("websocket", h.value)) return error.MissingHeader;
|
|
found |= FOUND_UPGRADE;
|
|
} else if (std.ascii.eqlIgnoreCase(h.key, "sec-websocket-version")) {
|
|
if (h.value.len != 2 or h.value[0] != '1' or h.value[1] != '3') return error.MissingHeader;
|
|
found |= FOUND_VERSION;
|
|
} else if (std.ascii.eqlIgnoreCase(h.key, "connection")) {
|
|
if (std.ascii.indexOfIgnoreCase(h.value, "upgrade") == null) return error.MissingHeader;
|
|
found |= FOUND_CONNECTION;
|
|
} else if (std.ascii.eqlIgnoreCase(h.key, "sec-websocket-key")) {
|
|
key = h.value;
|
|
found |= FOUND_KEY;
|
|
} else if (std.ascii.eqlIgnoreCase(h.key, "origin")) {
|
|
// Only a browser sends Origin, and a browser has no business
|
|
// driving CDP/BiDi: it's cross-origin to us by definition.
|
|
log.warn(.serve, "rejected websocket origin", .{ .origin = h.value[0..@min(h.value.len, 64)] });
|
|
return error.ForbiddenOrigin;
|
|
} else if (std.ascii.eqlIgnoreCase(h.key, "host")) {
|
|
// Defense in depth against DNS rebinding: only an IP literal can
|
|
// legitimately reach us (no name resolution involved). The one
|
|
// name we accept is `localhost:<port>`, which browsers hardwire
|
|
// to loopback without a lookup.
|
|
if (!std.mem.startsWith(u8, h.value, "localhost:")) {
|
|
_ = std.Io.net.IpAddress.parseLiteral(h.value) catch {
|
|
log.warn(.serve, "rejected websocket host", .{ .host = h.value[0..@min(h.value.len, 64)] });
|
|
return error.ForbiddenHost;
|
|
};
|
|
}
|
|
}
|
|
}
|
|
if (found != FOUND_ALL) {
|
|
return error.MissingHeader;
|
|
}
|
|
|
|
var sha: [20]u8 = undefined;
|
|
var hasher = std.crypto.hash.Sha1.init(.{});
|
|
hasher.update(key);
|
|
hasher.update("258EAFA5-E914-47DA-95CA-C5AB0DC85B11");
|
|
hasher.final(&sha);
|
|
_ = std.base64.standard.Encoder.encode(out, &sha);
|
|
return out;
|
|
}
|
|
|
|
const testing = @import("../testing.zig");
|
|
|
|
test "http: the read buffer grows with the request and gives the space back" {
|
|
var pair: [2]posix.socket_t = undefined;
|
|
if (std.c.socketpair(posix.AF.LOCAL, posix.SOCK.STREAM, 0, &pair) != 0) {
|
|
return error.SocketPairFailed;
|
|
}
|
|
defer sys_net.close(pair[0]);
|
|
defer sys_net.close(pair[1]);
|
|
|
|
const max = INITIAL_BUFFER_SIZE * 2;
|
|
var buffer = try Connection.Buffer.init(testing.allocator, max);
|
|
defer buffer.deinit();
|
|
|
|
// a connection commits the initial size, never the limit
|
|
try testing.expectEqual(INITIAL_BUFFER_SIZE, buffer.buf.len);
|
|
|
|
// a header declares no length, so the buffer doubles to take it
|
|
const filler = "a" ** max;
|
|
try sys_net.writeAll(pair[1], filler);
|
|
while (buffer.len < filler.len) {
|
|
_ = try buffer.read(pair[0]);
|
|
}
|
|
try testing.expectEqual(max, buffer.buf.len);
|
|
|
|
// and stops doubling at the limit
|
|
try sys_net.writeAll(pair[1], "a");
|
|
try testing.expectError(error.RequestTooLarge, buffer.read(pair[0]));
|
|
|
|
// the next request on this connection starts small again
|
|
buffer.reset();
|
|
try testing.expectEqual(0, buffer.len);
|
|
try testing.expectEqual(INITIAL_BUFFER_SIZE, buffer.buf.len);
|
|
}
|
|
|
|
test "http: a declared body is sized upfront" {
|
|
var pair: [2]posix.socket_t = undefined;
|
|
if (std.c.socketpair(posix.AF.LOCAL, posix.SOCK.STREAM, 0, &pair) != 0) {
|
|
return error.SocketPairFailed;
|
|
}
|
|
defer sys_net.close(pair[0]);
|
|
defer sys_net.close(pair[1]);
|
|
|
|
var buffer = try Connection.Buffer.init(testing.allocator, 1024 * 1024);
|
|
defer buffer.deinit();
|
|
|
|
var state: Connection.State = .header;
|
|
const body_len = INITIAL_BUFFER_SIZE * 4;
|
|
var head_buf: [64]u8 = undefined;
|
|
const head = try std.fmt.bufPrint(&head_buf, "POST /session HTTP/1.1\r\nContent-Length: {d}\r\n\r\n", .{body_len});
|
|
try sys_net.writeAll(pair[1], head);
|
|
|
|
// the header alone is enough to know how much room the body needs
|
|
const needed = switch (try state.parseHeader(testing.allocator, try buffer.read(pair[0]))) {
|
|
.complete => return error.UnexpectedlyComplete,
|
|
.need => |n| n,
|
|
};
|
|
try testing.expectEqual(head.len + body_len, needed);
|
|
try buffer.ensureCapacity(needed);
|
|
try testing.expectEqual(needed, buffer.buf.len);
|
|
|
|
// a body that can't fit is rejected without reading any of it
|
|
try testing.expectError(error.RequestTooLarge, buffer.ensureCapacity(buffer.max + 1));
|
|
}
|