tickets
All repositories: gitoria
112.4 KB
// hl:http1 plugin — HTTP/1.1 server with epoll + I/O thread pool// Compiled to libhttp1.so, loaded by runtime via dlopen//// Architecture:// Acceptor thread (epoll on server fd) → PARK (epoll) → I/O thread pool (TLS + parse)// → request queue → Hybriel main thread//// An IDLE connection never occupies an I/O worker (mission 084). A worker's read() is// blocking, so a connection handed straight to the pool pins a thread until bytes arrive;// with the default `threads = 4`, four idle keep-alive sockets (one open browser tab is// already several) starved every later request. Instead every connection that is not// known to have readable bytes is PARKED in the parker thread's epoll, and only enters// conn_queue when it is actually readable (or hung up). See ParkedConns below.//// Exports:// hl_http1_create_server(port, host, cert_path, key_path, threads) → iterator of request objects// hl_http1_listen(port) → same with defaults (backward compat)//// Each request object: method, path, query (object), headers (object), body, respond (handle),// remoteAddress, bytes (the body as a Bytes, ticket #88)// Respond handle: call("send", status_code, body [, content_type]) → writes responseconst std = @import("std");const api = @import("plugin_api");const http = @import("http_common");const ws = @import("ws_common");const HlValue = api.HlValue;const HlObject = api.HlObject;const HlField = api.HlField;const HlIterator = api.HlIterator;const HlHandle = api.HlHandle;const HlString = api.HlString;const linux = std.os.linux;const posix = std.posix;const c = std.c;const PthreadMutex = c.pthread_mutex_t;const PthreadCond = c.pthread_cond_t;fn mutexInit(m: *PthreadMutex) void {// Static initializer is sufficient; explicit init not exposed in this std.c._ = m;}fn mutexLock(m: *PthreadMutex) void {_ = c.pthread_mutex_lock(m);}fn mutexUnlock(m: *PthreadMutex) void {_ = c.pthread_mutex_unlock(m);}fn condInit(cnd: *PthreadCond) void {_ = cnd;}fn condSignal(cnd: *PthreadCond) void {_ = c.pthread_cond_signal(cnd);}fn condBroadcast(cnd: *PthreadCond) void {_ = c.pthread_cond_broadcast(cnd);}fn condWait(cnd: *PthreadCond, m: *PthreadMutex) void {_ = c.pthread_cond_wait(cnd, m);}// Use the SMP (thread-safe, production) allocator rather than the debug// GeneralPurposeAllocator: this plugin allocates per-request across the acceptor,// I/O-worker and interpreter threads, and the debug allocator's safety bookkeeping// (canaries/quarantine) segfaulted in its own free() path under request churn// (mission 027 flakiness). smp_allocator is thread-safe and has no debug tripwires.const allocator = std.heap.smp_allocator;// Use direct syscall for stderr writes — std.debug.print uses std.Progress// which has ABI-incompatible global state when loaded as a plugin into a// binary compiled with a different Zig version.fn logMsg(msg: []const u8) void {_ = linux.write(2, msg.ptr, msg.len);}fn logFmt(comptime fmt: []const u8, args: anytype) void {var buf: [512]u8 = undefined;const s = std.fmt.bufPrint(&buf, fmt, args) catch return;logMsg(s);}// =========================================================================// SSE (Server-Sent Events) push channel — mission 031, ADDRESSED in 080// A registry of connected event-stream clients. `sse_start` (on the respond// handle) upgrades a GET request into a long-lived text/event-stream, assigns// the connection a stable id and registers it; `hl_http1_sse_broadcast` writes// a `data:` frame to every registered client and `hl_http1_sse_send` to ONE of// them, dropping any that error (peer closed). Touched only from the// interpreter (main) thread, but guarded by a mutex for safety.//// Mission 080 (D26): the id is a MONOTONIC counter, never the fd — a closed// connection's fd is recycled by the kernel within milliseconds, so an fd-keyed// push would eventually land on a stranger's socket. `sse_start` returns the id// (it used to return the raw fd, which no caller used), and the framework sends// it to the browser as the `__hlHello` event so the client can name itself.// =========================================================================// Mission 084: a subscription is reaped PROACTIVELY, not on the first failed push.// A closed tab used to stay in this registry until something happened to be pushed to// it — and because a scoped push (§9.5) sends nothing to a client whose needs did not// change, "something" could be never. Two mechanisms, both in the reaper thread://// • EPOLLRDHUP on every registered fd — a closed tab sends FIN, which fires// immediately and costs no traffic at all. This is the primary detector.// • a periodic SSE COMMENT heartbeat (`:\n\n`, which EventSource ignores) — catches// a peer that vanished WITHOUT a FIN (killed machine, dropped NAT entry), which no// amount of epolling can see, and keeps intermediaries from timing the stream out.//// A heartbeat is also why one failed write is enough to declare death here: the probe// runs repeatedly, so the "first write after FIN succeeds, second gets EPIPE" TCP// behaviour just means the peer is reaped one tick later.const MSG_NOSIGNAL: u32 = 0x4000;const SseConn = struct { id: u32, fd: i32 };var sse_mutex: PthreadMutex = c.PTHREAD_MUTEX_INITIALIZER;var sse_subs: std.ArrayListUnmanaged(SseConn) = .empty;var sse_next_id: u32 = 1;/// how often the reaper writes its comment heartbeat / probes for silent deathconst SSE_HEARTBEAT_MS: i64 = 5_000;var sse_epoll_fd: i32 = -1;var sse_wake_fd: i32 = -1;var sse_reaper_thread: ?std.Thread = null;var sse_reaper_running = std.atomic.Value(bool).init(false);/// Start the reaper thread + its epoll. Idempotent; called under sse_mutex from the/// first sseRegister(), so a server that never opens a stream never spawns it.fn sseReaperEnsureLocked() void {if (sse_reaper_running.load(.acquire)) return;const ep_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);const ep: i32 = @bitCast(@as(u32, @truncate(ep_rc)));if (ep < 0) return;const ef_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);const ef: i32 = @bitCast(@as(u32, @truncate(ef_rc)));if (ef < 0) {_ = linux.close(ep);return;}var wake_ev = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .fd = ef } };_ = linux.epoll_ctl(ep, linux.EPOLL.CTL_ADD, ef, &wake_ev);sse_epoll_fd = ep;sse_wake_fd = ef;sse_reaper_running.store(true, .release);sse_reaper_thread = std.Thread.spawn(.{}, sseReaperLoop, .{}) catch {sse_reaper_running.store(false, .release);_ = linux.close(ep);_ = linux.close(ef);sse_epoll_fd = -1;sse_wake_fd = -1;return;};}/// Drop subscription at index `i`: epoll DEL, close, remove. Caller holds sse_mutex.fn sseDropLocked(i: usize) void {const fd = sse_subs.items[i].fd;if (sse_epoll_fd >= 0) _ = linux.epoll_ctl(sse_epoll_fd, linux.EPOLL.CTL_DEL, fd, null);_ = linux.close(fd);_ = sse_subs.swapRemove(i);}/// The reaper: EPOLLRDHUP wakes it the instant a tab closes; the 1s timeout paces the/// heartbeat. Nothing here pushes application data, so a reap needs no traffic from the/// app at all — which is the whole point (mission 080's gap).fn sseReaperLoop() void {var events: [64]linux.epoll_event = undefined;var last_beat: i64 = monotonicMs();while (sse_reaper_running.load(.acquire)) {const n_rc = linux.epoll_wait(sse_epoll_fd, &events, events.len, 1000);const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));if (n > 0) {for (events[0..@intCast(n)]) |ev| {if (ev.data.fd == sse_wake_fd) {var drain: u64 = 0;_ = linux.read(sse_wake_fd, @ptrCast(&drain), @sizeOf(u64));continue;}// A subscriber never SENDS on its stream, so any readable/hangup event// means the peer went away (or is misbehaving) — either way it is dead.mutexLock(&sse_mutex);for (sse_subs.items, 0..) |conn, i| {if (conn.fd == ev.data.fd) {sseDropLocked(i);break;}}mutexUnlock(&sse_mutex);}}const now = monotonicMs();if (now - last_beat >= SSE_HEARTBEAT_MS) {last_beat = now;mutexLock(&sse_mutex);var i: usize = 0;while (i < sse_subs.items.len) {// ":\n\n" is an SSE comment — EventSource ignores it, so this is a pure// liveness probe that never reaches an `onmessage` handler.if (!sseRawWrite(sse_subs.items[i].fd, ":\n\n")) {sseDropLocked(i);continue;}i += 1;}mutexUnlock(&sse_mutex);}}}// Write via sendto with MSG_NOSIGNAL so a dead peer yields EPIPE instead of// killing the process with SIGPIPE. Returns false on any short/failed write.fn sseRawWrite(fd: i32, data: []const u8) bool {var written: usize = 0;while (written < data.len) {const rc = linux.sendto(fd, data[written..].ptr, data.len - written, MSG_NOSIGNAL, null, 0);const n: isize = @bitCast(rc);if (n <= 0) return false;written += @intCast(n);}return true;}fn sseRegister(fd: i32) u32 {mutexLock(&sse_mutex);defer mutexUnlock(&sse_mutex);const id = sse_next_id;sse_next_id += 1;sse_subs.append(allocator, .{ .id = id, .fd = fd }) catch return 0;// Watch it for hangup from now on — see the reaper comment above.sseReaperEnsureLocked();if (sse_epoll_fd >= 0) {var ev = linux.epoll_event{.events = linux.EPOLL.RDHUP | linux.EPOLL.IN,.data = .{ .fd = fd },};_ = linux.epoll_ctl(sse_epoll_fd, linux.EPOLL.CTL_ADD, fd, &ev);}return id;}// One `data:` frame on one socket. Caller holds sse_mutex.fn sseWriteFrameLocked(fd: i32, msg: []const u8) bool {var ok = sseRawWrite(fd, "data: ");if (ok) ok = sseRawWrite(fd, msg);if (ok) ok = sseRawWrite(fd, "\n\n");return ok;}// Broadcast one SSE message to all subscribers. `msg` should be a single line// (callers send single-line JSON). Dead sockets are closed and removed.// Returns the number of clients successfully written to.fn sseBroadcast(msg: []const u8) usize {mutexLock(&sse_mutex);defer mutexUnlock(&sse_mutex);var count: usize = 0;var i: usize = 0;while (i < sse_subs.items.len) {if (!sseWriteFrameLocked(sse_subs.items[i].fd, msg)) {sseDropLocked(i);continue;}count += 1;i += 1;}return count;}// __native("http1.sse_count") → how many SSE subscriptions are currently LIVE.// The observable the proactive reaper exists to keep honest (mission 084): it must fall// when a client goes away, with nothing being pushed.export fn hl_http1_sse_count(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {_ = argc;_ = argv;mutexLock(&sse_mutex);defer mutexUnlock(&sse_mutex);return api.makeNumber(@floatFromInt(sse_subs.items.len));}// __native("http1.sse_alive", connId) → 1 if that subscription is still registered.// Lets the framework prune its own per-connection bookkeeping without pushing anything.export fn hl_http1_sse_alive(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 1 or argv[0].type != .hl_number) return api.makeNumber(0);const id: u32 = @intFromFloat(argv[0].data.number);mutexLock(&sse_mutex);defer mutexUnlock(&sse_mutex);for (sse_subs.items) |conn| {if (conn.id == id) return api.makeNumber(1);}return api.makeNumber(0);}// __native("http1.sse_broadcast", jsonString) → number of clients pushed to.export fn hl_http1_sse_broadcast(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 1 or argv[0].type != .hl_string) return api.makeNumber(0);const msg = argv[0].data.string.ptr[0..argv[0].data.string.len];const n = sseBroadcast(msg);return api.makeNumber(@floatFromInt(n));}// __native("http1.sse_send", connId, jsonString) → 1 pushed / 0 gone.// Mission 080 (D26): the ADDRESSED half of the channel — the dependency-scoped// push sends a different payload to each client, so it cannot use broadcast.export fn hl_http1_sse_send(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 2 or argv[0].type != .hl_number or argv[1].type != .hl_string) return api.makeNumber(0);const id: u32 = @intFromFloat(argv[0].data.number);const msg = argv[1].data.string.ptr[0..argv[1].data.string.len];mutexLock(&sse_mutex);defer mutexUnlock(&sse_mutex);for (sse_subs.items, 0..) |conn, i| {if (conn.id != id) continue;if (!sseWriteFrameLocked(conn.fd, msg)) {sseDropLocked(i);return api.makeNumber(0);}return api.makeNumber(1);}return api.makeNumber(0);}// =========================================================================// WebSocket engine (mission 064, decisions D8/D9)//// Framing lives in the shared core (plugins/http/ws_common.zig); this engine// owns the h1-specific parts: the Upgrade/101 handshake (done inline in the// I/O worker, see readAndParseRequest) and the socket lifecycle after it.//// One GLOBAL engine per plugin (mirrors the SSE registry): a single reader// thread epolls all upgraded sockets with level-triggered EPOLLIN only —// EPOLLOUT is never armed, so an idle connection never wakes the loop (the// busy-spin trap found in the http2 TLS path). Complete messages become// WsEvent entries that the interpreter drains via the hl_http1_ws_events// iterator (registered on the event loop like the request iterator).//// Outbound writes (send/broadcast/ping/close + pong replies) happen under// ws_mutex from either the interpreter thread or the reader thread. Sockets// are non-blocking; a slow consumer gets bounded EAGAIN retries (~100ms)// and is dropped rather than buffered (no outbound queue, no EPOLLOUT).// WS upgrades are plain-HTTP only for now (TLS WS = future work, like SSE).// =========================================================================const WsEventKind = enum(u8) { connect, message, close, pong };const WsEvent = struct {kind: WsEventKind,id: u32,data: ?[]u8 = null, // allocated payload (message data / close reason)is_binary: bool = false,code: u16 = 0, // close code// Mission 093: the UPGRADE REQUEST's `Cookie` header, verbatim, carried on the// `connect` event only. A WebSocket handshake is an ordinary HTTP request, so// the browser sends the session cookie with it automatically — this is the one// moment the socket can be attributed to whoever loaded the page, and after the// 101 the request (and its headers) is freed. Allocated; freed with the event.cookie: ?[]u8 = null,// Ticket #74: the upgrade request's `Host` header, verbatim, on `connect` only —// the address the browser dialled, which a page served per subdomain is rendered// for. Allocated; freed with the event.host: ?[]u8 = null,};const WsClient = struct {id: u32,fd: i32,decoder: ws.Decoder,closing: bool = false, // server sent close, awaiting peer echo// liveness (mission 091): when this socket last produced a frame, and when// the sweep's ping went out (0 = none outstanding). A TCP connection whose// peer vanished without a FIN stays writable indefinitely, so "still open"// is not evidence of a live peer — the pong is.last_seen_ms: i64 = 0,ping_at_ms: i64 = 0,};var ws_mutex: PthreadMutex = c.PTHREAD_MUTEX_INITIALIZER;var ws_clients: std.ArrayListUnmanaged(*WsClient) = .empty;var ws_event_queue: std.ArrayListUnmanaged(WsEvent) = .empty;var ws_epoll_fd: i32 = -1;var ws_wake_fd: i32 = -1;/// THE EVENT LOOP'S BELL for the WS event queue (mission 256), and a different fd/// from `ws_wake_fd` above — that one wakes the reader THREAD's own epoll, this/// one wakes the interpreter loop. Without it an app that turns sockets on has a/// source with no fd, and ONE such source puts the whole loop back on the 1ms/// poll (loop_wait.Waiter.observe) — so the server would busy-poll for as long as/// WebSockets were enabled.var ws_loop_wake_fd: i32 = -1;var ws_thread: ?std.Thread = null;var ws_running = std.atomic.Value(bool).init(false);var ws_enabled = std.atomic.Value(bool).init(false);var ws_next_id: u32 = 1;// --- liveness sweep (mission 091) ----------------------------------------// The reader thread already wakes every 500ms (the epoll timeout), so the sweep// costs no timer and no new thread: on every tick it pings sockets that have// gone quiet and drops the ones whose pong is overdue. The drop enqueues the// ordinary `close` event, so every consumer above (hl:web's subscription// registry included) prunes through the path it already had — nothing upstream// learns a new concept, and nothing at the hl level needs a timer.//// Operator knobs, read once at engine start; the defaults are production// values, the tests shrink them.const WS_PING_MS_DEFAULT: i64 = 15000; // quiet this long -> askconst WS_PONG_MS_DEFAULT: i64 = 10000; // no answer in this long -> gonevar ws_ping_ms: i64 = WS_PING_MS_DEFAULT;var ws_pong_ms: i64 = WS_PONG_MS_DEFAULT;/// MONOTONIC milliseconds — a timeout measured against the wall clock would fire/// early or never after an NTP step.fn wsNowMs() i64 {var ts: linux.timespec = undefined;_ = linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts);return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);}fn wsEnvMs(name: [:0]const u8, fallback: i64) i64 {const raw = std.mem.span(std.c.getenv(name) orelse return fallback);const n = std.fmt.parseInt(i64, std.mem.trim(u8, raw, " \t"), 10) catch return fallback;if (n <= 0) return fallback;return n;}const EAGAIN_ERR: isize = 11;const EINTR_ERR: isize = 4;/// WriteFn callback over a raw fd (context = fd stuffed into the pointer)./// Non-blocking socket: bounded EAGAIN retries, then give up (caller drops).fn wsFdWrite(ctx: ?*anyopaque, data: []const u8) bool {const fd: i32 = @intCast(@intFromPtr(ctx));var written: usize = 0;var retries: u32 = 0;while (written < data.len) {const rc = linux.sendto(fd, data[written..].ptr, data.len - written, MSG_NOSIGNAL, null, 0);const n: isize = @bitCast(rc);if (n > 0) {written += @intCast(n);retries = 0;continue;}const e = -n;if (e == EAGAIN_ERR) {retries += 1;if (retries > 100) return false; // ~100ms of backpressure → dropconst req = linux.timespec{ .sec = 0, .nsec = 1_000_000 }; // 1ms_ = linux.nanosleep(&req, null);continue;}if (e == EINTR_ERR) continue;return false;}return true;}fn wsFdCtx(fd: i32) ?*anyopaque {return @ptrFromInt(@as(usize, @intCast(fd)));}/// Start the reader thread + epoll instance. Idempotent. Called from the/// interpreter thread (hl_http1_ws_events); also flips ws_enabled so the I/O/// workers start honoring Upgrade requests.fn wsEnsureStarted() bool {mutexLock(&ws_mutex);defer mutexUnlock(&ws_mutex);if (ws_running.load(.acquire)) return true;ws_ping_ms = wsEnvMs("HL_WS_PING_MS", WS_PING_MS_DEFAULT);ws_pong_ms = wsEnvMs("HL_WS_PONG_TIMEOUT_MS", WS_PONG_MS_DEFAULT);const ep_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);const ep: i32 = @bitCast(@as(u32, @truncate(ep_rc)));if (ep < 0) return false;const ef_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);const ef: i32 = @bitCast(@as(u32, @truncate(ef_rc)));if (ef < 0) {_ = linux.close(ep);return false;}var wake_ev = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .fd = ef } };_ = linux.epoll_ctl(ep, linux.EPOLL.CTL_ADD, ef, &wake_ev);ws_epoll_fd = ep;ws_wake_fd = ef;if (ws_loop_wake_fd < 0) ws_loop_wake_fd = http.makeWakeFd(); // mission 256ws_running.store(true, .release);ws_thread = std.Thread.spawn(.{}, wsReaderLoop, .{}) catch {ws_running.store(false, .release);_ = linux.close(ep);_ = linux.close(ef);ws_epoll_fd = -1;ws_wake_fd = -1;return false;};ws_enabled.store(true, .release);return true;}fn wsWake() void {if (ws_wake_fd >= 0) {const one: u64 = 1;_ = linux.write(ws_wake_fd, @ptrCast(&one), @sizeOf(u64));}}/// Hand an upgraded socket to the engine. Called from an I/O worker thread/// right after the 101 was written. `leftover` = bytes the client sent after/// the handshake that were already consumed into the header buffer. `cookie` is/// the handshake's `Cookie` header (mission 093) — copied here, because the/// parsed request is freed the moment this returns. `host` is its `Host` header/// (ticket #74), copied for the same reason.fn wsRegisterClient(fd: i32, leftover: []const u8, cookie: []const u8, host: []const u8) void {// Non-blocking for the reader loopconst flags_rc = linux.fcntl(fd, linux.F.GETFL, @as(usize, 0));const flags_i: isize = @bitCast(flags_rc);if (flags_i >= 0) {var oflags: linux.O = @bitCast(@as(u32, @truncate(flags_rc)));oflags.NONBLOCK = true;_ = linux.fcntl(fd, linux.F.SETFL, @as(usize, @as(u32, @bitCast(oflags))));}mutexLock(&ws_mutex);const client = allocator.create(WsClient) catch {mutexUnlock(&ws_mutex);_ = linux.close(fd);return;};client.* = .{.id = ws_next_id,.fd = fd,.decoder = ws.Decoder.init(allocator),.last_seen_ms = wsNowMs(),};ws_next_id += 1;ws_clients.append(allocator, client) catch {allocator.destroy(client);mutexUnlock(&ws_mutex);_ = linux.close(fd);return;};const cookie_copy: ?[]u8 = if (cookie.len > 0) (allocator.dupe(u8, cookie) catch null) else null;const host_copy: ?[]u8 = if (host.len > 0) (allocator.dupe(u8, host) catch null) else null;ws_event_queue.append(allocator, .{ .kind = .connect, .id = client.id, .cookie = cookie_copy, .host = host_copy }) catch {if (cookie_copy) |cc| allocator.free(cc);if (host_copy) |hc| allocator.free(hc);};if (leftover.len > 0) client.decoder.feed(leftover) catch {};mutexUnlock(&ws_mutex);// Level-triggered EPOLLIN only (never EPOLLOUT): pending socket data fires// immediately, idle connections cost nothing.var ev = linux.epoll_event{ .events = linux.EPOLL.IN | linux.EPOLL.RDHUP, .data = .{ .fd = fd } };_ = linux.epoll_ctl(ws_epoll_fd, linux.EPOLL.CTL_ADD, fd, &ev);if (leftover.len > 0) wsWake(); // decoder-buffered bytes won't fire EPOLLIN// This runs on an I/O WORKER thread, not the reader thread, and it queued a// `connect` — so it rings the event loop itself (mission 256).http.ringWake(ws_loop_wake_fd);}fn wsFindByFdLocked(fd: i32) ?usize {for (ws_clients.items, 0..) |cl, i| {if (cl.fd == fd) return i;}return null;}fn wsFindByIdLocked(id: u32) ?usize {for (ws_clients.items, 0..) |cl, i| {if (cl.id == id) return i;}return null;}/// Remove client at index: epoll DEL, close fd, free state. Mutex held.fn wsRemoveLocked(idx: usize) void {const client = ws_clients.items[idx];_ = linux.epoll_ctl(ws_epoll_fd, linux.EPOLL.CTL_DEL, client.fd, null);_ = linux.close(client.fd);client.decoder.deinit();_ = ws_clients.swapRemove(idx);allocator.destroy(client);}/// Drain the client's decoder; enqueue events, auto-reply pings, run the/// close handshake. Returns true if the client was removed. Mutex held.fn wsDrainDecoderLocked(idx: usize) bool {const client = ws_clients.items[idx];while (true) {const maybe_ev = client.decoder.next() catch {// OOM mid-decode — drop the connectionws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1011 }) catch {};wsRemoveLocked(idx);return true;};const ev = maybe_ev orelse return false;// Any complete frame is proof of life, and it answers an outstanding// sweep ping whatever its opcode — a peer that is talking is not dead.client.last_seen_ms = wsNowMs();client.ping_at_ms = 0;switch (ev) {.text => |p| ws_event_queue.append(allocator, .{ .kind = .message, .id = client.id, .data = p }) catch allocator.free(p),.binary => |p| ws_event_queue.append(allocator, .{ .kind = .message, .id = client.id, .data = p, .is_binary = true }) catch allocator.free(p),.ping => |p| {_ = ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .pong, p);allocator.free(p);},.pong => |p| {allocator.free(p);ws_event_queue.append(allocator, .{ .kind = .pong, .id = client.id }) catch {};},.close => |cl| {if (!client.closing) {const echo_code = if (cl.code == ws.CLOSE_NO_STATUS) ws.CLOSE_NORMAL else cl.code;_ = ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), echo_code, "");}ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = cl.code, .data = cl.reason }) catch allocator.free(cl.reason);wsRemoveLocked(idx);return true;},.protocol_error => |code| {_ = ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), code, "");ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = code }) catch {};wsRemoveLocked(idx);return true;},}}}/// Read all available bytes from the client socket into its decoder, then/// drain. Returns true if the client was removed. Mutex held.fn wsServiceClientLocked(idx: usize) bool {const client = ws_clients.items[idx];var buf: [16384]u8 = undefined;while (true) {const rc = linux.read(client.fd, &buf, buf.len);const n: isize = @bitCast(rc);if (n > 0) {client.decoder.feed(buf[0..@intCast(n)]) catch {};continue;}if (n == 0) {// Peer closed without a close frame → 1006 abnormal closurews_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};wsRemoveLocked(idx);return true;}const e = -n;if (e == EINTR_ERR) continue;if (e == EAGAIN_ERR) break; // all available data consumedws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};wsRemoveLocked(idx);return true;}return wsDrainDecoderLocked(idx);}/// One liveness pass over every upgraded socket. Quiet for longer than/// ws_ping_ms → send a protocol ping (RFC 6455 §5.5.2; every browser answers it/// at the protocol layer, so nothing application-side is involved). A ping that/// stands unanswered for ws_pong_ms → the peer is gone however open the socket/// looks: enqueue the SAME `close` event a real FIN would have produced and drop/// the connection. That is the whole mechanism — no separate timer, no new/// thread, no concept added above this file.fn wsSweep() void {const now = wsNowMs();mutexLock(&ws_mutex);defer mutexUnlock(&ws_mutex);var i: usize = 0;while (i < ws_clients.items.len) {const client = ws_clients.items[i];if (client.ping_at_ms != 0) {if (now - client.ping_at_ms > ws_pong_ms) {ws_event_queue.append(allocator, .{.kind = .close,.id = client.id,.code = ws.CLOSE_GOING_AWAY,}) catch {};wsRemoveLocked(i);continue;}} else if (now - client.last_seen_ms > ws_ping_ms) {if (ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .ping, "")) {client.ping_at_ms = now;} else {// the write itself failed — this one needs no grace periodws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};wsRemoveLocked(i);continue;}}i += 1;}}fn wsReaderLoop() void {var events: [64]linux.epoll_event = undefined;while (ws_running.load(.acquire)) {// RING THE EVENT LOOP'S BELL AT THE END OF EVERY ITERATION (mission 256),// unconditionally and without taking `ws_mutex`. Fifteen places in this// file push onto `ws_event_queue`; a bell is coalesced and a spurious one// is harmless (loop_wait.zig: "only a signal that is never sent at all// could ever be a bug"), so one ring per pass covers all of them and can// never deadlock against a path that still holds the lock. The idle cost// is two empty loop rounds a second — this epoll has a 500ms timeout.defer http.ringWake(ws_loop_wake_fd);const n_rc = linux.epoll_wait(ws_epoll_fd, &events, events.len, 500);const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));// The 500ms epoll timeout IS the sweep's clock: an idle loop still ticks,// and a busy one sweeps just as often because the check is on wall time.if (n <= 0) {wsSweep();continue;}for (events[0..@intCast(n)]) |ev| {if (ev.data.fd == ws_wake_fd) {var drain: u64 = 0;_ = linux.read(ws_wake_fd, @ptrCast(&drain), @sizeOf(u64));if (!ws_running.load(.acquire)) return;// Drain decoder-buffered data (e.g. handshake leftover)mutexLock(&ws_mutex);var i: usize = 0;while (i < ws_clients.items.len) {if (!wsDrainDecoderLocked(i)) i += 1;}mutexUnlock(&ws_mutex);continue;}mutexLock(&ws_mutex);if (wsFindByFdLocked(ev.data.fd)) |idx| {_ = wsServiceClientLocked(idx);}mutexUnlock(&ws_mutex);}wsSweep();}}/// Stop the engine: join the reader thread, close all sockets, free queues./// Runs at interpreter teardown via the events-iterator deinit — BEFORE the/// runtime dlcloses this .so, so the thread never outlives its code.fn wsShutdown() void {if (!ws_running.swap(false, .acq_rel)) return;wsWake();if (ws_thread) |t| {t.join();ws_thread = null;}mutexLock(&ws_mutex);for (ws_clients.items) |client| {_ = ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), ws.CLOSE_GOING_AWAY, "");_ = linux.close(client.fd);client.decoder.deinit();allocator.destroy(client);}ws_clients.deinit(allocator);ws_clients = .empty;for (ws_event_queue.items) |*ev| {if (ev.data) |d| allocator.free(d);if (ev.cookie) |ck| allocator.free(ck);if (ev.host) |hc| allocator.free(hc);}ws_event_queue.deinit(allocator);ws_event_queue = .empty;mutexUnlock(&ws_mutex);if (ws_epoll_fd >= 0) _ = linux.close(ws_epoll_fd);if (ws_wake_fd >= 0) _ = linux.close(ws_wake_fd);ws_epoll_fd = -1;ws_wake_fd = -1;// `ws_loop_wake_fd` is deliberately NOT closed: the event loop may still hold// it in its epoll set, and a closed fd number gets REUSED — the loop would// then be watching whatever opened next. One eventfd per process, kept for// the process's life and reused if the engine restarts (mission 256).ws_enabled.store(false, .release);}// --- Interpreter-facing exports ------------------------------------------const ws_kind_names = [_][]const u8{ "connect", "message", "close", "pong" };fn wsEventObjDeinit(obj: *HlObject) callconv(.c) void {// fields[2] is "data", fields[5] "cookie" and fields[6] "host" — allocated// iff non-empty (empty = the static "").const s = obj.fields[2].value.data.string;if (s.len > 0) allocator.free(s.ptr[0..s.len]);const ck = obj.fields[5].value.data.string;if (ck.len > 0) allocator.free(ck.ptr[0..ck.len]);const hs = obj.fields[6].value.data.string;if (hs.len > 0) allocator.free(hs.ptr[0..hs.len]);allocator.free(obj.fields[0..obj.field_count]);allocator.destroy(obj);}/// try_next over the WS event queue → { kind, id, data, binary, code, cookie, host }/// or null. `cookie` is the handshake's Cookie header and `host` its Host header,/// both non-empty only on `connect` (mission 093, ticket #74).fn wsEventsTryNext(ctx: ?*anyopaque) callconv(.c) HlValue {_ = ctx;mutexLock(&ws_mutex);if (ws_event_queue.items.len == 0) {mutexUnlock(&ws_mutex);return api.makeNull();}const ev = ws_event_queue.orderedRemove(0);mutexUnlock(&ws_mutex);const fields = allocator.alloc(HlField, 7) catch {if (ev.data) |d| allocator.free(d);if (ev.cookie) |ck| allocator.free(ck);if (ev.host) |hc| allocator.free(hc);return api.makeNull();};const data_slice: []const u8 = if (ev.data) |d| d else "";const cookie_slice: []const u8 = if (ev.cookie) |ck| ck else "";const host_slice: []const u8 = if (ev.host) |hc| hc else "";fields[0] = .{ .key = http.hlStr("kind"), .value = api.makeString(ws_kind_names[@intFromEnum(ev.kind)]) };fields[1] = .{ .key = http.hlStr("id"), .value = api.makeNumber(@floatFromInt(ev.id)) };fields[2] = .{ .key = http.hlStr("data"), .value = api.makeString(data_slice) };fields[3] = .{ .key = http.hlStr("binary"), .value = api.makeBool(ev.is_binary) };fields[4] = .{ .key = http.hlStr("code"), .value = api.makeNumber(@floatFromInt(ev.code)) };fields[5] = .{ .key = http.hlStr("cookie"), .value = api.makeString(cookie_slice) };fields[6] = .{ .key = http.hlStr("host"), .value = api.makeString(host_slice) };const obj = allocator.create(HlObject) catch {if (ev.data) |d| allocator.free(d);if (ev.cookie) |ck| allocator.free(ck);if (ev.host) |hc| allocator.free(hc);allocator.free(fields);return api.makeNull();};obj.* = .{ .fields = fields.ptr, .field_count = 7, .deinit_fn = &wsEventObjDeinit };return api.makeObject(obj);}fn wsEventsIterDeinit(ctx: ?*anyopaque) callconv(.c) void {_ = ctx;wsShutdown();}/// __native("http1.ws_events") → event iterator; starting it enables WS/// upgrades on all plain-HTTP hl:http1 servers in this process.export fn hl_http1_ws_events(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {_ = argc;_ = argv;if (!wsEnsureStarted()) return api.makeNull();const iter = allocator.create(HlIterator) catch return api.makeNull();iter.* = .{.context = null,.next_fn = &wsEventsTryNext, // non-blocking either way — event loop only.deinit_fn = &wsEventsIterDeinit,.try_next_fn = &wsEventsTryNext,.wake_fd = ws_loop_wake_fd, // mission 256 — set by wsEnsureStarted above};return api.makeIterator(iter);}// --- cookie-grade random (mission 093) ------------------------------------// A session id is the ONLY thing standing between a stranger and someone else's// session, so it may not come from a seeded PRNG: `hl:math`'s random() is a// clock-seeded xoshiro, and a few of its outputs reveal its state, which would// make every other session's id derivable from one's own. This reads the// kernel's CSPRNG directly. It lives in hl:http1 because the session cookie is// an HTTP artifact and this plugin is the one that parses and sets it; there is// no other CSPRNG at the hl: level yet (named as a gap in mission 093's report).var token_buf: [128]u8 = undefined;/// __native("http1.random_token", n) → n lowercase hex chars (default 32, max 128)export fn hl_http1_random_token(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {var want: usize = 32;if (argc >= 1 and argv[0].type == .hl_number) {const n = argv[0].data.number;if (n >= 1 and n <= 128) want = @intFromFloat(n);}var raw: [64]u8 = undefined;const need = (want + 1) / 2;if (linux.getrandom(&raw, need, 0) != need) return api.makeNull();const hex = "0123456789abcdef";var i: usize = 0;while (i < want) : (i += 1) {const byte = raw[i / 2];const nib: u8 = if (i % 2 == 0) (byte >> 4) else (byte & 0x0f);token_buf[i] = hex[nib];}return api.makeString(token_buf[0..want]);}/// __native("http1.ws_send", id, text[, binaryFlag]) → 1 sent / 0 goneexport fn hl_http1_ws_send(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 2 or argv[0].type != .hl_number or argv[1].type != .hl_string) return api.makeNumber(0);const id: u32 = @intFromFloat(argv[0].data.number);const text = argv[1].data.string.ptr[0..argv[1].data.string.len];const opcode: ws.Opcode = if (argc >= 3 and argv[2].type == .hl_bool and argv[2].data.boolean) .binary else .text;mutexLock(&ws_mutex);defer mutexUnlock(&ws_mutex);const idx = wsFindByIdLocked(id) orelse return api.makeNumber(0);const client = ws_clients.items[idx];if (!ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), opcode, text)) {ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};wsRemoveLocked(idx);return api.makeNumber(0);}return api.makeNumber(1);}/// __native("http1.ws_broadcast", text) → number of clients writtenexport fn hl_http1_ws_broadcast(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 1 or argv[0].type != .hl_string) return api.makeNumber(0);const text = argv[0].data.string.ptr[0..argv[0].data.string.len];mutexLock(&ws_mutex);defer mutexUnlock(&ws_mutex);var count: usize = 0;var i: usize = 0;while (i < ws_clients.items.len) {const client = ws_clients.items[i];if (ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .text, text)) {count += 1;i += 1;} else {ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};wsRemoveLocked(i);}}return api.makeNumber(@floatFromInt(count));}/// __native("http1.ws_ping", id) → 1 sent / 0 goneexport fn hl_http1_ws_ping(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 1 or argv[0].type != .hl_number) return api.makeNumber(0);const id: u32 = @intFromFloat(argv[0].data.number);mutexLock(&ws_mutex);defer mutexUnlock(&ws_mutex);const idx = wsFindByIdLocked(id) orelse return api.makeNumber(0);const client = ws_clients.items[idx];if (!ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .ping, "")) {wsRemoveLocked(idx);return api.makeNumber(0);}return api.makeNumber(1);}/// __native("http1.ws_close", id[, code]) → 1 initiated / 0 gone./// Sends the close frame and waits for the peer echo (reader completes it).export fn hl_http1_ws_close(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 1 or argv[0].type != .hl_number) return api.makeNumber(0);const id: u32 = @intFromFloat(argv[0].data.number);const code: u16 = if (argc >= 2 and argv[1].type == .hl_number) @intFromFloat(argv[1].data.number) else ws.CLOSE_NORMAL;mutexLock(&ws_mutex);defer mutexUnlock(&ws_mutex);const idx = wsFindByIdLocked(id) orelse return api.makeNumber(0);const client = ws_clients.items[idx];if (!ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), code, "")) {ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};wsRemoveLocked(idx);return api.makeNumber(0);}client.closing = true;return api.makeNumber(1);}// =========================================================================// TLS context — OpenSSL via dlopen (optional, no compile-time dep)// =========================================================================const c_dlfcn = @cImport({@cInclude("dlfcn.h");});const SSL_CTX = opaque {};const SSL = opaque {};const SSL_METHOD = opaque {};// OpenSSL function pointer typesconst SSL_library_init_fn = *const fn () callconv(.c) c_int;const SSL_load_error_strings_fn = *const fn () callconv(.c) void;const TLS_server_method_fn = *const fn () callconv(.c) ?*const SSL_METHOD;const SSL_CTX_new_fn = *const fn (?*const SSL_METHOD) callconv(.c) ?*SSL_CTX;const SSL_CTX_free_fn = *const fn (?*SSL_CTX) callconv(.c) void;const SSL_CTX_use_certificate_chain_file_fn = *const fn (?*SSL_CTX, [*:0]const u8) callconv(.c) c_int;const SSL_CTX_use_PrivateKey_file_fn = *const fn (?*SSL_CTX, [*:0]const u8, c_int) callconv(.c) c_int;const SSL_new_fn = *const fn (?*SSL_CTX) callconv(.c) ?*SSL;const SSL_free_fn = *const fn (?*SSL) callconv(.c) void;const SSL_set_fd_fn = *const fn (?*SSL, c_int) callconv(.c) c_int;const SSL_accept_fn = *const fn (?*SSL) callconv(.c) c_int;const SSL_read_fn = *const fn (?*SSL, [*]u8, c_int) callconv(.c) c_int;const SSL_write_fn = *const fn (?*SSL, [*]const u8, c_int) callconv(.c) c_int;const SSL_shutdown_fn = *const fn (?*SSL) callconv(.c) c_int;const SSL_get_error_fn = *const fn (?*SSL, c_int) callconv(.c) c_int;const OPENSSL_init_ssl_fn = *const fn (u64, ?*anyopaque) callconv(.c) c_int;const SSL_FILETYPE_PEM: c_int = 1;const TlsContext = struct {ssl_ctx: ?*SSL_CTX = null,lib_ssl: ?*anyopaque = null,lib_crypto: ?*anyopaque = null,// Function pointersfn_ssl_ctx_new: ?SSL_CTX_new_fn = null,fn_ssl_ctx_free: ?SSL_CTX_free_fn = null,fn_ssl_ctx_use_cert: ?SSL_CTX_use_certificate_chain_file_fn = null,fn_ssl_ctx_use_key: ?SSL_CTX_use_PrivateKey_file_fn = null,fn_ssl_new: ?SSL_new_fn = null,fn_ssl_free: ?SSL_free_fn = null,fn_ssl_set_fd: ?SSL_set_fd_fn = null,fn_ssl_accept: ?SSL_accept_fn = null,fn_ssl_read: ?SSL_read_fn = null,fn_ssl_write: ?SSL_write_fn = null,fn_ssl_shutdown: ?SSL_shutdown_fn = null,fn_ssl_get_error: ?SSL_get_error_fn = null,fn loadSym(lib: ?*anyopaque, comptime T: type, name: [*:0]const u8) ?T {const sym = c_dlfcn.dlsym(lib, name) orelse return null;return @ptrCast(sym);}fn init(cert_path: []const u8, key_path: []const u8) ?TlsContext {var ctx = TlsContext{};// Try loading libssl and libcryptoconst ssl_paths = [_][*:0]const u8{ "libssl.so.3", "libssl.so.1.1", "libssl.so" };const crypto_paths = [_][*:0]const u8{ "libcrypto.so.3", "libcrypto.so.1.1", "libcrypto.so" };for (ssl_paths) |path| {ctx.lib_ssl = c_dlfcn.dlopen(path, c_dlfcn.RTLD_NOW | c_dlfcn.RTLD_LOCAL);if (ctx.lib_ssl != null) break;}if (ctx.lib_ssl == null) {logMsg("http1: TLS: failed to load libssl.so\n");return null;}for (crypto_paths) |path| {ctx.lib_crypto = c_dlfcn.dlopen(path, c_dlfcn.RTLD_NOW | c_dlfcn.RTLD_LOCAL);if (ctx.lib_crypto != null) break;}if (ctx.lib_crypto == null) {logMsg("http1: TLS: failed to load libcrypto.so\n");_ = c_dlfcn.dlclose(ctx.lib_ssl);return null;}// Load function pointers// Try OPENSSL_init_ssl first (OpenSSL 1.1+), fall back to SSL_library_initif (loadSym(ctx.lib_ssl, OPENSSL_init_ssl_fn, "OPENSSL_init_ssl")) |init_fn| {_ = init_fn(0, null);} else if (loadSym(ctx.lib_ssl, SSL_library_init_fn, "SSL_library_init")) |lib_init| {_ = lib_init();if (loadSym(ctx.lib_ssl, SSL_load_error_strings_fn, "SSL_load_error_strings")) |load_err| {load_err();}}const method_fn = loadSym(ctx.lib_ssl, TLS_server_method_fn, "TLS_server_method") orelse {logMsg("http1: TLS: TLS_server_method not found\n");ctx.deinit();return null;};ctx.fn_ssl_ctx_new = loadSym(ctx.lib_ssl, SSL_CTX_new_fn, "SSL_CTX_new");ctx.fn_ssl_ctx_free = loadSym(ctx.lib_ssl, SSL_CTX_free_fn, "SSL_CTX_free");ctx.fn_ssl_ctx_use_cert = loadSym(ctx.lib_ssl, SSL_CTX_use_certificate_chain_file_fn, "SSL_CTX_use_certificate_chain_file");ctx.fn_ssl_ctx_use_key = loadSym(ctx.lib_ssl, SSL_CTX_use_PrivateKey_file_fn, "SSL_CTX_use_PrivateKey_file");ctx.fn_ssl_new = loadSym(ctx.lib_ssl, SSL_new_fn, "SSL_new");ctx.fn_ssl_free = loadSym(ctx.lib_ssl, SSL_free_fn, "SSL_free");ctx.fn_ssl_set_fd = loadSym(ctx.lib_ssl, SSL_set_fd_fn, "SSL_set_fd");ctx.fn_ssl_accept = loadSym(ctx.lib_ssl, SSL_accept_fn, "SSL_accept");ctx.fn_ssl_read = loadSym(ctx.lib_ssl, SSL_read_fn, "SSL_read");ctx.fn_ssl_write = loadSym(ctx.lib_ssl, SSL_write_fn, "SSL_write");ctx.fn_ssl_shutdown = loadSym(ctx.lib_ssl, SSL_shutdown_fn, "SSL_shutdown");ctx.fn_ssl_get_error = loadSym(ctx.lib_ssl, SSL_get_error_fn, "SSL_get_error");if (ctx.fn_ssl_ctx_new == null or ctx.fn_ssl_new == null orctx.fn_ssl_set_fd == null or ctx.fn_ssl_accept == null orctx.fn_ssl_read == null or ctx.fn_ssl_write == null){logMsg("http1: TLS: missing required SSL symbols\n");ctx.deinit();return null;}// Create SSL_CTXconst method = method_fn();ctx.ssl_ctx = ctx.fn_ssl_ctx_new.?(method);if (ctx.ssl_ctx == null) {logMsg("http1: TLS: SSL_CTX_new failed\n");ctx.deinit();return null;}// Load cert and keyconst cert_z = allocator.dupeZ(u8, cert_path) catch {ctx.deinit();return null;};defer allocator.free(cert_z);const key_z = allocator.dupeZ(u8, key_path) catch {ctx.deinit();return null;};defer allocator.free(key_z);if (ctx.fn_ssl_ctx_use_cert) |use_cert| {if (use_cert(ctx.ssl_ctx, cert_z.ptr) != 1) {logFmt("http1: TLS: failed to load certificate: {s}\n", .{cert_path});ctx.deinit();return null;}}if (ctx.fn_ssl_ctx_use_key) |use_key| {if (use_key(ctx.ssl_ctx, key_z.ptr, SSL_FILETYPE_PEM) != 1) {logFmt("http1: TLS: failed to load private key: {s}\n", .{key_path});ctx.deinit();return null;}}logMsg("http1: TLS initialized\n");return ctx;}fn wrapConnection(self: *const TlsContext, fd: i32) ?*SSL {const ssl = self.fn_ssl_new.?(self.ssl_ctx);if (ssl == null) return null;_ = self.fn_ssl_set_fd.?(ssl, fd);const ret = self.fn_ssl_accept.?(ssl);if (ret != 1) {self.fn_ssl_free.?(ssl);return null;}return ssl;}fn sslRead(self: *const TlsContext, ssl: *SSL, buf: []u8) isize {const ret = self.fn_ssl_read.?(ssl, buf.ptr, @intCast(@min(buf.len, std.math.maxInt(c_int))));if (ret <= 0) return 0;return @intCast(ret);}fn sslWrite(self: *const TlsContext, ssl: *SSL, data: []const u8) isize {var written: usize = 0;while (written < data.len) {const chunk_len: c_int = @intCast(@min(data.len - written, std.math.maxInt(c_int)));const ret = self.fn_ssl_write.?(ssl, data[written..].ptr, chunk_len);if (ret <= 0) return @intCast(written);written += @intCast(ret);}return @intCast(written);}fn sslShutdown(self: *const TlsContext, ssl: *SSL) void {_ = self.fn_ssl_shutdown.?(ssl);self.fn_ssl_free.?(ssl);}fn deinit(self: *TlsContext) void {if (self.ssl_ctx != null) {if (self.fn_ssl_ctx_free) |free_fn| {free_fn(self.ssl_ctx);}self.ssl_ctx = null;}if (self.lib_ssl != null) {_ = c_dlfcn.dlclose(self.lib_ssl);self.lib_ssl = null;}if (self.lib_crypto != null) {_ = c_dlfcn.dlclose(self.lib_crypto);self.lib_crypto = null;}}};// =========================================================================// MPSC Request Queue — thread-safe, blocks on dequeue// =========================================================================const ParsedRequest = struct {method: []const u8,path: []const u8, // URL-decodedquery_params: []http.QueryParam, // parsed key-value pairsquery_raw: []const u8, // raw query stringheaders: std.ArrayListUnmanaged(HeaderPair),body: []const u8,client_fd: i32,ssl: ?*SSL, // null if plain HTTPkeep_alive: bool,};const HeaderPair = struct {key: []const u8,value: []const u8,};const RequestQueue = struct {queue: std.ArrayListUnmanaged(ParsedRequest),mutex: PthreadMutex,condvar: PthreadCond,shutdown: bool,/// THE EVENT LOOP'S BELL (mission 256). The condvar above wakes a `for (req of/// server)` consumer blocked in `dequeue`; the interpreter's event loop is/// NOT that consumer — it polls `tryDequeue` from a loop that also watches/// every other source, so it cannot block on a condvar belonging to one of/// them. An eventfd it can: `native/src/loop_wait.zig` epolls this fd, and an/// idle server costs nothing instead of 1000 empty poll rounds a second.loop_wake_fd: i32,fn init() RequestQueue {var self: RequestQueue = .{.queue = .empty,.mutex = c.PTHREAD_MUTEX_INITIALIZER,.condvar = c.PTHREAD_COND_INITIALIZER,.shutdown = false,.loop_wake_fd = http.makeWakeFd(),};mutexInit(&self.mutex);condInit(&self.condvar);return self;}fn enqueue(self: *RequestQueue, req: ParsedRequest) void {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);self.queue.append(allocator, req) catch return;condSignal(&self.condvar);// Rung INSIDE the lock, so the fd's counter is already up by the time the// request is visible to `tryDequeue`. A loop that is between its last poll// and its next `epoll_wait` therefore finds the level raised and returns// at once instead of sleeping on a queue that has work in it.http.ringWake(self.loop_wake_fd);}fn dequeue(self: *RequestQueue) ?ParsedRequest {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);while (self.queue.items.len == 0 and !self.shutdown) {condWait(&self.condvar, &self.mutex);}if (self.shutdown and self.queue.items.len == 0) return null;return self.queue.orderedRemove(0);}/// Non-blocking: returns null immediately if queue is emptyfn tryDequeue(self: *RequestQueue) ?ParsedRequest {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);if (self.queue.items.len == 0) return null;return self.queue.orderedRemove(0);}fn signalShutdown(self: *RequestQueue) void {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);self.shutdown = true;condBroadcast(&self.condvar);// The event loop is told too — it is blocked on this fd and its `while`// condition (has every source retired?) can only be re-read on the way out.http.ringWake(self.loop_wake_fd);}fn deinit(self: *RequestQueue) void {// Free any remaining queued requestsfor (self.queue.items) |*req| {freeRequest(req);}self.queue.deinit(allocator);if (self.loop_wake_fd >= 0) {_ = linux.close(self.loop_wake_fd);self.loop_wake_fd = -1;}_ = c.pthread_cond_destroy(&self.condvar);_ = c.pthread_mutex_destroy(&self.mutex);}};fn freeRequest(req: *ParsedRequest) void {allocator.free(req.method);allocator.free(req.path);allocator.free(req.query_raw);for (req.query_params) |param| {allocator.free(param.key);allocator.free(param.value);}allocator.free(req.query_params);for (req.headers.items) |hdr| {allocator.free(hdr.key);allocator.free(hdr.value);}req.headers.deinit(allocator);if (req.body.len > 0) allocator.free(req.body);}// =========================================================================// Connection — per-connection state tracked by I/O workers// =========================================================================const Connection = struct {fd: i32,ssl: ?*SSL,keep_alive: bool,request_count: u32,max_requests: u32,fn read(self: *const Connection, tls: ?*const TlsContext, buf: []u8) usize {if (self.ssl) |ssl_ptr| {if (tls) |t| {const n = t.sslRead(ssl_ptr, buf);if (n <= 0) return 0;return @intCast(n);}return 0;}return sysRead(self.fd, buf);}fn write(self: *const Connection, tls: ?*const TlsContext, data: []const u8) void {if (self.ssl) |ssl_ptr| {if (tls) |t| {_ = t.sslWrite(ssl_ptr, data);return;}}sysWriteAll(self.fd, data);}fn close(self: *Connection, tls: ?*const TlsContext) void {if (self.ssl) |ssl_ptr| {if (tls) |t| {t.sslShutdown(ssl_ptr);}self.ssl = null;}_ = linux.close(self.fd);}};// =========================================================================// ServerCore — owns server socket, epoll, threads, queue, optional TLS// =========================================================================const MAX_KEEPALIVE_REQUESTS: u32 = 100;const MAX_IO_THREADS: u32 = 16;const ServerCore = struct {port: u16 = 0,server_fd: i32,epoll_fd: i32,shutdown_fd: i32, // eventfd for shutdown signaltls: ?TlsContext,request_queue: RequestQueue,running: std.atomic.Value(bool),acceptor_thread: ?std.Thread,io_threads: []std.Thread,num_threads: u32,// Connection tracking for I/O threadsconn_queue: ConnectionQueue,/// idle connections wait here instead of inside a blocking worker read (mission 084)parked: ParkedConns,fn create(port: u16, host_str: []const u8, cert_path: []const u8, key_path: []const u8, num_threads: u32) ?*ServerCore {const core = allocator.create(ServerCore) catch return null;core.* = .{.server_fd = -1,.epoll_fd = -1,.shutdown_fd = -1,.tls = null,.request_queue = RequestQueue.init(),.running = std.atomic.Value(bool).init(false),.acceptor_thread = null,.io_threads = &[_]std.Thread{},.num_threads = @min(num_threads, MAX_IO_THREADS),.conn_queue = ConnectionQueue.init(),.parked = ParkedConns.init(),};// Create TCP socket. NONBLOCK IS NOT DECORATION (mission 260): the acceptor// drains with `while (true) accept4(...)` until the call fails, and on a// BLOCKING listener that last call does not fail — it sleeps in the kernel// (`inet_csk_accept`) until the next client arrives. The thread then never// re-reads `core.running`, so `shutdown()`'s `join(acceptor)` waits forever// and `close()` never returns. Measured: after ONE connection the process was// down to two threads, main in `__futex_wait` (the join) and the acceptor in// `inet_csk_accept`; each extra TCP connect added exactly one fd and put the// acceptor straight back into `inet_csk_accept`. hl:http2's listener// (`http2.zig:926`) has always carried NONBLOCK — http1 was the outlier.// The accept flags below apply to the ACCEPTED socket, never to this one.const sock_rc = linux.socket(linux.AF.INET, linux.SOCK.STREAM | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK, 0);const server_fd: i32 = @bitCast(@as(u32, @truncate(sock_rc)));if (server_fd < 0) {logMsg("http1: socket() failed\n");allocator.destroy(core);return null;}core.server_fd = server_fd;// SO_REUSEADDRconst one: i32 = 1;_ = linux.setsockopt(server_fd, linux.SOL.SOCKET, linux.SO.REUSEADDR, @ptrCast(&one), @sizeOf(i32));// Parse host addressvar addr_val: u32 = 0; // INADDR_ANYif (host_str.len > 0 and !std.mem.eql(u8, host_str, "0.0.0.0")) {addr_val = parseIPv4(host_str) orelse 0;}// Bindconst addr = linux.sockaddr.in{.port = std.mem.nativeToBig(u16, port),.addr = addr_val,};const bind_rc = linux.bind(server_fd, @ptrCast(&addr), @sizeOf(linux.sockaddr.in));const bind_err: i32 = @bitCast(@as(u32, @truncate(bind_rc)));if (bind_err < 0) {// The prefix is LOAD-BEARING: `tests/browser/fixtures.mjs` breaks its// readiness wait on `bind() failed on port N` (mission 285). The errno// and the holder are appended to it, never in front of it.var why: [320]u8 = undefined;logFmt("http1: bind() failed on port {d}: {s}\n", .{ port, http.bindFailureDetail(&why, bind_rc, port, false) });core.destroy();return null;}const listen_rc = linux.listen(server_fd, 128);const listen_err: i32 = @bitCast(@as(u32, @truncate(listen_rc)));if (listen_err < 0) {logMsg("http1: listen() failed\n");core.destroy();return null;}// Create eventfd for shutdown signalingconst efd_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);const shutdown_fd: i32 = @bitCast(@as(u32, @truncate(efd_rc)));if (shutdown_fd < 0) {logMsg("http1: eventfd() failed\n");core.destroy();return null;}core.shutdown_fd = shutdown_fd;// Create epoll instanceconst epoll_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);const epoll_fd: i32 = @bitCast(@as(u32, @truncate(epoll_rc)));if (epoll_fd < 0) {logMsg("http1: epoll_create1() failed\n");core.destroy();return null;}core.epoll_fd = epoll_fd;// Add server_fd to epollvar ev = linux.epoll_event{.events = linux.EPOLL.IN,.data = .{ .fd = server_fd },};_ = linux.epoll_ctl(epoll_fd, linux.EPOLL.CTL_ADD, server_fd, &ev);// Add shutdown_fd to epollvar shutdown_ev = linux.epoll_event{.events = linux.EPOLL.IN,.data = .{ .fd = shutdown_fd },};_ = linux.epoll_ctl(epoll_fd, linux.EPOLL.CTL_ADD, shutdown_fd, &shutdown_ev);// Initialize TLS if cert+key providedif (cert_path.len > 0 and key_path.len > 0) {core.tls = TlsContext.init(cert_path, key_path);if (core.tls == null) {logMsg("http1: TLS initialization failed, falling back to plain HTTP\n");}}logFmt("http1: listening on :{d}{s}\n", .{ port, if (core.tls != null) " (TLS)" else "" });// Start threadscore.running.store(true, .release);// Parker thread first: the acceptor parks into it from its very first accept.// If it cannot start we log and keep going — connections then go straight to the// pool, which is the pre-084 (starvable) behaviour rather than an outage.if (!core.parked.start(core)) {logMsg("http1: parker thread unavailable — idle keep-alive connections will hold I/O workers\n");}// Allocate I/O threadsconst threads = allocator.alloc(std.Thread, core.num_threads) catch {core.destroy();return null;};core.io_threads = threads;for (0..core.num_threads) |i| {core.io_threads[i] = std.Thread.spawn(.{}, ioWorker, .{core}) catch {logFmt("http1: failed to spawn I/O thread {d}\n", .{i});core.num_threads = @intCast(i);core.io_threads = core.io_threads[0..i];break;};}// Start acceptor threadcore.acceptor_thread = std.Thread.spawn(.{}, acceptorLoop, .{core}) catch {logMsg("http1: failed to spawn acceptor thread\n");core.shutdown();core.destroy();return null;};core.port = port;registerCore(core);return core;}fn shutdown(self: *ServerCore) void {unregisterCore(self);if (!self.running.swap(false, .acq_rel)) return;// Signal shutdown via eventfdif (self.shutdown_fd >= 0) {const val: u64 = 1;_ = linux.write(self.shutdown_fd, @ptrCast(&val), @sizeOf(u64));}// Wake up the request queue so dequeue() unblocksself.request_queue.signalShutdown();// Stop parking before the workers: the parker must not push new work into a// queue whose consumers are being torn down.self.parked.stop();// Signal connection queue to wake I/O workersself.conn_queue.signalShutdown();// Join acceptor threadif (self.acceptor_thread) |t| {t.join();self.acceptor_thread = null;}// Join I/O threadsfor (self.io_threads) |t| {t.join();}}fn destroy(self: *ServerCore) void {self.shutdown();if (self.epoll_fd >= 0) _ = linux.close(self.epoll_fd);if (self.shutdown_fd >= 0) _ = linux.close(self.shutdown_fd);if (self.server_fd >= 0) _ = linux.close(self.server_fd);if (self.tls) |*tls| tls.deinit();// Close every still-parked idle connection, then drain conn_queueself.parked.deinit(if (self.tls) |*t| t else null);self.conn_queue.deinit(if (self.tls) |*t| t else null);self.request_queue.deinit();if (self.io_threads.len > 0) allocator.free(self.io_threads);allocator.destroy(self);}};// =========================================================================// Connection Queue — MPSC queue for acceptor → I/O workers// =========================================================================const ConnectionQueue = struct {queue: std.ArrayListUnmanaged(Connection),mutex: PthreadMutex,condvar: PthreadCond,shutdown: bool,fn init() ConnectionQueue {var self: ConnectionQueue = .{.queue = .empty,.mutex = c.PTHREAD_MUTEX_INITIALIZER,.condvar = c.PTHREAD_COND_INITIALIZER,.shutdown = false,};mutexInit(&self.mutex);condInit(&self.condvar);return self;}fn enqueue(self: *ConnectionQueue, conn: Connection) void {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);self.queue.append(allocator, conn) catch return;condSignal(&self.condvar);}fn dequeue(self: *ConnectionQueue) ?Connection {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);while (self.queue.items.len == 0 and !self.shutdown) {condWait(&self.condvar, &self.mutex);}if (self.shutdown and self.queue.items.len == 0) return null;return self.queue.orderedRemove(0);}fn signalShutdown(self: *ConnectionQueue) void {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);self.shutdown = true;condBroadcast(&self.condvar);}/// Close everything still queued and stay a VALID, EMPTY queue — `deinit`/// leaves the list `undefined`, which is only safe on a core that is being/// freed, and a closed core deliberately is not (see `hl_http1_close`).fn closeAll(self: *ConnectionQueue, tls: ?*TlsContext) void {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);for (self.queue.items) |*conn| conn.close(tls);self.queue.clearRetainingCapacity();}fn deinit(self: *ConnectionQueue, tls: ?*TlsContext) void {for (self.queue.items) |*conn| {conn.close(tls);}self.queue.deinit(allocator);}};// =========================================================================// ParkedConns — idle connections wait HERE, not in an I/O worker (mission 084)// =========================================================================//// The starvation this fixes, measured: a `threads = 4` server, four keep-alive// connections that have gone quiet, and every subsequent request times out — each idle// socket sat inside a worker's blocking read(). Parking inverts that: an idle connection// costs one epoll registration and ZERO threads, and a worker only ever picks up a// connection that already has bytes waiting (or has hung up, which it reads as EOF and// closes). Connections are parked from the acceptor (a fresh socket may be silent — a// pre-connected browser socket routinely is) and after every keep-alive response.//// TLS caveat: an ESTABLISHED TLS connection is never parked. OpenSSL may hold already-// decrypted plaintext in its own buffer, which epoll on the raw fd cannot see, so parking// it could hang a live request. `ssl == null` covers all plain HTTP plus the pre-handshake// TLS socket (the ClientHello does arrive on the raw fd), which is what the park path takes./// How long a parked, silent connection is kept before it is closed. It costs no thread,/// only an fd, so this is generous compared to the old in-worker 30s SO_RCVTIMEO.const PARK_IDLE_TIMEOUT_MS: i64 = 60_000;/// Milliseconds on CLOCK_MONOTONIC. This zig's `std.time` exposes no timestamp function,/// and monotonic is the right clock anyway — a wall-clock step must not expire a live/// connection early or keep a dead one parked.fn monotonicMs() i64 {var ts: linux.timespec = undefined;if (linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts) != 0) return 0;return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);}const ParkedConn = struct {conn: Connection,/// monotonic ms after which this silent connection is closeddeadline_ms: i64,};const ParkedConns = struct {epoll_fd: i32,wake_fd: i32,thread: ?std.Thread,running: std.atomic.Value(bool),mutex: PthreadMutex,/// fd → parked connection. Guarded by `mutex`; the parker thread is the only/// consumer, park() the only producer, so a plain map is enough.map: std.AutoHashMapUnmanaged(i32, ParkedConn),fn init() ParkedConns {var self: ParkedConns = .{.epoll_fd = -1,.wake_fd = -1,.thread = null,.running = std.atomic.Value(bool).init(false),.mutex = c.PTHREAD_MUTEX_INITIALIZER,.map = .empty,};mutexInit(&self.mutex);return self;}/// Create the epoll instance + wake eventfd and spawn the parker thread./// Returns false if the kernel objects could not be made — the caller then falls/// back to handing connections straight to the pool (old behaviour, still correct,/// just starvable).fn start(self: *ParkedConns, core: *ServerCore) bool {const ep_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);const ep: i32 = @bitCast(@as(u32, @truncate(ep_rc)));if (ep < 0) return false;const ef_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);const ef: i32 = @bitCast(@as(u32, @truncate(ef_rc)));if (ef < 0) {_ = linux.close(ep);return false;}var wake_ev = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .fd = ef } };_ = linux.epoll_ctl(ep, linux.EPOLL.CTL_ADD, ef, &wake_ev);self.epoll_fd = ep;self.wake_fd = ef;self.running.store(true, .release);self.thread = std.Thread.spawn(.{}, parkerLoop, .{core}) catch {self.running.store(false, .release);_ = linux.close(ep);_ = linux.close(ef);self.epoll_fd = -1;self.wake_fd = -1;return false;};return true;}/// Hand a connection to the parker. The map insert happens BEFORE the epoll ADD so/// the parker can never see a readable fd it has no entry for.fn park(self: *ParkedConns, conn: Connection) bool {if (self.epoll_fd < 0 or !self.running.load(.acquire)) return false;mutexLock(&self.mutex);self.map.put(allocator, conn.fd, .{.conn = conn,.deadline_ms = monotonicMs() + PARK_IDLE_TIMEOUT_MS,}) catch {mutexUnlock(&self.mutex);return false;};var ev = linux.epoll_event{.events = linux.EPOLL.IN | linux.EPOLL.RDHUP,.data = .{ .fd = conn.fd },};const rc = linux.epoll_ctl(self.epoll_fd, linux.EPOLL.CTL_ADD, conn.fd, &ev);const err: i32 = @bitCast(@as(u32, @truncate(rc)));if (err < 0) {_ = self.map.remove(conn.fd);mutexUnlock(&self.mutex);return false;}mutexUnlock(&self.mutex);return true;}/// Take a parked connection off the epoll set. Returns it if we still owned it.fn take(self: *ParkedConns, fd: i32) ?Connection {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);const entry = self.map.fetchRemove(fd) orelse return null;_ = linux.epoll_ctl(self.epoll_fd, linux.EPOLL.CTL_DEL, fd, null);return entry.value.conn;}/// Close every parked connection whose silence outlived PARK_IDLE_TIMEOUT_MS.fn sweepExpired(self: *ParkedConns, tls: ?*TlsContext) void {const now = monotonicMs();mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);var expired: [64]i32 = undefined;var n: usize = 0;var it = self.map.iterator();while (it.next()) |kv| {if (kv.value_ptr.deadline_ms <= now) {if (n == expired.len) break;expired[n] = kv.key_ptr.*;n += 1;}}for (expired[0..n]) |fd| {if (self.map.fetchRemove(fd)) |e| {_ = linux.epoll_ctl(self.epoll_fd, linux.EPOLL.CTL_DEL, fd, null);var conn = e.value.conn;conn.close(tls);}}}fn stop(self: *ParkedConns) void {if (!self.running.swap(false, .acq_rel)) return;if (self.wake_fd >= 0) {const val: u64 = 1;_ = linux.write(self.wake_fd, @ptrCast(&val), @sizeOf(u64));}if (self.thread) |t| {t.join();self.thread = null;}}/// Close every parked connection and stay a VALID, EMPTY map. Unlike `deinit`/// this keeps the epoll/wake fds and the struct usable — a CLOSED core is not/// a freed one (`hl_http1_close`), and `park()` already refuses once `running`/// is false, so the emptied map simply stays empty.fn closeAll(self: *ParkedConns, tls: ?*TlsContext) void {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);var it = self.map.iterator();while (it.next()) |kv| {var conn = kv.value_ptr.conn;conn.close(tls);}self.map.clearRetainingCapacity();}fn deinit(self: *ParkedConns, tls: ?*TlsContext) void {self.stop();mutexLock(&self.mutex);var it = self.map.iterator();while (it.next()) |kv| {var conn = kv.value_ptr.conn;conn.close(tls);}self.map.deinit(allocator);self.map = .empty;mutexUnlock(&self.mutex);if (self.epoll_fd >= 0) _ = linux.close(self.epoll_fd);if (self.wake_fd >= 0) _ = linux.close(self.wake_fd);self.epoll_fd = -1;self.wake_fd = -1;}};/// The parker thread: waits for a parked connection to become READABLE and only then/// hands it to an I/O worker. The 1s epoll timeout doubles as the idle-sweep tick.fn parkerLoop(core: *ServerCore) void {var events: [64]linux.epoll_event = undefined;while (core.parked.running.load(.acquire)) {const n_rc = linux.epoll_wait(core.parked.epoll_fd, &events, events.len, 1000);const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));if (n > 0) {for (events[0..@intCast(n)]) |ev| {if (ev.data.fd == core.parked.wake_fd) {var drain: u64 = 0;_ = linux.read(core.parked.wake_fd, @ptrCast(&drain), @sizeOf(u64));continue;}// Readable, hung up or errored — all three are a worker's job: it either// parses the request or reads EOF and closes.if (core.parked.take(ev.data.fd)) |conn| {if (core.running.load(.acquire)) {core.conn_queue.enqueue(conn);} else {var dead = conn;dead.close(if (core.tls) |*t| t else null);}}}}core.parked.sweepExpired(if (core.tls) |*t| t else null);}}/// Park `conn` if we can, otherwise hand it straight to the pool. Every enqueue site/// that is NOT known to have bytes waiting goes through here.fn parkOrEnqueue(core: *ServerCore, conn: Connection) void {// An established TLS connection may hold decrypted bytes epoll cannot see — see the// ParkedConns header comment. Those go straight to a worker, as before.if (conn.ssl == null and core.parked.park(conn)) return;core.conn_queue.enqueue(conn);}// =========================================================================// Acceptor thread — epoll loop accepting new connections// =========================================================================fn acceptorLoop(core: *ServerCore) void {var events: [64]linux.epoll_event = undefined;while (core.running.load(.acquire)) {const n_rc = linux.epoll_wait(core.epoll_fd, &events, events.len, 1000);const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));if (n < 0) continue;if (n == 0) continue;for (events[0..@intCast(n)]) |ev| {if (ev.data.fd == core.shutdown_fd) {return; // shutdown signaled}if (ev.data.fd == core.server_fd) {// Accept all pending connections — the listener is NONBLOCK, so the// drain ends on EAGAIN instead of sleeping inside accept4().while (true) {var client_addr: linux.sockaddr.in = undefined;var addr_len: u32 = @sizeOf(linux.sockaddr.in);const accept_rc = linux.accept4(core.server_fd, @ptrCast(&client_addr), &addr_len, linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK);const client_fd: i32 = @bitCast(@as(u32, @truncate(accept_rc)));if (client_fd < 0) break;// A shutdown that landed mid-drain: the parker is stopped and the// conn_queue has no consumers left, so hand this socket to nobody —// close it and leave, rather than leaking the fd into a dead queue.if (!core.running.load(.acquire)) {_ = linux.close(client_fd);return;}// Set back to blocking for I/O workers (simpler read/write)const flags_rc = linux.fcntl(client_fd, linux.F.GETFL, @as(usize, 0));const flags_i: isize = @bitCast(flags_rc);if (flags_i >= 0) {var oflags: linux.O = @bitCast(@as(u32, @truncate(flags_rc)));oflags.NONBLOCK = false;_ = linux.fcntl(client_fd, linux.F.SETFL, @as(usize, @as(u32, @bitCast(oflags))));}// PARK, don't hand to a worker: a just-accepted socket has no bytes// yet (browsers routinely pre-open connections and send nothing), and// a worker blocking on it is exactly the starvation this replaces.parkOrEnqueue(core, .{.fd = client_fd,.ssl = null,.keep_alive = true,.request_count = 0,.max_requests = MAX_KEEPALIVE_REQUESTS,});}}}}}// =========================================================================// I/O Worker thread — TLS handshake + read/parse HTTP → enqueue request// =========================================================================fn ioWorker(core: *ServerCore) void {while (core.running.load(.acquire)) {var conn = core.conn_queue.dequeue() orelse return;// TLS handshake if needed (only on first request for this connection)if (core.tls != null and conn.ssl == null and conn.request_count == 0) {conn.ssl = core.tls.?.wrapConnection(conn.fd);if (conn.ssl == null) {_ = linux.close(conn.fd);continue;}}// Bound the worker's blocking read on EVERY connection, not just keep-alive ones.// A parked connection only reaches a worker once it is readable, so this is now a// backstop against a client that trickles (or stops mid-)headers rather than the// idle-timeout mechanism it used to be — that job belongs to PARK_IDLE_TIMEOUT_MS.const tv = linux.timeval{ .sec = 30, .usec = 0 };_ = linux.setsockopt(conn.fd, linux.SOL.SOCKET, linux.SO.RCVTIMEO, @ptrCast(&tv), @sizeOf(linux.timeval));// Read and parse HTTP requestswitch (readAndParseRequest(&conn, core)) {.request => |req| core.request_queue.enqueue(req),.upgraded => {}, // socket handed to the WebSocket engine.dead => conn.close(if (core.tls) |*t| t else null),}}}const ReadOutcome = union(enum) {request: ParsedRequest,upgraded, // WS handshake done — the WS engine owns the fd nowdead, // read failed / handshake rejected — caller closes};fn readAndParseRequest(conn: *Connection, core: *ServerCore) ReadOutcome {const tls = if (core.tls) |*t| t else null;// Read request headers (up to 64KB)var buf = allocator.alloc(u8, 65536) catch return .dead;defer allocator.free(buf);var total: usize = 0;var hdr_end: ?usize = null;while (total < buf.len) {const n = conn.read(tls, buf[total..]);if (n == 0) break;total += n;// Scan for \r\n\r\nconst start = if (total > n + 3) total - n - 3 else 0;if (std.mem.indexOf(u8, buf[start..total], "\r\n\r\n")) |pos| {hdr_end = start + pos;break;}}if (total == 0) return .dead;const headers_end = hdr_end orelse total;const header_section = buf[0..headers_end];const body_offset = if (hdr_end) |e| e + 4 else total;// Parse request line: METHOD PATH HTTP/1.xconst line_end = std.mem.indexOf(u8, header_section, "\r\n") orelse headers_end;const request_line = header_section[0..line_end];var method: []const u8 = "GET";var full_path: []const u8 = "/";if (std.mem.indexOfScalar(u8, request_line, ' ')) |sp1| {method = request_line[0..sp1];const rest = request_line[sp1 + 1 ..];if (std.mem.indexOfScalar(u8, rest, ' ')) |sp2| {full_path = rest[0..sp2];
Only the first lines are shown.
Branches
- mainmain branch
Latest commits
- 2ab91ee9tickets: gate checks rows appear once (session sync); re-vendor to ff51cf46 stopped on hybriel#122, stays 837fe120mre
- e01c2b1dtickets#24: installable app (manifest, service worker, offline list), own icon; gate waits for the hello's pongmre
- 752fbb7fdeploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
- 38bdd5e4deploy.sh: never send .git or .gitignore to Byrodinmre
- f12fa1bcState of 2026-09-27, before the move to gitoriamre