gitoriaLog in with ident

calendar

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commit3b2a4cd03b2a4cd0calendar mission 001 (3/4): code order — let only where reassigned (147 dropped incl. components/month-view; 110 left: 67 reassigned, 43 loop-bound); gates 78/0 (3 of 4 runs; 1 known 'next month' flake) + 18/0, live-data run = step 2mre3b2a4cd0/plugins/fetch/fetch.zig

57.9 KB

  1. // hl:fetch — the HTTP(S) CLIENT, as a server-realm plugin (mission 135).
  2. //
  3. // Shaped on the web `fetch` the creator already knows: a method, a url, headers and
  4. // a body go out; a status, headers and a body come back. What is NOT web-shaped is
  5. // where the wait happens, and that is the one design decision in this file:
  6. //
  7. // `fetch.start` NEVER WAITS. It hands the request to a thread of its own and
  8. // returns a TICKET (an HlHandle) immediately. `ticket.result()` is the wait.
  9. //
  10. // That split is what makes `concurrent [ fetchStart(a), fetchStart(b), … ]` put N
  11. // requests on the wire AT ONCE: a `concurrent` arm runs to completion before its
  12. // sibling starts (nothing in the interpreter yields today), so an arm that blocked
  13. // on the socket would serialize the whole group. An arm that only starts a thread
  14. // does not. `fetch()` in server.hl is `fetchStart(…).response()` — start plus wait,
  15. // the familiar one-liner, for code that has nothing else to do.
  16. //
  17. // THE EVENT LOOP is never blocked by a fetch in flight: no loop thread is inside a
  18. // socket call, and a completed request rings the mission-125 BELL — `fetch.completions`
  19. // is an ordinary loop source with a `wake_fd`, so an `on fetched(ev)` handler is woken
  20. // by the same `epoll_wait` that wakes an http1 request. (The wait in `result()` is a
  21. // wait like any other plugin call's: it blocks the FIBER that asked. Starting the
  22. // fetches early and reading them later is what keeps a server serving, and the
  23. // completions source is how a server never has to wait at all.)
  24. //
  25. // TLS is `plugins/http/tls_common.zig` — the same module hl:http1 and hl:http2 use,
  26. // extended with a client half in this mission. Verification is ON: system trust
  27. // store plus hostname checking, with `caFile` naming an EXTRA anchor for a test
  28. // fixture's self-signed certificate. There is no verify-off switch at all.
  29. //
  30. // HTTP/1.1 only. An h2 client is a later slice: it needs ALPN on the client side and
  31. // a multiplexed connection pool, neither of which this file has. `connection: close`
  32. // on every request, so there is no pool to keep coherent either.
  33. const std = @import("std");
  34. const api = @import("plugin_api");
  35. const http = @import("http_common");
  36. const tls = @import("tls_common");
  37. const HlValue = api.HlValue;
  38. const HlField = api.HlField;
  39. const HlObject = api.HlObject;
  40. const HlIterator = api.HlIterator;
  41. const HlHandle = api.HlHandle;
  42. const c = @cImport({
  43. @cInclude("netdb.h");
  44. @cInclude("sys/socket.h");
  45. @cInclude("sys/time.h");
  46. @cInclude("unistd.h");
  47. @cInclude("netinet/in.h");
  48. @cInclude("netinet/tcp.h");
  49. });
  50. // requests run on their own threads: the plugins' allocator (plugin_api.zig)
  51. const allocator = api.allocator;
  52. const linux = std.os.linux;
  53. const libc = std.c;
  54. // pthread rather than std.Thread.Mutex, exactly as hl:http1 does: this is a `.so`
  55. // loaded into a binary built by a possibly different Zig, and the C primitives are
  56. // the ones whose ABI is fixed by the platform rather than by the standard library.
  57. const PthreadMutex = libc.pthread_mutex_t;
  58. const PthreadCond = libc.pthread_cond_t;
  59. fn mutexLock(m: *PthreadMutex) void {
  60. _ = libc.pthread_mutex_lock(m);
  61. }
  62. fn mutexUnlock(m: *PthreadMutex) void {
  63. _ = libc.pthread_mutex_unlock(m);
  64. }
  65. fn condBroadcast(cnd: *PthreadCond) void {
  66. _ = libc.pthread_cond_broadcast(cnd);
  67. }
  68. fn condWait(cnd: *PthreadCond, m: *PthreadMutex) void {
  69. _ = libc.pthread_cond_wait(cnd, m);
  70. }
  71. fn logMsg(msg: []const u8) void {
  72. _ = linux.write(2, msg.ptr, msg.len);
  73. }
  74. fn logFmt(comptime fmt: []const u8, args: anytype) void {
  75. var buf: [512]u8 = undefined;
  76. const s = std.fmt.bufPrint(&buf, fmt, args) catch return;
  77. logMsg(s);
  78. }
  79. // ── Limits, all of them documented defaults rather than hard walls ───────────
  80. /// No answer at all within this many milliseconds is a failed fetch. Overridable
  81. /// per call with `timeoutMs`; it covers the WHOLE exchange (DNS, connect,
  82. /// handshake, request, response, and every redirect hop), not one syscall.
  83. const DEFAULT_TIMEOUT_MS: i64 = 30_000;
  84. /// How many 3xx hops are followed before the fetch fails with "too many redirects".
  85. /// Overridable with `maxRedirects`; 0 means the 3xx is RETURNED as the response.
  86. const DEFAULT_MAX_REDIRECTS: u32 = 5;
  87. /// A response body larger than this fails rather than growing the process without
  88. /// bound — a client that pulls other people's APIs must not be a memory bomb.
  89. const MAX_BODY: usize = 32 * 1024 * 1024;
  90. /// O_NONBLOCK on x86-64 Linux. Spelled here because <fcntl.h> cannot be imported
  91. /// (see tcpConnect); it is an ABI constant, not a libc detail.
  92. const O_NONBLOCK: usize = 0o4000;
  93. const Header = struct {
  94. name: []u8,
  95. value: []u8,
  96. };
  97. fn freeHeaders(list: []Header) void {
  98. for (list) |h| {
  99. allocator.free(h.name);
  100. allocator.free(h.value);
  101. }
  102. allocator.free(list);
  103. }
  104. // ── The request, as the ticket carries it ────────────────────────────────────
  105. const Request = struct {
  106. method: []u8,
  107. url: []u8,
  108. headers: []Header,
  109. body: []u8,
  110. has_body: bool,
  111. json_body: bool,
  112. timeout_ms: i64,
  113. max_redirects: u32,
  114. ca_file: ?[]u8,
  115. tag: []u8,
  116. /// STREAMING (the LLM-token mode): deliver the body INCREMENTALLY as chunk
  117. /// events on the completions source instead of accumulating it. The final
  118. /// completion then carries status/headers and an EMPTY body — the chunks
  119. /// were the body. Requires an armed completions source; `result()` still
  120. /// works and yields the empty-body summary when the stream ends.
  121. stream: bool = false,
  122. fn deinit(self: *Request) void {
  123. allocator.free(self.method);
  124. allocator.free(self.url);
  125. freeHeaders(self.headers);
  126. allocator.free(self.body);
  127. if (self.ca_file) |ca| allocator.free(ca);
  128. allocator.free(self.tag);
  129. }
  130. };
  131. // ── The ticket: one in-flight (or finished) request ──────────────────────────
  132. const Ticket = struct {
  133. mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER,
  134. cond: PthreadCond = libc.PTHREAD_COND_INITIALIZER,
  135. done: bool = false,
  136. thread: ?std.Thread = null,
  137. req: Request,
  138. status: u16 = 0,
  139. final_url: []u8 = &.{},
  140. headers: []Header = &.{},
  141. body: []u8 = &.{},
  142. /// Non-null = the fetch failed and this is the located message the handle
  143. /// hands back as an `hl_error`.
  144. err: ?[]u8 = null,
  145. /// ABORT (the LLM-interrupt primitive): the socket currently carrying this
  146. /// request, visible under the ticket mutex so `fetchAbort(tag)` can
  147. /// shutdown() it from the interpreter thread — the read unblocks, the peer
  148. /// sees the disconnect and stops generating. -1 = none in flight.
  149. live_fd: i32 = -1,
  150. aborted: bool = false,
  151. /// INSTANCE EVENTS (stream mode): a streaming ticket carries its OWN frame
  152. /// queue and bell, created at start — before the worker exists — so no chunk
  153. /// can outrun the arming. `handle.events()` wraps them into a loop source;
  154. /// pending_fetch.hl registers it and re-emits `chunk`/`done` at itself.
  155. /// -1 = not a streaming ticket (completions go the global route, if armed).
  156. ev_queue: std.ArrayListUnmanaged(Completion) = .empty,
  157. ev_wake_fd: i32 = -1,
  158. fn wait(self: *Ticket) void {
  159. mutexLock(&self.mutex);
  160. while (!self.done) condWait(&self.cond, &self.mutex);
  161. mutexUnlock(&self.mutex);
  162. }
  163. fn isDone(self: *Ticket) bool {
  164. mutexLock(&self.mutex);
  165. defer mutexUnlock(&self.mutex);
  166. return self.done;
  167. }
  168. fn setLiveFd(self: *Ticket, fd: i32) bool {
  169. mutexLock(&self.mutex);
  170. defer mutexUnlock(&self.mutex);
  171. if (self.aborted) return false; // aborted before the connect landed
  172. self.live_fd = fd;
  173. return true;
  174. }
  175. fn clearLiveFd(self: *Ticket) void {
  176. mutexLock(&self.mutex);
  177. self.live_fd = -1;
  178. mutexUnlock(&self.mutex);
  179. }
  180. /// The abort itself: mark, and shutdown() any in-flight socket so the
  181. /// worker's blocking read returns NOW. Closing is still the worker's job
  182. /// (its defer owns the fd); shutdown only cuts the conversation.
  183. fn abort(self: *Ticket) void {
  184. mutexLock(&self.mutex);
  185. self.aborted = true;
  186. if (self.live_fd >= 0) _ = c.shutdown(self.live_fd, 2); // SHUT_RDWR
  187. mutexUnlock(&self.mutex);
  188. }
  189. fn finish(self: *Ticket) void {
  190. mutexLock(&self.mutex);
  191. self.done = true;
  192. condBroadcast(&self.cond);
  193. mutexUnlock(&self.mutex);
  194. }
  195. fn fail(self: *Ticket, comptime fmt: []const u8, args: anytype) void {
  196. self.err = std.fmt.allocPrint(allocator, fmt, args) catch null;
  197. }
  198. };
  199. // ── URL ──────────────────────────────────────────────────────────────────────
  200. const Url = struct {
  201. https: bool,
  202. host: []const u8,
  203. port: u16,
  204. /// path + query, always starting with '/'
  205. target: []const u8,
  206. };
  207. /// `scheme://host[:port][/path][?query]`. Only http and https; a URL with
  208. /// credentials or an IPv6 literal is refused rather than half-understood.
  209. fn parseUrl(url: []const u8) ?Url {
  210. var rest: []const u8 = undefined;
  211. var https = false;
  212. if (std.mem.startsWith(u8, url, "http://")) {
  213. rest = url["http://".len..];
  214. } else if (std.mem.startsWith(u8, url, "https://")) {
  215. rest = url["https://".len..];
  216. https = true;
  217. } else return null;
  218. const slash = std.mem.indexOfScalar(u8, rest, '/');
  219. const authority = if (slash) |s| rest[0..s] else rest;
  220. const target = if (slash) |s| rest[s..] else "/";
  221. if (authority.len == 0) return null;
  222. if (std.mem.indexOfScalar(u8, authority, '@') != null) return null;
  223. if (std.mem.indexOfScalar(u8, authority, '[') != null) return null;
  224. var host = authority;
  225. var port: u16 = if (https) 443 else 80;
  226. if (std.mem.lastIndexOfScalar(u8, authority, ':')) |colon| {
  227. host = authority[0..colon];
  228. port = std.fmt.parseInt(u16, authority[colon + 1 ..], 10) catch return null;
  229. }
  230. if (host.len == 0) return null;
  231. return .{ .https = https, .host = host, .port = port, .target = target };
  232. }
  233. /// Resolve a `Location` against the URL it came from: absolute stays, `/x` replaces
  234. /// the path, anything else is relative to the current directory of the path.
  235. fn resolveLocation(base: []const u8, loc: []const u8) ?[]u8 {
  236. if (loc.len == 0) return null;
  237. if (std.mem.startsWith(u8, loc, "http://") or std.mem.startsWith(u8, loc, "https://")) {
  238. return allocator.dupe(u8, loc) catch null;
  239. }
  240. const b = parseUrl(base) orelse return null;
  241. const scheme = if (b.https) "https" else "http";
  242. const default_port: u16 = if (b.https) 443 else 80;
  243. var authority_buf: [300]u8 = undefined;
  244. const authority = if (b.port == default_port)
  245. std.fmt.bufPrint(&authority_buf, "{s}", .{b.host}) catch return null
  246. else
  247. std.fmt.bufPrint(&authority_buf, "{s}:{d}", .{ b.host, b.port }) catch return null;
  248. if (loc[0] == '/') {
  249. return std.fmt.allocPrint(allocator, "{s}://{s}{s}", .{ scheme, authority, loc }) catch null;
  250. }
  251. // Relative: cut the base target back to its last '/'.
  252. const path_only = if (std.mem.indexOfScalar(u8, b.target, '?')) |q| b.target[0..q] else b.target;
  253. const cut = std.mem.lastIndexOfScalar(u8, path_only, '/') orelse 0;
  254. return std.fmt.allocPrint(allocator, "{s}://{s}{s}/{s}", .{ scheme, authority, path_only[0..cut], loc }) catch null;
  255. }
  256. // ── Socket plumbing ──────────────────────────────────────────────────────────
  257. fn nowMs() i64 {
  258. var ts: linux.timespec = undefined;
  259. _ = linux.clock_gettime(linux.CLOCK.MONOTONIC, &ts);
  260. return @as(i64, ts.sec) * 1000 + @divTrunc(@as(i64, ts.nsec), 1_000_000);
  261. }
  262. /// Connect to host:port, giving up after `budget_ms`. The connect is non-blocking
  263. /// plus a poll so the DEADLINE is honoured even when the peer never answers SYN;
  264. /// the socket is put back into blocking mode with SO_RCVTIMEO/SO_SNDTIMEO after,
  265. /// which is what bounds the read side.
  266. fn tcpConnect(host: []const u8, port: u16, budget_ms: i64, err_out: *[]const u8) ?i32 {
  267. if (budget_ms <= 0) {
  268. err_out.* = "timed out";
  269. return null;
  270. }
  271. var host_buf: [256]u8 = undefined;
  272. if (host.len >= host_buf.len) {
  273. err_out.* = "host name too long";
  274. return null;
  275. }
  276. @memcpy(host_buf[0..host.len], host);
  277. host_buf[host.len] = 0;
  278. var port_buf: [8]u8 = undefined;
  279. const port_str = std.fmt.bufPrintZ(&port_buf, "{d}", .{port}) catch {
  280. err_out.* = "bad port";
  281. return null;
  282. };
  283. var hints: c.struct_addrinfo = std.mem.zeroes(c.struct_addrinfo);
  284. hints.ai_family = c.AF_UNSPEC;
  285. hints.ai_socktype = c.SOCK_STREAM;
  286. hints.ai_protocol = c.IPPROTO_TCP;
  287. var res: ?*c.struct_addrinfo = null;
  288. if (c.getaddrinfo(@ptrCast(&host_buf), port_str.ptr, &hints, &res) != 0 or res == null) {
  289. err_out.* = "could not resolve host";
  290. return null;
  291. }
  292. defer c.freeaddrinfo(res);
  293. const deadline = nowMs() + budget_ms;
  294. var last: []const u8 = "connection failed";
  295. var it: ?*c.struct_addrinfo = res;
  296. while (it) |ai| : (it = ai.ai_next) {
  297. const remaining = deadline - nowMs();
  298. if (remaining <= 0) {
  299. err_out.* = "timed out";
  300. return null;
  301. }
  302. const fd_c = c.socket(ai.ai_family, ai.ai_socktype, ai.ai_protocol);
  303. if (fd_c < 0) {
  304. last = "could not create a socket";
  305. continue;
  306. }
  307. const fd: i32 = @intCast(fd_c);
  308. // fcntl straight from the kernel: glibc's <fcntl.h> does not survive
  309. // translate-c under _FORTIFY_SOURCE (its open() overloads are declared
  310. // with attribute error), and the flags are ABI constants either way.
  311. const flags = linux.fcntl(fd, linux.F.GETFL, 0);
  312. _ = linux.fcntl(fd, linux.F.SETFL, flags | @as(usize, O_NONBLOCK));
  313. var connected = false;
  314. if (c.connect(fd, ai.ai_addr, ai.ai_addrlen) == 0) {
  315. connected = true;
  316. } else {
  317. // linux.poll, not glibc's: <poll.h> is fortified too (see above).
  318. var pfd = [_]linux.pollfd{.{ .fd = fd, .events = linux.POLL.OUT, .revents = 0 }};
  319. const pr: isize = @bitCast(linux.poll(&pfd, 1, @intCast(@min(remaining, std.math.maxInt(i32)))));
  320. if (pr > 0) {
  321. var so_err: c_int = 0;
  322. var len: c.socklen_t = @sizeOf(c_int);
  323. _ = c.getsockopt(fd, c.SOL_SOCKET, c.SO_ERROR, &so_err, &len);
  324. if (so_err == 0) {
  325. connected = true;
  326. } else if (so_err == 111) {
  327. last = "connection refused";
  328. } else {
  329. last = "connection failed";
  330. }
  331. } else if (pr == 0) {
  332. _ = c.close(fd);
  333. err_out.* = "timed out";
  334. return null;
  335. } else {
  336. last = "connection failed";
  337. }
  338. }
  339. if (!connected) {
  340. _ = c.close(fd);
  341. continue;
  342. }
  343. _ = linux.fcntl(fd, linux.F.SETFL, flags);
  344. const left = @max(deadline - nowMs(), 1);
  345. var tv = c.struct_timeval{
  346. .tv_sec = @intCast(@divTrunc(left, 1000)),
  347. .tv_usec = @intCast(@rem(left, 1000) * 1000),
  348. };
  349. _ = c.setsockopt(fd, c.SOL_SOCKET, c.SO_RCVTIMEO, &tv, @sizeOf(c.struct_timeval));
  350. _ = c.setsockopt(fd, c.SOL_SOCKET, c.SO_SNDTIMEO, &tv, @sizeOf(c.struct_timeval));
  351. const one: c_int = 1;
  352. _ = c.setsockopt(fd, c.IPPROTO_TCP, c.TCP_NODELAY, &one, @sizeOf(c_int));
  353. return fd;
  354. }
  355. err_out.* = last;
  356. return null;
  357. }
  358. /// The one place either transport is read or written, so the HTTP exchange below
  359. /// is written once for http:// and https://.
  360. const Conn = struct {
  361. fd: i32,
  362. ssl: ?*tls.SSL = null,
  363. tls_ctx: ?*const tls.TlsContext = null,
  364. fn writeAll(self: *Conn, data: []const u8) bool {
  365. if (self.ssl) |ssl| {
  366. return self.tls_ctx.?.writeAll(ssl, data) == data.len;
  367. }
  368. var sent: usize = 0;
  369. while (sent < data.len) {
  370. const n = c.write(self.fd, data.ptr + sent, data.len - sent);
  371. if (n <= 0) return false;
  372. sent += @intCast(n);
  373. }
  374. return true;
  375. }
  376. /// >0 bytes, 0 = clean end of body, <0 = error (or a socket timeout, which the
  377. /// caller separates from an error by looking at the clock).
  378. fn read(self: *Conn, buf: []u8) isize {
  379. if (self.ssl) |ssl| {
  380. const n = self.tls_ctx.?.read(ssl, buf);
  381. if (n > 0) return n;
  382. const e = self.tls_ctx.?.getError(ssl, n);
  383. if (e == tls.SSL_ERROR_ZERO_RETURN) return 0;
  384. return -1;
  385. }
  386. return c.read(self.fd, buf.ptr, buf.len);
  387. }
  388. fn close(self: *Conn) void {
  389. if (self.ssl) |ssl| {
  390. self.tls_ctx.?.shutdownAndFree(ssl);
  391. self.ssl = null;
  392. }
  393. _ = c.close(self.fd);
  394. }
  395. };
  396. // ── The client TLS context, one per CA file ──────────────────────────────────
  397. //
  398. // An SSL_CTX is expensive (it reads the system trust store) and thread-safe, so it
  399. // is built once per distinct `caFile` and shared. The map is tiny by construction:
  400. // a program has a system-trust context and at most a handful of pinned fixtures.
  401. const CtxEntry = struct {
  402. key: []u8,
  403. ctx: tls.TlsContext,
  404. };
  405. var ctx_mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER;
  406. var ctx_list: std.ArrayListUnmanaged(CtxEntry) = .empty;
  407. fn clientCtx(ca_file: ?[]const u8, err_out: *[]const u8) ?*const tls.TlsContext {
  408. const key = ca_file orelse "";
  409. mutexLock(&ctx_mutex);
  410. defer mutexUnlock(&ctx_mutex);
  411. for (ctx_list.items) |*e| {
  412. if (std.mem.eql(u8, e.key, key)) return &e.ctx;
  413. }
  414. const ctx = tls.TlsContext.initClient(.{ .ca_file = ca_file, .tag = "fetch" }) orelse {
  415. err_out.* = "https is not available (no usable OpenSSL, or the caFile could not be loaded)";
  416. return null;
  417. };
  418. const key_copy = allocator.dupe(u8, key) catch {
  419. err_out.* = "out of memory";
  420. return null;
  421. };
  422. ctx_list.append(allocator, .{ .key = key_copy, .ctx = ctx }) catch {
  423. allocator.free(key_copy);
  424. err_out.* = "out of memory";
  425. return null;
  426. };
  427. return &ctx_list.items[ctx_list.items.len - 1].ctx;
  428. }
  429. // ── One HTTP exchange ────────────────────────────────────────────────────────
  430. const Exchange = struct {
  431. status: u16 = 0,
  432. headers: []Header = &.{},
  433. body: []u8 = &.{},
  434. fn deinit(self: *Exchange) void {
  435. freeHeaders(self.headers);
  436. allocator.free(self.body);
  437. }
  438. };
  439. fn headerValue(headers: []const Header, name: []const u8) ?[]const u8 {
  440. for (headers) |h| {
  441. if (std.ascii.eqlIgnoreCase(h.name, name)) return h.value;
  442. }
  443. return null;
  444. }
  445. fn hasHeader(headers: []const Header, name: []const u8) bool {
  446. return headerValue(headers, name) != null;
  447. }
  448. /// Send one request over one connection and read the whole response.
  449. /// `err_out` gets a short static reason; the caller adds the URL.
  450. /// The status code off a raw head, or null while unparseable.
  451. fn parseStatus(head: []const u8) ?u16 {
  452. var lines = std.mem.splitSequence(u8, head, "\r\n");
  453. const status_line = lines.next() orelse return null;
  454. var parts = std.mem.splitScalar(u8, status_line, ' ');
  455. _ = parts.next(); // HTTP/1.1
  456. const code = parts.next() orelse return null;
  457. return std.fmt.parseInt(u16, code, 10) catch null;
  458. }
  459. /// Head only → an Exchange with an EMPTY body (the chunks were the body).
  460. fn parseHeadOnly(head_raw: []const u8, err_out: *[]const u8) ?Exchange {
  461. // reuse the full parser on a synthetic complete response: the head as-is
  462. // (its transfer-encoding stripped so no body is expected) + no body.
  463. var synth: std.ArrayListUnmanaged(u8) = .empty;
  464. defer synth.deinit(allocator);
  465. var lines = std.mem.splitSequence(u8, head_raw[0 .. head_raw.len - 4], "\r\n");
  466. while (lines.next()) |line| {
  467. if (line.len > 18 and std.ascii.eqlIgnoreCase(line[0..18], "transfer-encoding:")) continue;
  468. if (line.len > 15 and std.ascii.eqlIgnoreCase(line[0..15], "content-length:")) continue;
  469. if (!add(&synth, line)) return oom(err_out);
  470. if (!add(&synth, "\r\n")) return oom(err_out);
  471. }
  472. if (!add(&synth, "\r\n")) return oom(err_out);
  473. return parseResponse(synth.items, err_out);
  474. }
  475. /// STREAM the body: post every decoded piece as a chunk frame until the peer
  476. /// closes (`connection: close` is on every request, so close IS the end).
  477. /// Chunked transfer is decoded incrementally; anything else streams as-is.
  478. fn streamBody(t: *Ticket, conn: *Conn, raw: *std.ArrayListUnmanaged(u8), he: usize, deadline: i64, err_out: *[]const u8) ?Exchange {
  479. defer raw.deinit(allocator);
  480. const chunked = blk: {
  481. if (findHeaderIn(raw.items[0..he], "transfer-encoding")) |v| {
  482. break :blk std.ascii.indexOfIgnoreCase(v, "chunked") != null;
  483. }
  484. break :blk false;
  485. };
  486. // chunked-decoder state, carried across reads
  487. var pending: std.ArrayListUnmanaged(u8) = .empty;
  488. defer pending.deinit(allocator);
  489. var remaining: usize = 0; // bytes left of the current chunk's data
  490. var done = false;
  491. // seed with whatever body bytes arrived along with the head
  492. if (raw.items.len > he) {
  493. pending.appendSlice(allocator, raw.items[he..]) catch return oom(err_out);
  494. }
  495. var buf: [16 * 1024]u8 = undefined;
  496. while (true) {
  497. // decode + post what we hold
  498. if (!chunked) {
  499. if (pending.items.len > 0) {
  500. postChunk(t, pending.items);
  501. pending.clearRetainingCapacity();
  502. }
  503. } else while (!done) {
  504. if (remaining > 0) {
  505. const take = @min(remaining, pending.items.len);
  506. if (take == 0) break;
  507. postChunk(t, pending.items[0..take]);
  508. std.mem.copyForwards(u8, pending.items, pending.items[take..]);
  509. pending.shrinkRetainingCapacity(pending.items.len - take);
  510. remaining -= take;
  511. if (remaining == 0) {
  512. // the CRLF after the chunk data
  513. if (pending.items.len < 2) break;
  514. std.mem.copyForwards(u8, pending.items, pending.items[2..]);
  515. pending.shrinkRetainingCapacity(pending.items.len - 2);
  516. }
  517. continue;
  518. }
  519. const nl = std.mem.indexOf(u8, pending.items, "\r\n") orelse break;
  520. const size_line = pending.items[0..nl];
  521. const semi = std.mem.indexOfScalar(u8, size_line, ';') orelse size_line.len;
  522. const n = std.fmt.parseInt(usize, std.mem.trim(u8, size_line[0..semi], " \t"), 16) catch {
  523. err_out.* = "the peer sent a malformed chunked body";
  524. return null;
  525. };
  526. std.mem.copyForwards(u8, pending.items, pending.items[nl + 2 ..]);
  527. pending.shrinkRetainingCapacity(pending.items.len - nl - 2);
  528. if (n == 0) {
  529. done = true;
  530. break;
  531. }
  532. remaining = n;
  533. }
  534. if (done) break;
  535. const n = conn.read(&buf);
  536. if (n == 0) break; // close IS the end (identity), or a truncated chunked stream
  537. if (n < 0) {
  538. if (nowMs() >= deadline) {
  539. err_out.* = "timed out";
  540. return null;
  541. }
  542. break;
  543. }
  544. pending.appendSlice(allocator, buf[0..@intCast(n)]) catch return oom(err_out);
  545. }
  546. return parseHeadOnly(raw.items[0..he], err_out);
  547. }
  548. fn exchange(t: *Ticket, url: Url, method: []const u8, body: []const u8, send_body: bool, deadline: i64, err_out: *[]const u8) ?Exchange {
  549. const budget = deadline - nowMs();
  550. const fd = tcpConnect(url.host, url.port, budget, err_out) orelse return null;
  551. if (!t.setLiveFd(fd)) {
  552. _ = c.close(fd);
  553. err_out.* = "aborted";
  554. return null;
  555. }
  556. var conn = Conn{ .fd = fd };
  557. if (url.https) {
  558. const ctx = clientCtx(t.req.ca_file, err_out) orelse {
  559. _ = c.close(fd);
  560. return null;
  561. };
  562. // 137's shared client half logs the specific failure itself; the caller
  563. // gets the located summary (certificate not trusted, or not a TLS server)
  564. const ssl = ctx.connect(fd, url.host, "fetch") orelse {
  565. _ = c.close(fd);
  566. err_out.* = "TLS handshake failed (certificate not trusted, or not a TLS server)";
  567. return null;
  568. };
  569. conn.ssl = ssl;
  570. conn.tls_ctx = ctx;
  571. }
  572. defer conn.close();
  573. defer t.clearLiveFd();
  574. // --- request ---
  575. var req_buf: std.ArrayListUnmanaged(u8) = .empty;
  576. defer req_buf.deinit(allocator);
  577. const default_port: u16 = if (url.https) 443 else 80;
  578. var ok = true;
  579. ok = ok and addFmt(&req_buf, "{s} {s} HTTP/1.1\r\n", .{ method, url.target });
  580. if (url.port == default_port) {
  581. ok = ok and addFmt(&req_buf, "host: {s}\r\n", .{url.host});
  582. } else {
  583. ok = ok and addFmt(&req_buf, "host: {s}:{d}\r\n", .{ url.host, url.port });
  584. }
  585. // `connection: close` on every request: there is no connection pool, and a
  586. // closed connection is also the unambiguous end of a body with no length.
  587. ok = ok and add(&req_buf, "connection: close\r\n");
  588. if (!hasHeader(t.req.headers, "user-agent")) {
  589. ok = ok and add(&req_buf, "user-agent: hybriel-fetch/1\r\n");
  590. }
  591. if (!hasHeader(t.req.headers, "accept")) {
  592. ok = ok and add(&req_buf, "accept: */*\r\n");
  593. }
  594. if (send_body) {
  595. ok = ok and addFmt(&req_buf, "content-length: {d}\r\n", .{body.len});
  596. if (t.req.json_body and !hasHeader(t.req.headers, "content-type")) {
  597. ok = ok and add(&req_buf, "content-type: application/json\r\n");
  598. }
  599. }
  600. for (t.req.headers) |h| {
  601. // A caller header may not smuggle a second request in (CRLF injection).
  602. if (std.mem.indexOfAny(u8, h.name, "\r\n") != null or std.mem.indexOfAny(u8, h.value, "\r\n") != null) {
  603. err_out.* = "a header name or value contains a line break";
  604. return null;
  605. }
  606. ok = ok and addFmt(&req_buf, "{s}: {s}\r\n", .{ h.name, h.value });
  607. }
  608. ok = ok and add(&req_buf, "\r\n");
  609. if (send_body and body.len > 0) {
  610. ok = ok and add(&req_buf, body);
  611. }
  612. if (!ok) return oom(err_out);
  613. if (!conn.writeAll(req_buf.items)) {
  614. err_out.* = if (nowMs() >= deadline) "timed out" else "could not send the request";
  615. return null;
  616. }
  617. // --- response ---
  618. var raw: std.ArrayListUnmanaged(u8) = .empty;
  619. errdefer raw.deinit(allocator);
  620. var head_end: ?usize = null;
  621. var buf: [16 * 1024]u8 = undefined;
  622. while (true) {
  623. if (head_end == null) {
  624. if (std.mem.indexOf(u8, raw.items, "\r\n\r\n")) |idx| head_end = idx + 4;
  625. }
  626. if (head_end) |he| {
  627. // STREAMING (LLM tokens): the head is in — if this is the final
  628. // answer (not a redirect the caller follows), hand the rest of the
  629. // body over CHUNK BY CHUNK as it arrives instead of accumulating.
  630. if (t.req.stream) {
  631. const status_ok = parseStatus(raw.items[0..he]);
  632. if (status_ok != null and !isRedirect(status_ok.?)) {
  633. return streamBody(t, &conn, &raw, he, deadline, err_out);
  634. }
  635. }
  636. // Once the head is in, stop as soon as the body is provably complete.
  637. if (bodyComplete(raw.items, he)) break;
  638. }
  639. if (raw.items.len > MAX_BODY) {
  640. raw.deinit(allocator);
  641. err_out.* = "the response is larger than the 32 MiB limit";
  642. return null;
  643. }
  644. const n = conn.read(&buf);
  645. if (n == 0) break;
  646. if (n < 0) {
  647. if (nowMs() >= deadline) {
  648. raw.deinit(allocator);
  649. err_out.* = "timed out";
  650. return null;
  651. }
  652. // A TLS peer that just closes without close_notify, or a plain socket
  653. // reset after a complete answer: the parse below decides whether what
  654. // arrived is a whole response.
  655. break;
  656. }
  657. raw.appendSlice(allocator, buf[0..@intCast(n)]) catch {
  658. raw.deinit(allocator);
  659. return oom(err_out);
  660. };
  661. }
  662. const out = parseResponse(raw.items, err_out) orelse {
  663. raw.deinit(allocator);
  664. return null;
  665. };
  666. raw.deinit(allocator);
  667. return out;
  668. }
  669. fn add(buf: *std.ArrayListUnmanaged(u8), data: []const u8) bool {
  670. buf.appendSlice(allocator, data) catch return false;
  671. return true;
  672. }
  673. /// One request line or header, formatted. A header that does not fit 2 KiB is not
  674. /// a header anyone meant to send.
  675. fn addFmt(buf: *std.ArrayListUnmanaged(u8), comptime fmt: []const u8, args: anytype) bool {
  676. var tmp: [2048]u8 = undefined;
  677. const s = std.fmt.bufPrint(&tmp, fmt, args) catch return false;
  678. return add(buf, s);
  679. }
  680. fn oom(err_out: *[]const u8) ?Exchange {
  681. err_out.* = "out of memory";
  682. return null;
  683. }
  684. /// Is everything the response promised already in `raw`? Content-Length is exact,
  685. /// chunked ends at the zero chunk, and anything else ends when the peer closes.
  686. fn bodyComplete(raw: []const u8, head_end: usize) bool {
  687. const head = raw[0..head_end];
  688. if (findHeaderIn(head, "content-length")) |v| {
  689. const want = std.fmt.parseInt(usize, std.mem.trim(u8, v, " \t"), 10) catch return false;
  690. return raw.len - head_end >= want;
  691. }
  692. if (findHeaderIn(head, "transfer-encoding")) |v| {
  693. if (std.ascii.indexOfIgnoreCase(v, "chunked") != null) {
  694. return std.mem.indexOf(u8, raw[head_end..], "\r\n0\r\n") != null or
  695. std.mem.startsWith(u8, raw[head_end..], "0\r\n");
  696. }
  697. }
  698. return false;
  699. }
  700. fn findHeaderIn(head: []const u8, name: []const u8) ?[]const u8 {
  701. var lines = std.mem.splitSequence(u8, head, "\r\n");
  702. _ = lines.next(); // status line
  703. while (lines.next()) |line| {
  704. if (line.len == 0) break;
  705. const colon = std.mem.indexOfScalar(u8, line, ':') orelse continue;
  706. if (std.ascii.eqlIgnoreCase(std.mem.trim(u8, line[0..colon], " \t"), name)) {
  707. return std.mem.trim(u8, line[colon + 1 ..], " \t");
  708. }
  709. }
  710. return null;
  711. }
  712. fn parseResponse(raw: []const u8, err_out: *[]const u8) ?Exchange {
  713. const he = std.mem.indexOf(u8, raw, "\r\n\r\n") orelse {
  714. err_out.* = "the peer sent no complete HTTP response";
  715. return null;
  716. };
  717. const head = raw[0..he];
  718. const body_raw = raw[he + 4 ..];
  719. var lines = std.mem.splitSequence(u8, head, "\r\n");
  720. const status_line = lines.next() orelse {
  721. err_out.* = "the peer sent no status line";
  722. return null;
  723. };
  724. // "HTTP/1.1 200 OK"
  725. const sp = std.mem.indexOfScalar(u8, status_line, ' ') orelse {
  726. err_out.* = "the peer sent a malformed status line";
  727. return null;
  728. };
  729. const after = status_line[sp + 1 ..];
  730. const code_end = std.mem.indexOfScalar(u8, after, ' ') orelse after.len;
  731. const status = std.fmt.parseInt(u16, std.mem.trim(u8, after[0..code_end], " \t"), 10) catch {
  732. err_out.* = "the peer sent a malformed status line";
  733. return null;
  734. };
  735. var headers: std.ArrayListUnmanaged(Header) = .empty;
  736. errdefer {
  737. for (headers.items) |h| {
  738. allocator.free(h.name);
  739. allocator.free(h.value);
  740. }
  741. headers.deinit(allocator);
  742. }
  743. var chunked = false;
  744. while (lines.next()) |line| {
  745. if (line.len == 0) continue;
  746. const colon = std.mem.indexOfScalar(u8, line, ':') orelse continue;
  747. const name = std.mem.trim(u8, line[0..colon], " \t");
  748. const value = std.mem.trim(u8, line[colon + 1 ..], " \t");
  749. if (name.len == 0) continue;
  750. const name_copy = allocator.dupe(u8, name) catch return oom(err_out);
  751. http.toLowercase(name_copy);
  752. const value_copy = allocator.dupe(u8, value) catch {
  753. allocator.free(name_copy);
  754. return oom(err_out);
  755. };
  756. if (std.mem.eql(u8, name_copy, "transfer-encoding") and
  757. std.ascii.indexOfIgnoreCase(value_copy, "chunked") != null) chunked = true;
  758. headers.append(allocator, .{ .name = name_copy, .value = value_copy }) catch {
  759. allocator.free(name_copy);
  760. allocator.free(value_copy);
  761. return oom(err_out);
  762. };
  763. }
  764. var body: []u8 = undefined;
  765. if (chunked) {
  766. body = dechunk(body_raw) orelse {
  767. for (headers.items) |h| {
  768. allocator.free(h.name);
  769. allocator.free(h.value);
  770. }
  771. headers.deinit(allocator);
  772. err_out.* = "the peer sent a malformed chunked body";
  773. return null;
  774. };
  775. } else if (findHeaderIn(head, "content-length")) |v| {
  776. const want = std.fmt.parseInt(usize, v, 10) catch body_raw.len;
  777. const take = @min(want, body_raw.len);
  778. body = allocator.dupe(u8, body_raw[0..take]) catch return oom(err_out);
  779. } else {
  780. body = allocator.dupe(u8, body_raw) catch return oom(err_out);
  781. }
  782. const headers_slice = headers.toOwnedSlice(allocator) catch return oom(err_out);
  783. return .{ .status = status, .headers = headers_slice, .body = body };
  784. }
  785. fn dechunk(input: []const u8) ?[]u8 {
  786. var out: std.ArrayListUnmanaged(u8) = .empty;
  787. var i: usize = 0;
  788. while (true) {
  789. const line_end = std.mem.indexOfPos(u8, input, i, "\r\n") orelse {
  790. out.deinit(allocator);
  791. return null;
  792. };
  793. var size_txt = input[i..line_end];
  794. if (std.mem.indexOfScalar(u8, size_txt, ';')) |semi| size_txt = size_txt[0..semi];
  795. const size = std.fmt.parseInt(usize, std.mem.trim(u8, size_txt, " \t"), 16) catch {
  796. out.deinit(allocator);
  797. return null;
  798. };
  799. i = line_end + 2;
  800. if (size == 0) return out.toOwnedSlice(allocator) catch null;
  801. if (i + size > input.len) {
  802. out.deinit(allocator);
  803. return null;
  804. }
  805. out.appendSlice(allocator, input[i .. i + size]) catch {
  806. out.deinit(allocator);
  807. return null;
  808. };
  809. i += size + 2; // skip the chunk's trailing CRLF
  810. if (i > input.len) return out.toOwnedSlice(allocator) catch null;
  811. }
  812. }
  813. // ── The worker: one thread per fetch, redirects and all ───────────────────────
  814. fn isRedirect(status: u16) bool {
  815. return status == 301 or status == 302 or status == 303 or status == 307 or status == 308;
  816. }
  817. // ── The ACTIVE registry: every in-flight ticket, so `fetchAbort(tag)` can find
  818. // its socket. Added at spawn, removed at finish — a finished ticket is owned by
  819. // its handle alone and an abort of it is a no-op.
  820. var active_mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER;
  821. var active_tickets: std.ArrayListUnmanaged(*Ticket) = .empty;
  822. fn activeAdd(t: *Ticket) void {
  823. mutexLock(&active_mutex);
  824. active_tickets.append(allocator, t) catch {};
  825. mutexUnlock(&active_mutex);
  826. }
  827. fn activeRemove(t: *Ticket) void {
  828. mutexLock(&active_mutex);
  829. for (active_tickets.items, 0..) |a, i| {
  830. if (a == t) {
  831. _ = active_tickets.swapRemove(i);
  832. break;
  833. }
  834. }
  835. mutexUnlock(&active_mutex);
  836. }
  837. fn workerMain(t: *Ticket) void {
  838. runFetch(t);
  839. activeRemove(t);
  840. // an aborted request reports ABORTED whatever the read error looked like —
  841. // the app asked for the interrupt, the message says the interrupt happened
  842. if (t.aborted) {
  843. if (t.err) |e| allocator.free(e);
  844. t.err = std.fmt.allocPrint(allocator, "hl:fetch: aborted", .{}) catch null;
  845. }
  846. t.finish();
  847. postCompletion(t);
  848. }
  849. fn runFetch(t: *Ticket) void {
  850. const deadline = nowMs() + t.req.timeout_ms;
  851. var current = allocator.dupe(u8, t.req.url) catch {
  852. t.fail("hl:fetch: out of memory", .{});
  853. return;
  854. };
  855. var method = allocator.dupe(u8, t.req.method) catch {
  856. allocator.free(current);
  857. t.fail("hl:fetch: out of memory", .{});
  858. return;
  859. };
  860. var send_body = t.req.has_body;
  861. var hops: u32 = 0;
  862. while (true) {
  863. const url = parseUrl(current) orelse {
  864. t.fail("hl:fetch: '{s}' is not an http:// or https:// URL", .{current});
  865. allocator.free(current);
  866. allocator.free(method);
  867. return;
  868. };
  869. var reason: []const u8 = "connection failed";
  870. var ex = exchange(t, url, method, t.req.body, send_body, deadline, &reason) orelse {
  871. t.fail("hl:fetch {s} {s}: {s}", .{ method, current, reason });
  872. allocator.free(current);
  873. allocator.free(method);
  874. return;
  875. };
  876. if (isRedirect(ex.status) and hops < t.req.max_redirects) {
  877. if (headerValue(ex.headers, "location")) |loc| {
  878. const next = resolveLocation(current, loc) orelse {
  879. t.fail("hl:fetch {s}: redirect to an unusable Location '{s}'", .{ current, loc });
  880. ex.deinit();
  881. allocator.free(current);
  882. allocator.free(method);
  883. return;
  884. };
  885. // 303 always becomes a GET; 301/302 become one for anything that
  886. // is not already GET/HEAD — what every browser does, and what a
  887. // server that answers a POST with "see the result over there"
  888. // means. 307/308 keep the method AND the body by definition.
  889. if (ex.status == 303 or
  890. ((ex.status == 301 or ex.status == 302) and
  891. !std.mem.eql(u8, method, "GET") and !std.mem.eql(u8, method, "HEAD")))
  892. {
  893. allocator.free(method);
  894. method = allocator.dupe(u8, "GET") catch {
  895. allocator.free(next);
  896. ex.deinit();
  897. allocator.free(current);
  898. t.fail("hl:fetch: out of memory", .{});
  899. return;
  900. };
  901. send_body = false;
  902. }
  903. ex.deinit();
  904. allocator.free(current);
  905. current = next;
  906. hops += 1;
  907. continue;
  908. }
  909. }
  910. if (isRedirect(ex.status) and t.req.max_redirects > 0 and hops >= t.req.max_redirects) {
  911. t.fail("hl:fetch {s}: more than {d} redirects", .{ t.req.url, t.req.max_redirects });
  912. ex.deinit();
  913. allocator.free(current);
  914. allocator.free(method);
  915. return;
  916. }
  917. t.status = ex.status;
  918. t.headers = ex.headers;
  919. t.body = ex.body;
  920. t.final_url = current;
  921. allocator.free(method);
  922. return;
  923. }
  924. }
  925. // ── The value shapes handed back to Hybriel ──────────────────────────────────
  926. fn hlStr(s: []const u8) api.HlString {
  927. return .{ .ptr = s.ptr, .len = s.len };
  928. }
  929. /// Only the SHELLS are freed: every string in a `result()` object points into the
  930. /// ticket, which outlives the call (the loader copies what it keeps before this
  931. /// runs — the plugin ABI's ownership rule).
  932. fn respShellDeinit(obj: *HlObject) callconv(.c) void {
  933. const hdrs = obj.fields[3].value;
  934. if (hdrs.type == .hl_object) {
  935. const inner = hdrs.data.object;
  936. allocator.free(inner.fields[0..inner.field_count]);
  937. allocator.destroy(inner);
  938. }
  939. allocator.free(obj.fields[0..obj.field_count]);
  940. allocator.destroy(obj);
  941. }
  942. fn buildHeadersObject(headers: []const Header) ?*HlObject {
  943. const fields = allocator.alloc(HlField, headers.len) catch return null;
  944. for (headers, 0..) |h, i| {
  945. fields[i] = .{ .key = hlStr(h.name), .value = api.makeString(h.value) };
  946. }
  947. const obj = allocator.create(HlObject) catch {
  948. allocator.free(fields);
  949. return null;
  950. };
  951. obj.* = .{ .fields = fields.ptr, .field_count = headers.len, .deinit_fn = null };
  952. return obj;
  953. }
  954. /// `{ status, ok, url, headers, body, tag }` — the response as the .hl side sees
  955. /// it, hydrated into a FetchResponse by server.hl.
  956. fn buildResponse(t: *Ticket) HlValue {
  957. const headers_obj = buildHeadersObject(t.headers) orelse return api.makeError("hl:fetch: out of memory");
  958. const fields = allocator.alloc(HlField, 6) catch {
  959. allocator.free(headers_obj.fields[0..headers_obj.field_count]);
  960. allocator.destroy(headers_obj);
  961. return api.makeError("hl:fetch: out of memory");
  962. };
  963. fields[0] = .{ .key = hlStr("status"), .value = api.makeNumber(@floatFromInt(t.status)) };
  964. fields[1] = .{ .key = hlStr("ok"), .value = api.makeBool(t.status >= 200 and t.status < 300) };
  965. fields[2] = .{ .key = hlStr("url"), .value = api.makeString(t.final_url) };
  966. fields[3] = .{ .key = hlStr("headers"), .value = api.makeObject(headers_obj) };
  967. fields[4] = .{ .key = hlStr("body"), .value = api.makeString(t.body) };
  968. fields[5] = .{ .key = hlStr("tag"), .value = api.makeString(t.req.tag) };
  969. const obj = allocator.create(HlObject) catch {
  970. allocator.free(fields);
  971. allocator.free(headers_obj.fields[0..headers_obj.field_count]);
  972. allocator.destroy(headers_obj);
  973. return api.makeError("hl:fetch: out of memory");
  974. };
  975. obj.* = .{ .fields = fields.ptr, .field_count = 6, .deinit_fn = &respShellDeinit };
  976. return api.makeObject(obj);
  977. }
  978. // ── The ticket handle ────────────────────────────────────────────────────────
  979. fn ticketCall(ctx: ?*anyopaque, op: api.HlString, argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  980. _ = argc;
  981. _ = argv;
  982. const t: *Ticket = @ptrCast(@alignCast(ctx orelse return api.makeError("hl:fetch: dead ticket")));
  983. const name = op.ptr[0..op.len];
  984. if (std.mem.eql(u8, name, "done")) {
  985. return api.makeBool(t.isDone());
  986. }
  987. if (std.mem.eql(u8, name, "result")) {
  988. t.wait();
  989. if (t.err) |msg| return api.makeError(msg);
  990. return buildResponse(t);
  991. }
  992. if (std.mem.eql(u8, name, "events")) {
  993. if (t.ev_wake_fd < 0) return api.makeError("hl:fetch: events() is stream mode — start the fetch with stream = true");
  994. const iter = allocator.create(HlIterator) catch return api.makeError("hl:fetch: out of memory");
  995. iter.* = .{
  996. .context = @ptrCast(t),
  997. .next_fn = &eventsTryNext, // non-blocking either way — event loop only
  998. .deinit_fn = null,
  999. .try_next_fn = &eventsTryNext,
  1000. .wake_fd = t.ev_wake_fd,
  1001. };
  1002. return api.makeIterator(iter);
  1003. }
  1004. if (std.mem.eql(u8, name, "abort")) {
  1005. t.abort();
  1006. return api.makeBool(true);
  1007. }
  1008. return api.makeError("hl:fetch: a pending fetch has no method by that name");
  1009. }
  1010. fn ticketClose(ctx: ?*anyopaque) callconv(.c) void {
  1011. const t: *Ticket = @ptrCast(@alignCast(ctx orelse return));
  1012. // The worker writes into this ticket, so it is JOINED before anything is
  1013. // freed — a detached thread finishing after the free is a use-after-free the
  1014. // collector would surface as random corruption. A dropped, never-read fetch
  1015. // therefore costs at most its own timeout at collection time.
  1016. if (t.thread) |th| {
  1017. th.join();
  1018. t.thread = null;
  1019. }
  1020. t.req.deinit();
  1021. for (t.ev_queue.items) |comp| {
  1022. allocator.free(comp.tag);
  1023. allocator.free(comp.url);
  1024. allocator.free(comp.body);
  1025. allocator.free(comp.err);
  1026. if (comp.chunk) |ch| allocator.free(ch);
  1027. for (comp.headers) |h| {
  1028. allocator.free(h.name);
  1029. allocator.free(h.value);
  1030. }
  1031. if (comp.headers.len > 0) allocator.free(comp.headers);
  1032. }
  1033. t.ev_queue.deinit(allocator);
  1034. freeHeaders(t.headers);
  1035. allocator.free(t.body);
  1036. allocator.free(t.final_url);
  1037. if (t.err) |e| allocator.free(e);
  1038. allocator.destroy(t);
  1039. }
  1040. // ── Argument reading ─────────────────────────────────────────────────────────
  1041. fn fieldOf(v: HlValue, name: []const u8) ?HlValue {
  1042. if (v.type != .hl_object) return null;
  1043. const obj = v.data.object;
  1044. for (obj.fields[0..obj.field_count]) |f| {
  1045. if (std.mem.eql(u8, f.key.ptr[0..f.key.len], name)) return f.value;
  1046. }
  1047. return null;
  1048. }
  1049. fn dupString(v: ?HlValue, fallback: []const u8) []u8 {
  1050. if (v) |val| {
  1051. if (val.type == .hl_string) {
  1052. return allocator.dupe(u8, val.data.string.ptr[0..val.data.string.len]) catch
  1053. allocator.dupe(u8, fallback) catch unreachable;
  1054. }
  1055. }
  1056. return allocator.dupe(u8, fallback) catch unreachable;
  1057. }
  1058. fn numberOr(v: ?HlValue, fallback: i64) i64 {
  1059. if (v) |val| {
  1060. if (val.type == .hl_number) return @intFromFloat(val.data.number);
  1061. }
  1062. return fallback;
  1063. }
  1064. fn boolOr(v: ?HlValue, fallback: bool) bool {
  1065. if (v) |val| {
  1066. if (val.type == .hl_bool) return val.data.boolean;
  1067. }
  1068. return fallback;
  1069. }
  1070. fn readHeaders(v: ?HlValue) []Header {
  1071. const val = v orelse return &.{};
  1072. if (val.type != .hl_object) return &.{};
  1073. const obj = val.data.object;
  1074. var list: std.ArrayListUnmanaged(Header) = .empty;
  1075. for (obj.fields[0..obj.field_count]) |f| {
  1076. if (f.value.type != .hl_string) continue;
  1077. const name = allocator.dupe(u8, f.key.ptr[0..f.key.len]) catch continue;
  1078. const value = allocator.dupe(u8, f.value.data.string.ptr[0..f.value.data.string.len]) catch {
  1079. allocator.free(name);
  1080. continue;
  1081. };
  1082. list.append(allocator, .{ .name = name, .value = value }) catch {
  1083. allocator.free(name);
  1084. allocator.free(value);
  1085. };
  1086. }
  1087. return list.toOwnedSlice(allocator) catch &.{};
  1088. }
  1089. /// __native("fetch.start", url, options) → a ticket HANDLE, immediately.
  1090. /// Options: method, headers, body, jsonBody, timeoutMs, maxRedirects, caFile, tag.
  1091. export fn hl_fetch_start(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  1092. if (argc < 1 or argv[0].type != .hl_string) {
  1093. return api.makeError("hl:fetch: fetch() needs a URL string");
  1094. }
  1095. const url = argv[0].data.string.ptr[0..argv[0].data.string.len];
  1096. const opts: ?HlValue = if (argc > 1 and argv[1].type == .hl_object) argv[1] else null;
  1097. const method_raw = dupString(if (opts) |o| fieldOf(o, "method") else null, "GET");
  1098. for (method_raw) |*ch| {
  1099. if (ch.* >= 'a' and ch.* <= 'z') ch.* -= 32;
  1100. }
  1101. const body = dupString(if (opts) |o| fieldOf(o, "body") else null, "");
  1102. const has_body = blk: {
  1103. const bv = if (opts) |o| fieldOf(o, "body") else null;
  1104. if (bv) |v| break :blk v.type == .hl_string;
  1105. break :blk false;
  1106. };
  1107. const ca_val = if (opts) |o| fieldOf(o, "caFile") else null;
  1108. const ca_file: ?[]u8 = if (ca_val != null and ca_val.?.type == .hl_string)
  1109. allocator.dupe(u8, ca_val.?.data.string.ptr[0..ca_val.?.data.string.len]) catch null
  1110. else
  1111. null;
  1112. const t = allocator.create(Ticket) catch return api.makeError("hl:fetch: out of memory");
  1113. t.* = .{
  1114. .req = .{
  1115. .method = method_raw,
  1116. .url = allocator.dupe(u8, url) catch return api.makeError("hl:fetch: out of memory"),
  1117. .headers = readHeaders(if (opts) |o| fieldOf(o, "headers") else null),
  1118. .body = body,
  1119. .has_body = has_body,
  1120. .json_body = boolOr(if (opts) |o| fieldOf(o, "jsonBody") else null, false),
  1121. .timeout_ms = @max(numberOr(if (opts) |o| fieldOf(o, "timeoutMs") else null, DEFAULT_TIMEOUT_MS), 1),
  1122. .max_redirects = @intCast(@max(numberOr(if (opts) |o| fieldOf(o, "maxRedirects") else null, DEFAULT_MAX_REDIRECTS), 0)),
  1123. .ca_file = ca_file,
  1124. .tag = dupString(if (opts) |o| fieldOf(o, "tag") else null, ""),
  1125. .stream = boolOr(if (opts) |o| fieldOf(o, "stream") else null, false),
  1126. },
  1127. };
  1128. if (t.req.stream) t.ev_wake_fd = http.makeWakeFd();
  1129. activeAdd(t);
  1130. t.thread = std.Thread.spawn(.{}, workerMain, .{t}) catch {
  1131. // No thread: run it here rather than lying about having started. The call
  1132. // then blocks, which is slower and still correct.
  1133. workerMain(t);
  1134. t.thread = null;
  1135. return makeTicketHandle(t);
  1136. };
  1137. return makeTicketHandle(t);
  1138. }
  1139. fn makeTicketHandle(t: *Ticket) HlValue {
  1140. const h = allocator.create(HlHandle) catch return api.makeError("hl:fetch: out of memory");
  1141. h.* = .{
  1142. .context = @ptrCast(t),
  1143. .call_fn = &ticketCall,
  1144. .close_fn = &ticketClose,
  1145. .type_name = hlStr("pending fetch"),
  1146. };
  1147. return api.makeHandle(h);
  1148. }
  1149. // ── The completions source: a loop source with the mission-125 bell ──────────
  1150. //
  1151. // `for (ev of completions())` / `eventloop.register` — the SAME shape hl:http1's
  1152. // request and ws-event sources have, so a finished fetch reaches an `on fetched(ev)`
  1153. // handler through the loop's own `epoll_wait` instead of anyone polling. The queue
  1154. // is filled by the worker threads and drained by the loop thread.
  1155. const Completion = struct {
  1156. aborted: bool = false,
  1157. tag: []u8,
  1158. url: []u8,
  1159. status: u16,
  1160. headers: []Header,
  1161. body: []u8,
  1162. err: []u8,
  1163. /// non-null = a STREAM CHUNK frame ({tag, url, chunk}), not a completion
  1164. chunk: ?[]u8 = null,
  1165. };
  1166. var comp_mutex: PthreadMutex = libc.PTHREAD_MUTEX_INITIALIZER;
  1167. var comp_queue: std.ArrayListUnmanaged(Completion) = .empty;
  1168. var comp_wake_fd: i32 = -1;
  1169. var comp_enabled = std.atomic.Value(bool).init(false);
  1170. /// Called by every worker as its last act. Nothing is copied — and no wake is
  1171. /// rung — unless a completions source actually exists, so a program that only
  1172. /// uses `result()` pays nothing for this.
  1173. /// One decoded piece of a STREAMING body → a `{ tag, url, chunk }` frame on the
  1174. /// completions source, bell rung. Dropped (with one log line) when no source is
  1175. /// armed — a stream nobody listens to is a programming error worth seeing.
  1176. fn postChunk(t: *Ticket, bytes: []const u8) void {
  1177. if (t.ev_wake_fd < 0 and !comp_enabled.load(.acquire)) {
  1178. logMsg("hl:fetch: stream chunk dropped — nobody is listening");
  1179. return;
  1180. }
  1181. const comp = Completion{
  1182. .tag = allocator.dupe(u8, t.req.tag) catch return,
  1183. .url = allocator.dupe(u8, t.req.url) catch return,
  1184. .status = 0,
  1185. .headers = &.{},
  1186. .body = allocator.dupe(u8, "") catch return,
  1187. .err = allocator.dupe(u8, "") catch return,
  1188. .chunk = allocator.dupe(u8, bytes) catch return,
  1189. };
  1190. postFrame(t, comp);
  1191. }
  1192. /// One frame to wherever this ticket's listener lives: the ticket's own source
  1193. /// (stream mode) or the global completions source.
  1194. fn postFrame(t: *Ticket, comp: Completion) void {
  1195. if (t.ev_wake_fd >= 0) {
  1196. mutexLock(&t.mutex);
  1197. t.ev_queue.append(allocator, comp) catch {};
  1198. mutexUnlock(&t.mutex);
  1199. http.ringWake(t.ev_wake_fd);
  1200. return;
  1201. }
  1202. mutexLock(&comp_mutex);
  1203. comp_queue.append(allocator, comp) catch {};
  1204. mutexUnlock(&comp_mutex);
  1205. http.ringWake(comp_wake_fd);
  1206. }
  1207. fn postCompletion(t: *Ticket) void {
  1208. if (t.ev_wake_fd < 0 and !comp_enabled.load(.acquire)) return;
  1209. var headers = allocator.alloc(Header, t.headers.len) catch return;
  1210. var built: usize = 0;
  1211. for (t.headers, 0..) |h, i| {
  1212. const n = allocator.dupe(u8, h.name) catch break;
  1213. const v = allocator.dupe(u8, h.value) catch {
  1214. allocator.free(n);
  1215. break;
  1216. };
  1217. headers[i] = .{ .name = n, .value = v };
  1218. built = i + 1;
  1219. }
  1220. headers = headers[0..built];
  1221. const comp = Completion{
  1222. .aborted = t.aborted,
  1223. .tag = allocator.dupe(u8, t.req.tag) catch "",
  1224. .url = allocator.dupe(u8, if (t.final_url.len > 0) t.final_url else t.req.url) catch "",
  1225. .status = t.status,
  1226. .headers = headers,
  1227. .body = allocator.dupe(u8, t.body) catch "",
  1228. .err = allocator.dupe(u8, if (t.err) |e| e else "") catch "",
  1229. };
  1230. // Ring generously (D-125): the loop treats the fd as a bell, so a wake for a
  1231. // queue somebody already drained costs one empty round and nothing else.
  1232. postFrame(t, comp);
  1233. }
  1234. /// Frees the SHELLS and the payload: a completion is a snapshot nobody else owns.
  1235. fn compObjDeinit(obj: *HlObject) callconv(.c) void {
  1236. const fields = obj.fields[0..obj.field_count];
  1237. for (fields) |f| {
  1238. if (f.value.type == .hl_string) {
  1239. const s = f.value.data.string;
  1240. if (s.len > 0) allocator.free(@constCast(s.ptr[0..s.len]));
  1241. } else if (f.value.type == .hl_object) {
  1242. const inner = f.value.data.object;
  1243. for (inner.fields[0..inner.field_count]) |hf| {
  1244. allocator.free(@constCast(hf.key.ptr[0..hf.key.len]));
  1245. const hv = hf.value.data.string;
  1246. allocator.free(@constCast(hv.ptr[0..hv.len]));
  1247. }
  1248. allocator.free(inner.fields[0..inner.field_count]);
  1249. allocator.destroy(inner);
  1250. }
  1251. }
  1252. allocator.free(fields);
  1253. allocator.destroy(obj);
  1254. }
  1255. fn completionsTryNext(ctx: ?*anyopaque) callconv(.c) HlValue {
  1256. _ = ctx;
  1257. mutexLock(&comp_mutex);
  1258. if (comp_queue.items.len == 0) {
  1259. mutexUnlock(&comp_mutex);
  1260. return api.makeNull();
  1261. }
  1262. const comp = comp_queue.orderedRemove(0);
  1263. mutexUnlock(&comp_mutex);
  1264. return frameToValue(comp);
  1265. }
  1266. /// A streaming ticket's own source: drains that ticket's queue and nothing else.
  1267. fn eventsTryNext(ctx: ?*anyopaque) callconv(.c) HlValue {
  1268. const t: *Ticket = @ptrCast(@alignCast(ctx orelse return api.makeNull()));
  1269. mutexLock(&t.mutex);
  1270. if (t.ev_queue.items.len == 0) {
  1271. mutexUnlock(&t.mutex);
  1272. return api.makeNull();
  1273. }
  1274. const comp = t.ev_queue.orderedRemove(0);
  1275. mutexUnlock(&t.mutex);
  1276. return frameToValue(comp);
  1277. }
  1278. fn frameToValue(comp: Completion) HlValue {
  1279. // a STREAM CHUNK frame: { tag, url, chunk } and nothing else
  1280. if (comp.chunk) |chunk| {
  1281. allocator.free(comp.body);
  1282. allocator.free(comp.err);
  1283. const cfields = allocator.alloc(HlField, 3) catch return api.makeNull();
  1284. cfields[0] = .{ .key = hlStr("tag"), .value = api.makeString(comp.tag) };
  1285. cfields[1] = .{ .key = hlStr("url"), .value = api.makeString(comp.url) };
  1286. cfields[2] = .{ .key = hlStr("chunk"), .value = api.makeString(chunk) };
  1287. const cobj = allocator.create(HlObject) catch {
  1288. allocator.free(cfields);
  1289. return api.makeNull();
  1290. };
  1291. cobj.* = .{ .fields = cfields.ptr, .field_count = 3, .deinit_fn = &compObjDeinit };
  1292. return api.makeObject(cobj);
  1293. }
  1294. const hdr_fields = allocator.alloc(HlField, comp.headers.len) catch return api.makeNull();
  1295. for (comp.headers, 0..) |h, i| {
  1296. hdr_fields[i] = .{ .key = hlStr(h.name), .value = api.makeString(h.value) };
  1297. }
  1298. const hdr_obj = allocator.create(HlObject) catch {
  1299. allocator.free(hdr_fields);
  1300. return api.makeNull();
  1301. };
  1302. hdr_obj.* = .{ .fields = hdr_fields.ptr, .field_count = comp.headers.len, .deinit_fn = null };
  1303. allocator.free(comp.headers);
  1304. const fields = allocator.alloc(HlField, 8) catch return api.makeNull();
  1305. fields[0] = .{ .key = hlStr("tag"), .value = api.makeString(comp.tag) };
  1306. fields[1] = .{ .key = hlStr("url"), .value = api.makeString(comp.url) };
  1307. fields[2] = .{ .key = hlStr("status"), .value = api.makeNumber(@floatFromInt(comp.status)) };
  1308. fields[3] = .{ .key = hlStr("ok"), .value = api.makeBool(comp.status >= 200 and comp.status < 300) };
  1309. fields[4] = .{ .key = hlStr("body"), .value = api.makeString(comp.body) };
  1310. fields[5] = .{ .key = hlStr("error"), .value = api.makeString(comp.err) };
  1311. fields[6] = .{ .key = hlStr("headers"), .value = api.makeObject(hdr_obj) };
  1312. fields[7] = .{ .key = hlStr("aborted"), .value = api.makeBool(comp.aborted) };
  1313. const obj = allocator.create(HlObject) catch {
  1314. allocator.free(fields);
  1315. return api.makeNull();
  1316. };
  1317. obj.* = .{ .fields = fields.ptr, .field_count = 8, .deinit_fn = &compObjDeinit };
  1318. return api.makeObject(obj);
  1319. }
  1320. /// __native("fetch.abort", tag) — THE INTERRUPT (the LLM stop button): every
  1321. /// in-flight fetch whose `tag` matches is aborted by shutting its socket down.
  1322. /// The peer sees the disconnect (llama-server / OpenRouter stop generating and
  1323. /// billing), the worker's read unblocks, and the FINAL completion frame arrives
  1324. /// with `aborted = true` and error "hl:fetch: aborted". Returns how many were
  1325. /// aborted (0 = nothing in flight under that tag — already finished, or a typo).
  1326. export fn hl_fetch_abort(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  1327. if (argc < 1 or argv[0].type != .hl_string) {
  1328. return api.makeError("hl:fetch abort: pass the tag the fetch was started with");
  1329. }
  1330. const tag = argv[0].data.string.ptr[0..argv[0].data.string.len];
  1331. var n: u32 = 0;
  1332. mutexLock(&active_mutex);
  1333. for (active_tickets.items) |t| {
  1334. if (std.mem.eql(u8, t.req.tag, tag)) {
  1335. t.abort();
  1336. n += 1;
  1337. }
  1338. }
  1339. mutexUnlock(&active_mutex);
  1340. return api.makeNumber(@floatFromInt(n));
  1341. }
  1342. /// __native("fetch.completions") → the loop source. Calling it ARMS completion
  1343. /// delivery for every fetch started from here on.
  1344. export fn hl_fetch_completions(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  1345. _ = argc;
  1346. _ = argv;
  1347. mutexLock(&comp_mutex);
  1348. if (comp_wake_fd < 0) comp_wake_fd = http.makeWakeFd();
  1349. mutexUnlock(&comp_mutex);
  1350. comp_enabled.store(true, .release);
  1351. const iter = allocator.create(HlIterator) catch return api.makeError("hl:fetch: out of memory");
  1352. iter.* = .{
  1353. .context = null,
  1354. .next_fn = &completionsTryNext, // non-blocking either way — event loop only
  1355. .deinit_fn = null,
  1356. .try_next_fn = &completionsTryNext,
  1357. .wake_fd = comp_wake_fd,
  1358. };
  1359. return api.makeIterator(iter);
  1360. }

Branches

Latest commits

  • 3b2a4cd0calendar mission 001 (3/4): code order — let only where reassigned (147 dropped incl. components/month-view; 110 left: 67 reassigned, 43 loop-bound); gates 78/0 (3 of 4 runs; 1 known 'next month' flake) + 18/0, live-data run = step 2mre
  • f5ce6b1dcalendar mission 001 (2/4): code order — topics, map, thin faces: lib/util.hl, lib/events(-helpers).hl (+calendarView/settingsView/soonOf), lib/settings.hl, lib/users.hl (+userOfLoginCode, tagOf), lib/api(-helpers).hl; project.hl = map; login route takes &req/&sessions (failed-login reason now kept); gates 78/0 + 18/0mre
  • c6fbd011calendar mission 001 (1/4): code order — files moved: lib/events.hl, lib/users.hl, components/styles.hl (imports only); gates 78/0 + 18/0, live-data run identicalmre
  • d6c59290calendar: Hybriel master 06617221 (plugin allocators 3a781359 + 413f60e4, mpackdb 2cb7ae5e, http1 773de63e); gates 78/0 + 18/0mre
  • 7bd0337ccalendar: Hybriel master 190aa11d (fc838894 GC correctness, #126 closure scopes, #127); gates 78/0 + 18/0mre
  • ff41310ccalendar: Hybriel master 8efba065 (#126 memory, #48 lambda copy; audit: no & needed)mre
  • 14ba08c7antcolony#40: mission references point to the moved missionsmre
  • 99c73346antcolony#40: history (LOG.md), worker briefs (missions/) and reports moved here from antcolony, numbered per project; old numbers in antcolony docs/mission-map.mdmre
  • 76edaa62calendar: Hybriel master ff51cf46 (re-vendor round, static workaround removed)mre
  • 90a3fc2cdeploy.sh: back up live storage/.sessions/.env before every deploy (newest 5 kept)mre
  • 6722b72ddeploy.sh: never send .git or .gitignore to Byrodinmre
  • be099807State of 2026-09-27, before the move to gitoriamre