tracker
All repositories: gitoria
57.9 KB
// hl:fetch — the HTTP(S) CLIENT, as a server-realm plugin (mission 135).//// Shaped on the web `fetch` the creator already knows: a method, a url, headers and// a body go out; a status, headers and a body come back. What is NOT web-shaped is// where the wait happens, and that is the one design decision in this file://// `fetch.start` NEVER WAITS. It hands the request to a thread of its own and// returns a TICKET (an HlHandle) immediately. `ticket.result()` is the wait.//// That split is what makes `concurrent [ fetchStart(a), fetchStart(b), … ]` put N// requests on the wire AT ONCE: a `concurrent` arm runs to completion before its// sibling starts (nothing in the interpreter yields today), so an arm that blocked// on the socket would serialize the whole group. An arm that only starts a thread// does not. `fetch()` in server.hl is `fetchStart(…).response()` — start plus wait,// the familiar one-liner, for code that has nothing else to do.//// THE EVENT LOOP is never blocked by a fetch in flight: no loop thread is inside a// socket call, and a completed request rings the mission-125 BELL — `fetch.completions`// is an ordinary loop source with a `wake_fd`, so an `on fetched(ev)` handler is woken// by the same `epoll_wait` that wakes an http1 request. (The wait in `result()` is a// wait like any other plugin call's: it blocks the FIBER that asked. Starting the// fetches early and reading them later is what keeps a server serving, and the// completions source is how a server never has to wait at all.)//// TLS is `plugins/http/tls_common.zig` — the same module hl:http1 and hl:http2 use,// extended with a client half in this mission. Verification is ON: system trust// store plus hostname checking, with `caFile` naming an EXTRA anchor for a test// fixture's self-signed certificate. There is no verify-off switch at all.//// HTTP/1.1 only. An h2 client is a later slice: it needs ALPN on the client side and// a multiplexed connection pool, neither of which this file has. `connection: close`// on every request, so there is no pool to keep coherent either.const std = @import("std");const api = @import("plugin_api");const http = @import("http_common");const tls = @import("tls_common");const HlValue = api.HlValue;const HlField = api.HlField;const HlObject = api.HlObject;const HlIterator = api.HlIterator;const HlHandle = api.HlHandle;const c = @cImport({@cInclude("netdb.h");@cInclude("sys/socket.h");@cInclude("sys/time.h");@cInclude("unistd.h");@cInclude("netinet/in.h");@cInclude("netinet/tcp.h");});// requests run on their own threads: the plugins' allocator (plugin_api.zig)const allocator = api.allocator;const linux = std.os.linux;const libc = std.c;// pthread rather than std.Thread.Mutex, exactly as hl:http1 does: this is a `.so`// loaded into a binary built by a possibly different Zig, and the C primitives are// the ones whose ABI is fixed by the platform rather than by the standard library.const PthreadMutex = libc.pthread_mutex_t;const PthreadCond = libc.pthread_cond_t;fn mutexLock(m: *PthreadMutex) void {_ = libc.pthread_mutex_lock(m);}fn mutexUnlock(m: *PthreadMutex) void {_ = libc.pthread_mutex_unlock(m);}fn condBroadcast(cnd: *PthreadCond) void {_ = libc.pthread_cond_broadcast(cnd);}fn condWait(cnd: *PthreadCond, m: *PthreadMutex) void {_ = libc.pthread_cond_wait(cnd, m);}fn logMsg(msg: []const u8) void {_ = linux.write(2, msg.ptr, msg.len);}fn logFmt(comptime fmt: []const u8, args: anytype) void {var buf: [512]u8 = undefined;const s = std.fmt.bufPrint(&buf, fmt, args) catch return;logMsg(s);}// ── Limits, all of them documented defaults rather than hard walls ───────────/// No answer at all within this many milliseconds is a failed fetch. Overridable/// per call with `timeoutMs`; it covers the WHOLE exchange (DNS, connect,/// handshake, request, response, and every redirect hop), not one syscall.const DEFAULT_TIMEOUT_MS: i64 = 30_000;/// How many 3xx hops are followed before the fetch fails with "too many redirects"./// Overridable with `maxRedirects`; 0 means the 3xx is RETURNED as the response.const DEFAULT_MAX_REDIRECTS: u32 = 5;/// A response body larger than this fails rather than growing the process without/// bound — a client that pulls other people's APIs must not be a memory bomb.const MAX_BODY: usize = 32 * 1024 * 1024;/// O_NONBLOCK on x86-64 Linux. Spelled here because <fcntl.h> cannot be imported/// (see tcpConnect); it is an ABI constant, not a libc detail.const O_NONBLOCK: usize = 0o4000;const Header = struct {name: []u8,value: []u8,};fn freeHeaders(list: []Header) void {for (list) |h| {allocator.free(h.name);allocator.free(h.value);}allocator.free(list);}// ── The request, as the ticket carries it ────────────────────────────────────const Request = struct {method: []u8,url: []u8,headers: []Header,body: []u8,has_body: bool,json_body: bool,timeout_ms: i64,max_redirects: u32,ca_file: ?[]u8,tag: []u8,/// STREAMING (the LLM-token mode): deliver the body INCREMENTALLY as chunk/// events on the completions source instead of accumulating it. The final/// completion then carries status/headers and an EMPTY body — the chunks/// were the body. Requires an armed completions source; `result()` still/// works and yields the empty-body summary when the stream ends.stream: bool = false,fn deinit(self: *Request) void {allocator.free(self.method);allocator.free(self.url);freeHeaders(self.headers);allocator.free(self.body);if (self.ca_file) |ca| allocator.free(ca);allocator.free(self.tag);}};// ── The ticket: one in-flight (or finished) request ──────────────────────────const Ticket = struct {mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER,cond: PthreadCond = libc.PTHREAD_COND_INITIALIZER,done: bool = false,thread: ?std.Thread = null,req: Request,status: u16 = 0,final_url: []u8 = &.{},headers: []Header = &.{},body: []u8 = &.{},/// Non-null = the fetch failed and this is the located message the handle/// hands back as an `hl_error`.err: ?[]u8 = null,/// ABORT (the LLM-interrupt primitive): the socket currently carrying this/// request, visible under the ticket mutex so `fetchAbort(tag)` can/// shutdown() it from the interpreter thread — the read unblocks, the peer/// sees the disconnect and stops generating. -1 = none in flight.live_fd: i32 = -1,aborted: bool = false,/// INSTANCE EVENTS (stream mode): a streaming ticket carries its OWN frame/// queue and bell, created at start — before the worker exists — so no chunk/// can outrun the arming. `handle.events()` wraps them into a loop source;/// pending_fetch.hl registers it and re-emits `chunk`/`done` at itself./// -1 = not a streaming ticket (completions go the global route, if armed).ev_queue: std.ArrayListUnmanaged(Completion) = .empty,ev_wake_fd: i32 = -1,fn wait(self: *Ticket) void {mutexLock(&self.mutex);while (!self.done) condWait(&self.cond, &self.mutex);mutexUnlock(&self.mutex);}fn isDone(self: *Ticket) bool {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);return self.done;}fn setLiveFd(self: *Ticket, fd: i32) bool {mutexLock(&self.mutex);defer mutexUnlock(&self.mutex);if (self.aborted) return false; // aborted before the connect landedself.live_fd = fd;return true;}fn clearLiveFd(self: *Ticket) void {mutexLock(&self.mutex);self.live_fd = -1;mutexUnlock(&self.mutex);}/// The abort itself: mark, and shutdown() any in-flight socket so the/// worker's blocking read returns NOW. Closing is still the worker's job/// (its defer owns the fd); shutdown only cuts the conversation.fn abort(self: *Ticket) void {mutexLock(&self.mutex);self.aborted = true;if (self.live_fd >= 0) _ = c.shutdown(self.live_fd, 2); // SHUT_RDWRmutexUnlock(&self.mutex);}fn finish(self: *Ticket) void {mutexLock(&self.mutex);self.done = true;condBroadcast(&self.cond);mutexUnlock(&self.mutex);}fn fail(self: *Ticket, comptime fmt: []const u8, args: anytype) void {self.err = std.fmt.allocPrint(allocator, fmt, args) catch null;}};// ── URL ──────────────────────────────────────────────────────────────────────const Url = struct {https: bool,host: []const u8,port: u16,/// path + query, always starting with '/'target: []const u8,};/// `scheme://host[:port][/path][?query]`. Only http and https; a URL with/// credentials or an IPv6 literal is refused rather than half-understood.fn parseUrl(url: []const u8) ?Url {var rest: []const u8 = undefined;var https = false;if (std.mem.startsWith(u8, url, "http://")) {rest = url["http://".len..];} else if (std.mem.startsWith(u8, url, "https://")) {rest = url["https://".len..];https = true;} else return null;const slash = std.mem.indexOfScalar(u8, rest, '/');const authority = if (slash) |s| rest[0..s] else rest;const target = if (slash) |s| rest[s..] else "/";if (authority.len == 0) return null;if (std.mem.indexOfScalar(u8, authority, '@') != null) return null;if (std.mem.indexOfScalar(u8, authority, '[') != null) return null;var host = authority;var port: u16 = if (https) 443 else 80;if (std.mem.lastIndexOfScalar(u8, authority, ':')) |colon| {host = authority[0..colon];port = std.fmt.parseInt(u16, authority[colon + 1 ..], 10) catch return null;}if (host.len == 0) return null;return .{ .https = https, .host = host, .port = port, .target = target };}/// Resolve a `Location` against the URL it came from: absolute stays, `/x` replaces/// the path, anything else is relative to the current directory of the path.fn resolveLocation(base: []const u8, loc: []const u8) ?[]u8 {if (loc.len == 0) return null;if (std.mem.startsWith(u8, loc, "http://") or std.mem.startsWith(u8, loc, "https://")) {return allocator.dupe(u8, loc) catch null;}const b = parseUrl(base) orelse return null;const scheme = if (b.https) "https" else "http";const default_port: u16 = if (b.https) 443 else 80;var authority_buf: [300]u8 = undefined;const authority = if (b.port == default_port)std.fmt.bufPrint(&authority_buf, "{s}", .{b.host}) catch return nullelsestd.fmt.bufPrint(&authority_buf, "{s}:{d}", .{ b.host, b.port }) catch return null;if (loc[0] == '/') {return std.fmt.allocPrint(allocator, "{s}://{s}{s}", .{ scheme, authority, loc }) catch null;}// Relative: cut the base target back to its last '/'.const path_only = if (std.mem.indexOfScalar(u8, b.target, '?')) |q| b.target[0..q] else b.target;const cut = std.mem.lastIndexOfScalar(u8, path_only, '/') orelse 0;return std.fmt.allocPrint(allocator, "{s}://{s}{s}/{s}", .{ scheme, authority, path_only[0..cut], loc }) catch null;}// ── Socket plumbing ──────────────────────────────────────────────────────────fn nowMs() i64 {var ts: linux.timespec = undefined;_ = linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts);return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);}/// Connect to host:port, giving up after `budget_ms`. The connect is non-blocking/// plus a poll so the DEADLINE is honoured even when the peer never answers SYN;/// the socket is put back into blocking mode with SO_RCVTIMEO/SO_SNDTIMEO after,/// which is what bounds the read side.fn tcpConnect(host: []const u8, port: u16, budget_ms: i64, err_out: *[]const u8) ?i32 {if (budget_ms <= 0) {err_out.* = "timed out";return null;}var host_buf: [256]u8 = undefined;if (host.len >= host_buf.len) {err_out.* = "host name too long";return null;}@memcpy(host_buf[0..host.len], host);host_buf[host.len] = 0;var port_buf: [8]u8 = undefined;const port_str = std.fmt.bufPrintZ(&port_buf, "{d}", .{port}) catch {err_out.* = "bad port";return null;};var hints: c.struct_addrinfo = std.mem.zeroes(c.struct_addrinfo);hints.ai_family = c.AF_UNSPEC;hints.ai_socktype = c.SOCK_STREAM;hints.ai_protocol = c.IPPROTO_TCP;var res: ?*c.struct_addrinfo = null;if (c.getaddrinfo(@ptrCast(&host_buf), port_str.ptr, &hints, &res) != 0 or res == null) {err_out.* = "could not resolve host";return null;}defer c.freeaddrinfo(res);const deadline = nowMs() + budget_ms;var last: []const u8 = "connection failed";var it: ?*c.struct_addrinfo = res;while (it) |ai| : (it = ai.ai_next) {const remaining = deadline - nowMs();if (remaining <= 0) {err_out.* = "timed out";return null;}const fd_c = c.socket(ai.ai_family, ai.ai_socktype, ai.ai_protocol);if (fd_c < 0) {last = "could not create a socket";continue;}const fd: i32 = @intCast(fd_c);// fcntl straight from the kernel: glibc's <fcntl.h> does not survive// translate-c under _FORTIFY_SOURCE (its open() overloads are declared// with attribute error), and the flags are ABI constants either way.const flags = linux.fcntl(fd, linux.F.GETFL, 0);_ = linux.fcntl(fd, linux.F.SETFL, flags | @as(usize, O_NONBLOCK));var connected = false;if (c.connect(fd, ai.ai_addr, ai.ai_addrlen) == 0) {connected = true;} else {// linux.poll, not glibc's: <poll.h> is fortified too (see above).var pfd = [_]linux.pollfd{.{ .fd = fd, .events = linux.POLL.OUT, .revents = 0 }};const pr: isize = @bitCast(linux.poll(&pfd, 1, @intCast(@min(remaining, std.math.maxInt(i32)))));if (pr > 0) {var so_err: c_int = 0;var len: c.socklen_t = @sizeOf(c_int);_ = c.getsockopt(fd, c.SOL_SOCKET, c.SO_ERROR, &so_err, &len);if (so_err == 0) {connected = true;} else if (so_err == 111) {last = "connection refused";} else {last = "connection failed";}} else if (pr == 0) {_ = c.close(fd);err_out.* = "timed out";return null;} else {last = "connection failed";}}if (!connected) {_ = c.close(fd);continue;}_ = linux.fcntl(fd, linux.F.SETFL, flags);const left = @max(deadline - nowMs(), 1);var tv = c.struct_timeval{.tv_sec = @intCast(@divTrunc(left, 1000)),.tv_usec = @intCast(@rem(left, 1000) * 1000),};_ = c.setsockopt(fd, c.SOL_SOCKET, c.SO_RCVTIMEO, &tv, @sizeOf(c.struct_timeval));_ = c.setsockopt(fd, c.SOL_SOCKET, c.SO_SNDTIMEO, &tv, @sizeOf(c.struct_timeval));const one: c_int = 1;_ = c.setsockopt(fd, c.IPPROTO_TCP, c.TCP_NODELAY, &one, @sizeOf(c_int));return fd;}err_out.* = last;return null;}/// The one place either transport is read or written, so the HTTP exchange below/// is written once for http:// and https://.const Conn = struct {fd: i32,ssl: ?*tls.SSL = null,tls_ctx: ?*const tls.TlsContext = null,fn writeAll(self: *Conn, data: []const u8) bool {if (self.ssl) |ssl| {return self.tls_ctx.?.writeAll(ssl, data) == data.len;}var sent: usize = 0;while (sent < data.len) {const n = c.write(self.fd, data.ptr + sent, data.len - sent);if (n <= 0) return false;sent += @intCast(n);}return true;}/// >0 bytes, 0 = clean end of body, <0 = error (or a socket timeout, which the/// caller separates from an error by looking at the clock).fn read(self: *Conn, buf: []u8) isize {if (self.ssl) |ssl| {const n = self.tls_ctx.?.read(ssl, buf);if (n > 0) return n;const e = self.tls_ctx.?.getError(ssl, n);if (e == tls.SSL_ERROR_ZERO_RETURN) return 0;return -1;}return c.read(self.fd, buf.ptr, buf.len);}fn close(self: *Conn) void {if (self.ssl) |ssl| {self.tls_ctx.?.shutdownAndFree(ssl);self.ssl = null;}_ = c.close(self.fd);}};// ── The client TLS context, one per CA file ──────────────────────────────────//// An SSL_CTX is expensive (it reads the system trust store) and thread-safe, so it// is built once per distinct `caFile` and shared. The map is tiny by construction:// a program has a system-trust context and at most a handful of pinned fixtures.const CtxEntry = struct {key: []u8,ctx: tls.TlsContext,};var ctx_mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER;var ctx_list: std.ArrayListUnmanaged(CtxEntry) = .empty;fn clientCtx(ca_file: ?[]const u8, err_out: *[]const u8) ?*const tls.TlsContext {const key = ca_file orelse "";mutexLock(&ctx_mutex);defer mutexUnlock(&ctx_mutex);for (ctx_list.items) |*e| {if (std.mem.eql(u8, e.key, key)) return &e.ctx;}const ctx = tls.TlsContext.initClient(.{ .ca_file = ca_file, .tag = "fetch" }) orelse {err_out.* = "https is not available (no usable OpenSSL, or the caFile could not be loaded)";return null;};const key_copy = allocator.dupe(u8, key) catch {err_out.* = "out of memory";return null;};ctx_list.append(allocator, .{ .key = key_copy, .ctx = ctx }) catch {allocator.free(key_copy);err_out.* = "out of memory";return null;};return &ctx_list.items[ctx_list.items.len - 1].ctx;}// ── One HTTP exchange ────────────────────────────────────────────────────────const Exchange = struct {status: u16 = 0,headers: []Header = &.{},body: []u8 = &.{},fn deinit(self: *Exchange) void {freeHeaders(self.headers);allocator.free(self.body);}};fn headerValue(headers: []const Header, name: []const u8) ?[]const u8 {for (headers) |h| {if (std.ascii.eqlIgnoreCase(h.name, name)) return h.value;}return null;}fn hasHeader(headers: []const Header, name: []const u8) bool {return headerValue(headers, name) != null;}/// Send one request over one connection and read the whole response./// `err_out` gets a short static reason; the caller adds the URL./// The status code off a raw head, or null while unparseable.fn parseStatus(head: []const u8) ?u16 {var lines = std.mem.splitSequence(u8, head, "\r\n");const status_line = lines.next() orelse return null;var parts = std.mem.splitScalar(u8, status_line, ' ');_ = parts.next(); // HTTP/1.1const code = parts.next() orelse return null;return std.fmt.parseInt(u16, code, 10) catch null;}/// Head only → an Exchange with an EMPTY body (the chunks were the body).fn parseHeadOnly(head_raw: []const u8, err_out: *[]const u8) ?Exchange {// reuse the full parser on a synthetic complete response: the head as-is// (its transfer-encoding stripped so no body is expected) + no body.var synth: std.ArrayListUnmanaged(u8) = .empty;defer synth.deinit(allocator);var lines = std.mem.splitSequence(u8, head_raw[0 .. head_raw.len - 4], "\r\n");while (lines.next()) |line| {if (line.len > 18 and std.ascii.eqlIgnoreCase(line[0..18], "transfer-encoding:")) continue;if (line.len > 15 and std.ascii.eqlIgnoreCase(line[0..15], "content-length:")) continue;if (!add(&synth, line)) return oom(err_out);if (!add(&synth, "\r\n")) return oom(err_out);}if (!add(&synth, "\r\n")) return oom(err_out);return parseResponse(synth.items, err_out);}/// STREAM the body: post every decoded piece as a chunk frame until the peer/// closes (`connection: close` is on every request, so close IS the end)./// Chunked transfer is decoded incrementally; anything else streams as-is.fn streamBody(t: *Ticket, conn: *Conn, raw: *std.ArrayListUnmanaged(u8), he: usize, deadline: i64, err_out: *[]const u8) ?Exchange {defer raw.deinit(allocator);const chunked = blk: {if (findHeaderIn(raw.items[0..he], "transfer-encoding")) |v| {break :blk std.ascii.indexOfIgnoreCase(v, "chunked") != null;}break :blk false;};// chunked-decoder state, carried across readsvar pending: std.ArrayListUnmanaged(u8) = .empty;defer pending.deinit(allocator);var remaining: usize = 0; // bytes left of the current chunk's datavar done = false;// seed with whatever body bytes arrived along with the headif (raw.items.len > he) {pending.appendSlice(allocator, raw.items[he..]) catch return oom(err_out);}var buf: [16 * 1024]u8 = undefined;while (true) {// decode + post what we holdif (!chunked) {if (pending.items.len > 0) {postChunk(t, pending.items);pending.clearRetainingCapacity();}} else while (!done) {if (remaining > 0) {const take = @min(remaining, pending.items.len);if (take == 0) break;postChunk(t, pending.items[0..take]);std.mem.copyForwards(u8, pending.items, pending.items[take..]);pending.shrinkRetainingCapacity(pending.items.len - take);remaining -= take;if (remaining == 0) {// the CRLF after the chunk dataif (pending.items.len < 2) break;std.mem.copyForwards(u8, pending.items, pending.items[2..]);pending.shrinkRetainingCapacity(pending.items.len - 2);}continue;}const nl = std.mem.indexOf(u8, pending.items, "\r\n") orelse break;const size_line = pending.items[0..nl];const semi = std.mem.indexOfScalar(u8, size_line, ';') orelse size_line.len;const n = std.fmt.parseInt(usize, std.mem.trim(u8, size_line[0..semi], " \t"), 16) catch {err_out.* = "the peer sent a malformed chunked body";return null;};std.mem.copyForwards(u8, pending.items, pending.items[nl + 2 ..]);pending.shrinkRetainingCapacity(pending.items.len - nl - 2);if (n == 0) {done = true;break;}remaining = n;}if (done) break;const n = conn.read(&buf);if (n == 0) break; // close IS the end (identity), or a truncated chunked streamif (n < 0) {if (nowMs() >= deadline) {err_out.* = "timed out";return null;}break;}pending.appendSlice(allocator, buf[0..@intCast(n)]) catch return oom(err_out);}return parseHeadOnly(raw.items[0..he], err_out);}fn exchange(t: *Ticket, url: Url, method: []const u8, body: []const u8, send_body: bool, deadline: i64, err_out: *[]const u8) ?Exchange {const budget = deadline - nowMs();const fd = tcpConnect(url.host, url.port, budget, err_out) orelse return null;if (!t.setLiveFd(fd)) {_ = c.close(fd);err_out.* = "aborted";return null;}var conn = Conn{ .fd = fd };if (url.https) {const ctx = clientCtx(t.req.ca_file, err_out) orelse {_ = c.close(fd);return null;};// 137's shared client half logs the specific failure itself; the caller// gets the located summary (certificate not trusted, or not a TLS server)const ssl = ctx.connect(fd, url.host, "fetch") orelse {_ = c.close(fd);err_out.* = "TLS handshake failed (certificate not trusted, or not a TLS server)";return null;};conn.ssl = ssl;conn.tls_ctx = ctx;}defer conn.close();defer t.clearLiveFd();// --- request ---var req_buf: std.ArrayListUnmanaged(u8) = .empty;defer req_buf.deinit(allocator);const default_port: u16 = if (url.https) 443 else 80;var ok = true;ok = ok and addFmt(&req_buf, "{s} {s} HTTP/1.1\r\n", .{ method, url.target });if (url.port == default_port) {ok = ok and addFmt(&req_buf, "host: {s}\r\n", .{url.host});} else {ok = ok and addFmt(&req_buf, "host: {s}:{d}\r\n", .{ url.host, url.port });}// `connection: close` on every request: there is no connection pool, and a// closed connection is also the unambiguous end of a body with no length.ok = ok and add(&req_buf, "connection: close\r\n");if (!hasHeader(t.req.headers, "user-agent")) {ok = ok and add(&req_buf, "user-agent: hybriel-fetch/1\r\n");}if (!hasHeader(t.req.headers, "accept")) {ok = ok and add(&req_buf, "accept: */*\r\n");}if (send_body) {ok = ok and addFmt(&req_buf, "content-length: {d}\r\n", .{body.len});if (t.req.json_body and !hasHeader(t.req.headers, "content-type")) {ok = ok and add(&req_buf, "content-type: application/json\r\n");}}for (t.req.headers) |h| {// A caller header may not smuggle a second request in (CRLF injection).if (std.mem.indexOfAny(u8, h.name, "\r\n") != null or std.mem.indexOfAny(u8, h.value, "\r\n") != null) {err_out.* = "a header name or value contains a line break";return null;}ok = ok and addFmt(&req_buf, "{s}: {s}\r\n", .{ h.name, h.value });}ok = ok and add(&req_buf, "\r\n");if (send_body and body.len > 0) {ok = ok and add(&req_buf, body);}if (!ok) return oom(err_out);if (!conn.writeAll(req_buf.items)) {err_out.* = if (nowMs() >= deadline) "timed out" else "could not send the request";return null;}// --- response ---var raw: std.ArrayListUnmanaged(u8) = .empty;errdefer raw.deinit(allocator);var head_end: ?usize = null;var buf: [16 * 1024]u8 = undefined;while (true) {if (head_end == null) {if (std.mem.indexOf(u8, raw.items, "\r\n\r\n")) |idx| head_end = idx + 4;}if (head_end) |he| {// STREAMING (LLM tokens): the head is in — if this is the final// answer (not a redirect the caller follows), hand the rest of the// body over CHUNK BY CHUNK as it arrives instead of accumulating.if (t.req.stream) {const status_ok = parseStatus(raw.items[0..he]);if (status_ok != null and !isRedirect(status_ok.?)) {return streamBody(t, &conn, &raw, he, deadline, err_out);}}// Once the head is in, stop as soon as the body is provably complete.if (bodyComplete(raw.items, he)) break;}if (raw.items.len > MAX_BODY) {raw.deinit(allocator);err_out.* = "the response is larger than the 32 MiB limit";return null;}const n = conn.read(&buf);if (n == 0) break;if (n < 0) {if (nowMs() >= deadline) {raw.deinit(allocator);err_out.* = "timed out";return null;}// A TLS peer that just closes without close_notify, or a plain socket// reset after a complete answer: the parse below decides whether what// arrived is a whole response.break;}raw.appendSlice(allocator, buf[0..@intCast(n)]) catch {raw.deinit(allocator);return oom(err_out);};}const out = parseResponse(raw.items, err_out) orelse {raw.deinit(allocator);return null;};raw.deinit(allocator);return out;}fn add(buf: *std.ArrayListUnmanaged(u8), data: []const u8) bool {buf.appendSlice(allocator, data) catch return false;return true;}/// One request line or header, formatted. A header that does not fit 2 KiB is not/// a header anyone meant to send.fn addFmt(buf: *std.ArrayListUnmanaged(u8), comptime fmt: []const u8, args: anytype) bool {var tmp: [2048]u8 = undefined;const s = std.fmt.bufPrint(&tmp, fmt, args) catch return false;return add(buf, s);}fn oom(err_out: *[]const u8) ?Exchange {err_out.* = "out of memory";return null;}/// Is everything the response promised already in `raw`? Content-Length is exact,/// chunked ends at the zero chunk, and anything else ends when the peer closes.fn bodyComplete(raw: []const u8, head_end: usize) bool {const head = raw[0..head_end];if (findHeaderIn(head, "content-length")) |v| {const want = std.fmt.parseInt(usize, std.mem.trim(u8, v, " \t"), 10) catch return false;return raw.len - head_end >= want;}if (findHeaderIn(head, "transfer-encoding")) |v| {if (std.ascii.indexOfIgnoreCase(v, "chunked") != null) {return std.mem.indexOf(u8, raw[head_end..], "\r\n0\r\n") != null orstd.mem.startsWith(u8, raw[head_end..], "0\r\n");}}return false;}fn findHeaderIn(head: []const u8, name: []const u8) ?[]const u8 {var lines = std.mem.splitSequence(u8, head, "\r\n");_ = lines.next(); // status linewhile (lines.next()) |line| {if (line.len == 0) break;const colon = std.mem.indexOfScalar(u8, line, ':') orelse continue;if (std.ascii.eqlIgnoreCase(std.mem.trim(u8, line[0..colon], " \t"), name)) {return std.mem.trim(u8, line[colon + 1 ..], " \t");}}return null;}fn parseResponse(raw: []const u8, err_out: *[]const u8) ?Exchange {const he = std.mem.indexOf(u8, raw, "\r\n\r\n") orelse {err_out.* = "the peer sent no complete HTTP response";return null;};const head = raw[0..he];const body_raw = raw[he + 4 ..];var lines = std.mem.splitSequence(u8, head, "\r\n");const status_line = lines.next() orelse {err_out.* = "the peer sent no status line";return null;};// "HTTP/1.1 200 OK"const sp = std.mem.indexOfScalar(u8, status_line, ' ') orelse {err_out.* = "the peer sent a malformed status line";return null;};const after = status_line[sp + 1 ..];const code_end = std.mem.indexOfScalar(u8, after, ' ') orelse after.len;const status = std.fmt.parseInt(u16, std.mem.trim(u8, after[0..code_end], " \t"), 10) catch {err_out.* = "the peer sent a malformed status line";return null;};var headers: std.ArrayListUnmanaged(Header) = .empty;errdefer {for (headers.items) |h| {allocator.free(h.name);allocator.free(h.value);}headers.deinit(allocator);}var chunked = false;while (lines.next()) |line| {if (line.len == 0) continue;const colon = std.mem.indexOfScalar(u8, line, ':') orelse continue;const name = std.mem.trim(u8, line[0..colon], " \t");const value = std.mem.trim(u8, line[colon + 1 ..], " \t");if (name.len == 0) continue;const name_copy = allocator.dupe(u8, name) catch return oom(err_out);http.toLowercase(name_copy);const value_copy = allocator.dupe(u8, value) catch {allocator.free(name_copy);return oom(err_out);};if (std.mem.eql(u8, name_copy, "transfer-encoding") andstd.ascii.indexOfIgnoreCase(value_copy, "chunked") != null) chunked = true;headers.append(allocator, .{ .name = name_copy, .value = value_copy }) catch {allocator.free(name_copy);allocator.free(value_copy);return oom(err_out);};}var body: []u8 = undefined;if (chunked) {body = dechunk(body_raw) orelse {for (headers.items) |h| {allocator.free(h.name);allocator.free(h.value);}headers.deinit(allocator);err_out.* = "the peer sent a malformed chunked body";return null;};} else if (findHeaderIn(head, "content-length")) |v| {const want = std.fmt.parseInt(usize, v, 10) catch body_raw.len;const take = @min(want, body_raw.len);body = allocator.dupe(u8, body_raw[0..take]) catch return oom(err_out);} else {body = allocator.dupe(u8, body_raw) catch return oom(err_out);}const headers_slice = headers.toOwnedSlice(allocator) catch return oom(err_out);return .{ .status = status, .headers = headers_slice, .body = body };}fn dechunk(input: []const u8) ?[]u8 {var out: std.ArrayListUnmanaged(u8) = .empty;var i: usize = 0;while (true) {const line_end = std.mem.indexOfPos(u8, input, i, "\r\n") orelse {out.deinit(allocator);return null;};var size_txt = input[i..line_end];if (std.mem.indexOfScalar(u8, size_txt, ';')) |semi| size_txt = size_txt[0..semi];const size = std.fmt.parseInt(usize, std.mem.trim(u8, size_txt, " \t"), 16) catch {out.deinit(allocator);return null;};i = line_end + 2;if (size == 0) return out.toOwnedSlice(allocator) catch null;if (i + size > input.len) {out.deinit(allocator);return null;}out.appendSlice(allocator, input[i .. i + size]) catch {out.deinit(allocator);return null;};i += size + 2; // skip the chunk's trailing CRLFif (i > input.len) return out.toOwnedSlice(allocator) catch null;}}// ── The worker: one thread per fetch, redirects and all ───────────────────────fn isRedirect(status: u16) bool {return status == 301 or status == 302 or status == 303 or status == 307 or status == 308;}// ── The ACTIVE registry: every in-flight ticket, so `fetchAbort(tag)` can find// its socket. Added at spawn, removed at finish — a finished ticket is owned by// its handle alone and an abort of it is a no-op.var active_mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER;var active_tickets: std.ArrayListUnmanaged(*Ticket) = .empty;fn activeAdd(t: *Ticket) void {mutexLock(&active_mutex);active_tickets.append(allocator, t) catch {};mutexUnlock(&active_mutex);}fn activeRemove(t: *Ticket) void {mutexLock(&active_mutex);for (active_tickets.items, 0..) |a, i| {if (a == t) {_ = active_tickets.swapRemove(i);break;}}mutexUnlock(&active_mutex);}fn workerMain(t: *Ticket) void {runFetch(t);activeRemove(t);// an aborted request reports ABORTED whatever the read error looked like —// the app asked for the interrupt, the message says the interrupt happenedif (t.aborted) {if (t.err) |e| allocator.free(e);t.err = std.fmt.allocPrint(allocator, "hl:fetch: aborted", .{}) catch null;}t.finish();postCompletion(t);}fn runFetch(t: *Ticket) void {const deadline = nowMs() + t.req.timeout_ms;var current = allocator.dupe(u8, t.req.url) catch {t.fail("hl:fetch: out of memory", .{});return;};var method = allocator.dupe(u8, t.req.method) catch {allocator.free(current);t.fail("hl:fetch: out of memory", .{});return;};var send_body = t.req.has_body;var hops: u32 = 0;while (true) {const url = parseUrl(current) orelse {t.fail("hl:fetch: '{s}' is not an http:// or https:// URL", .{current});allocator.free(current);allocator.free(method);return;};var reason: []const u8 = "connection failed";var ex = exchange(t, url, method, t.req.body, send_body, deadline, &reason) orelse {t.fail("hl:fetch {s} {s}: {s}", .{ method, current, reason });allocator.free(current);allocator.free(method);return;};if (isRedirect(ex.status) and hops < t.req.max_redirects) {if (headerValue(ex.headers, "location")) |loc| {const next = resolveLocation(current, loc) orelse {t.fail("hl:fetch {s}: redirect to an unusable Location '{s}'", .{ current, loc });ex.deinit();allocator.free(current);allocator.free(method);return;};// 303 always becomes a GET; 301/302 become one for anything that// is not already GET/HEAD — what every browser does, and what a// server that answers a POST with "see the result over there"// means. 307/308 keep the method AND the body by definition.if (ex.status == 303 or((ex.status == 301 or ex.status == 302) and!std.mem.eql(u8, method, "GET") and !std.mem.eql(u8, method, "HEAD"))){allocator.free(method);method = allocator.dupe(u8, "GET") catch {allocator.free(next);ex.deinit();allocator.free(current);t.fail("hl:fetch: out of memory", .{});return;};send_body = false;}ex.deinit();allocator.free(current);current = next;hops += 1;continue;}}if (isRedirect(ex.status) and t.req.max_redirects > 0 and hops >= t.req.max_redirects) {t.fail("hl:fetch {s}: more than {d} redirects", .{ t.req.url, t.req.max_redirects });ex.deinit();allocator.free(current);allocator.free(method);return;}t.status = ex.status;t.headers = ex.headers;t.body = ex.body;t.final_url = current;allocator.free(method);return;}}// ── The value shapes handed back to Hybriel ──────────────────────────────────fn hlStr(s: []const u8) api.HlString {return .{ .ptr = s.ptr, .len = s.len };}/// Only the SHELLS are freed: every string in a `result()` object points into the/// ticket, which outlives the call (the loader copies what it keeps before this/// runs — the plugin ABI's ownership rule).fn respShellDeinit(obj: *HlObject) callconv(.c) void {const hdrs = obj.fields[3].value;if (hdrs.type == .hl_object) {const inner = hdrs.data.object;allocator.free(inner.fields[0..inner.field_count]);allocator.destroy(inner);}allocator.free(obj.fields[0..obj.field_count]);allocator.destroy(obj);}fn buildHeadersObject(headers: []const Header) ?*HlObject {const fields = allocator.alloc(HlField, headers.len) catch return null;for (headers, 0..) |h, i| {fields[i] = .{ .key = hlStr(h.name), .value = api.makeString(h.value) };}const obj = allocator.create(HlObject) catch {allocator.free(fields);return null;};obj.* = .{ .fields = fields.ptr, .field_count = headers.len, .deinit_fn = null };return obj;}/// `{ status, ok, url, headers, body, tag }` — the response as the .hl side sees/// it, hydrated into a FetchResponse by server.hl.fn buildResponse(t: *Ticket) HlValue {const headers_obj = buildHeadersObject(t.headers) orelse return api.makeError("hl:fetch: out of memory");const fields = allocator.alloc(HlField, 6) catch {allocator.free(headers_obj.fields[0..headers_obj.field_count]);allocator.destroy(headers_obj);return api.makeError("hl:fetch: out of memory");};fields[0] = .{ .key = hlStr("status"), .value = api.makeNumber(@floatFromInt(t.status)) };fields[1] = .{ .key = hlStr("ok"), .value = api.makeBool(t.status >= 200 and t.status < 300) };fields[2] = .{ .key = hlStr("url"), .value = api.makeString(t.final_url) };fields[3] = .{ .key = hlStr("headers"), .value = api.makeObject(headers_obj) };fields[4] = .{ .key = hlStr("body"), .value = api.makeString(t.body) };fields[5] = .{ .key = hlStr("tag"), .value = api.makeString(t.req.tag) };const obj = allocator.create(HlObject) catch {allocator.free(fields);allocator.free(headers_obj.fields[0..headers_obj.field_count]);allocator.destroy(headers_obj);return api.makeError("hl:fetch: out of memory");};obj.* = .{ .fields = fields.ptr, .field_count = 6, .deinit_fn = &respShellDeinit };return api.makeObject(obj);}// ── The ticket handle ────────────────────────────────────────────────────────fn ticketCall(ctx: ?*anyopaque, op: api.HlString, argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {_ = argc;_ = argv;const t: *Ticket = @ptrCast(@alignCast(ctx orelse return api.makeError("hl:fetch: dead ticket")));const name = op.ptr[0..op.len];if (std.mem.eql(u8, name, "done")) {return api.makeBool(t.isDone());}if (std.mem.eql(u8, name, "result")) {t.wait();if (t.err) |msg| return api.makeError(msg);return buildResponse(t);}if (std.mem.eql(u8, name, "events")) {if (t.ev_wake_fd < 0) return api.makeError("hl:fetch: events() is stream mode — start the fetch with stream = true");const iter = allocator.create(HlIterator) catch return api.makeError("hl:fetch: out of memory");iter.* = .{.context = @ptrCast(t),.next_fn = &eventsTryNext, // non-blocking either way — event loop only.deinit_fn = null,.try_next_fn = &eventsTryNext,.wake_fd = t.ev_wake_fd,};return api.makeIterator(iter);}if (std.mem.eql(u8, name, "abort")) {t.abort();return api.makeBool(true);}return api.makeError("hl:fetch: a pending fetch has no method by that name");}fn ticketClose(ctx: ?*anyopaque) callconv(.c) void {const t: *Ticket = @ptrCast(@alignCast(ctx orelse return));// The worker writes into this ticket, so it is JOINED before anything is// freed — a detached thread finishing after the free is a use-after-free the// collector would surface as random corruption. A dropped, never-read fetch// therefore costs at most its own timeout at collection time.if (t.thread) |th| {th.join();t.thread = null;}t.req.deinit();for (t.ev_queue.items) |comp| {allocator.free(comp.tag);allocator.free(comp.url);allocator.free(comp.body);allocator.free(comp.err);if (comp.chunk) |ch| allocator.free(ch);for (comp.headers) |h| {allocator.free(h.name);allocator.free(h.value);}if (comp.headers.len > 0) allocator.free(comp.headers);}t.ev_queue.deinit(allocator);freeHeaders(t.headers);allocator.free(t.body);allocator.free(t.final_url);if (t.err) |e| allocator.free(e);allocator.destroy(t);}// ── Argument reading ─────────────────────────────────────────────────────────fn fieldOf(v: HlValue, name: []const u8) ?HlValue {if (v.type != .hl_object) return null;const obj = v.data.object;for (obj.fields[0..obj.field_count]) |f| {if (std.mem.eql(u8, f.key.ptr[0..f.key.len], name)) return f.value;}return null;}fn dupString(v: ?HlValue, fallback: []const u8) []u8 {if (v) |val| {if (val.type == .hl_string) {return allocator.dupe(u8, val.data.string.ptr[0..val.data.string.len]) catchallocator.dupe(u8, fallback) catch unreachable;}}return allocator.dupe(u8, fallback) catch unreachable;}fn numberOr(v: ?HlValue, fallback: i64) i64 {if (v) |val| {if (val.type == .hl_number) return @intFromFloat(val.data.number);}return fallback;}fn boolOr(v: ?HlValue, fallback: bool) bool {if (v) |val| {if (val.type == .hl_bool) return val.data.boolean;}return fallback;}fn readHeaders(v: ?HlValue) []Header {const val = v orelse return &.{};if (val.type != .hl_object) return &.{};const obj = val.data.object;var list: std.ArrayListUnmanaged(Header) = .empty;for (obj.fields[0..obj.field_count]) |f| {if (f.value.type != .hl_string) continue;const name = allocator.dupe(u8, f.key.ptr[0..f.key.len]) catch continue;const value = allocator.dupe(u8, f.value.data.string.ptr[0..f.value.data.string.len]) catch {allocator.free(name);continue;};list.append(allocator, .{ .name = name, .value = value }) catch {allocator.free(name);allocator.free(value);};}return list.toOwnedSlice(allocator) catch &.{};}/// __native("fetch.start", url, options) → a ticket HANDLE, immediately./// Options: method, headers, body, jsonBody, timeoutMs, maxRedirects, caFile, tag.export fn hl_fetch_start(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 1 or argv[0].type != .hl_string) {return api.makeError("hl:fetch: fetch() needs a URL string");}const url = argv[0].data.string.ptr[0..argv[0].data.string.len];const opts: ?HlValue = if (argc > 1 and argv[1].type == .hl_object) argv[1] else null;const method_raw = dupString(if (opts) |o| fieldOf(o, "method") else null, "GET");for (method_raw) |*ch| {if (ch.* >= 'a' and ch.* <= 'z') ch.* -= 32;}const body = dupString(if (opts) |o| fieldOf(o, "body") else null, "");const has_body = blk: {const bv = if (opts) |o| fieldOf(o, "body") else null;if (bv) |v| break :blk v.type == .hl_string;break :blk false;};const ca_val = if (opts) |o| fieldOf(o, "caFile") else null;const ca_file: ?[]u8 = if (ca_val != null and ca_val.?.type == .hl_string)allocator.dupe(u8, ca_val.?.data.string.ptr[0..ca_val.?.data.string.len]) catch nullelsenull;const t = allocator.create(Ticket) catch return api.makeError("hl:fetch: out of memory");t.* = .{.req = .{.method = method_raw,.url = allocator.dupe(u8, url) catch return api.makeError("hl:fetch: out of memory"),.headers = readHeaders(if (opts) |o| fieldOf(o, "headers") else null),.body = body,.has_body = has_body,.json_body = boolOr(if (opts) |o| fieldOf(o, "jsonBody") else null, false),.timeout_ms = @max(numberOr(if (opts) |o| fieldOf(o, "timeoutMs") else null, DEFAULT_TIMEOUT_MS), 1),.max_redirects = @intCast(@max(numberOr(if (opts) |o| fieldOf(o, "maxRedirects") else null, DEFAULT_MAX_REDIRECTS), 0)),.ca_file = ca_file,.tag = dupString(if (opts) |o| fieldOf(o, "tag") else null, ""),.stream = boolOr(if (opts) |o| fieldOf(o, "stream") else null, false),},};if (t.req.stream) t.ev_wake_fd = http.makeWakeFd();activeAdd(t);t.thread = std.Thread.spawn(.{}, workerMain, .{t}) catch {// No thread: run it here rather than lying about having started. The call// then blocks, which is slower and still correct.workerMain(t);t.thread = null;return makeTicketHandle(t);};return makeTicketHandle(t);}fn makeTicketHandle(t: *Ticket) HlValue {const h = allocator.create(HlHandle) catch return api.makeError("hl:fetch: out of memory");h.* = .{.context = @ptrCast(t),.call_fn = &ticketCall,.close_fn = &ticketClose,.type_name = hlStr("pending fetch"),};return api.makeHandle(h);}// ── The completions source: a loop source with the mission-125 bell ──────────//// `for (ev of completions())` / `eventloop.register` — the SAME shape hl:http1's// request and ws-event sources have, so a finished fetch reaches an `on fetched(ev)`// handler through the loop's own `epoll_wait` instead of anyone polling. The queue// is filled by the worker threads and drained by the loop thread.const Completion = struct {aborted: bool = false,tag: []u8,url: []u8,status: u16,headers: []Header,body: []u8,err: []u8,/// non-null = a STREAM CHUNK frame ({tag, url, chunk}), not a completionchunk: ?[]u8 = null,};var comp_mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER;var comp_queue: std.ArrayListUnmanaged(Completion) = .empty;var comp_wake_fd: i32 = -1;var comp_enabled = std.atomic.Value(bool).init(false);/// Called by every worker as its last act. Nothing is copied — and no wake is/// rung — unless a completions source actually exists, so a program that only/// uses `result()` pays nothing for this./// One decoded piece of a STREAMING body → a `{ tag, url, chunk }` frame on the/// completions source, bell rung. Dropped (with one log line) when no source is/// armed — a stream nobody listens to is a programming error worth seeing.fn postChunk(t: *Ticket, bytes: []const u8) void {if (t.ev_wake_fd < 0 and !comp_enabled.load(.acquire)) {logMsg("hl:fetch: stream chunk dropped — nobody is listening");return;}const comp = Completion{.tag = allocator.dupe(u8, t.req.tag) catch return,.url = allocator.dupe(u8, t.req.url) catch return,.status = 0,.headers = &.{},.body = allocator.dupe(u8, "") catch return,.err = allocator.dupe(u8, "") catch return,.chunk = allocator.dupe(u8, bytes) catch return,};postFrame(t, comp);}/// One frame to wherever this ticket's listener lives: the ticket's own source/// (stream mode) or the global completions source.fn postFrame(t: *Ticket, comp: Completion) void {if (t.ev_wake_fd >= 0) {mutexLock(&t.mutex);t.ev_queue.append(allocator, comp) catch {};mutexUnlock(&t.mutex);http.ringWake(t.ev_wake_fd);return;}mutexLock(&comp_mutex);comp_queue.append(allocator, comp) catch {};mutexUnlock(&comp_mutex);http.ringWake(comp_wake_fd);}fn postCompletion(t: *Ticket) void {if (t.ev_wake_fd < 0 and !comp_enabled.load(.acquire)) return;var headers = allocator.alloc(Header, t.headers.len) catch return;var built: usize = 0;for (t.headers, 0..) |h, i| {const n = allocator.dupe(u8, h.name) catch break;const v = allocator.dupe(u8, h.value) catch {allocator.free(n);break;};headers[i] = .{ .name = n, .value = v };built = i + 1;}headers = headers[0..built];const comp = Completion{.aborted = t.aborted,.tag = allocator.dupe(u8, t.req.tag) catch "",.url = allocator.dupe(u8, if (t.final_url.len > 0) t.final_url else t.req.url) catch "",.status = t.status,.headers = headers,.body = allocator.dupe(u8, t.body) catch "",.err = allocator.dupe(u8, if (t.err) |e| e else "") catch "",};// Ring generously (D-125): the loop treats the fd as a bell, so a wake for a// queue somebody already drained costs one empty round and nothing else.postFrame(t, comp);}/// Frees the SHELLS and the payload: a completion is a snapshot nobody else owns.fn compObjDeinit(obj: *HlObject) callconv(.c) void {const fields = obj.fields[0..obj.field_count];for (fields) |f| {if (f.value.type == .hl_string) {const s = f.value.data.string;if (s.len > 0) allocator.free(@constCast(s.ptr[0..s.len]));} else if (f.value.type == .hl_object) {const inner = f.value.data.object;for (inner.fields[0..inner.field_count]) |hf| {allocator.free(@constCast(hf.key.ptr[0..hf.key.len]));const hv = hf.value.data.string;allocator.free(@constCast(hv.ptr[0..hv.len]));}allocator.free(inner.fields[0..inner.field_count]);allocator.destroy(inner);}}allocator.free(fields);allocator.destroy(obj);}fn completionsTryNext(ctx: ?*anyopaque) callconv(.c) HlValue {_ = ctx;mutexLock(&comp_mutex);if (comp_queue.items.len == 0) {mutexUnlock(&comp_mutex);return api.makeNull();}const comp = comp_queue.orderedRemove(0);mutexUnlock(&comp_mutex);return frameToValue(comp);}/// A streaming ticket's own source: drains that ticket's queue and nothing else.fn eventsTryNext(ctx: ?*anyopaque) callconv(.c) HlValue {const t: *Ticket = @ptrCast(@alignCast(ctx orelse return api.makeNull()));mutexLock(&t.mutex);if (t.ev_queue.items.len == 0) {mutexUnlock(&t.mutex);return api.makeNull();}const comp = t.ev_queue.orderedRemove(0);mutexUnlock(&t.mutex);return frameToValue(comp);}fn frameToValue(comp: Completion) HlValue {// a STREAM CHUNK frame: { tag, url, chunk } and nothing elseif (comp.chunk) |chunk| {allocator.free(comp.body);allocator.free(comp.err);const cfields = allocator.alloc(HlField, 3) catch return api.makeNull();cfields[0] = .{ .key = hlStr("tag"), .value = api.makeString(comp.tag) };cfields[1] = .{ .key = hlStr("url"), .value = api.makeString(comp.url) };cfields[2] = .{ .key = hlStr("chunk"), .value = api.makeString(chunk) };const cobj = allocator.create(HlObject) catch {allocator.free(cfields);return api.makeNull();};cobj.* = .{ .fields = cfields.ptr, .field_count = 3, .deinit_fn = &compObjDeinit };return api.makeObject(cobj);}const hdr_fields = allocator.alloc(HlField, comp.headers.len) catch return api.makeNull();for (comp.headers, 0..) |h, i| {hdr_fields[i] = .{ .key = hlStr(h.name), .value = api.makeString(h.value) };}const hdr_obj = allocator.create(HlObject) catch {allocator.free(hdr_fields);return api.makeNull();};hdr_obj.* = .{ .fields = hdr_fields.ptr, .field_count = comp.headers.len, .deinit_fn = null };allocator.free(comp.headers);const fields = allocator.alloc(HlField, 8) catch return api.makeNull();fields[0] = .{ .key = hlStr("tag"), .value = api.makeString(comp.tag) };fields[1] = .{ .key = hlStr("url"), .value = api.makeString(comp.url) };fields[2] = .{ .key = hlStr("status"), .value = api.makeNumber(@floatFromInt(comp.status)) };fields[3] = .{ .key = hlStr("ok"), .value = api.makeBool(comp.status >= 200 and comp.status < 300) };fields[4] = .{ .key = hlStr("body"), .value = api.makeString(comp.body) };fields[5] = .{ .key = hlStr("error"), .value = api.makeString(comp.err) };fields[6] = .{ .key = hlStr("headers"), .value = api.makeObject(hdr_obj) };fields[7] = .{ .key = hlStr("aborted"), .value = api.makeBool(comp.aborted) };const obj = allocator.create(HlObject) catch {allocator.free(fields);return api.makeNull();};obj.* = .{ .fields = fields.ptr, .field_count = 8, .deinit_fn = &compObjDeinit };return api.makeObject(obj);}/// __native("fetch.abort", tag) — THE INTERRUPT (the LLM stop button): every/// in-flight fetch whose `tag` matches is aborted by shutting its socket down./// The peer sees the disconnect (llama-server / OpenRouter stop generating and/// billing), the worker's read unblocks, and the FINAL completion frame arrives/// with `aborted = true` and error "hl:fetch: aborted". Returns how many were/// aborted (0 = nothing in flight under that tag — already finished, or a typo).export fn hl_fetch_abort(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {if (argc < 1 or argv[0].type != .hl_string) {return api.makeError("hl:fetch abort: pass the tag the fetch was started with");}const tag = argv[0].data.string.ptr[0..argv[0].data.string.len];var n: u32 = 0;mutexLock(&active_mutex);for (active_tickets.items) |t| {if (std.mem.eql(u8, t.req.tag, tag)) {t.abort();n += 1;}}mutexUnlock(&active_mutex);return api.makeNumber(@floatFromInt(n));}/// __native("fetch.completions") → the loop source. Calling it ARMS completion/// delivery for every fetch started from here on.export fn hl_fetch_completions(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {_ = argc;_ = argv;mutexLock(&comp_mutex);if (comp_wake_fd < 0) comp_wake_fd = http.makeWakeFd();mutexUnlock(&comp_mutex);comp_enabled.store(true, .release);const iter = allocator.create(HlIterator) catch return api.makeError("hl:fetch: out of memory");iter.* = .{.context = null,.next_fn = &completionsTryNext, // non-blocking either way — event loop only.deinit_fn = null,.try_next_fn = &completionsTryNext,.wake_fd = comp_wake_fd,};return api.makeIterator(iter);}
Branches
- mainmain branch
Latest commits
- 93dfb0batracker#40 (mission 034): the curated franchises — data/franchises.json (17 franchises, 31 timelines, 285 TMDB titles, movies + series, in-universe/release order, 12 TMDB collections attached); lib/franchiseseed.hl + jobs.hl franchiseSeedTick (last start job, imports missing titles via details.hl importWithCredits = the search's Add, one per step paced, then one build; franchiseseed.db: editor changes win, the creator's same-name franchise adopted / timeline left alone, 404 remembered, resumable, idempotent); timeline heads 'N titles · in-universe order' (orderKind) and wrap on a phone; series pages show the widget; new gate tests/franchiseseed.mjs (5th in deploy.sh), the others run with TRACKER_FRANCHISE_SEED=0; gates 365/0, 32/0, 52/0, 24/0, check-theme 0; real copy 196 imported, 0 failed, 7 min, restart unchanged=31mre
- 96ba683adeploy.sh: a gate without a 'passed,' line (check-theme) no longer ends the scriptmre
- eb3b9205tracker: report 031mre
- 9b5d2e89tracker mission 031: README (What it does, Test: four gates + the #32 checks, Files: theme/, new pages), STATUS (real copy, A/B load, how to repeat, open points), LOGmre
- 39950e4ctracker#32 (mission 031): the WorldAPI theme (theme/ vendored verbatim from layouts.worldapi.org 85b5654; styles.hl inherits it: accent green-dark, type colours 1-6; own base/header rules, row lines, genre-pill and inverted-button frames removed, the season foldable keeps its line; check-theme 21 -> 0, 4th deploy gate; main actions class primary) and the #32 header (theme AppHeader/MainMenu/UserMenu/Sidebar/ContentFirst: desktop brand, search, Series|Shows|Movies|Genres|People, user icon with Unwatched..Settings, Logout; signed out the ident selector, phone the iD icon dropdown; phone menu in the sidebar overlay; marked entry by :has); /find -> /search/<q>, /genres, /people(/<letter>), /settings; main { ContentFirst { slot } } works around the hl:web one-line slot bug; gates 365/0, 32/0, 52/0, check-theme 0mre
- a386dc92tracker: reports 029 + 030mre
- 71e0fd7dtracker missions 029 + 030: README (What it does, Files, gate count), STATUS (real-copy numbers, how to repeat, open points), LOGmre
- d36ea6eatracker#34 + #35 (mission 030): Follow directly under the poster, as wide as the poster (show.hl, styles.hl); the status pill next to a series' title — TVmaze's status (new tvmazeStatus, stored by the sync's TVmaze merge) else TMDB's, TVmaze Ended + TMDB Canceled = Canceled, inverted (filled, dark text, no border), green running / yellow pending / red canceled / muted ended (shows.hl statusOf); the daily delta asks TVmaze's status of an unfollowed series TVmaze's change list names (dailysync.hl syncRunStep, sync.hl syncTvmazeStatus); the status backfill after the details repair (backfill.hl, jobs.hl statusTick; resumable, 550 ms per TVmaze request); gates 354/0, 32/0, 52/0mre
- 7d7d4487tracker#33 (mission 029): reduced titles — every title TMDB's details never went through this app (no detailsAt, no tmdbSync) is incomplete (shows.hl isIncomplete; the old tracker's migrated rows passed #26's test: 5,697 non-adult on the live copy, 691 series without seasons); the repair job does the visibly reduced first (shows.hl missingParts), the page completes one on open; a title TMDB has no poster for (The Remaining) shows the placeholder; tools/count-incomplete.hl; gate fixtures stand for synced titles (tmdbSync), tests/seed-reduced.hl + #33 checks; gates 347/0, 32/0, 52/0mre
- 661c2592tracker: report 028mre
- 27c916fatracker mission 028: README ("Code order", the new file map), STATUS (counts before/after, tests, how to repeat, open), LOGmre
- d924f398tracker mission 028: comments name the new files (sync.hl, dailysync.hl, backfill.hl, credits.hl, jobs.hl, images.hl …); tools/ref-params.py + tools/lambda-audit.py also scan lib/ (they globbed the root only), lambda-audit counts a plain `x = p` alias like `let x = p`mre
- 2e89b968tracker mission 028 (code order) 5/5 let: `let` only where a variable is reassigned — 667 never-reassigned lets became plain declarations (project.hl, lib/, components/, tools/, tests/); kept: 264 in loop bodies (a plain declaration there is 'Cannot reassign' on the 2nd pass), 234 reassigned, 27 whose name is also a member/outer/free name (a plain write would rebind it); tools/let-audit.py decides and fixes (README 'Code order'); tests/realdata-m028.{sh,mjs} = the page-output diff on a real copy; gates 342/0, 32/0, 52/0, real-copy pages identicalmre
- 54796ff2tracker mission 028 (code order) 4/5 thin faces + last copies: the show page's check/follow faces call lib/watches.hl toggleWatched / toggleSeasonWatched (seasonAllWatched moved there) and lib/follows.hl toggleFollowed; both logins (header selector face, /login/callback) share lib/users.hl userOfCode; todayStr/listOf copies in components and the export readers copied into tools/migrate.hl + tools/old-short-ids.hl now once (lib/util.hl, lib/export.hl); gates 342/0, 32/0, 52/0; old-short-ids output byte-identical, migrate output identicalmre
- 06b078e3tracker mission 028 (code order) 3/5 project.hl is the map: config, routes, wiring and a feature → file index (914 → 258 lines); the background jobs (daily sync run, backfills, details repair, credits job, merge, short ids, collection seed) moved unchanged into lib/jobs.hl (a class: their state is reassigned every step, a static cannot be; one instance made after the server), the login callback into lib/users.hl, poster/photo serving into lib/images.hl, the /shows/<slug> rule into lib/shows.hl showsMovedPath; route handlers are thin wrappers; gates 342/0, 32/0, 52/0, real-copy pages identicalmre
- 94716fd2tracker mission 028 (code order) 2/5 util + topics: lib/util.hl holds envOr, storageDir, postersDir, profilesDir, newId, hexDigits, todayStr, dateOr, textOr, hasId, listOr, firstOf, sortDesc once (were copied into up to 5 files); tmdbsync.hl split into tmdb.hl (TMDB/TVmaze requests), sync.hl (one title's sync), sync-helpers.hl, backfill.hl; details.hl split into details.hl, credits.hl, credits-helpers.hl (isIncomplete to shows.hl); search-helpers.hl (words, query, ranking, slugs); collections.hl (the TMDB collection seed, out of franchises.hl); deltasync.hl renamed dailysync.hl; no behaviour change: gates 342/0, 32/0, 52/0, real-copy pages identicalmre
- 186079b0tracker mission 028 (code order) 1/5 move: every root .hl except project.hl into lib/ (styles.hl into components/), import paths only; gates 342/0, 32/0, 52/0; real-copy pages identicalmre
- 4f47f181tracker: report 027mre
- dc1d4be4tracker mission 027: Hybriel master 06617221 vendored (plugin allocator fixes 3a781359 + 413f60e4); real copy RSS through first-start jobs + 400 loads flat ~2.55 GB (190aa11d 2.3 -> 5.6 GB), page times <= 1.1x; gates 342/0, 32/0, 52/0mre
- 84e1b3e1tracker: reports 025 + 026mre