gitoriaLog in with ident

tracker

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commiteb3b9205eb3b9205tracker: report 031mreeb3b9205/plugins/smtp/smtp.zig

51.8 KB

  1. // hl:smtp plugin — an SMTP submission client (libsmtp.so, dlopen'd by the runtime)
  2. //
  3. // Mission 137. Sending mail, and nothing else: there is no receiving (IMAP/POP),
  4. // no DKIM signing and no queue-with-retry here — see plugins/smtp/README.md for
  5. // why each of those is somewhere else's job.
  6. //
  7. // SHAPE
  8. // hl_smtp_open(host, port, options) → Number mailer id
  9. // hl_smtp_send(id, message) → Number job id, or an hl_error REFUSAL
  10. // hl_smtp_results(id) → the result source (iterator + wake fd)
  11. // hl_smtp_close(id) → null
  12. //
  13. // THE SPLIT BETWEEN THE THREADS is the whole design:
  14. //
  15. // • `send()` runs on the Hybriel thread and does everything that can FAIL
  16. // LOCALLY — validating the addresses, refusing a header injection, rendering
  17. // the RFC 5322 message. Those are located errors at the caller's own line,
  18. // which they can only be if they happen before the job leaves.
  19. // • the WORKER thread does the conversation (connect, EHLO, STARTTLS, AUTH,
  20. // MAIL/RCPT/DATA, QUIT). Nothing on the Hybriel thread waits for it: the
  21. // result goes onto a queue and the queue's eventfd is rung, which is mission
  22. // 125's bell — the event loop is parked in one `epoll_wait` over every
  23. // source's fd and comes back with the answer.
  24. //
  25. // So a server that sends a mail keeps serving while the mail is in flight, and a
  26. // script that only sends mail drains the same source with `for (r of m.results())`.
  27. // One source, the two consumptions the runtime already gives every plugin source
  28. // (`hl:http1`'s server iterator has been both since mission 012/052).
  29. //
  30. // TLS rides plugins/http/tls_common.zig — the SAME OpenSSL layer the servers use,
  31. // through the client half mission 137 added to it. There is no second TLS here.
  32. const std = @import("std");
  33. const api = @import("plugin_api");
  34. const http = @import("http_common");
  35. const tls = @import("tls_common");
  36. const HlValue = api.HlValue;
  37. const HlObject = api.HlObject;
  38. const HlField = api.HlField;
  39. const HlIterator = api.HlIterator;
  40. const HlString = api.HlString;
  41. const linux = std.os.linux;
  42. const c = std.c;
  43. const PthreadMutex = c.pthread_mutex_t;
  44. const PthreadCond = c.pthread_cond_t;
  45. // per-mailer state is touched by two threads: the plugins' allocator (plugin_api.zig)
  46. const allocator = api.allocator;
  47. fn mutexLock(m: *PthreadMutex) void {
  48. _ = c.pthread_mutex_lock(m);
  49. }
  50. fn mutexUnlock(m: *PthreadMutex) void {
  51. _ = c.pthread_mutex_unlock(m);
  52. }
  53. fn condSignal(cnd: *PthreadCond) void {
  54. _ = c.pthread_cond_signal(cnd);
  55. }
  56. fn condBroadcast(cnd: *PthreadCond) void {
  57. _ = c.pthread_cond_broadcast(cnd);
  58. }
  59. fn condWait(cnd: *PthreadCond, m: *PthreadMutex) void {
  60. _ = c.pthread_cond_wait(cnd, m);
  61. }
  62. /// Direct syscall for stderr — std.debug.print pulls in std.Progress, whose
  63. /// global state is ABI-incompatible when a .so is loaded into a differently
  64. /// built binary (the same reason tls_common.zig has its own).
  65. fn logMsg(msg: []const u8) void {
  66. _ = linux.write(2, msg.ptr, msg.len);
  67. }
  68. // ── REFUSALS ────────────────────────────────────────────────────────────────
  69. //
  70. // An hl_error returned at the top level of a plugin call becomes a LOCATED
  71. // runtime error at the `__native` call site and is offered to `on Error(e)`
  72. // first (plugin_loader / native_bridge). That is the channel every local refusal
  73. // here uses, so `m.send(…)` with a newline in the subject points at the line the
  74. // author wrote and is catchable exactly like any other plugin failure.
  75. //
  76. // The text is rendered into a module buffer: a plugin call never yields, and the
  77. // runtime dupes the string into its own tracker the instant the call returns.
  78. var err_buf: [1024]u8 = undefined;
  79. fn refuse(comptime fmt: []const u8, args: anytype) HlValue {
  80. const s = std.fmt.bufPrint(&err_buf, fmt, args) catch "hl:smtp refused the message";
  81. return api.makeError(s);
  82. }
  83. // ── ARGUMENT READING ────────────────────────────────────────────────────────
  84. fn fieldOf(v: HlValue, key: []const u8) ?HlValue {
  85. if (v.type != .hl_object) return null;
  86. const obj = v.data.object;
  87. for (obj.fields[0..obj.field_count]) |f| {
  88. if (std.mem.eql(u8, f.key.ptr[0..f.key.len], key)) return f.value;
  89. }
  90. return null;
  91. }
  92. fn strField(v: HlValue, key: []const u8) ?[]const u8 {
  93. const f = fieldOf(v, key) orelse return null;
  94. if (f.type != .hl_string) return null;
  95. return f.data.string.ptr[0..f.data.string.len];
  96. }
  97. fn boolField(v: HlValue, key: []const u8, dflt: bool) bool {
  98. const f = fieldOf(v, key) orelse return dflt;
  99. return switch (f.type) {
  100. .hl_bool => f.data.boolean,
  101. else => dflt,
  102. };
  103. }
  104. fn numField(v: HlValue, key: []const u8, dflt: f64) f64 {
  105. const f = fieldOf(v, key) orelse return dflt;
  106. if (f.type != .hl_number) return dflt;
  107. return f.data.number;
  108. }
  109. /// THE SENDER'S ADDRESS, under either name.
  110. ///
  111. /// `from` is a RESERVED WORD in Hybriel (`import X from './x.hl'`), so
  112. /// `{ from = "a@b" }` does not parse and `msg.from` is not a member access.
  113. /// Rather than force every caller to write the quoted-key form, both spellings
  114. /// are read here: `"from" = …` (the quoted key, which parses fine and is what
  115. /// mission 137's brief names) and `sender = …` (the bare-identifier form, for a
  116. /// caller who would rather not quote). Same field, one meaning; the quoted one
  117. /// wins if a message somehow carries both.
  118. fn senderField(v: HlValue) ?[]const u8 {
  119. if (strField(v, "from")) |s| return s;
  120. return strField(v, "sender");
  121. }
  122. fn dupOpt(s: ?[]const u8) ?[]u8 {
  123. const src = s orelse return null;
  124. if (src.len == 0) return null;
  125. return allocator.dupe(u8, src) catch null;
  126. }
  127. // ── CONFIGURATION ───────────────────────────────────────────────────────────
  128. const TlsMode = enum { starttls, implicit, none };
  129. const AuthMode = enum { auto, plain, login, none };
  130. // ── THE JOB AND ITS RESULT ──────────────────────────────────────────────────
  131. const Job = struct {
  132. id: u32,
  133. from: []u8,
  134. /// Envelope recipients — every RCPT TO the conversation issues.
  135. rcpt: [][]u8,
  136. /// The rendered RFC 5322 message: CRLF line endings, dot-stuffed, WITHOUT
  137. /// the terminating ".". Built on the Hybriel thread by `send`.
  138. data: []u8,
  139. fn deinit(self: *Job) void {
  140. allocator.free(self.from);
  141. for (self.rcpt) |r| allocator.free(r);
  142. allocator.free(self.rcpt);
  143. allocator.free(self.data);
  144. }
  145. };
  146. const Result = struct {
  147. id: u32,
  148. ok: bool,
  149. /// The last SMTP reply code the conversation saw (0 = never got one).
  150. code: u16,
  151. /// Which step answered: "connect", "greeting", "ehlo", "starttls", "auth",
  152. /// "mail", "rcpt", "data", "body", "quit", "done".
  153. stage: []const u8,
  154. /// The server's own text, or this plugin's reason for giving up.
  155. message: []u8,
  156. /// Did the conversation run inside TLS by the time DATA was sent?
  157. secure: bool,
  158. rcpt: [][]u8,
  159. };
  160. // ── THE MAILER ──────────────────────────────────────────────────────────────
  161. const Mailer = struct {
  162. id: u32,
  163. host: []u8,
  164. port: u16,
  165. user: ?[]u8,
  166. pass: ?[]u8,
  167. from: []u8,
  168. helo: []u8,
  169. tls_mode: TlsMode,
  170. auth_mode: AuthMode,
  171. ca_file: ?[]u8,
  172. verify: bool,
  173. /// Send AUTH over a connection that never became TLS. OFF by default: the
  174. /// credentials would be on the wire in the clear, and a default that leaks
  175. /// them is not a default, it is a bug with a setting.
  176. allow_insecure_auth: bool,
  177. timeout_ms: u32,
  178. mutex: PthreadMutex,
  179. /// The worker parks here for the next job; the drain parks here for the next
  180. /// result. One lock covers both queues and `outstanding`, which is what makes
  181. /// "no results and nothing outstanding" a single consistent observation.
  182. condvar: PthreadCond,
  183. jobs: std.ArrayListUnmanaged(Job) = .empty,
  184. results: std.ArrayListUnmanaged(Result) = .empty,
  185. /// Jobs accepted and not yet answered. The blocking drain ends when this is
  186. /// zero and the result queue is empty — that is what lets a script that only
  187. /// sends mail terminate.
  188. outstanding: usize = 0,
  189. next_job: u32 = 1,
  190. stop: bool = false,
  191. /// Mission 125's bell for the event loop. Rung whenever a result becomes
  192. /// visible; over-ringing is free, a missed ring is the only bug there is.
  193. wake_fd: i32 = -1,
  194. worker: ?std.Thread = null,
  195. /// The source object is made ONCE: two iterators over one queue would each
  196. /// steal half the results.
  197. source: ?*HlIterator = null,
  198. registered: bool = false,
  199. fn enqueue(self: *Mailer, job: Job) void {
  200. mutexLock(&self.mutex);
  201. defer mutexUnlock(&self.mutex);
  202. self.jobs.append(allocator, job) catch return;
  203. self.outstanding += 1;
  204. condBroadcast(&self.condvar);
  205. }
  206. fn takeJob(self: *Mailer) ?Job {
  207. mutexLock(&self.mutex);
  208. defer mutexUnlock(&self.mutex);
  209. while (self.jobs.items.len == 0 and !self.stop) {
  210. condWait(&self.condvar, &self.mutex);
  211. }
  212. if (self.jobs.items.len == 0) return null;
  213. return self.jobs.orderedRemove(0);
  214. }
  215. fn publish(self: *Mailer, result: Result) void {
  216. mutexLock(&self.mutex);
  217. self.results.append(allocator, result) catch {};
  218. if (self.outstanding > 0) self.outstanding -= 1;
  219. condBroadcast(&self.condvar);
  220. // Rung while the append is still under the lock — the loop drains the
  221. // eventfd counter only immediately AFTER a wake, never before blocking,
  222. // so a bell rung any time after that drain leaves the level up and the
  223. // next epoll_wait returns at once (native/src/loop_wait.zig).
  224. http.ringWake(self.wake_fd);
  225. mutexUnlock(&self.mutex);
  226. }
  227. fn tryTake(self: *Mailer) ?Result {
  228. mutexLock(&self.mutex);
  229. defer mutexUnlock(&self.mutex);
  230. if (self.results.items.len == 0) return null;
  231. return self.results.orderedRemove(0);
  232. }
  233. /// Blocking take for the drain form. `null` means DONE, not "nothing yet":
  234. /// every accepted job has been answered and its answer taken.
  235. fn take(self: *Mailer) ?Result {
  236. mutexLock(&self.mutex);
  237. defer mutexUnlock(&self.mutex);
  238. while (self.results.items.len == 0 and self.outstanding > 0 and !self.stop) {
  239. condWait(&self.condvar, &self.mutex);
  240. }
  241. if (self.results.items.len == 0) return null;
  242. return self.results.orderedRemove(0);
  243. }
  244. };
  245. var mailers: std.ArrayListUnmanaged(*Mailer) = .empty;
  246. var mailers_mutex: PthreadMutex = c.PTHREAD_MUTEX_INITIALIZER;
  247. var next_mailer_id: u32 = 1;
  248. fn mailerById(id: u32) ?*Mailer {
  249. mutexLock(&mailers_mutex);
  250. defer mutexUnlock(&mailers_mutex);
  251. for (mailers.items) |m| {
  252. if (m.id == id) return m;
  253. }
  254. return null;
  255. }
  256. // ── HEADER INJECTION IS A REFUSAL, NOT AN ESCAPE ────────────────────────────
  257. //
  258. // A CR or an LF inside a `to`, a `from` or a `subject` ends the header and
  259. // starts a new one — that is the whole of SMTP header injection, and it is how a
  260. // contact form becomes an open relay (`subject = "hi\r\nBcc: everyone@…"`).
  261. // Stripping the characters silently would deliver a mail the author did not
  262. // write; escaping them is not defined for a header field. So it is a located
  263. // refusal, at the caller's line, naming the field and the character.
  264. fn checkHeaderValue(field: []const u8, value: []const u8) ?HlValue {
  265. for (value, 0..) |ch, i| {
  266. if (ch == '\r' or ch == '\n') {
  267. return refuse(
  268. "hl:smtp refused the message: {s} contains a {s} at byte {d} — a line break in a header field is how mail headers are injected, so it is never sent and never stripped",
  269. .{ field, if (ch == '\r') "carriage return" else "line feed", i },
  270. );
  271. }
  272. if (ch == 0) {
  273. return refuse(
  274. "hl:smtp refused the message: {s} contains a NUL byte at byte {d}",
  275. .{ field, i },
  276. );
  277. }
  278. }
  279. return null;
  280. }
  281. /// An envelope address additionally may not carry the delimiters the envelope
  282. /// itself is made of: `MAIL FROM:<a@b>` is a command, and a `<`, `>` or a space
  283. /// inside the address would rewrite it.
  284. fn checkAddress(field: []const u8, value: []const u8) ?HlValue {
  285. if (checkHeaderValue(field, value)) |e| return e;
  286. if (value.len == 0) {
  287. return refuse("hl:smtp refused the message: {s} is empty", .{field});
  288. }
  289. for (value) |ch| {
  290. if (ch == '<' or ch == '>' or ch == ' ' or ch == '\t') {
  291. return refuse(
  292. "hl:smtp refused the message: the address '{s}' in {s} contains '{c}' — an envelope address is written between angle brackets and may not carry them, or whitespace",
  293. .{ value, field, ch },
  294. );
  295. }
  296. }
  297. if (std.mem.indexOfScalar(u8, value, '@') == null) {
  298. return refuse(
  299. "hl:smtp refused the message: the address '{s}' in {s} has no '@'",
  300. .{ value, field },
  301. );
  302. }
  303. return null;
  304. }
  305. // ── MESSAGE RENDERING ───────────────────────────────────────────────────────
  306. const b64 = std.base64.standard.Encoder;
  307. fn b64Alloc(src: []const u8) ![]u8 {
  308. const out = try allocator.alloc(u8, b64.calcSize(src.len));
  309. _ = b64.encode(out, src);
  310. return out;
  311. }
  312. fn isAscii(s: []const u8) bool {
  313. for (s) |ch| {
  314. if (ch >= 0x80) return false;
  315. }
  316. return true;
  317. }
  318. const day_names = [_][]const u8{ "Thu", "Fri", "Sat", "Sun", "Mon", "Tue", "Wed" };
  319. const month_names = [_][]const u8{ "Jan", "Feb", "Mar", "Apr", "May", "Jun", "Jul", "Aug", "Sep", "Oct", "Nov", "Dec" };
  320. /// RFC 5322 §3.3 date-time, always +0000. A mail without a Date is not a mail —
  321. /// receivers add their own and it reads as a forgery.
  322. fn writeDate(w: *std.ArrayListUnmanaged(u8)) !void {
  323. var ts: linux.timespec = undefined;
  324. _ = linux.clock_gettime(linux.CLOCK.REALTIME, &ts);
  325. const secs: u64 = @intCast(@max(@as(i64, ts.sec), 0));
  326. const es = std.time.epoch.EpochSeconds{ .secs = secs };
  327. const yd = es.getEpochDay().calculateYearDay();
  328. const md = yd.calculateMonthDay();
  329. const ds = es.getDaySeconds();
  330. // 1970-01-01 was a Thursday, which is why day_names starts there.
  331. const dow = day_names[@intCast(es.getEpochDay().day % 7)];
  332. var buf: [64]u8 = undefined;
  333. const line = try std.fmt.bufPrint(&buf, "Date: {s}, {d:0>2} {s} {d} {d:0>2}:{d:0>2}:{d:0>2} +0000\r\n", .{
  334. dow,
  335. @as(u32, md.day_index) + 1,
  336. month_names[md.month.numeric() - 1],
  337. yd.year,
  338. ds.getHoursIntoDay(),
  339. ds.getMinutesIntoHour(),
  340. ds.getSecondsIntoMinute(),
  341. });
  342. try w.appendSlice(allocator, line);
  343. }
  344. fn randomHex(out: []u8) void {
  345. var raw: [32]u8 = undefined;
  346. const need = (out.len + 1) / 2;
  347. if (linux.getrandom(&raw, need, 0) != need) {
  348. @memset(out, '0');
  349. return;
  350. }
  351. const hex = "0123456789abcdef";
  352. for (out, 0..) |*ch, i| {
  353. const byte = raw[i / 2];
  354. const nib: u8 = if (i % 2 == 0) (byte >> 4) else (byte & 0x0f);
  355. ch.* = hex[nib];
  356. }
  357. }
  358. /// A header value that is not pure ASCII becomes an RFC 2047 encoded-word.
  359. /// Writing raw UTF-8 into a header is only legal with SMTPUTF8 negotiated, which
  360. /// this client does not do — so an umlaut in a subject would arrive as mojibake.
  361. fn appendHeader(w: *std.ArrayListUnmanaged(u8), name: []const u8, value: []const u8) !void {
  362. try w.appendSlice(allocator, name);
  363. try w.appendSlice(allocator, ": ");
  364. if (isAscii(value)) {
  365. try w.appendSlice(allocator, value);
  366. } else {
  367. const enc = try b64Alloc(value);
  368. defer allocator.free(enc);
  369. try w.appendSlice(allocator, "=?UTF-8?B?");
  370. try w.appendSlice(allocator, enc);
  371. try w.appendSlice(allocator, "?=");
  372. }
  373. try w.appendSlice(allocator, "\r\n");
  374. }
  375. /// CRLF DISCIPLINE AND DOT-STUFFING, in one pass over the assembled message.
  376. ///
  377. /// SMTP's line terminator is CRLF and its end-of-data marker is a line holding
  378. /// exactly ".", so a body line that legitimately begins with "." must be sent as
  379. /// ".." (RFC 5321 §4.5.2) or it truncates the mail — and a lone LF, which every
  380. /// text a program builds is full of, is not a line ending at all on the wire.
  381. fn appendDotStuffed(w: *std.ArrayListUnmanaged(u8), body: []const u8) !void {
  382. var at_line_start = true;
  383. var i: usize = 0;
  384. while (i < body.len) : (i += 1) {
  385. const ch = body[i];
  386. if (ch == '\r') {
  387. // Normalise CR and CRLF alike to one CRLF; a bare CR is a line
  388. // ending in nothing that talks SMTP.
  389. if (i + 1 < body.len and body[i + 1] == '\n') i += 1;
  390. try w.appendSlice(allocator, "\r\n");
  391. at_line_start = true;
  392. continue;
  393. }
  394. if (ch == '\n') {
  395. try w.appendSlice(allocator, "\r\n");
  396. at_line_start = true;
  397. continue;
  398. }
  399. if (at_line_start and ch == '.') try w.append(allocator, '.');
  400. try w.append(allocator, ch);
  401. at_line_start = false;
  402. }
  403. if (!at_line_start) try w.appendSlice(allocator, "\r\n");
  404. }
  405. // ── THE CONNECTION ──────────────────────────────────────────────────────────
  406. const Conn = struct {
  407. fd: i32 = -1,
  408. ssl: ?*tls.SSL = null,
  409. ctx: ?tls.TlsContext = null,
  410. buf: [4096]u8 = undefined,
  411. len: usize = 0,
  412. pos: usize = 0,
  413. fn writeAll(self: *Conn, bytes: []const u8) bool {
  414. if (self.ssl) |ssl| {
  415. const ctx = &(self.ctx orelse return false);
  416. return ctx.writeAll(ssl, bytes) == bytes.len;
  417. }
  418. var off: usize = 0;
  419. while (off < bytes.len) {
  420. const rc = linux.write(self.fd, bytes.ptr + off, bytes.len - off);
  421. const n: isize = @bitCast(rc);
  422. if (n <= 0) return false;
  423. off += @intCast(n);
  424. }
  425. return true;
  426. }
  427. fn fill(self: *Conn) bool {
  428. if (self.ssl) |ssl| {
  429. const ctx = &(self.ctx orelse return false);
  430. const n = ctx.read(ssl, self.buf[0..]);
  431. if (n <= 0) return false;
  432. self.len = @intCast(n);
  433. self.pos = 0;
  434. return true;
  435. }
  436. const rc = linux.read(self.fd, &self.buf, self.buf.len);
  437. const n: isize = @bitCast(rc);
  438. if (n <= 0) return false;
  439. self.len = @intCast(n);
  440. self.pos = 0;
  441. return true;
  442. }
  443. /// One CRLF-terminated reply line, without its terminator, appended into
  444. /// `out` (which is cleared first). `false` = the peer went away or timed out.
  445. fn readLine(self: *Conn, out: *std.ArrayListUnmanaged(u8)) bool {
  446. out.clearRetainingCapacity();
  447. while (true) {
  448. if (self.pos >= self.len) {
  449. if (!self.fill()) return false;
  450. }
  451. const ch = self.buf[self.pos];
  452. self.pos += 1;
  453. if (ch == '\n') {
  454. if (out.items.len > 0 and out.items[out.items.len - 1] == '\r') {
  455. _ = out.pop();
  456. }
  457. return true;
  458. }
  459. out.append(allocator, ch) catch return false;
  460. if (out.items.len > 65536) return false;
  461. }
  462. }
  463. fn close(self: *Conn) void {
  464. if (self.ssl) |ssl| {
  465. if (self.ctx) |*ctx| ctx.shutdownAndFree(ssl);
  466. self.ssl = null;
  467. }
  468. if (self.ctx) |*ctx| {
  469. ctx.deinit();
  470. self.ctx = null;
  471. }
  472. if (self.fd >= 0) {
  473. _ = linux.close(self.fd);
  474. self.fd = -1;
  475. }
  476. }
  477. };
  478. /// A complete SMTP reply: the code, and every line's text joined with "; ".
  479. const Reply = struct {
  480. code: u16 = 0,
  481. text: std.ArrayListUnmanaged(u8) = .empty,
  482. fn deinit(self: *Reply) void {
  483. self.text.deinit(allocator);
  484. }
  485. };
  486. /// Read a reply, following RFC 5321 §4.2's multiline form (`250-…` continues,
  487. /// `250 …` ends). Returns false if the connection died mid-reply.
  488. fn readReply(conn: *Conn, reply: *Reply) bool {
  489. reply.text.clearRetainingCapacity();
  490. reply.code = 0;
  491. var line: std.ArrayListUnmanaged(u8) = .empty;
  492. defer line.deinit(allocator);
  493. var first = true;
  494. while (true) {
  495. if (!conn.readLine(&line)) return false;
  496. const s = line.items;
  497. if (s.len < 3) return false;
  498. const code = std.fmt.parseInt(u16, s[0..3], 10) catch return false;
  499. if (first) {
  500. reply.code = code;
  501. first = false;
  502. }
  503. const rest = if (s.len > 4) s[4..] else "";
  504. if (reply.text.items.len > 0) reply.text.appendSlice(allocator, "; ") catch {};
  505. reply.text.appendSlice(allocator, rest) catch {};
  506. if (s.len == 3 or s[3] != '-') return true;
  507. }
  508. }
  509. fn command(conn: *Conn, reply: *Reply, parts: []const []const u8) bool {
  510. for (parts) |p| {
  511. if (!conn.writeAll(p)) return false;
  512. }
  513. if (!conn.writeAll("\r\n")) return false;
  514. return readReply(conn, reply);
  515. }
  516. /// Does the EHLO capability list advertise `name`? Case-insensitive, and
  517. /// anchored on a word boundary so "AUTH" does not match "AUTHOR".
  518. fn advertises(caps: []const u8, name: []const u8) bool {
  519. var i: usize = 0;
  520. while (i + name.len <= caps.len) : (i += 1) {
  521. if (std.ascii.eqlIgnoreCase(caps[i .. i + name.len], name)) {
  522. const before_ok = i == 0 or !std.ascii.isAlphanumeric(caps[i - 1]);
  523. const after = i + name.len;
  524. const after_ok = after >= caps.len or !std.ascii.isAlphanumeric(caps[after]);
  525. if (before_ok and after_ok) return true;
  526. }
  527. }
  528. return false;
  529. }
  530. // ── DIALLING ────────────────────────────────────────────────────────────────
  531. /// Connect to host:port with `timeout_ms` bounding BOTH the connect and every
  532. /// later read/write. A blocking connect(2) ignores SO_SNDTIMEO, so the connect
  533. /// is made non-blocking and polled; the socket goes back to blocking afterwards
  534. /// and the two socket timeouts carry the rest of the conversation.
  535. fn dial(host: []const u8, port: u16, timeout_ms: u32, reason: *std.ArrayListUnmanaged(u8)) ?i32 {
  536. // libc's resolver, not a hand-rolled one: /etc/hosts, /etc/resolv.conf,
  537. // nsswitch and IPv6 are the host's configuration, and a mail client that
  538. // answers differently from every other program on the machine is a bug.
  539. var host_z_buf: [256]u8 = undefined;
  540. if (host.len >= host_z_buf.len) {
  541. reason.appendSlice(allocator, "the host name is too long") catch {};
  542. return null;
  543. }
  544. @memcpy(host_z_buf[0..host.len], host);
  545. host_z_buf[host.len] = 0;
  546. var port_buf: [8]u8 = undefined;
  547. const port_s = std.fmt.bufPrintZ(&port_buf, "{d}", .{port}) catch {
  548. reason.appendSlice(allocator, "bad port") catch {};
  549. return null;
  550. };
  551. const hints = c.addrinfo{
  552. .flags = .{},
  553. .family = c.AF.UNSPEC,
  554. .socktype = c.SOCK.STREAM,
  555. .protocol = 0,
  556. .addrlen = 0,
  557. .canonname = null,
  558. .addr = null,
  559. .next = null,
  560. };
  561. var res: ?*c.addrinfo = null;
  562. const rc_ai = c.getaddrinfo(@ptrCast(&host_z_buf), port_s.ptr, &hints, &res);
  563. if (rc_ai != @as(c.EAI, @enumFromInt(0)) or res == null) {
  564. reason.appendSlice(allocator, "could not resolve the host") catch {};
  565. return null;
  566. }
  567. defer c.freeaddrinfo(res.?);
  568. var last: []const u8 = "connection refused";
  569. var cursor: ?*c.addrinfo = res;
  570. while (cursor) |ai| : (cursor = ai.next) {
  571. const sa = ai.addr orelse continue;
  572. const rc = linux.socket(@intCast(ai.family), linux.SOCK.STREAM | linux.SOCK.CLOEXEC, 0);
  573. const fd: i32 = @intCast(@as(isize, @bitCast(rc)));
  574. if (fd < 0) continue;
  575. const flags_rc = linux.fcntl(fd, linux.F.GETFL, @as(usize, 0));
  576. var oflags: linux.O = @bitCast(@as(u32, @truncate(flags_rc)));
  577. oflags.NONBLOCK = true;
  578. _ = linux.fcntl(fd, linux.F.SETFL, @as(usize, @as(u32, @bitCast(oflags))));
  579. const crc = linux.connect(fd, @ptrCast(sa), ai.addrlen);
  580. const cerr: isize = @bitCast(crc);
  581. var connected = cerr == 0;
  582. if (!connected) {
  583. // A negative return is -errno; INPROGRESS is the non-blocking
  584. // connect's ordinary answer, not a failure.
  585. const e: linux.E = @enumFromInt(@as(u16, @intCast(-cerr)));
  586. if (e != .INPROGRESS and e != .INTR) {
  587. last = "connection refused";
  588. _ = linux.close(fd);
  589. continue;
  590. }
  591. var pfd = [_]linux.pollfd{.{ .fd = fd, .events = linux.POLL.OUT, .revents = 0 }};
  592. const prc = linux.poll(&pfd, 1, @intCast(timeout_ms));
  593. const pn: isize = @bitCast(prc);
  594. if (pn <= 0) {
  595. last = "the connection attempt timed out";
  596. _ = linux.close(fd);
  597. continue;
  598. }
  599. var soerr: i32 = 0;
  600. var slen: linux.socklen_t = @sizeOf(i32);
  601. _ = linux.getsockopt(fd, linux.SOL.SOCKET, linux.SO.ERROR, @ptrCast(&soerr), &slen);
  602. if (soerr != 0) {
  603. last = "connection refused";
  604. _ = linux.close(fd);
  605. continue;
  606. }
  607. connected = true;
  608. }
  609. oflags.NONBLOCK = false;
  610. _ = linux.fcntl(fd, linux.F.SETFL, @as(usize, @as(u32, @bitCast(oflags))));
  611. const tv = linux.timeval{
  612. .sec = @intCast(timeout_ms / 1000),
  613. .usec = @intCast((timeout_ms % 1000) * 1000),
  614. };
  615. _ = linux.setsockopt(fd, linux.SOL.SOCKET, linux.SO.RCVTIMEO, @ptrCast(&tv), @sizeOf(linux.timeval));
  616. _ = linux.setsockopt(fd, linux.SOL.SOCKET, linux.SO.SNDTIMEO, @ptrCast(&tv), @sizeOf(linux.timeval));
  617. return fd;
  618. }
  619. reason.appendSlice(allocator, last) catch {};
  620. return null;
  621. }
  622. // ── THE CONVERSATION ────────────────────────────────────────────────────────
  623. const Outcome = struct {
  624. ok: bool,
  625. code: u16,
  626. stage: []const u8,
  627. message: []u8,
  628. secure: bool,
  629. };
  630. fn own(text: []const u8) []u8 {
  631. return allocator.dupe(u8, text) catch allocator.alloc(u8, 0) catch unreachable;
  632. }
  633. fn fail(stage: []const u8, code: u16, secure: bool, text: []const u8) Outcome {
  634. return .{ .ok = false, .code = code, .stage = stage, .message = own(text), .secure = secure };
  635. }
  636. fn runJob(m: *Mailer, job: *const Job) Outcome {
  637. var conn = Conn{};
  638. defer conn.close();
  639. var reason: std.ArrayListUnmanaged(u8) = .empty;
  640. defer reason.deinit(allocator);
  641. const fd = dial(m.host, m.port, m.timeout_ms, &reason) orelse
  642. return fail("connect", 0, false, reason.items);
  643. conn.fd = fd;
  644. var secure = false;
  645. if (m.tls_mode == .implicit) {
  646. const ctx = tls.TlsContext.initClient(.{
  647. .ca_file = if (m.ca_file) |ca| ca else null,
  648. .verify = m.verify,
  649. .tag = "smtp",
  650. }) orelse return fail("connect", 0, false, "TLS is unavailable (no usable libssl on this host, or the CA file could not be loaded)");
  651. conn.ctx = ctx;
  652. conn.ssl = conn.ctx.?.connect(fd, m.host, "smtp") orelse
  653. return fail("connect", 0, false, "the implicit TLS handshake failed (see the smtp: lines on stderr)");
  654. secure = true;
  655. }
  656. var reply = Reply{};
  657. defer reply.deinit();
  658. if (!readReply(&conn, &reply)) return fail("greeting", 0, secure, "the server sent no greeting");
  659. if (reply.code != 220) return fail("greeting", reply.code, secure, reply.text.items);
  660. // EHLO, and keep the capability list — STARTTLS and AUTH are both read off it.
  661. var caps: std.ArrayListUnmanaged(u8) = .empty;
  662. defer caps.deinit(allocator);
  663. if (!command(&conn, &reply, &.{ "EHLO ", m.helo })) return fail("ehlo", 0, secure, "the connection closed during EHLO");
  664. if (reply.code != 250) return fail("ehlo", reply.code, secure, reply.text.items);
  665. caps.appendSlice(allocator, reply.text.items) catch {};
  666. if (m.tls_mode == .starttls and !secure) {
  667. if (!advertises(caps.items, "STARTTLS")) {
  668. return fail("starttls", 0, secure, "the server does not advertise STARTTLS (set tls = \"none\" to accept a plaintext session deliberately)");
  669. }
  670. if (!command(&conn, &reply, &.{"STARTTLS"})) return fail("starttls", 0, secure, "the connection closed during STARTTLS");
  671. if (reply.code != 220) return fail("starttls", reply.code, secure, reply.text.items);
  672. const ctx = tls.TlsContext.initClient(.{
  673. .ca_file = if (m.ca_file) |ca| ca else null,
  674. .verify = m.verify,
  675. .tag = "smtp",
  676. }) orelse return fail("starttls", 0, secure, "TLS is unavailable (no usable libssl on this host, or the CA file could not be loaded)");
  677. conn.ctx = ctx;
  678. // The SAME socket, mid-protocol: this is what STARTTLS is, and the
  679. // client half of tls_common.zig takes it as it stands.
  680. conn.ssl = conn.ctx.?.connect(fd, m.host, "smtp") orelse
  681. return fail("starttls", 0, secure, "the STARTTLS handshake failed (see the smtp: lines on stderr)");
  682. secure = true;
  683. conn.len = 0;
  684. conn.pos = 0;
  685. // RFC 3207 §4.2: everything learned before the upgrade is discarded and
  686. // EHLO is re-issued inside TLS — a capability list from the plaintext
  687. // half is exactly what an attacker in the middle would have written.
  688. caps.clearRetainingCapacity();
  689. if (!command(&conn, &reply, &.{ "EHLO ", m.helo })) return fail("ehlo", 0, secure, "the connection closed during the post-STARTTLS EHLO");
  690. if (reply.code != 250) return fail("ehlo", reply.code, secure, reply.text.items);
  691. caps.appendSlice(allocator, reply.text.items) catch {};
  692. }
  693. // AUTH
  694. if (m.user != null and m.pass != null and m.auth_mode != .none) {
  695. if (!secure and !m.allow_insecure_auth) {
  696. return fail("auth", 0, secure, "refusing to send credentials over a plaintext connection — use tls = \"starttls\" (the default) or \"implicit\", or set allowInsecureAuth = true to say the clear text is intended");
  697. }
  698. const wants_login = switch (m.auth_mode) {
  699. .login => true,
  700. .plain => false,
  701. else => !advertises(caps.items, "PLAIN") and advertises(caps.items, "LOGIN"),
  702. };
  703. const outcome = if (wants_login) authLogin(m, &conn, &reply, secure) else authPlain(m, &conn, &reply, secure);
  704. if (outcome) |o| return o;
  705. }
  706. if (!command(&conn, &reply, &.{ "MAIL FROM:<", job.from, ">" })) return fail("mail", 0, secure, "the connection closed during MAIL FROM");
  707. if (reply.code != 250) return fail("mail", reply.code, secure, reply.text.items);
  708. for (job.rcpt) |r| {
  709. if (!command(&conn, &reply, &.{ "RCPT TO:<", r, ">" })) return fail("rcpt", 0, secure, "the connection closed during RCPT TO");
  710. // 251 = "will forward"; anything else non-2xx is a rejected recipient,
  711. // and one rejected recipient fails the whole send rather than silently
  712. // delivering to the rest (the app decides what to do about it).
  713. if (reply.code != 250 and reply.code != 251) return fail("rcpt", reply.code, secure, reply.text.items);
  714. }
  715. if (!command(&conn, &reply, &.{"DATA"})) return fail("data", 0, secure, "the connection closed during DATA");
  716. if (reply.code != 354) return fail("data", reply.code, secure, reply.text.items);
  717. if (!conn.writeAll(job.data)) return fail("body", 0, secure, "the connection closed while the message was being written");
  718. if (!conn.writeAll(".\r\n")) return fail("body", 0, secure, "the connection closed at the end-of-data marker");
  719. if (!readReply(&conn, &reply)) return fail("body", 0, secure, "the server never acknowledged the message");
  720. if (reply.code != 250) return fail("body", reply.code, secure, reply.text.items);
  721. const accepted = own(reply.text.items);
  722. // QUIT is courtesy: the message is accepted the moment DATA answered 250, so
  723. // a server that drops the socket instead of saying 221 has not lost it.
  724. _ = command(&conn, &reply, &.{"QUIT"});
  725. return .{ .ok = true, .code = 250, .stage = "done", .message = accepted, .secure = secure };
  726. }
  727. /// AUTH PLAIN — RFC 4616: base64 of "\0user\0pass", in the command itself.
  728. fn authPlain(m: *Mailer, conn: *Conn, reply: *Reply, secure: bool) ?Outcome {
  729. const user = m.user.?;
  730. const pass = m.pass.?;
  731. const raw = allocator.alloc(u8, 2 + user.len + pass.len) catch return fail("auth", 0, secure, "out of memory");
  732. defer allocator.free(raw);
  733. raw[0] = 0;
  734. @memcpy(raw[1 .. 1 + user.len], user);
  735. raw[1 + user.len] = 0;
  736. @memcpy(raw[2 + user.len ..], pass);
  737. const enc = b64Alloc(raw) catch return fail("auth", 0, secure, "out of memory");
  738. defer allocator.free(enc);
  739. if (!command(conn, reply, &.{ "AUTH PLAIN ", enc })) return fail("auth", 0, secure, "the connection closed during AUTH PLAIN");
  740. if (reply.code != 235) return fail("auth", reply.code, secure, reply.text.items);
  741. return null;
  742. }
  743. /// AUTH LOGIN — the de-facto challenge/response form: base64 username, then
  744. /// base64 password, each answering a 334.
  745. fn authLogin(m: *Mailer, conn: *Conn, reply: *Reply, secure: bool) ?Outcome {
  746. if (!command(conn, reply, &.{"AUTH LOGIN"})) return fail("auth", 0, secure, "the connection closed during AUTH LOGIN");
  747. if (reply.code != 334) return fail("auth", reply.code, secure, reply.text.items);
  748. const user_enc = b64Alloc(m.user.?) catch return fail("auth", 0, secure, "out of memory");
  749. defer allocator.free(user_enc);
  750. if (!command(conn, reply, &.{user_enc})) return fail("auth", 0, secure, "the connection closed after the username");
  751. if (reply.code != 334) return fail("auth", reply.code, secure, reply.text.items);
  752. const pass_enc = b64Alloc(m.pass.?) catch return fail("auth", 0, secure, "out of memory");
  753. defer allocator.free(pass_enc);
  754. if (!command(conn, reply, &.{pass_enc})) return fail("auth", 0, secure, "the connection closed after the password");
  755. if (reply.code != 235) return fail("auth", reply.code, secure, reply.text.items);
  756. return null;
  757. }
  758. fn workerLoop(m: *Mailer) void {
  759. while (true) {
  760. var job = m.takeJob() orelse break;
  761. const outcome = runJob(m, &job);
  762. const rcpt: [][]u8 = allocator.alloc([]u8, job.rcpt.len) catch &.{};
  763. for (job.rcpt, 0..) |r, i| {
  764. if (i < rcpt.len) rcpt[i] = allocator.dupe(u8, r) catch own("");
  765. }
  766. m.publish(.{
  767. .id = job.id,
  768. .ok = outcome.ok,
  769. .code = outcome.code,
  770. .stage = outcome.stage,
  771. .message = outcome.message,
  772. .secure = outcome.secure,
  773. .rcpt = rcpt,
  774. });
  775. job.deinit();
  776. }
  777. }
  778. // ── THE RESULT SOURCE ───────────────────────────────────────────────────────
  779. /// The loader converts a nested `hl_object` by recursing and running THAT
  780. /// object's own `deinit_fn` before it runs this one, so the `to` list is already
  781. /// gone by the time this is called — touching it here is a double free (measured,
  782. /// mission 137). Only the string values this object owns are freed here.
  783. fn resultObjDeinit(obj: *HlObject) callconv(.c) void {
  784. const fields = obj.fields[0..obj.field_count];
  785. for (fields) |f| {
  786. if (f.value.type == .hl_string) {
  787. allocator.free(@constCast(f.value.data.string.ptr[0..f.value.data.string.len]));
  788. }
  789. }
  790. allocator.free(fields);
  791. allocator.destroy(obj);
  792. }
  793. fn rcptObjDeinit(obj: *HlObject) callconv(.c) void {
  794. const fields = obj.fields[0..obj.field_count];
  795. for (fields) |f| {
  796. allocator.free(@constCast(f.key.ptr[0..f.key.len]));
  797. if (f.value.type == .hl_string) {
  798. allocator.free(@constCast(f.value.data.string.ptr[0..f.value.data.string.len]));
  799. }
  800. }
  801. allocator.free(fields);
  802. allocator.destroy(obj);
  803. }
  804. /// Contiguous "0".."n-1" keys is how the loader recognises an ordered hybrid, so
  805. /// `res.to` arrives in Hybriel as a list however many recipients there were.
  806. fn rcptValue(rcpt: [][]u8) HlValue {
  807. const fields = allocator.alloc(HlField, rcpt.len) catch return api.makeNull();
  808. var built: usize = 0;
  809. for (rcpt, 0..) |r, i| {
  810. var key_buf: [24]u8 = undefined;
  811. const key_src = std.fmt.bufPrint(&key_buf, "{d}", .{i}) catch break;
  812. const key = allocator.dupe(u8, key_src) catch break;
  813. fields[built] = .{ .key = .{ .ptr = key.ptr, .len = key.len }, .value = api.makeString(r) };
  814. built += 1;
  815. }
  816. const obj = allocator.create(HlObject) catch return api.makeNull();
  817. obj.* = .{ .fields = fields.ptr, .field_count = built, .deinit_fn = &rcptObjDeinit };
  818. return api.makeObject(obj);
  819. }
  820. fn resultValue(r: Result) HlValue {
  821. const fields = allocator.alloc(HlField, 7) catch return api.makeNull();
  822. fields[0] = .{ .key = http.hlStr("id"), .value = api.makeNumber(@floatFromInt(r.id)) };
  823. fields[1] = .{ .key = http.hlStr("ok"), .value = api.makeBool(r.ok) };
  824. fields[2] = .{ .key = http.hlStr("code"), .value = api.makeNumber(@floatFromInt(r.code)) };
  825. // `stage` names a compile-time constant, so it is DUPED here: the deinit
  826. // below frees every string field it finds, and handing it a `.rodata`
  827. // pointer is a segfault (measured, mission 137). One rule for all fields
  828. // beats a per-field exception nobody will remember.
  829. fields[3] = .{ .key = http.hlStr("stage"), .value = api.makeString(own(r.stage)) };
  830. fields[4] = .{ .key = http.hlStr("message"), .value = api.makeString(r.message) };
  831. fields[5] = .{ .key = http.hlStr("secure"), .value = api.makeBool(r.secure) };
  832. fields[6] = .{ .key = http.hlStr("to"), .value = rcptValue(r.rcpt) };
  833. const obj = allocator.create(HlObject) catch return api.makeNull();
  834. obj.* = .{ .fields = fields.ptr, .field_count = 7, .deinit_fn = &resultObjDeinit };
  835. allocator.free(r.rcpt);
  836. return api.makeObject(obj);
  837. }
  838. fn resultsTryNext(ctx: ?*anyopaque) callconv(.c) HlValue {
  839. const m: *Mailer = @ptrCast(@alignCast(ctx orelse return api.makeNull()));
  840. const r = m.tryTake() orelse return api.makeNull();
  841. return resultValue(r);
  842. }
  843. fn resultsNext(ctx: ?*anyopaque) callconv(.c) HlValue {
  844. const m: *Mailer = @ptrCast(@alignCast(ctx orelse return api.makeNull()));
  845. const r = m.take() orelse return api.makeNull();
  846. return resultValue(r);
  847. }
  848. fn resultsDeinit(_: ?*anyopaque) callconv(.c) void {}
  849. // ── EXPORTS ─────────────────────────────────────────────────────────────────
  850. /// __native("smtp.open", host, port, options) → Number mailer id
  851. export fn hl_smtp_open(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  852. if (argc < 2 or argv[0].type != .hl_string or argv[1].type != .hl_number) {
  853. return refuse("hl:smtp open: host must be a String and port a Number", .{});
  854. }
  855. const host = argv[0].data.string.ptr[0..argv[0].data.string.len];
  856. const port_f = argv[1].data.number;
  857. if (port_f < 1 or port_f > 65535) {
  858. return refuse("hl:smtp open: {d} is not a port number", .{port_f});
  859. }
  860. const opts: HlValue = if (argc >= 3) argv[2] else api.makeNull();
  861. const tls_name = strField(opts, "tls") orelse "starttls";
  862. const tls_mode: TlsMode = if (std.ascii.eqlIgnoreCase(tls_name, "starttls"))
  863. .starttls
  864. else if (std.ascii.eqlIgnoreCase(tls_name, "implicit"))
  865. .implicit
  866. else if (std.ascii.eqlIgnoreCase(tls_name, "none"))
  867. .none
  868. else
  869. return refuse("hl:smtp open: tls = \"{s}\" is not one of \"starttls\", \"implicit\", \"none\"", .{tls_name});
  870. const auth_name = strField(opts, "auth") orelse "auto";
  871. const auth_mode: AuthMode = if (std.ascii.eqlIgnoreCase(auth_name, "auto"))
  872. .auto
  873. else if (std.ascii.eqlIgnoreCase(auth_name, "plain"))
  874. .plain
  875. else if (std.ascii.eqlIgnoreCase(auth_name, "login"))
  876. .login
  877. else if (std.ascii.eqlIgnoreCase(auth_name, "none"))
  878. .none
  879. else
  880. return refuse("hl:smtp open: auth = \"{s}\" is not one of \"auto\", \"plain\", \"login\", \"none\"", .{auth_name});
  881. // A MAILER WITHOUT A SENDER OPENS (ticket #13): an app that never sends
  882. // must not abort at load for want of one. A message that ends up with no
  883. // sender at all is refused at send(), located at that call.
  884. const from = senderField(opts) orelse "";
  885. if (from.len > 0) {
  886. if (checkAddress("the mailer's \"from\" (or sender)", from)) |e| return e;
  887. }
  888. const helo = strField(opts, "helo") orelse "localhost";
  889. if (checkHeaderValue("helo", helo)) |e| return e;
  890. const m = allocator.create(Mailer) catch return refuse("hl:smtp open: out of memory", .{});
  891. m.* = .{
  892. .id = 0,
  893. .host = allocator.dupe(u8, host) catch return refuse("hl:smtp open: out of memory", .{}),
  894. .port = @intFromFloat(port_f),
  895. .user = dupOpt(strField(opts, "user")),
  896. .pass = dupOpt(strField(opts, "pass")),
  897. .from = allocator.dupe(u8, from) catch return refuse("hl:smtp open: out of memory", .{}),
  898. .helo = allocator.dupe(u8, helo) catch return refuse("hl:smtp open: out of memory", .{}),
  899. .tls_mode = tls_mode,
  900. .auth_mode = auth_mode,
  901. .ca_file = dupOpt(strField(opts, "caFile")),
  902. .verify = boolField(opts, "verify", true),
  903. .allow_insecure_auth = boolField(opts, "allowInsecureAuth", false),
  904. .timeout_ms = @intFromFloat(@max(1000, numField(opts, "timeout", 30000))),
  905. .mutex = c.PTHREAD_MUTEX_INITIALIZER,
  906. .condvar = c.PTHREAD_COND_INITIALIZER,
  907. .wake_fd = http.makeWakeFd(),
  908. };
  909. mutexLock(&mailers_mutex);
  910. m.id = next_mailer_id;
  911. next_mailer_id += 1;
  912. mailers.append(allocator, m) catch {};
  913. mutexUnlock(&mailers_mutex);
  914. m.worker = std.Thread.spawn(.{}, workerLoop, .{m}) catch {
  915. return refuse("hl:smtp open: could not start the sender thread", .{});
  916. };
  917. return api.makeNumber(@floatFromInt(m.id));
  918. }
  919. /// __native("smtp.send", id, message) → Number job id, or a located refusal.
  920. ///
  921. /// Everything this function does happens on the HYBRIEL thread, before the job
  922. /// exists: that is deliberate, because a refusal is only located if it is raised
  923. /// at the call site and a worker thread has no call site.
  924. export fn hl_smtp_send(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  925. if (argc < 2 or argv[0].type != .hl_number) {
  926. return refuse("hl:smtp send: expected a mailer id and a message", .{});
  927. }
  928. const m = mailerById(@intFromFloat(argv[0].data.number)) orelse
  929. return refuse("hl:smtp send: this Mailer is closed", .{});
  930. const msg = argv[1];
  931. if (msg.type != .hl_object) {
  932. return refuse("hl:smtp send: the message must be a hybrid with to, subject and text (or html)", .{});
  933. }
  934. const from = senderField(msg) orelse m.from;
  935. if (from.len == 0) {
  936. return refuse("hl:smtp refused the message: it has no \"from\" and neither has the mailer — give one in the message or in the Mailer's options", .{});
  937. }
  938. if (checkAddress("\"from\"", from)) |e| return e;
  939. // `to` is one address or a list of them. A single String is the common case
  940. // and a list is the same message with more RCPT TO lines — never several
  941. // messages, so every recipient sees the same To: header.
  942. var rcpt: std.ArrayListUnmanaged([]u8) = .empty;
  943. var rcpt_failed = false;
  944. defer if (rcpt_failed) {
  945. for (rcpt.items) |r| allocator.free(r);
  946. rcpt.deinit(allocator);
  947. };
  948. const to_field = fieldOf(msg, "to") orelse return refuse("hl:smtp send: the message has no 'to'", .{});
  949. switch (to_field.type) {
  950. .hl_string => {
  951. const one = to_field.data.string.ptr[0..to_field.data.string.len];
  952. if (checkAddress("to", one)) |e| {
  953. rcpt_failed = true;
  954. return e;
  955. }
  956. rcpt.append(allocator, allocator.dupe(u8, one) catch "") catch {};
  957. },
  958. .hl_object => {
  959. const list = to_field.data.object;
  960. for (list.fields[0..list.field_count]) |f| {
  961. if (f.value.type != .hl_string) continue;
  962. const one = f.value.data.string.ptr[0..f.value.data.string.len];
  963. if (checkAddress("to", one)) |e| {
  964. rcpt_failed = true;
  965. return e;
  966. }
  967. rcpt.append(allocator, allocator.dupe(u8, one) catch "") catch {};
  968. }
  969. },
  970. else => return refuse("hl:smtp send: 'to' must be a String or a list of Strings", .{}),
  971. }
  972. if (rcpt.items.len == 0) {
  973. rcpt_failed = true;
  974. return refuse("hl:smtp send: 'to' names no recipient", .{});
  975. }
  976. const subject = strField(msg, "subject") orelse "";
  977. if (checkHeaderValue("subject", subject)) |e| {
  978. rcpt_failed = true;
  979. return e;
  980. }
  981. const text = strField(msg, "text");
  982. const html = strField(msg, "html");
  983. if (text == null and html == null) {
  984. rcpt_failed = true;
  985. return refuse("hl:smtp send: the message has neither 'text' nor 'html'", .{});
  986. }
  987. var body: std.ArrayListUnmanaged(u8) = .empty;
  988. var body_failed = false;
  989. defer if (body_failed) body.deinit(allocator);
  990. renderMessage(&body, m, from, rcpt.items, subject, text, html) catch {
  991. body_failed = true;
  992. rcpt_failed = true;
  993. return refuse("hl:smtp send: out of memory rendering the message", .{});
  994. };
  995. mutexLock(&m.mutex);
  996. const job_id = m.next_job;
  997. m.next_job += 1;
  998. mutexUnlock(&m.mutex);
  999. m.enqueue(.{
  1000. .id = job_id,
  1001. .from = allocator.dupe(u8, from) catch "",
  1002. .rcpt = rcpt.toOwnedSlice(allocator) catch &[_][]u8{},
  1003. .data = body.toOwnedSlice(allocator) catch &[_]u8{},
  1004. });
  1005. return api.makeNumber(@floatFromInt(job_id));
  1006. }
  1007. fn renderMessage(
  1008. out: *std.ArrayListUnmanaged(u8),
  1009. m: *Mailer,
  1010. from: []const u8,
  1011. rcpt: []const []u8,
  1012. subject: []const u8,
  1013. text: ?[]const u8,
  1014. html: ?[]const u8,
  1015. ) !void {
  1016. var head: std.ArrayListUnmanaged(u8) = .empty;
  1017. defer head.deinit(allocator);
  1018. try appendHeader(&head, "From", from);
  1019. var to_join: std.ArrayListUnmanaged(u8) = .empty;
  1020. defer to_join.deinit(allocator);
  1021. for (rcpt, 0..) |r, i| {
  1022. if (i > 0) try to_join.appendSlice(allocator, ", ");
  1023. try to_join.appendSlice(allocator, r);
  1024. }
  1025. try appendHeader(&head, "To", to_join.items);
  1026. try appendHeader(&head, "Subject", subject);
  1027. try writeDate(&head);
  1028. var mid: [32]u8 = undefined;
  1029. randomHex(&mid);
  1030. try head.appendSlice(allocator, "Message-ID: <");
  1031. try head.appendSlice(allocator, &mid);
  1032. try head.append(allocator, '@');
  1033. try head.appendSlice(allocator, m.helo);
  1034. try head.appendSlice(allocator, ">\r\n");
  1035. try head.appendSlice(allocator, "MIME-Version: 1.0\r\n");
  1036. if (text != null and html != null) {
  1037. var boundary: [24]u8 = undefined;
  1038. randomHex(&boundary);
  1039. try head.appendSlice(allocator, "Content-Type: multipart/alternative; boundary=\"hl-");
  1040. try head.appendSlice(allocator, &boundary);
  1041. try head.appendSlice(allocator, "\"\r\n\r\n");
  1042. // The plain part first: RFC 2046 §5.1.4 orders alternatives worst-first,
  1043. // so a reader that understands HTML picks the LAST one it can render.
  1044. try head.appendSlice(allocator, "--hl-");
  1045. try head.appendSlice(allocator, &boundary);
  1046. try head.appendSlice(allocator, "\r\nContent-Type: text/plain; charset=utf-8\r\nContent-Transfer-Encoding: 8bit\r\n\r\n");
  1047. try out.appendSlice(allocator, head.items);
  1048. try appendDotStuffed(out, text.?);
  1049. try out.appendSlice(allocator, "--hl-");
  1050. try out.appendSlice(allocator, &boundary);
  1051. try out.appendSlice(allocator, "\r\nContent-Type: text/html; charset=utf-8\r\nContent-Transfer-Encoding: 8bit\r\n\r\n");
  1052. try appendDotStuffed(out, html.?);
  1053. try out.appendSlice(allocator, "--hl-");
  1054. try out.appendSlice(allocator, &boundary);
  1055. try out.appendSlice(allocator, "--\r\n");
  1056. return;
  1057. }
  1058. const single = text orelse html.?;
  1059. const ctype = if (text != null) "text/plain" else "text/html";
  1060. try head.appendSlice(allocator, "Content-Type: ");
  1061. try head.appendSlice(allocator, ctype);
  1062. try head.appendSlice(allocator, "; charset=utf-8\r\nContent-Transfer-Encoding: 8bit\r\n\r\n");
  1063. try out.appendSlice(allocator, head.items);
  1064. try appendDotStuffed(out, single);
  1065. }
  1066. /// __native("smtp.results", id) → the mailer's ONE result source.
  1067. export fn hl_smtp_results(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  1068. if (argc < 1 or argv[0].type != .hl_number) return api.makeNull();
  1069. const m = mailerById(@intFromFloat(argv[0].data.number)) orelse
  1070. return refuse("hl:smtp results: this Mailer is closed", .{});
  1071. if (m.source == null) {
  1072. const iter = allocator.create(HlIterator) catch return api.makeNull();
  1073. iter.* = .{
  1074. .context = @ptrCast(m),
  1075. .next_fn = &resultsNext,
  1076. .deinit_fn = &resultsDeinit,
  1077. .try_next_fn = &resultsTryNext,
  1078. .wake_fd = m.wake_fd,
  1079. };
  1080. m.source = iter;
  1081. }
  1082. return api.makeIterator(m.source.?);
  1083. }
  1084. /// __native("smtp.mark_registered", id) → refuses a SECOND consumer.
  1085. /// The drain (`for (r of m.results())`) and the event loop take results off the
  1086. /// same queue, so a program doing both would see each result exactly once, in
  1087. /// one of two places, at random. Saying so is cheaper than debugging it.
  1088. export fn hl_smtp_mark_registered(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  1089. if (argc < 1 or argv[0].type != .hl_number) return api.makeNull();
  1090. const m = mailerById(@intFromFloat(argv[0].data.number)) orelse
  1091. return refuse("hl:smtp deliver: this Mailer is closed", .{});
  1092. if (m.registered) {
  1093. return refuse("hl:smtp deliver: this Mailer already delivers its results to the event loop", .{});
  1094. }
  1095. m.registered = true;
  1096. return api.makeNull();
  1097. }
  1098. /// __native("smtp.pending", id) → how many sends have not been answered yet.
  1099. export fn hl_smtp_pending(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  1100. if (argc < 1 or argv[0].type != .hl_number) return api.makeNumber(0);
  1101. const m = mailerById(@intFromFloat(argv[0].data.number)) orelse return api.makeNumber(0);
  1102. mutexLock(&m.mutex);
  1103. defer mutexUnlock(&m.mutex);
  1104. return api.makeNumber(@floatFromInt(m.outstanding));
  1105. }
  1106. /// __native("smtp.close", id) — stop the worker. Queued jobs still in hand are
  1107. /// dropped; a job already in conversation finishes, because abandoning a socket
  1108. /// between DATA and its 250 is how a message gets delivered twice.
  1109. export fn hl_smtp_close(argc: u32, argv: [*]const HlValue) callconv(.c) HlValue {
  1110. if (argc < 1 or argv[0].type != .hl_number) return api.makeNull();
  1111. const m = mailerById(@intFromFloat(argv[0].data.number)) orelse return api.makeNull();
  1112. mutexLock(&m.mutex);
  1113. m.stop = true;
  1114. condBroadcast(&m.condvar);
  1115. mutexUnlock(&m.mutex);
  1116. http.ringWake(m.wake_fd);
  1117. if (m.worker) |t| {
  1118. t.join();
  1119. m.worker = null;
  1120. }
  1121. mutexLock(&mailers_mutex);
  1122. for (mailers.items, 0..) |candidate, i| {
  1123. if (candidate == m) {
  1124. _ = mailers.orderedRemove(i);
  1125. break;
  1126. }
  1127. }
  1128. mutexUnlock(&mailers_mutex);
  1129. return api.makeNull();
  1130. }

Branches

Latest commits

  • eb3b9205tracker: report 031mre
  • 9b5d2e89tracker mission 031: README (What it does, Test: four gates + the #32 checks, Files: theme/, new pages), STATUS (real copy, A/B load, how to repeat, open points), LOGmre
  • 39950e4ctracker#32 (mission 031): the WorldAPI theme (theme/ vendored verbatim from layouts.worldapi.org 85b5654; styles.hl inherits it: accent green-dark, type colours 1-6; own base/header rules, row lines, genre-pill and inverted-button frames removed, the season foldable keeps its line; check-theme 21 -> 0, 4th deploy gate; main actions class primary) and the #32 header (theme AppHeader/MainMenu/UserMenu/Sidebar/ContentFirst: desktop brand, search, Series|Shows|Movies|Genres|People, user icon with Unwatched..Settings, Logout; signed out the ident selector, phone the iD icon dropdown; phone menu in the sidebar overlay; marked entry by :has); /find -> /search/<q>, /genres, /people(/<letter>), /settings; main { ContentFirst { slot } } works around the hl:web one-line slot bug; gates 365/0, 32/0, 52/0, check-theme 0mre
  • a386dc92tracker: reports 029 + 030mre
  • 71e0fd7dtracker missions 029 + 030: README (What it does, Files, gate count), STATUS (real-copy numbers, how to repeat, open points), LOGmre
  • d36ea6eatracker#34 + #35 (mission 030): Follow directly under the poster, as wide as the poster (show.hl, styles.hl); the status pill next to a series' title — TVmaze's status (new tvmazeStatus, stored by the sync's TVmaze merge) else TMDB's, TVmaze Ended + TMDB Canceled = Canceled, inverted (filled, dark text, no border), green running / yellow pending / red canceled / muted ended (shows.hl statusOf); the daily delta asks TVmaze's status of an unfollowed series TVmaze's change list names (dailysync.hl syncRunStep, sync.hl syncTvmazeStatus); the status backfill after the details repair (backfill.hl, jobs.hl statusTick; resumable, 550 ms per TVmaze request); gates 354/0, 32/0, 52/0mre
  • 7d7d4487tracker#33 (mission 029): reduced titles — every title TMDB's details never went through this app (no detailsAt, no tmdbSync) is incomplete (shows.hl isIncomplete; the old tracker's migrated rows passed #26's test: 5,697 non-adult on the live copy, 691 series without seasons); the repair job does the visibly reduced first (shows.hl missingParts), the page completes one on open; a title TMDB has no poster for (The Remaining) shows the placeholder; tools/count-incomplete.hl; gate fixtures stand for synced titles (tmdbSync), tests/seed-reduced.hl + #33 checks; gates 347/0, 32/0, 52/0mre
  • 661c2592tracker: report 028mre
  • 27c916fatracker mission 028: README ("Code order", the new file map), STATUS (counts before/after, tests, how to repeat, open), LOGmre
  • d924f398tracker mission 028: comments name the new files (sync.hl, dailysync.hl, backfill.hl, credits.hl, jobs.hl, images.hl …); tools/ref-params.py + tools/lambda-audit.py also scan lib/ (they globbed the root only), lambda-audit counts a plain `x = p` alias like `let x = p`mre
  • 2e89b968tracker mission 028 (code order) 5/5 let: `let` only where a variable is reassigned — 667 never-reassigned lets became plain declarations (project.hl, lib/, components/, tools/, tests/); kept: 264 in loop bodies (a plain declaration there is 'Cannot reassign' on the 2nd pass), 234 reassigned, 27 whose name is also a member/outer/free name (a plain write would rebind it); tools/let-audit.py decides and fixes (README 'Code order'); tests/realdata-m028.{sh,mjs} = the page-output diff on a real copy; gates 342/0, 32/0, 52/0, real-copy pages identicalmre
  • 54796ff2tracker mission 028 (code order) 4/5 thin faces + last copies: the show page's check/follow faces call lib/watches.hl toggleWatched / toggleSeasonWatched (seasonAllWatched moved there) and lib/follows.hl toggleFollowed; both logins (header selector face, /login/callback) share lib/users.hl userOfCode; todayStr/listOf copies in components and the export readers copied into tools/migrate.hl + tools/old-short-ids.hl now once (lib/util.hl, lib/export.hl); gates 342/0, 32/0, 52/0; old-short-ids output byte-identical, migrate output identicalmre
  • 06b078e3tracker mission 028 (code order) 3/5 project.hl is the map: config, routes, wiring and a feature → file index (914 → 258 lines); the background jobs (daily sync run, backfills, details repair, credits job, merge, short ids, collection seed) moved unchanged into lib/jobs.hl (a class: their state is reassigned every step, a static cannot be; one instance made after the server), the login callback into lib/users.hl, poster/photo serving into lib/images.hl, the /shows/<slug> rule into lib/shows.hl showsMovedPath; route handlers are thin wrappers; gates 342/0, 32/0, 52/0, real-copy pages identicalmre
  • 94716fd2tracker mission 028 (code order) 2/5 util + topics: lib/util.hl holds envOr, storageDir, postersDir, profilesDir, newId, hexDigits, todayStr, dateOr, textOr, hasId, listOr, firstOf, sortDesc once (were copied into up to 5 files); tmdbsync.hl split into tmdb.hl (TMDB/TVmaze requests), sync.hl (one title's sync), sync-helpers.hl, backfill.hl; details.hl split into details.hl, credits.hl, credits-helpers.hl (isIncomplete to shows.hl); search-helpers.hl (words, query, ranking, slugs); collections.hl (the TMDB collection seed, out of franchises.hl); deltasync.hl renamed dailysync.hl; no behaviour change: gates 342/0, 32/0, 52/0, real-copy pages identicalmre
  • 186079b0tracker mission 028 (code order) 1/5 move: every root .hl except project.hl into lib/ (styles.hl into components/), import paths only; gates 342/0, 32/0, 52/0; real-copy pages identicalmre
  • 4f47f181tracker: report 027mre
  • dc1d4be4tracker mission 027: Hybriel master 06617221 vendored (plugin allocator fixes 3a781359 + 413f60e4); real copy RSS through first-start jobs + 400 loads flat ~2.55 GB (190aa11d 2.3 -> 5.6 GB), page times <= 1.1x; gates 342/0, 32/0, 52/0mre
  • 84e1b3e1tracker: reports 025 + 026mre
  • dc40d859tracker#31 (mission 026): duplicate titles merged — the 68 type+tmdbId pairs held by 157 records were the old tracker's (all migrated); merge.hl repair job (own clock, before the TMDB jobs) keeps one keeper per title (follows/watches > old short id > oldest), moves follows, watches, seasons, cast, credits, timelines, tombstones the rest (mergedInto, never deleted), slugs + short ids 301 to the keeper; stray seasons merged into their listed twin (Reacher S3 watches) or linked when watched; search import re-checks before its put; deploy.sh waits up to 90 s for 200; real copy 68 -> 0 dup ids, az5b2 follows/watches equal; gates 342/0, 32/0, 52/0mre
  • 2667da05tracker#30 (mission 025): /my/ pages from slim cached title cards, episode rows and watch sets (after the jobs /my/series 1.8 s -> 0.06 s, /my/unwatched 4.1 -> 0.18 s); timeline page shows its name once; franchise widget under the poster/title; movies with TV leftovers (First Contact) go through the details repair; tools/count-tmdb-ids.hl; gates 327/0, 52/0, 32/0mre