gitoriaLog in with ident

tracker

All repositories: gitoria

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

Branches

Latest commits

  • f83571c3tracker#18 (mission 067): daily sync by change lists — TMDB /tv|movie/changes (since the stored day, paged) + TVmaze /updates/shows → only our changed titles (followed: full step, unfollowed: light step — changed seasons, no TVmaze), full walk on first run / gap > 14 days / failed list; show record refreshed (title, tmdbSummary, tagline, status, genres …; renamed titles re-indexed); summary = the creator's own text (page: summary > tmdbSummary > tvmazeSummary), one-time clear of copied summaries (9,647 on the live copy); gate 261, tests/realdata-018*.mjs, README + STATUSmre
  • 10bb3f93tracker#28 (mission 061): full cast (all seasons, main cast by episodes, guest stars) + crew (created by, directed by, written by, screenplay, story, music) — stored by the details completion, the daily sync, the search import (one details request) and a background credits job (resumes, RSS limit); show page collapsed after 20 with client-side Show all; showBySlug via a slug map; gate 249, tests/realdata-028.mjs, README + STATUSmre
  • 25a50bc4tracker#26 (mission 059): titles from a filmography are completed — on open (skeleton, step-wise face showComplete, no reload) and by the in-app details repair (resumes, TMDB-paced, series in parts); cast from TMDB credits; gate 231, tests/realdata-026.mjs, README + STATUSmre
  • 1704ec45tracker#17 (mission 057): season caret down/up, skeleton rows while a season loads, sessionless showSeasonEpisodes face (no page re-mount), client-only close; gate 214, tests/realdata-057.mjs, README + STATUSmre
  • f2fe3e36mission 056: README + STATUS (merge, fixes, Hybriel 8590df63, real-data check), tests/realdata-056.mjs, tools/check-public-slugs.hlmre
  • f40c250emission 056: re-vendor hybriel master 8590df63 (#121, #122); an adult title's page is Not found for non-followers; gate: leave the page before stopping the servermre
  • 2b7fdd6cmission 056: signed-out header one row on phones ("Log in", nowrap), backfill skips adult titles' posters, gate checksmre
  • 2c53d5efMerge branch 't16-person' (tracker#16 person pages) into main; filmography shows only public titles (054 adult flag), gate race fix (backfill start line)mre
  • c171227emission 054: hide adult/unknown titles from the public lists and the search; in-app adult-flag backfill (TMDB details + poster per title, resumes), gate + real-data proofmre
  • 139fafd8tracker#16: short bio (4 lines, click = all), real-data check script, README + STATUSmre
  • 93be9476tracker#16: person pages /person/<slug> with the filmography fetched from TMDB on the first visit (step by step), gatemre
  • 47a3cae6STATUS: mission 053 merge commit idsmre
  • dcc5eecaMerge branch 't14-search'mre
  • 03edc783Merge branch 't15-tvmaze'mre
  • 71b46345tracker#15: numbering check by date or title, placeholder titles in other languages, docs + real-data proofmre
  • 6bb2daf1tracker#13: homepage (tiles, intro, latest movies/shows), /shows, /movies/page/N, /my/movies; lists cached in memorymre
  • b8bd1157tracker#14: README + STATUS (search, real-data numbers, gate, merge notes)mre
  • 65c694a8tracker#14: search — header magnifier, /search/<text> (in-memory word-prefix index over titles + people), Fetch from web (TMDB search/multi, ours left out), Add = import via syncShow; gate +25 checks, real-data scriptmre
  • 34f2c15btracker#15: TVmaze merge in the sync (gaps only: new episodes/seasons, empty titles/air dates; numbering check), fake TVmaze episodes + gatemre
  • cbdc4ea7tracker#12: link icons TMDB/IMDb/TVDB/TVmaze; sync fills missing ids (TVmaze lookup); movies fetched via /movie/mre