gitoriaLog in with ident

antcolony

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Main branchmain3a4d0324antcolony#37: a too-long report gets up to 3 fix tries, finished work is never thrown away for lengthmremain/plugins/http1/http1.zig

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

Only the first lines are shown.

Branches

Latest commits

  • 3a4d0324antcolony#37: a too-long report gets up to 3 fix tries, finished work is never thrown away for lengthmre
  • a6af7883tracker: worker box sees calendar.worldapi.org (login to copy)mre
  • c613d26btemplates: bridges to external components (login.js for ident's selector) are allowed (creator 2026-09-27)mre
  • 9062978ctracker: worker box sees /media/STORAGE/projects/old-tracker read-only (tracker#2 source data)mre
  • 7f9660eeState of 2026-09-27, before the move to gitoriamre