gitoriaLog in with ident

antcolony

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commita6af7883a6af7883tracker: worker box sees calendar.worldapi.org (login to copy)mrea6af7883/plugins/fetch/fetch.zig

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

Branches

Latest commits

  • a6af7883tracker: worker box sees calendar.worldapi.org (login to copy)mre
  • c613d26btemplates: bridges to external components (login.js for ident's selector) are allowed (creator 2026-09-27)mre
  • 9062978ctracker: worker box sees /media/STORAGE/projects/old-tracker read-only (tracker#2 source data)mre
  • 7f9660eeState of 2026-09-27, before the move to gitoriamre