gitoriaLog in with ident

tickets

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commita75e0279a75e0279mission 010 (code order) 3/4: let only where a variable is reassigned or re-bound in a loop body (456 lets → plain declarations; Hybriel refuses a plain declaration inside a loop on its 2nd pass). gate 249/0, connect 60/0, real-data reads identical, a 50-step write sequence (API + faces) identical to the old codemrea75e0279/plugins/http1/http1.zig

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

Only the first lines are shown.

Branches

Latest commits

  • a75e0279mission 010 (code order) 3/4: let only where a variable is reassigned or re-bound in a loop body (456 lets → plain declarations; Hybriel refuses a plain declaration inside a loop on its 2nd pass). gate 249/0, connect 60/0, real-data reads identical, a 50-step write sequence (API + faces) identical to the old codemre
  • e9d5c618mission 010 (code order) 2/4: one lib/ file per topic — store.hl split into projects / tickets (+ relations) / events / tickets-helpers, util.hl shared helpers (env, storage dir, URLs, sorts, Vienna time), the function routes out of project.hl into lib/api.hl (thin; auth/filters/Accept in api-helpers.hl), invite + member-removal logic out of the faces/routes into invites.hl / tickets.hl; project.hl is the map. /login/callback gets req + the session store by reference. gate 249/0, connect 60/0, real-data reads identicalmre
  • 97e269b5mission 010 (code order) 1/4: .hl files out of the root — lib/ (store, users, connections, invites, migrate, markdown, mdview, import = ticketfile, util = localtime, jsoncheck, api-helpers = api), tools/import.hl, components/styles.hl; import paths only. gate 249/0, connect 60/0, real-data reads identicalmre
  • 38f9d10ftickets: Hybriel master 06617221 (plugin allocators 3a781359 + 413f60e4, mpackdb 2cb7ae5e, http1 773de63e); gate 249/0, connect 60/0mre
  • d3db6139tickets: Hybriel master 190aa11d (fc838894 GC correctness, #126 closure scopes, #127); gate 249/0, connect 60/0mre
  • bce182e3tickets: Hybriel master 7eea0d32 (#126 memory, #48 lambda copies its argument); migrate.hl lambdas take &logmre
  • 4137be0fantcolony#40: mission references point to the moved missionsmre
  • 9bfba36aantcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmre
  • c7bd2645tickets: Hybriel master 73267707 (#122 fixed); compactNow workaround removed (#110 covered)mre
  • 2ab91ee9tickets: gate checks rows appear once (session sync); re-vendor to ff51cf46 stopped on hybriel#122, stays 837fe120mre
  • e01c2b1dtickets#24: installable app (manifest, service worker, offline list), own icon; gate waits for the hello's pongmre
  • 752fbb7fdeploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
  • 38bdd5e4deploy.sh: never send .git or .gitignore to Byrodinmre
  • f12fa1bcState of 2026-09-27, before the move to gitoriamre