gitoriaLog in with ident

ident

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commitff805b9aff805b9aantcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmreff805b9a/plugins/http1/http1.zig

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

Only the first lines are shown.

Branches

Latest commits

  • ff805b9aantcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmre
  • 51a7bcdfident: Hybriel master 73267707 (#122); /code uses the new page() signature; pending address passed as parameter; once-checksmre
  • 836f644fident#24: installable app (manifest, service worker, data-free offline /start), own iconmre
  • 8bebbbf2deploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
  • cc063ea2deploy.sh: never send .git or .gitignore to Byrodinmre
  • 81b15b7bState of 2026-09-27, before the move to gitoriamre