gitoriaLog in with ident

gitoria

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commit4a2d71254a2d7125initial commitmre4a2d7125/plugins/http1/http1.zig

106.7 KB

  1. // hl:http1 plugin — HTTP/1.1 server with epoll + I/O thread pool
  2. // Compiled to libhttp1.so, loaded by runtime via dlopen
  3. //
  4. // Architecture:
  5. // Acceptor thread (epoll on server fd) → PARK (epoll) → I/O thread pool (TLS + parse)
  6. // → request queue → Hybriel main thread
  7. //
  8. // An IDLE connection never occupies an I/O worker (mission 084). A worker's read() is
  9. // blocking, so a connection handed straight to the pool pins a thread until bytes arrive;
  10. // with the default `threads = 4`, four idle keep-alive sockets (one open browser tab is
  11. // already several) starved every later request. Instead every connection that is not
  12. // known to have readable bytes is PARKED in the parker thread's epoll, and only enters
  13. // conn_queue when it is actually readable (or hung up). See ParkedConns below.
  14. //
  15. // Exports:
  16. // hl_http1_create_server(port, host, cert_path, key_path, threads) → iterator of request objects
  17. // hl_http1_listen(port) → same with defaults (backward compat)
  18. //
  19. // Each request object: method, path, query (object), headers (object), body, respond (handle)
  20. // Respond handle: call("send", status_code, body [, content_type]) → writes response
  21. const std = @import("std");
  22. const api = @import("plugin_api");
  23. const http = @import("http_common");
  24. const ws = @import("ws_common");
  25. const HlValue = api.HlValue;
  26. const HlObject = api.HlObject;
  27. const HlField = api.HlField;
  28. const HlIterator = api.HlIterator;
  29. const HlHandle = api.HlHandle;
  30. const HlString = api.HlString;
  31. const linux = std.os.linux;
  32. const posix = std.posix;
  33. const c = std.c;
  34. const PthreadMutex = c.pthread_mutex_t;
  35. const PthreadCond = c.pthread_cond_t;
  36. fn mutexInit(m: *PthreadMutex) void {
  37. // Static initializer is sufficient; explicit init not exposed in this std.c.
  38. _ = m;
  39. }
  40. fn mutexLock(m: *PthreadMutex) void {
  41. _ = c.pthread_mutex_lock(m);
  42. }
  43. fn mutexUnlock(m: *PthreadMutex) void {
  44. _ = c.pthread_mutex_unlock(m);
  45. }
  46. fn condInit(cnd: *PthreadCond) void {
  47. _ = cnd;
  48. }
  49. fn condSignal(cnd: *PthreadCond) void {
  50. _ = c.pthread_cond_signal(cnd);
  51. }
  52. fn condBroadcast(cnd: *PthreadCond) void {
  53. _ = c.pthread_cond_broadcast(cnd);
  54. }
  55. fn condWait(cnd: *PthreadCond, m: *PthreadMutex) void {
  56. _ = c.pthread_cond_wait(cnd, m);
  57. }
  58. // Use the SMP (thread-safe, production) allocator rather than the debug
  59. // GeneralPurposeAllocator: this plugin allocates per-request across the acceptor,
  60. // I/O-worker and interpreter threads, and the debug allocator's safety bookkeeping
  61. // (canaries/quarantine) segfaulted in its own free() path under request churn
  62. // (mission 027 flakiness). smp_allocator is thread-safe and has no debug tripwires.
  63. const allocator = std.heap.smp_allocator;
  64. // Use direct syscall for stderr writes — std.debug.print uses std.Progress
  65. // which has ABI-incompatible global state when loaded as a plugin into a
  66. // binary compiled with a different Zig version.
  67. fn logMsg(msg: []const u8) void {
  68. _ = linux.write(2, msg.ptr, msg.len);
  69. }
  70. fn logFmt(comptime fmt: []const u8, args: anytype) void {
  71. var buf: [512]u8 = undefined;
  72. const s = std.fmt.bufPrint(&buf, fmt, args) catch return;
  73. logMsg(s);
  74. }
  75. // =========================================================================
  76. // SSE (Server-Sent Events) push channel — mission 031, ADDRESSED in 080
  77. // A registry of connected event-stream clients. `sse_start` (on the respond
  78. // handle) upgrades a GET request into a long-lived text/event-stream, assigns
  79. // the connection a stable id and registers it; `hl_http1_sse_broadcast` writes
  80. // a `data:` frame to every registered client and `hl_http1_sse_send` to ONE of
  81. // them, dropping any that error (peer closed). Touched only from the
  82. // interpreter (main) thread, but guarded by a mutex for safety.
  83. //
  84. // Mission 080 (D26): the id is a MONOTONIC counter, never the fd — a closed
  85. // connection's fd is recycled by the kernel within milliseconds, so an fd-keyed
  86. // push would eventually land on a stranger's socket. `sse_start` returns the id
  87. // (it used to return the raw fd, which no caller used), and the framework sends
  88. // it to the browser as the `__hlHello` event so the client can name itself.
  89. // =========================================================================
  90. // Mission 084: a subscription is reaped PROACTIVELY, not on the first failed push.
  91. // A closed tab used to stay in this registry until something happened to be pushed to
  92. // it — and because a scoped push (§9.5) sends nothing to a client whose needs did not
  93. // change, "something" could be never. Two mechanisms, both in the reaper thread:
  94. //
  95. // • EPOLLRDHUP on every registered fd — a closed tab sends FIN, which fires
  96. // immediately and costs no traffic at all. This is the primary detector.
  97. // • a periodic SSE COMMENT heartbeat (`:\n\n`, which EventSource ignores) — catches
  98. // a peer that vanished WITHOUT a FIN (killed machine, dropped NAT entry), which no
  99. // amount of epolling can see, and keeps intermediaries from timing the stream out.
  100. //
  101. // A heartbeat is also why one failed write is enough to declare death here: the probe
  102. // runs repeatedly, so the "first write after FIN succeeds, second gets EPIPE" TCP
  103. // behaviour just means the peer is reaped one tick later.
  104. const MSG_NOSIGNAL: u32 = 0x4000;
  105. const SseConn = struct { id: u32, fd: i32 };
  106. var sse_mutex: PthreadMutex = c.PTHREAD_MUTEX_INITIALIZER;
  107. var sse_subs: std.ArrayListUnmanaged(SseConn) = .empty;
  108. var sse_next_id: u32 = 1;
  109. /// how often the reaper writes its comment heartbeat / probes for silent death
  110. const SSE_HEARTBEAT_MS: i64 = 5_000;
  111. var sse_epoll_fd: i32 = -1;
  112. var sse_wake_fd: i32 = -1;
  113. var sse_reaper_thread: ?std.Thread = null;
  114. var sse_reaper_running = std.atomic.Value(bool).init(false);
  115. /// Start the reaper thread + its epoll. Idempotent; called under sse_mutex from the
  116. /// first sseRegister(), so a server that never opens a stream never spawns it.
  117. fn sseReaperEnsureLocked() void {
  118. if (sse_reaper_running.load(.acquire)) return;
  119. const ep_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);
  120. const ep: i32 = @bitCast(@as(u32, @truncate(ep_rc)));
  121. if (ep < 0) return;
  122. const ef_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);
  123. const ef: i32 = @bitCast(@as(u32, @truncate(ef_rc)));
  124. if (ef < 0) {
  125. _ = linux.close(ep);
  126. return;
  127. }
  128. var wake_ev = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .fd = ef } };
  129. _ = linux.epoll_ctl(ep, linux.EPOLL.CTL_ADD, ef, &wake_ev);
  130. sse_epoll_fd = ep;
  131. sse_wake_fd = ef;
  132. sse_reaper_running.store(true, .release);
  133. sse_reaper_thread = std.Thread.spawn(.{}, sseReaperLoop, .{}) catch {
  134. sse_reaper_running.store(false, .release);
  135. _ = linux.close(ep);
  136. _ = linux.close(ef);
  137. sse_epoll_fd = -1;
  138. sse_wake_fd = -1;
  139. return;
  140. };
  141. }
  142. /// Drop subscription at index `i`: epoll DEL, close, remove. Caller holds sse_mutex.
  143. fn sseDropLocked(i: usize) void {
  144. const fd = sse_subs.items[i].fd;
  145. if (sse_epoll_fd >= 0) _ = linux.epoll_ctl(sse_epoll_fd, linux.EPOLL.CTL_DEL, fd, null);
  146. _ = linux.close(fd);
  147. _ = sse_subs.swapRemove(i);
  148. }
  149. /// The reaper: EPOLLRDHUP wakes it the instant a tab closes; the 1s timeout paces the
  150. /// heartbeat. Nothing here pushes application data, so a reap needs no traffic from the
  151. /// app at all — which is the whole point (mission 080's gap).
  152. fn sseReaperLoop() void {
  153. var events: [64]linux.epoll_event = undefined;
  154. var last_beat: i64 = monotonicMs();
  155. while (sse_reaper_running.load(.acquire)) {
  156. const n_rc = linux.epoll_wait(sse_epoll_fd, &events, events.len, 1000);
  157. const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));
  158. if (n > 0) {
  159. for (events[0..@intCast(n)]) |ev| {
  160. if (ev.data.fd == sse_wake_fd) {
  161. var drain: u64 = 0;
  162. _ = linux.read(sse_wake_fd, @ptrCast(&drain), @sizeOf(u64));
  163. continue;
  164. }
  165. // A subscriber never SENDS on its stream, so any readable/hangup event
  166. // means the peer went away (or is misbehaving) — either way it is dead.
  167. mutexLock(&sse_mutex);
  168. for (sse_subs.items, 0..) |conn, i| {
  169. if (conn.fd == ev.data.fd) {
  170. sseDropLocked(i);
  171. break;
  172. }
  173. }
  174. mutexUnlock(&sse_mutex);
  175. }
  176. }
  177. const now = monotonicMs();
  178. if (now - last_beat >= SSE_HEARTBEAT_MS) {
  179. last_beat = now;
  180. mutexLock(&sse_mutex);
  181. var i: usize = 0;
  182. while (i < sse_subs.items.len) {
  183. // ":\n\n" is an SSE comment — EventSource ignores it, so this is a pure
  184. // liveness probe that never reaches an `onmessage` handler.
  185. if (!sseRawWrite(sse_subs.items[i].fd, ":\n\n")) {
  186. sseDropLocked(i);
  187. continue;
  188. }
  189. i += 1;
  190. }
  191. mutexUnlock(&sse_mutex);
  192. }
  193. }
  194. }
  195. // Write via sendto with MSG_NOSIGNAL so a dead peer yields EPIPE instead of
  196. // killing the process with SIGPIPE. Returns false on any short/failed write.
  197. fn sseRawWrite(fd: i32, data: []const u8) bool {
  198. var written: usize = 0;
  199. while (written < data.len) {
  200. const rc = linux.sendto(fd, data[written..].ptr, data.len - written, MSG_NOSIGNAL, null, 0);
  201. const n: isize = @bitCast(rc);
  202. if (n <= 0) return false;
  203. written += @intCast(n);
  204. }
  205. return true;
  206. }
  207. fn sseRegister(fd: i32) u32 {
  208. mutexLock(&sse_mutex);
  209. defer mutexUnlock(&sse_mutex);
  210. const id = sse_next_id;
  211. sse_next_id += 1;
  212. sse_subs.append(allocator, .{ .id = id, .fd = fd }) catch return 0;
  213. // Watch it for hangup from now on — see the reaper comment above.
  214. sseReaperEnsureLocked();
  215. if (sse_epoll_fd >= 0) {
  216. var ev = linux.epoll_event{
  217. .events = linux.EPOLL.RDHUP | linux.EPOLL.IN,
  218. .data = .{ .fd = fd },
  219. };
  220. _ = linux.epoll_ctl(sse_epoll_fd, linux.EPOLL.CTL_ADD, fd, &ev);
  221. }
  222. return id;
  223. }
  224. // One `data:` frame on one socket. Caller holds sse_mutex.
  225. fn sseWriteFrameLocked(fd: i32, msg: []const u8) bool {
  226. var ok = sseRawWrite(fd, "data: ");
  227. if (ok) ok = sseRawWrite(fd, msg);
  228. if (ok) ok = sseRawWrite(fd, "\n\n");
  229. return ok;
  230. }
  231. // Broadcast one SSE message to all subscribers. `msg` should be a single line
  232. // (callers send single-line JSON). Dead sockets are closed and removed.
  233. // Returns the number of clients successfully written to.
  234. fn sseBroadcast(msg: []const u8) usize {
  235. mutexLock(&sse_mutex);
  236. defer mutexUnlock(&sse_mutex);
  237. var count: usize = 0;
  238. var i: usize = 0;
  239. while (i < sse_subs.items.len) {
  240. if (!sseWriteFrameLocked(sse_subs.items[i].fd, msg)) {
  241. sseDropLocked(i);
  242. continue;
  243. }
  244. count += 1;
  245. i += 1;
  246. }
  247. return count;
  248. }
  249. // __native("http1.sse_count") → how many SSE subscriptions are currently LIVE.
  250. // The observable the proactive reaper exists to keep honest (mission 084): it must fall
  251. // when a client goes away, with nothing being pushed.
  252. export fn hl_http1_sse_count(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  253. _ = argc;
  254. _ = argv;
  255. mutexLock(&sse_mutex);
  256. defer mutexUnlock(&sse_mutex);
  257. return api.makeNumber(@floatFromInt(sse_subs.items.len));
  258. }
  259. // __native("http1.sse_alive", connId) → 1 if that subscription is still registered.
  260. // Lets the framework prune its own per-connection bookkeeping without pushing anything.
  261. export fn hl_http1_sse_alive(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  262. if (argc < 1 or argv[0].type != .hl_number) return api.makeNumber(0);
  263. const id: u32 = @intFromFloat(argv[0].data.number);
  264. mutexLock(&sse_mutex);
  265. defer mutexUnlock(&sse_mutex);
  266. for (sse_subs.items) |conn| {
  267. if (conn.id == id) return api.makeNumber(1);
  268. }
  269. return api.makeNumber(0);
  270. }
  271. // __native("http1.sse_broadcast", jsonString) → number of clients pushed to.
  272. export fn hl_http1_sse_broadcast(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  273. if (argc < 1 or argv[0].type != .hl_string) return api.makeNumber(0);
  274. const msg = argv[0].data.string.ptr[0..argv[0].data.string.len];
  275. const n = sseBroadcast(msg);
  276. return api.makeNumber(@floatFromInt(n));
  277. }
  278. // __native("http1.sse_send", connId, jsonString) → 1 pushed / 0 gone.
  279. // Mission 080 (D26): the ADDRESSED half of the channel — the dependency-scoped
  280. // push sends a different payload to each client, so it cannot use broadcast.
  281. export fn hl_http1_sse_send(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  282. if (argc < 2 or argv[0].type != .hl_number or argv[1].type != .hl_string) return api.makeNumber(0);
  283. const id: u32 = @intFromFloat(argv[0].data.number);
  284. const msg = argv[1].data.string.ptr[0..argv[1].data.string.len];
  285. mutexLock(&sse_mutex);
  286. defer mutexUnlock(&sse_mutex);
  287. for (sse_subs.items, 0..) |conn, i| {
  288. if (conn.id != id) continue;
  289. if (!sseWriteFrameLocked(conn.fd, msg)) {
  290. sseDropLocked(i);
  291. return api.makeNumber(0);
  292. }
  293. return api.makeNumber(1);
  294. }
  295. return api.makeNumber(0);
  296. }
  297. // =========================================================================
  298. // WebSocket engine (mission 064, decisions D8/D9)
  299. //
  300. // Framing lives in the shared core (plugins/http/ws_common.zig); this engine
  301. // owns the h1-specific parts: the Upgrade/101 handshake (done inline in the
  302. // I/O worker, see readAndParseRequest) and the socket lifecycle after it.
  303. //
  304. // One GLOBAL engine per plugin (mirrors the SSE registry): a single reader
  305. // thread epolls all upgraded sockets with level-triggered EPOLLIN only —
  306. // EPOLLOUT is never armed, so an idle connection never wakes the loop (the
  307. // busy-spin trap found in the http2 TLS path). Complete messages become
  308. // WsEvent entries that the interpreter drains via the hl_http1_ws_events
  309. // iterator (registered on the event loop like the request iterator).
  310. //
  311. // Outbound writes (send/broadcast/ping/close + pong replies) happen under
  312. // ws_mutex from either the interpreter thread or the reader thread. Sockets
  313. // are non-blocking; a slow consumer gets bounded EAGAIN retries (~100ms)
  314. // and is dropped rather than buffered (no outbound queue, no EPOLLOUT).
  315. // WS upgrades are plain-HTTP only for now (TLS WS = future work, like SSE).
  316. // =========================================================================
  317. const WsEventKind = enum(u8) { connect, message, close, pong };
  318. const WsEvent = struct {
  319. kind: WsEventKind,
  320. id: u32,
  321. data: ?[]u8 = null, // allocated payload (message data / close reason)
  322. is_binary: bool = false,
  323. code: u16 = 0, // close code
  324. // Mission 093: the UPGRADE REQUEST's `Cookie` header, verbatim, carried on the
  325. // `connect` event only. A WebSocket handshake is an ordinary HTTP request, so
  326. // the browser sends the session cookie with it automatically — this is the one
  327. // moment the socket can be attributed to whoever loaded the page, and after the
  328. // 101 the request (and its headers) is freed. Allocated; freed with the event.
  329. cookie: ?[]u8 = null,
  330. // Ticket #74: the upgrade request's `Host` header, verbatim, on `connect` only —
  331. // the address the browser dialled, which a page served per subdomain is rendered
  332. // for. Allocated; freed with the event.
  333. host: ?[]u8 = null,
  334. };
  335. const WsClient = struct {
  336. id: u32,
  337. fd: i32,
  338. decoder: ws.Decoder,
  339. closing: bool = false, // server sent close, awaiting peer echo
  340. // liveness (mission 091): when this socket last produced a frame, and when
  341. // the sweep's ping went out (0 = none outstanding). A TCP connection whose
  342. // peer vanished without a FIN stays writable indefinitely, so "still open"
  343. // is not evidence of a live peer — the pong is.
  344. last_seen_ms: i64 = 0,
  345. ping_at_ms: i64 = 0,
  346. };
  347. var ws_mutex: PthreadMutex = c.PTHREAD_MUTEX_INITIALIZER;
  348. var ws_clients: std.ArrayListUnmanaged(*WsClient) = .empty;
  349. var ws_event_queue: std.ArrayListUnmanaged(WsEvent) = .empty;
  350. var ws_epoll_fd: i32 = -1;
  351. var ws_wake_fd: i32 = -1;
  352. /// THE EVENT LOOP'S BELL for the WS event queue (mission 256), and a different fd
  353. /// from `ws_wake_fd` above — that one wakes the reader THREAD's own epoll, this
  354. /// one wakes the interpreter loop. Without it an app that turns sockets on has a
  355. /// source with no fd, and ONE such source puts the whole loop back on the 1ms
  356. /// poll (loop_wait.Waiter.observe) — so the server would busy-poll for as long as
  357. /// WebSockets were enabled.
  358. var ws_loop_wake_fd: i32 = -1;
  359. var ws_thread: ?std.Thread = null;
  360. var ws_running = std.atomic.Value(bool).init(false);
  361. var ws_enabled = std.atomic.Value(bool).init(false);
  362. var ws_next_id: u32 = 1;
  363. // --- liveness sweep (mission 091) ----------------------------------------
  364. // The reader thread already wakes every 500ms (the epoll timeout), so the sweep
  365. // costs no timer and no new thread: on every tick it pings sockets that have
  366. // gone quiet and drops the ones whose pong is overdue. The drop enqueues the
  367. // ordinary `close` event, so every consumer above (hl:web's subscription
  368. // registry included) prunes through the path it already had — nothing upstream
  369. // learns a new concept, and nothing at the hl level needs a timer.
  370. //
  371. // Operator knobs, read once at engine start; the defaults are production
  372. // values, the tests shrink them.
  373. const WS_PING_MS_DEFAULT: i64 = 15000; // quiet this long -> ask
  374. const WS_PONG_MS_DEFAULT: i64 = 10000; // no answer in this long -> gone
  375. var ws_ping_ms: i64 = WS_PING_MS_DEFAULT;
  376. var ws_pong_ms: i64 = WS_PONG_MS_DEFAULT;
  377. /// MONOTONIC milliseconds — a timeout measured against the wall clock would fire
  378. /// early or never after an NTP step.
  379. fn wsNowMs() i64 {
  380. var ts: linux.timespec = undefined;
  381. _ = linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts);
  382. return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);
  383. }
  384. fn wsEnvMs(name: [:0]const u8, fallback: i64) i64 {
  385. const raw = std.mem.span(std.c.getenv(name) orelse return fallback);
  386. const n = std.fmt.parseInt(i64, std.mem.trim(u8, raw, " \t"), 10) catch return fallback;
  387. if (n <= 0) return fallback;
  388. return n;
  389. }
  390. const EAGAIN_ERR: isize = 11;
  391. const EINTR_ERR: isize = 4;
  392. /// WriteFn callback over a raw fd (context = fd stuffed into the pointer).
  393. /// Non-blocking socket: bounded EAGAIN retries, then give up (caller drops).
  394. fn wsFdWrite(ctx: ?*anyopaque, data: []const u8) bool {
  395. const fd: i32 = @intCast(@intFromPtr(ctx));
  396. var written: usize = 0;
  397. var retries: u32 = 0;
  398. while (written < data.len) {
  399. const rc = linux.sendto(fd, data[written..].ptr, data.len - written, MSG_NOSIGNAL, null, 0);
  400. const n: isize = @bitCast(rc);
  401. if (n > 0) {
  402. written += @intCast(n);
  403. retries = 0;
  404. continue;
  405. }
  406. const e = -n;
  407. if (e == EAGAIN_ERR) {
  408. retries += 1;
  409. if (retries > 100) return false; // ~100ms of backpressure → drop
  410. const req = linux.timespec{ .sec = 0, .nsec = 1_000_000 }; // 1ms
  411. _ = linux.nanosleep(&req, null);
  412. continue;
  413. }
  414. if (e == EINTR_ERR) continue;
  415. return false;
  416. }
  417. return true;
  418. }
  419. fn wsFdCtx(fd: i32) ?*anyopaque {
  420. return @ptrFromInt(@as(usize, @intCast(fd)));
  421. }
  422. /// Start the reader thread + epoll instance. Idempotent. Called from the
  423. /// interpreter thread (hl_http1_ws_events); also flips ws_enabled so the I/O
  424. /// workers start honoring Upgrade requests.
  425. fn wsEnsureStarted() bool {
  426. mutexLock(&ws_mutex);
  427. defer mutexUnlock(&ws_mutex);
  428. if (ws_running.load(.acquire)) return true;
  429. ws_ping_ms = wsEnvMs("HL_WS_PING_MS", WS_PING_MS_DEFAULT);
  430. ws_pong_ms = wsEnvMs("HL_WS_PONG_TIMEOUT_MS", WS_PONG_MS_DEFAULT);
  431. const ep_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);
  432. const ep: i32 = @bitCast(@as(u32, @truncate(ep_rc)));
  433. if (ep < 0) return false;
  434. const ef_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);
  435. const ef: i32 = @bitCast(@as(u32, @truncate(ef_rc)));
  436. if (ef < 0) {
  437. _ = linux.close(ep);
  438. return false;
  439. }
  440. var wake_ev = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .fd = ef } };
  441. _ = linux.epoll_ctl(ep, linux.EPOLL.CTL_ADD, ef, &wake_ev);
  442. ws_epoll_fd = ep;
  443. ws_wake_fd = ef;
  444. if (ws_loop_wake_fd < 0) ws_loop_wake_fd = http.makeWakeFd(); // mission 256
  445. ws_running.store(true, .release);
  446. ws_thread = std.Thread.spawn(.{}, wsReaderLoop, .{}) catch {
  447. ws_running.store(false, .release);
  448. _ = linux.close(ep);
  449. _ = linux.close(ef);
  450. ws_epoll_fd = -1;
  451. ws_wake_fd = -1;
  452. return false;
  453. };
  454. ws_enabled.store(true, .release);
  455. return true;
  456. }
  457. fn wsWake() void {
  458. if (ws_wake_fd >= 0) {
  459. const one: u64 = 1;
  460. _ = linux.write(ws_wake_fd, @ptrCast(&one), @sizeOf(u64));
  461. }
  462. }
  463. /// Hand an upgraded socket to the engine. Called from an I/O worker thread
  464. /// right after the 101 was written. `leftover` = bytes the client sent after
  465. /// the handshake that were already consumed into the header buffer. `cookie` is
  466. /// the handshake's `Cookie` header (mission 093) — copied here, because the
  467. /// parsed request is freed the moment this returns. `host` is its `Host` header
  468. /// (ticket #74), copied for the same reason.
  469. fn wsRegisterClient(fd: i32, leftover: []const u8, cookie: []const u8, host: []const u8) void {
  470. // Non-blocking for the reader loop
  471. const flags_rc = linux.fcntl(fd, linux.F.GETFL, @as(usize, 0));
  472. const flags_i: isize = @bitCast(flags_rc);
  473. if (flags_i >= 0) {
  474. var oflags: linux.O = @bitCast(@as(u32, @truncate(flags_rc)));
  475. oflags.NONBLOCK = true;
  476. _ = linux.fcntl(fd, linux.F.SETFL, @as(usize, @as(u32, @bitCast(oflags))));
  477. }
  478. mutexLock(&ws_mutex);
  479. const client = allocator.create(WsClient) catch {
  480. mutexUnlock(&ws_mutex);
  481. _ = linux.close(fd);
  482. return;
  483. };
  484. client.* = .{
  485. .id = ws_next_id,
  486. .fd = fd,
  487. .decoder = ws.Decoder.init(allocator),
  488. .last_seen_ms = wsNowMs(),
  489. };
  490. ws_next_id += 1;
  491. ws_clients.append(allocator, client) catch {
  492. allocator.destroy(client);
  493. mutexUnlock(&ws_mutex);
  494. _ = linux.close(fd);
  495. return;
  496. };
  497. const cookie_copy: ?[]u8 = if (cookie.len > 0) (allocator.dupe(u8, cookie) catch null) else null;
  498. const host_copy: ?[]u8 = if (host.len > 0) (allocator.dupe(u8, host) catch null) else null;
  499. ws_event_queue.append(allocator, .{ .kind = .connect, .id = client.id, .cookie = cookie_copy, .host = host_copy }) catch {
  500. if (cookie_copy) |cc| allocator.free(cc);
  501. if (host_copy) |hc| allocator.free(hc);
  502. };
  503. if (leftover.len > 0) client.decoder.feed(leftover) catch {};
  504. mutexUnlock(&ws_mutex);
  505. // Level-triggered EPOLLIN only (never EPOLLOUT): pending socket data fires
  506. // immediately, idle connections cost nothing.
  507. var ev = linux.epoll_event{ .events = linux.EPOLL.IN | linux.EPOLL.RDHUP, .data = .{ .fd = fd } };
  508. _ = linux.epoll_ctl(ws_epoll_fd, linux.EPOLL.CTL_ADD, fd, &ev);
  509. if (leftover.len > 0) wsWake(); // decoder-buffered bytes won't fire EPOLLIN
  510. // This runs on an I/O WORKER thread, not the reader thread, and it queued a
  511. // `connect` — so it rings the event loop itself (mission 256).
  512. http.ringWake(ws_loop_wake_fd);
  513. }
  514. fn wsFindByFdLocked(fd: i32) ?usize {
  515. for (ws_clients.items, 0..) |cl, i| {
  516. if (cl.fd == fd) return i;
  517. }
  518. return null;
  519. }
  520. fn wsFindByIdLocked(id: u32) ?usize {
  521. for (ws_clients.items, 0..) |cl, i| {
  522. if (cl.id == id) return i;
  523. }
  524. return null;
  525. }
  526. /// Remove client at index: epoll DEL, close fd, free state. Mutex held.
  527. fn wsRemoveLocked(idx: usize) void {
  528. const client = ws_clients.items[idx];
  529. _ = linux.epoll_ctl(ws_epoll_fd, linux.EPOLL.CTL_DEL, client.fd, null);
  530. _ = linux.close(client.fd);
  531. client.decoder.deinit();
  532. _ = ws_clients.swapRemove(idx);
  533. allocator.destroy(client);
  534. }
  535. /// Drain the client's decoder; enqueue events, auto-reply pings, run the
  536. /// close handshake. Returns true if the client was removed. Mutex held.
  537. fn wsDrainDecoderLocked(idx: usize) bool {
  538. const client = ws_clients.items[idx];
  539. while (true) {
  540. const maybe_ev = client.decoder.next() catch {
  541. // OOM mid-decode — drop the connection
  542. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1011 }) catch {};
  543. wsRemoveLocked(idx);
  544. return true;
  545. };
  546. const ev = maybe_ev orelse return false;
  547. // Any complete frame is proof of life, and it answers an outstanding
  548. // sweep ping whatever its opcode — a peer that is talking is not dead.
  549. client.last_seen_ms = wsNowMs();
  550. client.ping_at_ms = 0;
  551. switch (ev) {
  552. .text => |p| ws_event_queue.append(allocator, .{ .kind = .message, .id = client.id, .data = p }) catch allocator.free(p),
  553. .binary => |p| ws_event_queue.append(allocator, .{ .kind = .message, .id = client.id, .data = p, .is_binary = true }) catch allocator.free(p),
  554. .ping => |p| {
  555. _ = ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .pong, p);
  556. allocator.free(p);
  557. },
  558. .pong => |p| {
  559. allocator.free(p);
  560. ws_event_queue.append(allocator, .{ .kind = .pong, .id = client.id }) catch {};
  561. },
  562. .close => |cl| {
  563. if (!client.closing) {
  564. const echo_code = if (cl.code == ws.CLOSE_NO_STATUS) ws.CLOSE_NORMAL else cl.code;
  565. _ = ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), echo_code, "");
  566. }
  567. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = cl.code, .data = cl.reason }) catch allocator.free(cl.reason);
  568. wsRemoveLocked(idx);
  569. return true;
  570. },
  571. .protocol_error => |code| {
  572. _ = ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), code, "");
  573. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = code }) catch {};
  574. wsRemoveLocked(idx);
  575. return true;
  576. },
  577. }
  578. }
  579. }
  580. /// Read all available bytes from the client socket into its decoder, then
  581. /// drain. Returns true if the client was removed. Mutex held.
  582. fn wsServiceClientLocked(idx: usize) bool {
  583. const client = ws_clients.items[idx];
  584. var buf: [16384]u8 = undefined;
  585. while (true) {
  586. const rc = linux.read(client.fd, &buf, buf.len);
  587. const n: isize = @bitCast(rc);
  588. if (n > 0) {
  589. client.decoder.feed(buf[0..@intCast(n)]) catch {};
  590. continue;
  591. }
  592. if (n == 0) {
  593. // Peer closed without a close frame → 1006 abnormal closure
  594. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};
  595. wsRemoveLocked(idx);
  596. return true;
  597. }
  598. const e = -n;
  599. if (e == EINTR_ERR) continue;
  600. if (e == EAGAIN_ERR) break; // all available data consumed
  601. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};
  602. wsRemoveLocked(idx);
  603. return true;
  604. }
  605. return wsDrainDecoderLocked(idx);
  606. }
  607. /// One liveness pass over every upgraded socket. Quiet for longer than
  608. /// ws_ping_ms → send a protocol ping (RFC 6455 §5.5.2; every browser answers it
  609. /// at the protocol layer, so nothing application-side is involved). A ping that
  610. /// stands unanswered for ws_pong_ms → the peer is gone however open the socket
  611. /// looks: enqueue the SAME `close` event a real FIN would have produced and drop
  612. /// the connection. That is the whole mechanism — no separate timer, no new
  613. /// thread, no concept added above this file.
  614. fn wsSweep() void {
  615. const now = wsNowMs();
  616. mutexLock(&ws_mutex);
  617. defer mutexUnlock(&ws_mutex);
  618. var i: usize = 0;
  619. while (i < ws_clients.items.len) {
  620. const client = ws_clients.items[i];
  621. if (client.ping_at_ms != 0) {
  622. if (now - client.ping_at_ms > ws_pong_ms) {
  623. ws_event_queue.append(allocator, .{
  624. .kind = .close,
  625. .id = client.id,
  626. .code = ws.CLOSE_GOING_AWAY,
  627. }) catch {};
  628. wsRemoveLocked(i);
  629. continue;
  630. }
  631. } else if (now - client.last_seen_ms > ws_ping_ms) {
  632. if (ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .ping, "")) {
  633. client.ping_at_ms = now;
  634. } else {
  635. // the write itself failed — this one needs no grace period
  636. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};
  637. wsRemoveLocked(i);
  638. continue;
  639. }
  640. }
  641. i += 1;
  642. }
  643. }
  644. fn wsReaderLoop() void {
  645. var events: [64]linux.epoll_event = undefined;
  646. while (ws_running.load(.acquire)) {
  647. // RING THE EVENT LOOP'S BELL AT THE END OF EVERY ITERATION (mission 256),
  648. // unconditionally and without taking `ws_mutex`. Fifteen places in this
  649. // file push onto `ws_event_queue`; a bell is coalesced and a spurious one
  650. // is harmless (loop_wait.zig: "only a signal that is never sent at all
  651. // could ever be a bug"), so one ring per pass covers all of them and can
  652. // never deadlock against a path that still holds the lock. The idle cost
  653. // is two empty loop rounds a second — this epoll has a 500ms timeout.
  654. defer http.ringWake(ws_loop_wake_fd);
  655. const n_rc = linux.epoll_wait(ws_epoll_fd, &events, events.len, 500);
  656. const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));
  657. // The 500ms epoll timeout IS the sweep's clock: an idle loop still ticks,
  658. // and a busy one sweeps just as often because the check is on wall time.
  659. if (n <= 0) {
  660. wsSweep();
  661. continue;
  662. }
  663. for (events[0..@intCast(n)]) |ev| {
  664. if (ev.data.fd == ws_wake_fd) {
  665. var drain: u64 = 0;
  666. _ = linux.read(ws_wake_fd, @ptrCast(&drain), @sizeOf(u64));
  667. if (!ws_running.load(.acquire)) return;
  668. // Drain decoder-buffered data (e.g. handshake leftover)
  669. mutexLock(&ws_mutex);
  670. var i: usize = 0;
  671. while (i < ws_clients.items.len) {
  672. if (!wsDrainDecoderLocked(i)) i += 1;
  673. }
  674. mutexUnlock(&ws_mutex);
  675. continue;
  676. }
  677. mutexLock(&ws_mutex);
  678. if (wsFindByFdLocked(ev.data.fd)) |idx| {
  679. _ = wsServiceClientLocked(idx);
  680. }
  681. mutexUnlock(&ws_mutex);
  682. }
  683. wsSweep();
  684. }
  685. }
  686. /// Stop the engine: join the reader thread, close all sockets, free queues.
  687. /// Runs at interpreter teardown via the events-iterator deinit — BEFORE the
  688. /// runtime dlcloses this .so, so the thread never outlives its code.
  689. fn wsShutdown() void {
  690. if (!ws_running.swap(false, .acq_rel)) return;
  691. wsWake();
  692. if (ws_thread) |t| {
  693. t.join();
  694. ws_thread = null;
  695. }
  696. mutexLock(&ws_mutex);
  697. for (ws_clients.items) |client| {
  698. _ = ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), ws.CLOSE_GOING_AWAY, "");
  699. _ = linux.close(client.fd);
  700. client.decoder.deinit();
  701. allocator.destroy(client);
  702. }
  703. ws_clients.deinit(allocator);
  704. ws_clients = .empty;
  705. for (ws_event_queue.items) |*ev| {
  706. if (ev.data) |d| allocator.free(d);
  707. if (ev.cookie) |ck| allocator.free(ck);
  708. if (ev.host) |hc| allocator.free(hc);
  709. }
  710. ws_event_queue.deinit(allocator);
  711. ws_event_queue = .empty;
  712. mutexUnlock(&ws_mutex);
  713. if (ws_epoll_fd >= 0) _ = linux.close(ws_epoll_fd);
  714. if (ws_wake_fd >= 0) _ = linux.close(ws_wake_fd);
  715. ws_epoll_fd = -1;
  716. ws_wake_fd = -1;
  717. // `ws_loop_wake_fd` is deliberately NOT closed: the event loop may still hold
  718. // it in its epoll set, and a closed fd number gets REUSED — the loop would
  719. // then be watching whatever opened next. One eventfd per process, kept for
  720. // the process's life and reused if the engine restarts (mission 256).
  721. ws_enabled.store(false, .release);
  722. }
  723. // --- Interpreter-facing exports ------------------------------------------
  724. const ws_kind_names = [_][]const u8{ "connect", "message", "close", "pong" };
  725. fn wsEventObjDeinit(obj: *HlObject) callconv(.c) void {
  726. // fields[2] is "data", fields[5] "cookie" and fields[6] "host" — allocated
  727. // iff non-empty (empty = the static "").
  728. const s = obj.fields[2].value.data.string;
  729. if (s.len > 0) allocator.free(s.ptr[0..s.len]);
  730. const ck = obj.fields[5].value.data.string;
  731. if (ck.len > 0) allocator.free(ck.ptr[0..ck.len]);
  732. const hs = obj.fields[6].value.data.string;
  733. if (hs.len > 0) allocator.free(hs.ptr[0..hs.len]);
  734. allocator.free(obj.fields[0..obj.field_count]);
  735. allocator.destroy(obj);
  736. }
  737. /// try_next over the WS event queue → { kind, id, data, binary, code, cookie, host }
  738. /// or null. `cookie` is the handshake's Cookie header and `host` its Host header,
  739. /// both non-empty only on `connect` (mission 093, ticket #74).
  740. fn wsEventsTryNext(ctx: ?*anyopaque) callconv(.c) HlValue {
  741. _ = ctx;
  742. mutexLock(&ws_mutex);
  743. if (ws_event_queue.items.len == 0) {
  744. mutexUnlock(&ws_mutex);
  745. return api.makeNull();
  746. }
  747. const ev = ws_event_queue.orderedRemove(0);
  748. mutexUnlock(&ws_mutex);
  749. const fields = allocator.alloc(HlField, 7) catch {
  750. if (ev.data) |d| allocator.free(d);
  751. if (ev.cookie) |ck| allocator.free(ck);
  752. if (ev.host) |hc| allocator.free(hc);
  753. return api.makeNull();
  754. };
  755. const data_slice: []const u8 = if (ev.data) |d| d else "";
  756. const cookie_slice: []const u8 = if (ev.cookie) |ck| ck else "";
  757. const host_slice: []const u8 = if (ev.host) |hc| hc else "";
  758. fields[0] = .{ .key = http.hlStr("kind"), .value = api.makeString(ws_kind_names[@intFromEnum(ev.kind)]) };
  759. fields[1] = .{ .key = http.hlStr("id"), .value = api.makeNumber(@floatFromInt(ev.id)) };
  760. fields[2] = .{ .key = http.hlStr("data"), .value = api.makeString(data_slice) };
  761. fields[3] = .{ .key = http.hlStr("binary"), .value = api.makeBool(ev.is_binary) };
  762. fields[4] = .{ .key = http.hlStr("code"), .value = api.makeNumber(@floatFromInt(ev.code)) };
  763. fields[5] = .{ .key = http.hlStr("cookie"), .value = api.makeString(cookie_slice) };
  764. fields[6] = .{ .key = http.hlStr("host"), .value = api.makeString(host_slice) };
  765. const obj = allocator.create(HlObject) catch {
  766. if (ev.data) |d| allocator.free(d);
  767. if (ev.cookie) |ck| allocator.free(ck);
  768. if (ev.host) |hc| allocator.free(hc);
  769. allocator.free(fields);
  770. return api.makeNull();
  771. };
  772. obj.* = .{ .fields = fields.ptr, .field_count = 7, .deinit_fn = &wsEventObjDeinit };
  773. return api.makeObject(obj);
  774. }
  775. fn wsEventsIterDeinit(ctx: ?*anyopaque) callconv(.c) void {
  776. _ = ctx;
  777. wsShutdown();
  778. }
  779. /// __native("http1.ws_events") → event iterator; starting it enables WS
  780. /// upgrades on all plain-HTTP hl:http1 servers in this process.
  781. export fn hl_http1_ws_events(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  782. _ = argc;
  783. _ = argv;
  784. if (!wsEnsureStarted()) return api.makeNull();
  785. const iter = allocator.create(HlIterator) catch return api.makeNull();
  786. iter.* = .{
  787. .context = null,
  788. .next_fn = &wsEventsTryNext, // non-blocking either way — event loop only
  789. .deinit_fn = &wsEventsIterDeinit,
  790. .try_next_fn = &wsEventsTryNext,
  791. .wake_fd = ws_loop_wake_fd, // mission 256 — set by wsEnsureStarted above
  792. };
  793. return api.makeIterator(iter);
  794. }
  795. // --- cookie-grade random (mission 093) ------------------------------------
  796. // A session id is the ONLY thing standing between a stranger and someone else's
  797. // session, so it may not come from a seeded PRNG: `hl:math`'s random() is a
  798. // clock-seeded xoshiro, and a few of its outputs reveal its state, which would
  799. // make every other session's id derivable from one's own. This reads the
  800. // kernel's CSPRNG directly. It lives in hl:http1 because the session cookie is
  801. // an HTTP artifact and this plugin is the one that parses and sets it; there is
  802. // no other CSPRNG at the hl: level yet (named as a gap in mission 093's report).
  803. var token_buf: [128]u8 = undefined;
  804. /// __native("http1.random_token", n) → n lowercase hex chars (default 32, max 128)
  805. export fn hl_http1_random_token(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  806. var want: usize = 32;
  807. if (argc >= 1 and argv[0].type == .hl_number) {
  808. const n = argv[0].data.number;
  809. if (n >= 1 and n <= 128) want = @intFromFloat(n);
  810. }
  811. var raw: [64]u8 = undefined;
  812. const need = (want + 1) / 2;
  813. if (linux.getrandom(&raw, need, 0) != need) return api.makeNull();
  814. const hex = "0123456789abcdef";
  815. var i: usize = 0;
  816. while (i < want) : (i += 1) {
  817. const byte = raw[i / 2];
  818. const nib: u8 = if (i % 2 == 0) (byte >> 4) else (byte & 0x0f);
  819. token_buf[i] = hex[nib];
  820. }
  821. return api.makeString(token_buf[0..want]);
  822. }
  823. /// __native("http1.ws_send", id, text[, binaryFlag]) → 1 sent / 0 gone
  824. export fn hl_http1_ws_send(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  825. if (argc < 2 or argv[0].type != .hl_number or argv[1].type != .hl_string) return api.makeNumber(0);
  826. const id: u32 = @intFromFloat(argv[0].data.number);
  827. const text = argv[1].data.string.ptr[0..argv[1].data.string.len];
  828. const opcode: ws.Opcode = if (argc >= 3 and argv[2].type == .hl_bool and argv[2].data.boolean) .binary else .text;
  829. mutexLock(&ws_mutex);
  830. defer mutexUnlock(&ws_mutex);
  831. const idx = wsFindByIdLocked(id) orelse return api.makeNumber(0);
  832. const client = ws_clients.items[idx];
  833. if (!ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), opcode, text)) {
  834. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};
  835. wsRemoveLocked(idx);
  836. return api.makeNumber(0);
  837. }
  838. return api.makeNumber(1);
  839. }
  840. /// __native("http1.ws_broadcast", text) → number of clients written
  841. export fn hl_http1_ws_broadcast(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  842. if (argc < 1 or argv[0].type != .hl_string) return api.makeNumber(0);
  843. const text = argv[0].data.string.ptr[0..argv[0].data.string.len];
  844. mutexLock(&ws_mutex);
  845. defer mutexUnlock(&ws_mutex);
  846. var count: usize = 0;
  847. var i: usize = 0;
  848. while (i < ws_clients.items.len) {
  849. const client = ws_clients.items[i];
  850. if (ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .text, text)) {
  851. count += 1;
  852. i += 1;
  853. } else {
  854. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};
  855. wsRemoveLocked(i);
  856. }
  857. }
  858. return api.makeNumber(@floatFromInt(count));
  859. }
  860. /// __native("http1.ws_ping", id) → 1 sent / 0 gone
  861. export fn hl_http1_ws_ping(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  862. if (argc < 1 or argv[0].type != .hl_number) return api.makeNumber(0);
  863. const id: u32 = @intFromFloat(argv[0].data.number);
  864. mutexLock(&ws_mutex);
  865. defer mutexUnlock(&ws_mutex);
  866. const idx = wsFindByIdLocked(id) orelse return api.makeNumber(0);
  867. const client = ws_clients.items[idx];
  868. if (!ws.writeFrame(&wsFdWrite, wsFdCtx(client.fd), .ping, "")) {
  869. wsRemoveLocked(idx);
  870. return api.makeNumber(0);
  871. }
  872. return api.makeNumber(1);
  873. }
  874. /// __native("http1.ws_close", id[, code]) → 1 initiated / 0 gone.
  875. /// Sends the close frame and waits for the peer echo (reader completes it).
  876. export fn hl_http1_ws_close(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  877. if (argc < 1 or argv[0].type != .hl_number) return api.makeNumber(0);
  878. const id: u32 = @intFromFloat(argv[0].data.number);
  879. const code: u16 = if (argc >= 2 and argv[1].type == .hl_number) @intFromFloat(argv[1].data.number) else ws.CLOSE_NORMAL;
  880. mutexLock(&ws_mutex);
  881. defer mutexUnlock(&ws_mutex);
  882. const idx = wsFindByIdLocked(id) orelse return api.makeNumber(0);
  883. const client = ws_clients.items[idx];
  884. if (!ws.writeClose(&wsFdWrite, wsFdCtx(client.fd), code, "")) {
  885. ws_event_queue.append(allocator, .{ .kind = .close, .id = client.id, .code = 1006 }) catch {};
  886. wsRemoveLocked(idx);
  887. return api.makeNumber(0);
  888. }
  889. client.closing = true;
  890. return api.makeNumber(1);
  891. }
  892. // =========================================================================
  893. // TLS context — OpenSSL via dlopen (optional, no compile-time dep)
  894. // =========================================================================
  895. const c_dlfcn = @cImport({
  896. @cInclude("dlfcn.h");
  897. });
  898. const SSL_CTX = opaque {};
  899. const SSL = opaque {};
  900. const SSL_METHOD = opaque {};
  901. // OpenSSL function pointer types
  902. const SSL_library_init_fn = *const fn () callconv(.c) c_int;
  903. const SSL_load_error_strings_fn = *const fn () callconv(.c) void;
  904. const TLS_server_method_fn = *const fn () callconv(.c) ?*const SSL_METHOD;
  905. const SSL_CTX_new_fn = *const fn (?*const SSL_METHOD) callconv(.c) ?*SSL_CTX;
  906. const SSL_CTX_free_fn = *const fn (?*SSL_CTX) callconv(.c) void;
  907. const SSL_CTX_use_certificate_chain_file_fn = *const fn (?*SSL_CTX, [*:0]const u8) callconv(.c) c_int;
  908. const SSL_CTX_use_PrivateKey_file_fn = *const fn (?*SSL_CTX, [*:0]const u8, c_int) callconv(.c) c_int;
  909. const SSL_new_fn = *const fn (?*SSL_CTX) callconv(.c) ?*SSL;
  910. const SSL_free_fn = *const fn (?*SSL) callconv(.c) void;
  911. const SSL_set_fd_fn = *const fn (?*SSL, c_int) callconv(.c) c_int;
  912. const SSL_accept_fn = *const fn (?*SSL) callconv(.c) c_int;
  913. const SSL_read_fn = *const fn (?*SSL, [*]u8, c_int) callconv(.c) c_int;
  914. const SSL_write_fn = *const fn (?*SSL, [*]const u8, c_int) callconv(.c) c_int;
  915. const SSL_shutdown_fn = *const fn (?*SSL) callconv(.c) c_int;
  916. const SSL_get_error_fn = *const fn (?*SSL, c_int) callconv(.c) c_int;
  917. const OPENSSL_init_ssl_fn = *const fn (u64, ?*anyopaque) callconv(.c) c_int;
  918. const SSL_FILETYPE_PEM: c_int = 1;
  919. const TlsContext = struct {
  920. ssl_ctx: ?*SSL_CTX = null,
  921. lib_ssl: ?*anyopaque = null,
  922. lib_crypto: ?*anyopaque = null,
  923. // Function pointers
  924. fn_ssl_ctx_new: ?SSL_CTX_new_fn = null,
  925. fn_ssl_ctx_free: ?SSL_CTX_free_fn = null,
  926. fn_ssl_ctx_use_cert: ?SSL_CTX_use_certificate_chain_file_fn = null,
  927. fn_ssl_ctx_use_key: ?SSL_CTX_use_PrivateKey_file_fn = null,
  928. fn_ssl_new: ?SSL_new_fn = null,
  929. fn_ssl_free: ?SSL_free_fn = null,
  930. fn_ssl_set_fd: ?SSL_set_fd_fn = null,
  931. fn_ssl_accept: ?SSL_accept_fn = null,
  932. fn_ssl_read: ?SSL_read_fn = null,
  933. fn_ssl_write: ?SSL_write_fn = null,
  934. fn_ssl_shutdown: ?SSL_shutdown_fn = null,
  935. fn_ssl_get_error: ?SSL_get_error_fn = null,
  936. fn loadSym(lib: ?*anyopaque, comptime T: type, name: [*:0]const u8) ?T {
  937. const sym = c_dlfcn.dlsym(lib, name) orelse return null;
  938. return @ptrCast(sym);
  939. }
  940. fn init(cert_path: []const u8, key_path: []const u8) ?TlsContext {
  941. var ctx = TlsContext{};
  942. // Try loading libssl and libcrypto
  943. const ssl_paths = [_][*:0]const u8{ "libssl.so.3", "libssl.so.1.1", "libssl.so" };
  944. const crypto_paths = [_][*:0]const u8{ "libcrypto.so.3", "libcrypto.so.1.1", "libcrypto.so" };
  945. for (ssl_paths) |path| {
  946. ctx.lib_ssl = c_dlfcn.dlopen(path, c_dlfcn.RTLD_NOW | c_dlfcn.RTLD_LOCAL);
  947. if (ctx.lib_ssl != null) break;
  948. }
  949. if (ctx.lib_ssl == null) {
  950. logMsg("http1: TLS: failed to load libssl.so\n");
  951. return null;
  952. }
  953. for (crypto_paths) |path| {
  954. ctx.lib_crypto = c_dlfcn.dlopen(path, c_dlfcn.RTLD_NOW | c_dlfcn.RTLD_LOCAL);
  955. if (ctx.lib_crypto != null) break;
  956. }
  957. if (ctx.lib_crypto == null) {
  958. logMsg("http1: TLS: failed to load libcrypto.so\n");
  959. _ = c_dlfcn.dlclose(ctx.lib_ssl);
  960. return null;
  961. }
  962. // Load function pointers
  963. // Try OPENSSL_init_ssl first (OpenSSL 1.1+), fall back to SSL_library_init
  964. if (loadSym(ctx.lib_ssl, OPENSSL_init_ssl_fn, "OPENSSL_init_ssl")) |init_fn| {
  965. _ = init_fn(0, null);
  966. } else if (loadSym(ctx.lib_ssl, SSL_library_init_fn, "SSL_library_init")) |lib_init| {
  967. _ = lib_init();
  968. if (loadSym(ctx.lib_ssl, SSL_load_error_strings_fn, "SSL_load_error_strings")) |load_err| {
  969. load_err();
  970. }
  971. }
  972. const method_fn = loadSym(ctx.lib_ssl, TLS_server_method_fn, "TLS_server_method") orelse {
  973. logMsg("http1: TLS: TLS_server_method not found\n");
  974. ctx.deinit();
  975. return null;
  976. };
  977. ctx.fn_ssl_ctx_new = loadSym(ctx.lib_ssl, SSL_CTX_new_fn, "SSL_CTX_new");
  978. ctx.fn_ssl_ctx_free = loadSym(ctx.lib_ssl, SSL_CTX_free_fn, "SSL_CTX_free");
  979. ctx.fn_ssl_ctx_use_cert = loadSym(ctx.lib_ssl, SSL_CTX_use_certificate_chain_file_fn, "SSL_CTX_use_certificate_chain_file");
  980. ctx.fn_ssl_ctx_use_key = loadSym(ctx.lib_ssl, SSL_CTX_use_PrivateKey_file_fn, "SSL_CTX_use_PrivateKey_file");
  981. ctx.fn_ssl_new = loadSym(ctx.lib_ssl, SSL_new_fn, "SSL_new");
  982. ctx.fn_ssl_free = loadSym(ctx.lib_ssl, SSL_free_fn, "SSL_free");
  983. ctx.fn_ssl_set_fd = loadSym(ctx.lib_ssl, SSL_set_fd_fn, "SSL_set_fd");
  984. ctx.fn_ssl_accept = loadSym(ctx.lib_ssl, SSL_accept_fn, "SSL_accept");
  985. ctx.fn_ssl_read = loadSym(ctx.lib_ssl, SSL_read_fn, "SSL_read");
  986. ctx.fn_ssl_write = loadSym(ctx.lib_ssl, SSL_write_fn, "SSL_write");
  987. ctx.fn_ssl_shutdown = loadSym(ctx.lib_ssl, SSL_shutdown_fn, "SSL_shutdown");
  988. ctx.fn_ssl_get_error = loadSym(ctx.lib_ssl, SSL_get_error_fn, "SSL_get_error");
  989. if (ctx.fn_ssl_ctx_new == null or ctx.fn_ssl_new == null or
  990. ctx.fn_ssl_set_fd == null or ctx.fn_ssl_accept == null or
  991. ctx.fn_ssl_read == null or ctx.fn_ssl_write == null)
  992. {
  993. logMsg("http1: TLS: missing required SSL symbols\n");
  994. ctx.deinit();
  995. return null;
  996. }
  997. // Create SSL_CTX
  998. const method = method_fn();
  999. ctx.ssl_ctx = ctx.fn_ssl_ctx_new.?(method);
  1000. if (ctx.ssl_ctx == null) {
  1001. logMsg("http1: TLS: SSL_CTX_new failed\n");
  1002. ctx.deinit();
  1003. return null;
  1004. }
  1005. // Load cert and key
  1006. const cert_z = allocator.dupeZ(u8, cert_path) catch {
  1007. ctx.deinit();
  1008. return null;
  1009. };
  1010. defer allocator.free(cert_z);
  1011. const key_z = allocator.dupeZ(u8, key_path) catch {
  1012. ctx.deinit();
  1013. return null;
  1014. };
  1015. defer allocator.free(key_z);
  1016. if (ctx.fn_ssl_ctx_use_cert) |use_cert| {
  1017. if (use_cert(ctx.ssl_ctx, cert_z.ptr) != 1) {
  1018. logFmt("http1: TLS: failed to load certificate: {s}\n", .{cert_path});
  1019. ctx.deinit();
  1020. return null;
  1021. }
  1022. }
  1023. if (ctx.fn_ssl_ctx_use_key) |use_key| {
  1024. if (use_key(ctx.ssl_ctx, key_z.ptr, SSL_FILETYPE_PEM) != 1) {
  1025. logFmt("http1: TLS: failed to load private key: {s}\n", .{key_path});
  1026. ctx.deinit();
  1027. return null;
  1028. }
  1029. }
  1030. logMsg("http1: TLS initialized\n");
  1031. return ctx;
  1032. }
  1033. fn wrapConnection(self: *const TlsContext, fd: i32) ?*SSL {
  1034. const ssl = self.fn_ssl_new.?(self.ssl_ctx);
  1035. if (ssl == null) return null;
  1036. _ = self.fn_ssl_set_fd.?(ssl, fd);
  1037. const ret = self.fn_ssl_accept.?(ssl);
  1038. if (ret != 1) {
  1039. self.fn_ssl_free.?(ssl);
  1040. return null;
  1041. }
  1042. return ssl;
  1043. }
  1044. fn sslRead(self: *const TlsContext, ssl: *SSL, buf: []u8) isize {
  1045. const ret = self.fn_ssl_read.?(ssl, buf.ptr, @intCast(@min(buf.len, std.math.maxInt(c_int))));
  1046. if (ret <= 0) return 0;
  1047. return @intCast(ret);
  1048. }
  1049. fn sslWrite(self: *const TlsContext, ssl: *SSL, data: []const u8) isize {
  1050. var written: usize = 0;
  1051. while (written < data.len) {
  1052. const chunk_len: c_int = @intCast(@min(data.len - written, std.math.maxInt(c_int)));
  1053. const ret = self.fn_ssl_write.?(ssl, data[written..].ptr, chunk_len);
  1054. if (ret <= 0) return @intCast(written);
  1055. written += @intCast(ret);
  1056. }
  1057. return @intCast(written);
  1058. }
  1059. fn sslShutdown(self: *const TlsContext, ssl: *SSL) void {
  1060. _ = self.fn_ssl_shutdown.?(ssl);
  1061. self.fn_ssl_free.?(ssl);
  1062. }
  1063. fn deinit(self: *TlsContext) void {
  1064. if (self.ssl_ctx != null) {
  1065. if (self.fn_ssl_ctx_free) |free_fn| {
  1066. free_fn(self.ssl_ctx);
  1067. }
  1068. self.ssl_ctx = null;
  1069. }
  1070. if (self.lib_ssl != null) {
  1071. _ = c_dlfcn.dlclose(self.lib_ssl);
  1072. self.lib_ssl = null;
  1073. }
  1074. if (self.lib_crypto != null) {
  1075. _ = c_dlfcn.dlclose(self.lib_crypto);
  1076. self.lib_crypto = null;
  1077. }
  1078. }
  1079. };
  1080. // =========================================================================
  1081. // MPSC Request Queue — thread-safe, blocks on dequeue
  1082. // =========================================================================
  1083. const ParsedRequest = struct {
  1084. method: []const u8,
  1085. path: []const u8, // URL-decoded
  1086. query_params: []http.QueryParam, // parsed key-value pairs
  1087. query_raw: []const u8, // raw query string
  1088. headers: std.ArrayListUnmanaged(HeaderPair),
  1089. body: []const u8,
  1090. client_fd: i32,
  1091. ssl: ?*SSL, // null if plain HTTP
  1092. keep_alive: bool,
  1093. };
  1094. const HeaderPair = struct {
  1095. key: []const u8,
  1096. value: []const u8,
  1097. };
  1098. const RequestQueue = struct {
  1099. queue: std.ArrayListUnmanaged(ParsedRequest),
  1100. mutex: PthreadMutex,
  1101. condvar: PthreadCond,
  1102. shutdown: bool,
  1103. /// THE EVENT LOOP'S BELL (mission 256). The condvar above wakes a `for (req of
  1104. /// server)` consumer blocked in `dequeue`; the interpreter's event loop is
  1105. /// NOT that consumer — it polls `tryDequeue` from a loop that also watches
  1106. /// every other source, so it cannot block on a condvar belonging to one of
  1107. /// them. An eventfd it can: `native/src/loop_wait.zig` epolls this fd, and an
  1108. /// idle server costs nothing instead of 1000 empty poll rounds a second.
  1109. loop_wake_fd: i32,
  1110. fn init() RequestQueue {
  1111. var self: RequestQueue = .{
  1112. .queue = .empty,
  1113. .mutex = c.PTHREAD_MUTEX_INITIALIZER,
  1114. .condvar = c.PTHREAD_COND_INITIALIZER,
  1115. .shutdown = false,
  1116. .loop_wake_fd = http.makeWakeFd(),
  1117. };
  1118. mutexInit(&self.mutex);
  1119. condInit(&self.condvar);
  1120. return self;
  1121. }
  1122. fn enqueue(self: *RequestQueue, req: ParsedRequest) void {
  1123. mutexLock(&self.mutex);
  1124. defer mutexUnlock(&self.mutex);
  1125. self.queue.append(allocator, req) catch return;
  1126. condSignal(&self.condvar);
  1127. // Rung INSIDE the lock, so the fd's counter is already up by the time the
  1128. // request is visible to `tryDequeue`. A loop that is between its last poll
  1129. // and its next `epoll_wait` therefore finds the level raised and returns
  1130. // at once instead of sleeping on a queue that has work in it.
  1131. http.ringWake(self.loop_wake_fd);
  1132. }
  1133. fn dequeue(self: *RequestQueue) ?ParsedRequest {
  1134. mutexLock(&self.mutex);
  1135. defer mutexUnlock(&self.mutex);
  1136. while (self.queue.items.len == 0 and !self.shutdown) {
  1137. condWait(&self.condvar, &self.mutex);
  1138. }
  1139. if (self.shutdown and self.queue.items.len == 0) return null;
  1140. return self.queue.orderedRemove(0);
  1141. }
  1142. /// Non-blocking: returns null immediately if queue is empty
  1143. fn tryDequeue(self: *RequestQueue) ?ParsedRequest {
  1144. mutexLock(&self.mutex);
  1145. defer mutexUnlock(&self.mutex);
  1146. if (self.queue.items.len == 0) return null;
  1147. return self.queue.orderedRemove(0);
  1148. }
  1149. fn signalShutdown(self: *RequestQueue) void {
  1150. mutexLock(&self.mutex);
  1151. defer mutexUnlock(&self.mutex);
  1152. self.shutdown = true;
  1153. condBroadcast(&self.condvar);
  1154. // The event loop is told too — it is blocked on this fd and its `while`
  1155. // condition (has every source retired?) can only be re-read on the way out.
  1156. http.ringWake(self.loop_wake_fd);
  1157. }
  1158. fn deinit(self: *RequestQueue) void {
  1159. // Free any remaining queued requests
  1160. for (self.queue.items) |*req| {
  1161. freeRequest(req);
  1162. }
  1163. self.queue.deinit(allocator);
  1164. if (self.loop_wake_fd >= 0) {
  1165. _ = linux.close(self.loop_wake_fd);
  1166. self.loop_wake_fd = -1;
  1167. }
  1168. _ = c.pthread_cond_destroy(&self.condvar);
  1169. _ = c.pthread_mutex_destroy(&self.mutex);
  1170. }
  1171. };
  1172. fn freeRequest(req: *ParsedRequest) void {
  1173. allocator.free(req.method);
  1174. allocator.free(req.path);
  1175. allocator.free(req.query_raw);
  1176. for (req.query_params) |param| {
  1177. allocator.free(param.key);
  1178. allocator.free(param.value);
  1179. }
  1180. allocator.free(req.query_params);
  1181. for (req.headers.items) |hdr| {
  1182. allocator.free(hdr.key);
  1183. allocator.free(hdr.value);
  1184. }
  1185. req.headers.deinit(allocator);
  1186. if (req.body.len > 0) allocator.free(req.body);
  1187. }
  1188. // =========================================================================
  1189. // Connection — per-connection state tracked by I/O workers
  1190. // =========================================================================
  1191. const Connection = struct {
  1192. fd: i32,
  1193. ssl: ?*SSL,
  1194. keep_alive: bool,
  1195. request_count: u32,
  1196. max_requests: u32,
  1197. fn read(self: *const Connection, tls: ?*const TlsContext, buf: []u8) usize {
  1198. if (self.ssl) |ssl_ptr| {
  1199. if (tls) |t| {
  1200. const n = t.sslRead(ssl_ptr, buf);
  1201. if (n <= 0) return 0;
  1202. return @intCast(n);
  1203. }
  1204. return 0;
  1205. }
  1206. return sysRead(self.fd, buf);
  1207. }
  1208. fn write(self: *const Connection, tls: ?*const TlsContext, data: []const u8) void {
  1209. if (self.ssl) |ssl_ptr| {
  1210. if (tls) |t| {
  1211. _ = t.sslWrite(ssl_ptr, data);
  1212. return;
  1213. }
  1214. }
  1215. sysWriteAll(self.fd, data);
  1216. }
  1217. fn close(self: *Connection, tls: ?*const TlsContext) void {
  1218. if (self.ssl) |ssl_ptr| {
  1219. if (tls) |t| {
  1220. t.sslShutdown(ssl_ptr);
  1221. }
  1222. self.ssl = null;
  1223. }
  1224. _ = linux.close(self.fd);
  1225. }
  1226. };
  1227. // =========================================================================
  1228. // ServerCore — owns server socket, epoll, threads, queue, optional TLS
  1229. // =========================================================================
  1230. const MAX_KEEPALIVE_REQUESTS: u32 = 100;
  1231. const MAX_IO_THREADS: u32 = 16;
  1232. const ServerCore = struct {
  1233. port: u16 = 0,
  1234. server_fd: i32,
  1235. epoll_fd: i32,
  1236. shutdown_fd: i32, // eventfd for shutdown signal
  1237. tls: ?TlsContext,
  1238. request_queue: RequestQueue,
  1239. running: std.atomic.Value(bool),
  1240. acceptor_thread: ?std.Thread,
  1241. io_threads: []std.Thread,
  1242. num_threads: u32,
  1243. // Connection tracking for I/O threads
  1244. conn_queue: ConnectionQueue,
  1245. /// idle connections wait here instead of inside a blocking worker read (mission 084)
  1246. parked: ParkedConns,
  1247. fn create(port: u16, host_str: []const u8, cert_path: []const u8, key_path: []const u8, num_threads: u32) ?*ServerCore {
  1248. const core = allocator.create(ServerCore) catch return null;
  1249. core.* = .{
  1250. .server_fd = -1,
  1251. .epoll_fd = -1,
  1252. .shutdown_fd = -1,
  1253. .tls = null,
  1254. .request_queue = RequestQueue.init(),
  1255. .running = std.atomic.Value(bool).init(false),
  1256. .acceptor_thread = null,
  1257. .io_threads = &[_]std.Thread{},
  1258. .num_threads = @min(num_threads, MAX_IO_THREADS),
  1259. .conn_queue = ConnectionQueue.init(),
  1260. .parked = ParkedConns.init(),
  1261. };
  1262. // Create TCP socket. NONBLOCK IS NOT DECORATION (mission 260): the acceptor
  1263. // drains with `while (true) accept4(...)` until the call fails, and on a
  1264. // BLOCKING listener that last call does not fail — it sleeps in the kernel
  1265. // (`inet_csk_accept`) until the next client arrives. The thread then never
  1266. // re-reads `core.running`, so `shutdown()`'s `join(acceptor)` waits forever
  1267. // and `close()` never returns. Measured: after ONE connection the process was
  1268. // down to two threads, main in `__futex_wait` (the join) and the acceptor in
  1269. // `inet_csk_accept`; each extra TCP connect added exactly one fd and put the
  1270. // acceptor straight back into `inet_csk_accept`. hl:http2's listener
  1271. // (`http2.zig:926`) has always carried NONBLOCK — http1 was the outlier.
  1272. // The accept flags below apply to the ACCEPTED socket, never to this one.
  1273. const sock_rc = linux.socket(linux.AF.INET, linux.SOCK.STREAM | linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK, 0);
  1274. const server_fd: i32 = @bitCast(@as(u32, @truncate(sock_rc)));
  1275. if (server_fd < 0) {
  1276. logMsg("http1: socket() failed\n");
  1277. allocator.destroy(core);
  1278. return null;
  1279. }
  1280. core.server_fd = server_fd;
  1281. // SO_REUSEADDR
  1282. const one: i32 = 1;
  1283. _ = linux.setsockopt(server_fd, linux.SOL.SOCKET, linux.SO.REUSEADDR, @ptrCast(&one), @sizeOf(i32));
  1284. // Parse host address
  1285. var addr_val: u32 = 0; // INADDR_ANY
  1286. if (host_str.len > 0 and !std.mem.eql(u8, host_str, "0.0.0.0")) {
  1287. addr_val = parseIPv4(host_str) orelse 0;
  1288. }
  1289. // Bind
  1290. const addr = linux.sockaddr.in{
  1291. .port = std.mem.nativeToBig(u16, port),
  1292. .addr = addr_val,
  1293. };
  1294. const bind_rc = linux.bind(server_fd, @ptrCast(&addr), @sizeOf(linux.sockaddr.in));
  1295. const bind_err: i32 = @bitCast(@as(u32, @truncate(bind_rc)));
  1296. if (bind_err < 0) {
  1297. // The prefix is LOAD-BEARING: `tests/browser/fixtures.mjs` breaks its
  1298. // readiness wait on `bind() failed on port N` (mission 285). The errno
  1299. // and the holder are appended to it, never in front of it.
  1300. var why: [320]u8 = undefined;
  1301. logFmt("http1: bind() failed on port {d}: {s}\n", .{ port, http.bindFailureDetail(&why, bind_rc, port, false) });
  1302. core.destroy();
  1303. return null;
  1304. }
  1305. const listen_rc = linux.listen(server_fd, 128);
  1306. const listen_err: i32 = @bitCast(@as(u32, @truncate(listen_rc)));
  1307. if (listen_err < 0) {
  1308. logMsg("http1: listen() failed\n");
  1309. core.destroy();
  1310. return null;
  1311. }
  1312. // Create eventfd for shutdown signaling
  1313. const efd_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);
  1314. const shutdown_fd: i32 = @bitCast(@as(u32, @truncate(efd_rc)));
  1315. if (shutdown_fd < 0) {
  1316. logMsg("http1: eventfd() failed\n");
  1317. core.destroy();
  1318. return null;
  1319. }
  1320. core.shutdown_fd = shutdown_fd;
  1321. // Create epoll instance
  1322. const epoll_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);
  1323. const epoll_fd: i32 = @bitCast(@as(u32, @truncate(epoll_rc)));
  1324. if (epoll_fd < 0) {
  1325. logMsg("http1: epoll_create1() failed\n");
  1326. core.destroy();
  1327. return null;
  1328. }
  1329. core.epoll_fd = epoll_fd;
  1330. // Add server_fd to epoll
  1331. var ev = linux.epoll_event{
  1332. .events = linux.EPOLL.IN,
  1333. .data = .{ .fd = server_fd },
  1334. };
  1335. _ = linux.epoll_ctl(epoll_fd, linux.EPOLL.CTL_ADD, server_fd, &ev);
  1336. // Add shutdown_fd to epoll
  1337. var shutdown_ev = linux.epoll_event{
  1338. .events = linux.EPOLL.IN,
  1339. .data = .{ .fd = shutdown_fd },
  1340. };
  1341. _ = linux.epoll_ctl(epoll_fd, linux.EPOLL.CTL_ADD, shutdown_fd, &shutdown_ev);
  1342. // Initialize TLS if cert+key provided
  1343. if (cert_path.len > 0 and key_path.len > 0) {
  1344. core.tls = TlsContext.init(cert_path, key_path);
  1345. if (core.tls == null) {
  1346. logMsg("http1: TLS initialization failed, falling back to plain HTTP\n");
  1347. }
  1348. }
  1349. logFmt("http1: listening on :{d}{s}\n", .{ port, if (core.tls != null) " (TLS)" else "" });
  1350. // Start threads
  1351. core.running.store(true, .release);
  1352. // Parker thread first: the acceptor parks into it from its very first accept.
  1353. // If it cannot start we log and keep going — connections then go straight to the
  1354. // pool, which is the pre-084 (starvable) behaviour rather than an outage.
  1355. if (!core.parked.start(core)) {
  1356. logMsg("http1: parker thread unavailable — idle keep-alive connections will hold I/O workers\n");
  1357. }
  1358. // Allocate I/O threads
  1359. const threads = allocator.alloc(std.Thread, core.num_threads) catch {
  1360. core.destroy();
  1361. return null;
  1362. };
  1363. core.io_threads = threads;
  1364. for (0..core.num_threads) |i| {
  1365. core.io_threads[i] = std.Thread.spawn(.{}, ioWorker, .{core}) catch {
  1366. logFmt("http1: failed to spawn I/O thread {d}\n", .{i});
  1367. core.num_threads = @intCast(i);
  1368. core.io_threads = core.io_threads[0..i];
  1369. break;
  1370. };
  1371. }
  1372. // Start acceptor thread
  1373. core.acceptor_thread = std.Thread.spawn(.{}, acceptorLoop, .{core}) catch {
  1374. logMsg("http1: failed to spawn acceptor thread\n");
  1375. core.shutdown();
  1376. core.destroy();
  1377. return null;
  1378. };
  1379. core.port = port;
  1380. registerCore(core);
  1381. return core;
  1382. }
  1383. fn shutdown(self: *ServerCore) void {
  1384. unregisterCore(self);
  1385. if (!self.running.swap(false, .acq_rel)) return;
  1386. // Signal shutdown via eventfd
  1387. if (self.shutdown_fd >= 0) {
  1388. const val: u64 = 1;
  1389. _ = linux.write(self.shutdown_fd, @ptrCast(&val), @sizeOf(u64));
  1390. }
  1391. // Wake up the request queue so dequeue() unblocks
  1392. self.request_queue.signalShutdown();
  1393. // Stop parking before the workers: the parker must not push new work into a
  1394. // queue whose consumers are being torn down.
  1395. self.parked.stop();
  1396. // Signal connection queue to wake I/O workers
  1397. self.conn_queue.signalShutdown();
  1398. // Join acceptor thread
  1399. if (self.acceptor_thread) |t| {
  1400. t.join();
  1401. self.acceptor_thread = null;
  1402. }
  1403. // Join I/O threads
  1404. for (self.io_threads) |t| {
  1405. t.join();
  1406. }
  1407. }
  1408. fn destroy(self: *ServerCore) void {
  1409. self.shutdown();
  1410. if (self.epoll_fd >= 0) _ = linux.close(self.epoll_fd);
  1411. if (self.shutdown_fd >= 0) _ = linux.close(self.shutdown_fd);
  1412. if (self.server_fd >= 0) _ = linux.close(self.server_fd);
  1413. if (self.tls) |*tls| tls.deinit();
  1414. // Close every still-parked idle connection, then drain conn_queue
  1415. self.parked.deinit(if (self.tls) |*t| t else null);
  1416. self.conn_queue.deinit(if (self.tls) |*t| t else null);
  1417. self.request_queue.deinit();
  1418. if (self.io_threads.len > 0) allocator.free(self.io_threads);
  1419. allocator.destroy(self);
  1420. }
  1421. };
  1422. // =========================================================================
  1423. // Connection Queue — MPSC queue for acceptor → I/O workers
  1424. // =========================================================================
  1425. const ConnectionQueue = struct {
  1426. queue: std.ArrayListUnmanaged(Connection),
  1427. mutex: PthreadMutex,
  1428. condvar: PthreadCond,
  1429. shutdown: bool,
  1430. fn init() ConnectionQueue {
  1431. var self: ConnectionQueue = .{
  1432. .queue = .empty,
  1433. .mutex = c.PTHREAD_MUTEX_INITIALIZER,
  1434. .condvar = c.PTHREAD_COND_INITIALIZER,
  1435. .shutdown = false,
  1436. };
  1437. mutexInit(&self.mutex);
  1438. condInit(&self.condvar);
  1439. return self;
  1440. }
  1441. fn enqueue(self: *ConnectionQueue, conn: Connection) void {
  1442. mutexLock(&self.mutex);
  1443. defer mutexUnlock(&self.mutex);
  1444. self.queue.append(allocator, conn) catch return;
  1445. condSignal(&self.condvar);
  1446. }
  1447. fn dequeue(self: *ConnectionQueue) ?Connection {
  1448. mutexLock(&self.mutex);
  1449. defer mutexUnlock(&self.mutex);
  1450. while (self.queue.items.len == 0 and !self.shutdown) {
  1451. condWait(&self.condvar, &self.mutex);
  1452. }
  1453. if (self.shutdown and self.queue.items.len == 0) return null;
  1454. return self.queue.orderedRemove(0);
  1455. }
  1456. fn signalShutdown(self: *ConnectionQueue) void {
  1457. mutexLock(&self.mutex);
  1458. defer mutexUnlock(&self.mutex);
  1459. self.shutdown = true;
  1460. condBroadcast(&self.condvar);
  1461. }
  1462. /// Close everything still queued and stay a VALID, EMPTY queue — `deinit`
  1463. /// leaves the list `undefined`, which is only safe on a core that is being
  1464. /// freed, and a closed core deliberately is not (see `hl_http1_close`).
  1465. fn closeAll(self: *ConnectionQueue, tls: ?*TlsContext) void {
  1466. mutexLock(&self.mutex);
  1467. defer mutexUnlock(&self.mutex);
  1468. for (self.queue.items) |*conn| conn.close(tls);
  1469. self.queue.clearRetainingCapacity();
  1470. }
  1471. fn deinit(self: *ConnectionQueue, tls: ?*TlsContext) void {
  1472. for (self.queue.items) |*conn| {
  1473. conn.close(tls);
  1474. }
  1475. self.queue.deinit(allocator);
  1476. }
  1477. };
  1478. // =========================================================================
  1479. // ParkedConns — idle connections wait HERE, not in an I/O worker (mission 084)
  1480. // =========================================================================
  1481. //
  1482. // The starvation this fixes, measured: a `threads = 4` server, four keep-alive
  1483. // connections that have gone quiet, and every subsequent request times out — each idle
  1484. // socket sat inside a worker's blocking read(). Parking inverts that: an idle connection
  1485. // costs one epoll registration and ZERO threads, and a worker only ever picks up a
  1486. // connection that already has bytes waiting (or has hung up, which it reads as EOF and
  1487. // closes). Connections are parked from the acceptor (a fresh socket may be silent — a
  1488. // pre-connected browser socket routinely is) and after every keep-alive response.
  1489. //
  1490. // TLS caveat: an ESTABLISHED TLS connection is never parked. OpenSSL may hold already-
  1491. // decrypted plaintext in its own buffer, which epoll on the raw fd cannot see, so parking
  1492. // it could hang a live request. `ssl == null` covers all plain HTTP plus the pre-handshake
  1493. // TLS socket (the ClientHello does arrive on the raw fd), which is what the park path takes.
  1494. /// How long a parked, silent connection is kept before it is closed. It costs no thread,
  1495. /// only an fd, so this is generous compared to the old in-worker 30s SO_RCVTIMEO.
  1496. const PARK_IDLE_TIMEOUT_MS: i64 = 60_000;
  1497. /// Milliseconds on CLOCK_MONOTONIC. This zig's `std.time` exposes no timestamp function,
  1498. /// and monotonic is the right clock anyway — a wall-clock step must not expire a live
  1499. /// connection early or keep a dead one parked.
  1500. fn monotonicMs() i64 {
  1501. var ts: linux.timespec = undefined;
  1502. if (linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts) != 0) return 0;
  1503. return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);
  1504. }
  1505. const ParkedConn = struct {
  1506. conn: Connection,
  1507. /// monotonic ms after which this silent connection is closed
  1508. deadline_ms: i64,
  1509. };
  1510. const ParkedConns = struct {
  1511. epoll_fd: i32,
  1512. wake_fd: i32,
  1513. thread: ?std.Thread,
  1514. running: std.atomic.Value(bool),
  1515. mutex: PthreadMutex,
  1516. /// fd → parked connection. Guarded by `mutex`; the parker thread is the only
  1517. /// consumer, park() the only producer, so a plain map is enough.
  1518. map: std.AutoHashMapUnmanaged(i32, ParkedConn),
  1519. fn init() ParkedConns {
  1520. var self: ParkedConns = .{
  1521. .epoll_fd = -1,
  1522. .wake_fd = -1,
  1523. .thread = null,
  1524. .running = std.atomic.Value(bool).init(false),
  1525. .mutex = c.PTHREAD_MUTEX_INITIALIZER,
  1526. .map = .empty,
  1527. };
  1528. mutexInit(&self.mutex);
  1529. return self;
  1530. }
  1531. /// Create the epoll instance + wake eventfd and spawn the parker thread.
  1532. /// Returns false if the kernel objects could not be made — the caller then falls
  1533. /// back to handing connections straight to the pool (old behaviour, still correct,
  1534. /// just starvable).
  1535. fn start(self: *ParkedConns, core: *ServerCore) bool {
  1536. const ep_rc = linux.epoll_create1(linux.EPOLL.CLOEXEC);
  1537. const ep: i32 = @bitCast(@as(u32, @truncate(ep_rc)));
  1538. if (ep < 0) return false;
  1539. const ef_rc = linux.eventfd(0, linux.EFD.CLOEXEC | linux.EFD.NONBLOCK);
  1540. const ef: i32 = @bitCast(@as(u32, @truncate(ef_rc)));
  1541. if (ef < 0) {
  1542. _ = linux.close(ep);
  1543. return false;
  1544. }
  1545. var wake_ev = linux.epoll_event{ .events = linux.EPOLL.IN, .data = .{ .fd = ef } };
  1546. _ = linux.epoll_ctl(ep, linux.EPOLL.CTL_ADD, ef, &wake_ev);
  1547. self.epoll_fd = ep;
  1548. self.wake_fd = ef;
  1549. self.running.store(true, .release);
  1550. self.thread = std.Thread.spawn(.{}, parkerLoop, .{core}) catch {
  1551. self.running.store(false, .release);
  1552. _ = linux.close(ep);
  1553. _ = linux.close(ef);
  1554. self.epoll_fd = -1;
  1555. self.wake_fd = -1;
  1556. return false;
  1557. };
  1558. return true;
  1559. }
  1560. /// Hand a connection to the parker. The map insert happens BEFORE the epoll ADD so
  1561. /// the parker can never see a readable fd it has no entry for.
  1562. fn park(self: *ParkedConns, conn: Connection) bool {
  1563. if (self.epoll_fd < 0 or !self.running.load(.acquire)) return false;
  1564. mutexLock(&self.mutex);
  1565. self.map.put(allocator, conn.fd, .{
  1566. .conn = conn,
  1567. .deadline_ms = monotonicMs() + PARK_IDLE_TIMEOUT_MS,
  1568. }) catch {
  1569. mutexUnlock(&self.mutex);
  1570. return false;
  1571. };
  1572. var ev = linux.epoll_event{
  1573. .events = linux.EPOLL.IN | linux.EPOLL.RDHUP,
  1574. .data = .{ .fd = conn.fd },
  1575. };
  1576. const rc = linux.epoll_ctl(self.epoll_fd, linux.EPOLL.CTL_ADD, conn.fd, &ev);
  1577. const err: i32 = @bitCast(@as(u32, @truncate(rc)));
  1578. if (err < 0) {
  1579. _ = self.map.remove(conn.fd);
  1580. mutexUnlock(&self.mutex);
  1581. return false;
  1582. }
  1583. mutexUnlock(&self.mutex);
  1584. return true;
  1585. }
  1586. /// Take a parked connection off the epoll set. Returns it if we still owned it.
  1587. fn take(self: *ParkedConns, fd: i32) ?Connection {
  1588. mutexLock(&self.mutex);
  1589. defer mutexUnlock(&self.mutex);
  1590. const entry = self.map.fetchRemove(fd) orelse return null;
  1591. _ = linux.epoll_ctl(self.epoll_fd, linux.EPOLL.CTL_DEL, fd, null);
  1592. return entry.value.conn;
  1593. }
  1594. /// Close every parked connection whose silence outlived PARK_IDLE_TIMEOUT_MS.
  1595. fn sweepExpired(self: *ParkedConns, tls: ?*TlsContext) void {
  1596. const now = monotonicMs();
  1597. mutexLock(&self.mutex);
  1598. defer mutexUnlock(&self.mutex);
  1599. var expired: [64]i32 = undefined;
  1600. var n: usize = 0;
  1601. var it = self.map.iterator();
  1602. while (it.next()) |kv| {
  1603. if (kv.value_ptr.deadline_ms <= now) {
  1604. if (n == expired.len) break;
  1605. expired[n] = kv.key_ptr.*;
  1606. n += 1;
  1607. }
  1608. }
  1609. for (expired[0..n]) |fd| {
  1610. if (self.map.fetchRemove(fd)) |e| {
  1611. _ = linux.epoll_ctl(self.epoll_fd, linux.EPOLL.CTL_DEL, fd, null);
  1612. var conn = e.value.conn;
  1613. conn.close(tls);
  1614. }
  1615. }
  1616. }
  1617. fn stop(self: *ParkedConns) void {
  1618. if (!self.running.swap(false, .acq_rel)) return;
  1619. if (self.wake_fd >= 0) {
  1620. const val: u64 = 1;
  1621. _ = linux.write(self.wake_fd, @ptrCast(&val), @sizeOf(u64));
  1622. }
  1623. if (self.thread) |t| {
  1624. t.join();
  1625. self.thread = null;
  1626. }
  1627. }
  1628. /// Close every parked connection and stay a VALID, EMPTY map. Unlike `deinit`
  1629. /// this keeps the epoll/wake fds and the struct usable — a CLOSED core is not
  1630. /// a freed one (`hl_http1_close`), and `park()` already refuses once `running`
  1631. /// is false, so the emptied map simply stays empty.
  1632. fn closeAll(self: *ParkedConns, tls: ?*TlsContext) void {
  1633. mutexLock(&self.mutex);
  1634. defer mutexUnlock(&self.mutex);
  1635. var it = self.map.iterator();
  1636. while (it.next()) |kv| {
  1637. var conn = kv.value_ptr.conn;
  1638. conn.close(tls);
  1639. }
  1640. self.map.clearRetainingCapacity();
  1641. }
  1642. fn deinit(self: *ParkedConns, tls: ?*TlsContext) void {
  1643. self.stop();
  1644. mutexLock(&self.mutex);
  1645. var it = self.map.iterator();
  1646. while (it.next()) |kv| {
  1647. var conn = kv.value_ptr.conn;
  1648. conn.close(tls);
  1649. }
  1650. self.map.deinit(allocator);
  1651. self.map = .empty;
  1652. mutexUnlock(&self.mutex);
  1653. if (self.epoll_fd >= 0) _ = linux.close(self.epoll_fd);
  1654. if (self.wake_fd >= 0) _ = linux.close(self.wake_fd);
  1655. self.epoll_fd = -1;
  1656. self.wake_fd = -1;
  1657. }
  1658. };
  1659. /// The parker thread: waits for a parked connection to become READABLE and only then
  1660. /// hands it to an I/O worker. The 1s epoll timeout doubles as the idle-sweep tick.
  1661. fn parkerLoop(core: *ServerCore) void {
  1662. var events: [64]linux.epoll_event = undefined;
  1663. while (core.parked.running.load(.acquire)) {
  1664. const n_rc = linux.epoll_wait(core.parked.epoll_fd, &events, events.len, 1000);
  1665. const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));
  1666. if (n > 0) {
  1667. for (events[0..@intCast(n)]) |ev| {
  1668. if (ev.data.fd == core.parked.wake_fd) {
  1669. var drain: u64 = 0;
  1670. _ = linux.read(core.parked.wake_fd, @ptrCast(&drain), @sizeOf(u64));
  1671. continue;
  1672. }
  1673. // Readable, hung up or errored — all three are a worker's job: it either
  1674. // parses the request or reads EOF and closes.
  1675. if (core.parked.take(ev.data.fd)) |conn| {
  1676. if (core.running.load(.acquire)) {
  1677. core.conn_queue.enqueue(conn);
  1678. } else {
  1679. var dead = conn;
  1680. dead.close(if (core.tls) |*t| t else null);
  1681. }
  1682. }
  1683. }
  1684. }
  1685. core.parked.sweepExpired(if (core.tls) |*t| t else null);
  1686. }
  1687. }
  1688. /// Park `conn` if we can, otherwise hand it straight to the pool. Every enqueue site
  1689. /// that is NOT known to have bytes waiting goes through here.
  1690. fn parkOrEnqueue(core: *ServerCore, conn: Connection) void {
  1691. // An established TLS connection may hold decrypted bytes epoll cannot see — see the
  1692. // ParkedConns header comment. Those go straight to a worker, as before.
  1693. if (conn.ssl == null and core.parked.park(conn)) return;
  1694. core.conn_queue.enqueue(conn);
  1695. }
  1696. // =========================================================================
  1697. // Acceptor thread — epoll loop accepting new connections
  1698. // =========================================================================
  1699. fn acceptorLoop(core: *ServerCore) void {
  1700. var events: [64]linux.epoll_event = undefined;
  1701. while (core.running.load(.acquire)) {
  1702. const n_rc = linux.epoll_wait(core.epoll_fd, &events, events.len, 1000);
  1703. const n: i32 = @bitCast(@as(u32, @truncate(n_rc)));
  1704. if (n < 0) continue;
  1705. if (n == 0) continue;
  1706. for (events[0..@intCast(n)]) |ev| {
  1707. if (ev.data.fd == core.shutdown_fd) {
  1708. return; // shutdown signaled
  1709. }
  1710. if (ev.data.fd == core.server_fd) {
  1711. // Accept all pending connections — the listener is NONBLOCK, so the
  1712. // drain ends on EAGAIN instead of sleeping inside accept4().
  1713. while (true) {
  1714. var client_addr: linux.sockaddr.in = undefined;
  1715. var addr_len: u32 = @sizeOf(linux.sockaddr.in);
  1716. const accept_rc = linux.accept4(core.server_fd, @ptrCast(&client_addr), &addr_len, linux.SOCK.CLOEXEC | linux.SOCK.NONBLOCK);
  1717. const client_fd: i32 = @bitCast(@as(u32, @truncate(accept_rc)));
  1718. if (client_fd < 0) break;
  1719. // A shutdown that landed mid-drain: the parker is stopped and the
  1720. // conn_queue has no consumers left, so hand this socket to nobody —
  1721. // close it and leave, rather than leaking the fd into a dead queue.
  1722. if (!core.running.load(.acquire)) {
  1723. _ = linux.close(client_fd);
  1724. return;
  1725. }
  1726. // Set back to blocking for I/O workers (simpler read/write)
  1727. const flags_rc = linux.fcntl(client_fd, linux.F.GETFL, @as(usize, 0));
  1728. const flags_i: isize = @bitCast(flags_rc);
  1729. if (flags_i >= 0) {
  1730. var oflags: linux.O = @bitCast(@as(u32, @truncate(flags_rc)));
  1731. oflags.NONBLOCK = false;
  1732. _ = linux.fcntl(client_fd, linux.F.SETFL, @as(usize, @as(u32, @bitCast(oflags))));
  1733. }
  1734. // PARK, don't hand to a worker: a just-accepted socket has no bytes
  1735. // yet (browsers routinely pre-open connections and send nothing), and
  1736. // a worker blocking on it is exactly the starvation this replaces.
  1737. parkOrEnqueue(core, .{
  1738. .fd = client_fd,
  1739. .ssl = null,
  1740. .keep_alive = true,
  1741. .request_count = 0,
  1742. .max_requests = MAX_KEEPALIVE_REQUESTS,
  1743. });
  1744. }
  1745. }
  1746. }
  1747. }
  1748. }
  1749. // =========================================================================
  1750. // I/O Worker thread — TLS handshake + read/parse HTTP → enqueue request
  1751. // =========================================================================
  1752. fn ioWorker(core: *ServerCore) void {
  1753. while (core.running.load(.acquire)) {
  1754. var conn = core.conn_queue.dequeue() orelse return;
  1755. // TLS handshake if needed (only on first request for this connection)
  1756. if (core.tls != null and conn.ssl == null and conn.request_count == 0) {
  1757. conn.ssl = core.tls.?.wrapConnection(conn.fd);
  1758. if (conn.ssl == null) {
  1759. _ = linux.close(conn.fd);
  1760. continue;
  1761. }
  1762. }
  1763. // Bound the worker's blocking read on EVERY connection, not just keep-alive ones.
  1764. // A parked connection only reaches a worker once it is readable, so this is now a
  1765. // backstop against a client that trickles (or stops mid-)headers rather than the
  1766. // idle-timeout mechanism it used to be — that job belongs to PARK_IDLE_TIMEOUT_MS.
  1767. const tv = linux.timeval{ .sec = 30, .usec = 0 };
  1768. _ = linux.setsockopt(conn.fd, linux.SOL.SOCKET, linux.SO.RCVTIMEO, @ptrCast(&tv), @sizeOf(linux.timeval));
  1769. // Read and parse HTTP request
  1770. switch (readAndParseRequest(&conn, core)) {
  1771. .request => |req| core.request_queue.enqueue(req),
  1772. .upgraded => {}, // socket handed to the WebSocket engine
  1773. .dead => conn.close(if (core.tls) |*t| t else null),
  1774. }
  1775. }
  1776. }
  1777. const ReadOutcome = union(enum) {
  1778. request: ParsedRequest,
  1779. upgraded, // WS handshake done — the WS engine owns the fd now
  1780. dead, // read failed / handshake rejected — caller closes
  1781. };
  1782. fn readAndParseRequest(conn: *Connection, core: *ServerCore) ReadOutcome {
  1783. const tls = if (core.tls) |*t| t else null;
  1784. // Read request headers (up to 64KB)
  1785. var buf = allocator.alloc(u8, 65536) catch return .dead;
  1786. defer allocator.free(buf);
  1787. var total: usize = 0;
  1788. var hdr_end: ?usize = null;
  1789. while (total < buf.len) {
  1790. const n = conn.read(tls, buf[total..]);
  1791. if (n == 0) break;
  1792. total += n;
  1793. // Scan for \r\n\r\n
  1794. const start = if (total > n + 3) total - n - 3 else 0;
  1795. if (std.mem.indexOf(u8, buf[start..total], "\r\n\r\n")) |pos| {
  1796. hdr_end = start + pos;
  1797. break;
  1798. }
  1799. }
  1800. if (total == 0) return .dead;
  1801. const headers_end = hdr_end orelse total;
  1802. const header_section = buf[0..headers_end];
  1803. const body_offset = if (hdr_end) |e| e + 4 else total;
  1804. // Parse request line: METHOD PATH HTTP/1.x
  1805. const line_end = std.mem.indexOf(u8, header_section, "\r\n") orelse headers_end;
  1806. const request_line = header_section[0..line_end];
  1807. var method: []const u8 = "GET";
  1808. var full_path: []const u8 = "/";
  1809. if (std.mem.indexOfScalar(u8, request_line, ' ')) |sp1| {
  1810. method = request_line[0..sp1];
  1811. const rest = request_line[sp1 + 1 ..];
  1812. if (std.mem.indexOfScalar(u8, rest, ' ')) |sp2| {
  1813. full_path = rest[0..sp2];
  1814. } else {

Only the first lines are shown.

Branches

Latest commits

  • 4a2d7125initial commitmre