gitoriaLog in with ident

tickets

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commit38bdd5e438bdd5e4deploy.sh: never send .git or .gitignore to Byrodinmre38bdd5e4/plugins/http1/http1.zig

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

Only the first lines are shown.

Branches

Latest commits

  • 38bdd5e4deploy.sh: never send .git or .gitignore to Byrodinmre
  • f12fa1bcState of 2026-09-27, before the move to gitoriamre