diff --git a/lith/mail-inbound.mjs b/lith/mail-inbound.mjs index 5df76cd89c..c1b6de35cc 100644 --- a/lith/mail-inbound.mjs +++ b/lith/mail-inbound.mjs @@ -17,6 +17,7 @@ // node mail-inbound.mjs # :25, Google relays only // node mail-inbound.mjs --port 2525 --open # any sender — local testing +import { mailTrace, recordMailEvent, flushMailEvents } from "../system/backend/mail-events.mjs"; import { SMTPServer } from "smtp-server"; import { simpleParser } from "mailparser"; import { BlockList } from "node:net"; @@ -161,20 +162,41 @@ export function makeLimiter({ perPairPerHour = 20, pushPerPair = 3, perMinute = const prune = (arr, span, t) => { while (arr.length && t - arr[0] > span) arr.shift(); }; - return { - take(pair) { - const t = now(); - prune(minute, 60_000, t); - if (minute.length >= perMinute) return { ok: false, code: 451, why: "The door is busy, try later" }; - const hits = pairs.get(pair) || []; + function takeMany(keys) { + const t = now(); + prune(minute, 60_000, t); + if (minute.length + keys.length > perMinute) return { ok: false, code: 451, reason: "rate_global", why: "The door is busy, try later" }; + const planned = new Map(); + for (const pair of keys) { + const hits = planned.get(pair) || [...(pairs.get(pair) || [])]; prune(hits, 3_600_000, t); - if (hits.length >= perPairPerHour) return { ok: false, code: 550, why: "Too many letters to this box this hour" }; + if (hits.length >= perPairPerHour) return { ok: false, code: 550, reason: "rate_pair", why: "Too many letters to this box this hour" }; hits.push(t); - pairs.set(pair, hits); - minute.push(t); - if (pairs.size > 5000) for (const [k, v] of pairs) if (!v.length || t - v[v.length - 1] > 3_600_000) pairs.delete(k); - return { ok: true, quiet: hits.length > pushPerPair }; - }, + planned.set(pair, hits); + } + for (const [key, hits] of planned) pairs.set(key, hits); + for (const key of keys) minute.push(t); + if (pairs.size > 5000) for (const [k, v] of pairs) if (!v.length || t - v[v.length - 1] > 3_600_000) pairs.delete(k); + let released = false; + return { + ok: true, + gates: keys.map((key) => ({ quiet: planned.get(key).length > pushPerPair })), + release() { + if (released) return; + released = true; + for (const key of keys) { + const hits = pairs.get(key); + const index = hits?.lastIndexOf(t) ?? -1; + if (index >= 0) hits.splice(index, 1); + const globalIndex = minute.lastIndexOf(t); + if (globalIndex >= 0) minute.splice(globalIndex, 1); + } + }, + }; + } + return { + takeMany, + take(pair) { const result = takeMany([pair]); return result.ok ? { ok: true, ...result.gates[0] } : result; }, }; } @@ -189,13 +211,31 @@ export function createInbound({ file, domains, relays = null, + getRelays = () => relays, open = false, tls = null, secret = null, limiter = makeLimiter(), log = console.log, + record = (fields) => recordMailEvent(null, fields, log), }) { + const event = (session, name, fields = {}) => { + session.mailTrace ||= mailTrace(); + record({ ...fields, event: name, trace: session.mailTrace, transport: "smtp-in" }); + }; + const sessions = new Map(); + const ignore = () => {}; const server = new SMTPServer({ + // smtp-server can refuse SIZE/protocol commands before our callbacks. + // Consume ONLY response status digits; discard its verbose protocol logs. + logger: { + trace: ignore, info: ignore, warn: ignore, error: ignore, fatal: ignore, + debug(meta, label, payload) { + if (meta?.tnx !== "send" || label !== "S:" || typeof payload !== "string") return; + const status = Number(payload.match(/^([45]\d\d)\b/)?.[1]); + if (status) event(sessions.get(meta.cid) || {}, "smtp_response", { status }); + }, + }, name: HOST, banner: "Amail — aesthetic.computer", size: MAX_SIZE, @@ -205,27 +245,46 @@ export function createInbound({ ...(tls || {}), onConnect(session, cb) { - if (open || isRelay(relays, session.remoteAddress)) return cb(); - log("mail.inbound.refused.relay"); - cb(new Error("Only Google Workspace delivers here")); + sessions.set(session.id, session); + if (open || isRelay(getRelays(), session.remoteAddress)) { event(session, "connected"); return cb(); } + event(session, "rejected", { reason: "relay", status: 554 }); + cb(Object.assign(new Error("Only Google Workspace delivers here"), { responseCode: 554 })); + }, + + onMailFrom(address, session, cb) { + // SMTP connections can carry several independent messages, including RSET. + session.amail = new Map(); + session.mailTrace = mailTrace(); + sessions.set(session.id, session); + event(session, "started"); + cb(); + }, + + onClose(session) { + // Includes protocol/SIZE rejections made by smtp-server before callbacks. + const status = Number(String(session.error || "").match(/^([45]\d\d)\b/)?.[1]); + event(session, "disconnected", { ...(status ? { status } : {}) }); + sessions.delete(session.id); }, async onRcptTo(address, session, cb) { const who = parseRecipient(address.address, domains); - if (!who) return cb(Object.assign(new Error("No such mailbox"), { responseCode: 550 })); + if (!who) { event(session, "rejected", { reason: "recipient", status: 550 }); return cb(Object.assign(new Error("No such mailbox"), { responseCode: 550 })); } try { const sub = await lookup(who.local); - if (!sub) return cb(Object.assign(new Error(`No handle ${who.local}`), { responseCode: 550 })); + if (!sub) { event(session, "rejected", { reason: "recipient", status: 550 }); return cb(Object.assign(new Error("No such mailbox"), { responseCode: 550 })); } session.amail = session.amail || new Map(); - session.amail.set(address.address, { sub, ...who }); + session.amail.set(sub, { sub, ...who }); + event(session, "recipient", { recipients: session.amail.size }); cb(); } catch (err) { - log("mail.inbound.lookup.error", mailErrorCode(err)); + event(session, "deferred", { reason: "lookup", status: 451, error: mailErrorCode(err) }); cb(Object.assign(new Error("Try again later"), { responseCode: 451 })); } }, async onData(stream, session, cb) { + let reservation; try { // Drain an oversized message without buffering it or passing it to // mailparser. SMTP's advertised SIZE alone doesn't bound parser memory. @@ -237,6 +296,7 @@ export function createInbound({ else chunks.length = 0; } if (size > MAX_SIZE || stream.sizeExceeded) { + event(session, "rejected", { reason: "wire_size", status: 552, bytes: size }); return cb(Object.assign(new Error("Letter too large"), { responseCode: 552 })); } const parsed = await simpleParser(Buffer.concat(chunks), { skipImageLinks: true }); @@ -244,13 +304,13 @@ export function createInbound({ // Only letters that came through OUR routing rule carry the stamp. if (secret && parsed.headers?.get?.("x-amail-route") !== secret) { - log("mail.inbound.refused.stamp"); + event(session, "rejected", { reason: "route", status: 550 }); return cb(Object.assign(new Error("Not our route"), { responseCode: 550 })); } const auth = authFrom(parsed); if (auth.dmarc === "fail") { - log("mail.inbound.refused.dmarc"); + event(session, "rejected", { reason: "dmarc", status: 550 }); return cb(Object.assign(new Error("Sender's domain disowns this letter"), { responseCode: 550 })); } @@ -258,32 +318,33 @@ export function createInbound({ // path. Validate the entire letter before filing it for any recipient. const attachments = incomingAttachments(parsed.attachments); - const results = []; - let refused = null; - for (const rcpt of session.amail?.values() || []) { - const gate = limiter.take(`${sender}→${rcpt.sub}`); - if (!gate.ok) { - refused = refused || gate; - log("mail.inbound.refused.rate"); - continue; - } - results.push(await file({ ...rcpt, parsed, attachments, auth, quiet: gate.quiet, remote: session.remoteAddress })); + const recipients = [...(session.amail?.values() || [])]; + // Check and reserve ALL recipients before storing ANY copy. A DATA 250 + // must never acknowledge a message whose recipients were silently skipped. + reservation = limiter.takeMany(recipients.map((rcpt) => `${sender}→${rcpt.sub}`)); + if (!reservation.ok) { + event(session, reservation.code === 451 ? "deferred" : "rejected", { reason: reservation.reason, status: reservation.code, recipients: recipients.length }); + return cb(Object.assign(new Error(reservation.why), { responseCode: reservation.code })); } - if (!results.length && refused) { - return cb(Object.assign(new Error(refused.why), { responseCode: refused.code })); + const results = []; + for (const [i, rcpt] of recipients.entries()) { + results.push(await file({ ...rcpt, parsed, attachments, auth, quiet: reservation.gates[i].quiet, trace: session.mailTrace, remote: session.remoteAddress })); } - log("mail.inbound.filed", { + event(session, "accepted", { + status: 250, bytes: size, attachments: attachments.length, recipients: results.length, duplicates: results.filter((r) => r.duplicate).length, - quiet: results.filter((r) => r.quiet).length, verified: auth.verified === true, }); cb(); } catch (err) { + // A temporary storage outage must not turn retries into quota bounces. + reservation?.release?.(); if (err.responseCode === 552) { + event(session, "rejected", { reason: "attachments", status: 552 }); return cb(Object.assign(new Error("At most 10 files and 8 MiB of attachments per letter"), { responseCode: 552 })); } - log("mail.inbound.file.error", mailErrorCode(err)); + event(session, "deferred", { reason: "storage", status: 451, error: mailErrorCode(err) }); cb(Object.assign(new Error("Could not file that letter"), { responseCode: 451 })); } }, @@ -304,6 +365,7 @@ async function main() { } = await import("../system/backend/mail.mjs"); const database = await connect(); + const record = (fields) => recordMailEvent(database, { transport: "smtp-in", ...fields }); let relays = null; if (!OPEN) { @@ -311,8 +373,9 @@ async function main() { setInterval(async () => { try { relays = await googleRelays(); + record({ event: "relays_refreshed" }); } catch (err) { - console.log("mail.inbound.relays.error", mailErrorCode(err)); + record({ event: "relays_failed", error: mailErrorCode(err) }); } }, 60 * 60 * 1000).unref(); } @@ -322,15 +385,17 @@ async function main() { if (!secret) console.log("🟡 AMAIL_ROUTE_SECRET is unset — any Google tenant's route would be accepted"); const server = createInbound({ domains: INBOUND_DOMAINS, - relays, + getRelays: () => relays, + record, open: OPEN, tls, secret, lookup: (local) => subFromAddress(local, database), - file: async ({ sub, parsed, attachments, reply, auth, quiet }) => { + file: async ({ sub, parsed, attachments, reply, auth, quiet, trace }) => { const sender = parsed.from?.value?.[0] || {}; return deliverFromOutside( { + trace, to: sub, fromEmail: (sender.address || "").toLowerCase(), fromName: sender.name || "", @@ -346,8 +411,9 @@ async function main() { }, }); - server.on("error", (err) => console.log("mail.inbound.smtp.error", mailErrorCode(err))); + server.on("error", (err) => record({ event: "smtp_error", error: mailErrorCode(err) })); server.listen(PORT, () => { + record({ event: "ready", tls: !!tls, routeSecret: !!secret, open: OPEN }); console.log( `📮 Amail inbound on :${PORT} for ${INBOUND_DOMAINS.join(", ")} — ` + `${OPEN ? "OPEN to any sender" : "Google relays only"}, ${tls ? "STARTTLS" : "no TLS (no cert yet)"}`, @@ -370,7 +436,7 @@ async function main() { } for (const signal of ["SIGINT", "SIGTERM"]) { - process.on(signal, () => server.close(() => process.exit(0))); + process.on(signal, () => server.close(async () => { await flushMailEvents(); process.exit(0); })); } } diff --git a/spec/mail-events-spec.mjs b/spec/mail-events-spec.mjs new file mode 100644 index 0000000000..2eb217d093 --- /dev/null +++ b/spec/mail-events-spec.mjs @@ -0,0 +1,105 @@ +// Real loopback SMTP; no credentials, production mailbox, or external delivery. +import assert from "node:assert/strict"; +import nodemailer from "nodemailer"; +import { BlockList } from "node:net"; +import { mailEvent, mailTrace, recordMailEvent, flushMailEvents, observeMail } from "../system/backend/mail-events.mjs"; +import { options, queryFor } from "../system/backend/mail-events-cli.mjs"; +import { createInbound, makeLimiter } from "../lith/mail-inbound.mjs"; + +const canary = "PRIVATE_MAIL_CANARY@example.invalid"; +const trace = mailTrace(); +const sanitized = mailEvent({ event: "stored", trace, transport: "internal", letterId: "a".repeat(24), subject: canary, text: canary, from: canary, filename: canary, error: canary, reason: canary, status: canary, attachments: 2 }); +assert.ok(!JSON.stringify(sanitized).includes(canary)); +assert.equal(sanitized.error, "UNKNOWN"); +assert.equal(mailEvent({ event: canary }), null); +assert.equal(mailEvent({ event: "failed", trace: canary }).trace, undefined); +const logs = [], saved = [], indexes = []; +const db = { db: { collection(name) { + assert.equal(name, "mail-events"); + return { createIndex: async (...args) => indexes.push(args), insertOne: async (row) => saved.push(row) }; +} } }; +for (let i = 0; i < 3; i++) recordMailEvent(db, { event: "stored", trace, subject: canary }, (line) => logs.push(line)); +await flushMailEvents(); +assert.equal(saved.length, 3); +assert.equal(indexes.length, 3, "concurrent events share index setup"); +assert.equal(indexes[0][1].expireAfterSeconds, 2592000); +assert.ok(!JSON.stringify([logs, saved]).includes(canary)); +recordMailEvent({ db: { collection() { throw new Error(canary); } } }, { event: "accepted", trace }, (line) => logs.push(line)); +await flushMailEvents(); +assert.equal(JSON.parse(logs.at(-1)).event, "telemetry_unavailable"); +assert.ok(!JSON.stringify(logs).includes(canary)); +// Telemetry outages don't delay a successful operation, even with hung storage. +let unblock; +const stuck = new Promise((resolve) => { unblock = resolve; }); +const hangingDB = { db: { collection: () => ({ createIndex: () => stuck, insertOne: async () => {} }) } }; +const oldLog = console.log; console.log = (line) => logs.push(line); +try { + assert.equal(await observeMail("internal", { trace }, hangingDB, async (_opts, event) => { event("stored"); return "delivered"; }), "delivered"); +} finally { console.log = oldLog; unblock(); await flushMailEvents(); } + +const limit = makeLimiter({ perPairPerHour: 1, perMinute: 100 }); +assert.equal(limit.take("full").ok, true); +assert.equal(limit.takeMany(["fresh", "full"]).ok, false); +assert.equal(limit.take("fresh").ok, true, "failed batch must not consume fresh quota"); +const globalLimit = makeLimiter({ perPairPerHour: 100, perMinute: 1 }); +assert.equal(globalLimit.takeMany(["one", "two"]).code, 451); +assert.equal(globalLimit.take("one").ok, true); + +const events = [], deliveries = []; +const limiter = makeLimiter({ perPairPerHour: 2, perMinute: 100 }); +const server = createInbound({ + open: true, domains: ["example.invalid"], limiter, + record: (fields) => events.push(mailEvent(fields)), + lookup: async (local) => ["a", "b", "full", "fresh", "broken"].includes(local) ? local : undefined, + file: async (letter) => { + if (letter.sub === "broken") throw new Error(canary); + deliveries.push({ to: letter.sub, trace: letter.trace }); return {}; + }, +}); +await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); +const transport = nodemailer.createTransport({ host: "127.0.0.1", port: server.server.address().port, secure: false, ignoreTLS: true, pool: true, maxConnections: 1, maxMessages: 100 }); +const send = (to, extra = {}) => transport.sendMail({ from: canary, to, subject: canary, text: canary, ...extra }); +try { + await send("a@example.invalid"); + await send("b@example.invalid"); + assert.deepEqual(deliveries.map((d) => d.to), ["a", "b"], "reused SMTP socket must not keep previous recipients"); + assert.equal(events.filter((e) => e.event === "connected").length, 1, "test actually reused a connection"); + assert.notEqual(deliveries[0].trace, deliveries[1].trace, "trace is per transaction"); + await send("full@example.invalid"); await send("full@example.invalid"); + const before = deliveries.length; + await assert.rejects(send(["fresh@example.invalid", "full@example.invalid"]), (err) => err.responseCode === 550); + assert.equal(deliveries.length, before, "mixed rate-limited delivery must not acknowledge/store partial batch"); + await send("fresh@example.invalid"); + await assert.rejects(send("unknown@example.invalid"), (err) => err.responseCode === 550); + for (let i = 0; i < 3; i++) await assert.rejects(send("broken@example.invalid"), (err) => err.responseCode === 451, "temporary failures release quota for retry"); + await assert.rejects(send("a@example.invalid", { headers: { "Authentication-Results": "mx; dmarc=fail" } }), (err) => err.responseCode === 550); + await assert.rejects(send("a@example.invalid", { attachments: Array.from({ length: 11 }, (_, i) => ({ filename: `${i}-${canary}`, content: "x" })) }), (err) => err.responseCode === 552); + // Over advertised MIME cap without trusting a client's SIZE declaration. + await assert.rejects(send("a@example.invalid", { text: "x".repeat(13 * 1024 * 1024) }), (err) => err.responseCode === 552); + for (const reason of ["rate_pair", "recipient", "storage", "dmarc", "attachments", "wire_size"]) assert.ok(events.some((e) => e.reason === reason), reason); + assert.ok(events.some((e) => e.event === "smtp_response" && e.status === 552)); + assert.ok(events.some((e) => e.event === "accepted" && e.status === 250)); + assert.ok(!JSON.stringify(events).includes(canary), "no SMTP payload in telemetry"); +} finally { + transport.close(); + await new Promise((resolve) => server.close(resolve)); +} +// A newly refreshed relay list is consulted on every connection. +let relayList = new BlockList(); +const relayEvents = []; +const relayServer = createInbound({ domains: [], lookup: async () => {}, file: async () => {}, getRelays: () => relayList, record: (fields) => relayEvents.push(fields) }); +let rejected; +relayServer.options.onConnect({ remoteAddress: "127.0.0.1" }, (err) => { rejected = err; }); +assert.equal(rejected.responseCode, 554); +relayList = new BlockList(); relayList.addAddress("127.0.0.1"); +relayServer.options.onConnect({ remoteAddress: "127.0.0.1" }, (err) => { rejected = err; }); +assert.equal(rejected, undefined); + +assert.throws(() => options(["--limit", "99999"])); +assert.throws(() => options(["--trace", canary])); +const opts = options(["--failures", "--since", "60", "--trace", trace, "--json"]); +const query = queryFor(opts, 3600000); +assert.equal(query.trace, trace); +assert.equal(query.when.$gte.getTime(), 0); +assert.ok(query.$or.length); +console.log("mail event spec passed: privacy, retention, outage isolation, atomic throttles, SMTP reuse/rejections, relay refresh and inspection filters"); diff --git a/spec/mail-media-spec.mjs b/spec/mail-media-spec.mjs index 4ce1bd2b61..b04a567d14 100644 --- a/spec/mail-media-spec.mjs +++ b/spec/mail-media-spec.mjs @@ -1,3 +1,15 @@ +import * as mailEvents from "../system/backend/mail-events.mjs"; +const eventMocks = { + ...mailEvents, + recordMailEvent: (_database, fields) => logs.push(mailEvents.mailEvent(fields)), + observeMail: async (transport, options, _database, operation) => { + const trace = options.trace || mailEvents.mailTrace(); + const event = (event, fields = {}) => logs.push(mailEvents.mailEvent({ ...fields, event, transport, trace })); + event("started"); + try { return await operation({ ...options, trace }, event); } + catch (error) { event("failed", { error: privacy.mailErrorCode(error) }); throw error; } + }, +}; // node --experimental-vm-modules spec/mail-media-spec.mjs // Real MIME/SMTP and production handlers; isolated in-memory mailbox/auth. import assert from 'node:assert/strict'; @@ -64,6 +76,7 @@ async function load(path, mocks) { } context.process = { env: {} }; const backend = await load('../system/backend/mail.mjs', { + "./mail-events.mjs": eventMocks, './authorization.mjs': { handleFor: async (sub) => sub, userIDFromHandleOrEmail: async () => 'recipient' }, './filter.mjs': { filter: (s) => s }, './shell.mjs': { shell: { log: (...args) => logs.push(args) } }, './mail-media.mjs': media, '../../shared/mail-privacy.mjs': privacy, @@ -74,6 +87,7 @@ const backend = await load('../system/backend/mail.mjs', { } }) } }, }); const api = await load('../system/netlify/functions/mail.mjs', { + "../../backend/mail-events.mjs": eventMocks, '../../backend/authorization.mjs': { authorize: async () => identity }, '../../backend/database.mjs': { connect: async () => database }, '../../backend/http.mjs': { respond }, '../../backend/mail.mjs': backend, diff --git a/spec/mail-privacy-spec.mjs b/spec/mail-privacy-spec.mjs index 4e9b2e7243..a436ab61fc 100644 --- a/spec/mail-privacy-spec.mjs +++ b/spec/mail-privacy-spec.mjs @@ -1,3 +1,15 @@ +import * as mailEvents from "../system/backend/mail-events.mjs"; +const eventMocks = { + ...mailEvents, + recordMailEvent: (_database, fields) => logs.push(mailEvents.mailEvent(fields)), + observeMail: async (transport, options, _database, operation) => { + const trace = options.trace || mailEvents.mailTrace(); + const event = (event, fields = {}) => logs.push(mailEvents.mailEvent({ ...fields, event, transport, trace })); + event("started"); + try { return await operation({ ...options, trace }, event); } + catch (error) { event("failed", { error: privacy.mailErrorCode(error) }); throw error; } + }, +}; // Run: node --experimental-vm-modules spec/mail-privacy-spec.mjs // Execute production modules with fake storage/transports; no credentials/network. import assert from "node:assert/strict"; @@ -34,15 +46,17 @@ async function load(path, mocks) { } const cleanJSON = (value) => JSON.parse(JSON.stringify(value)); const privateFree = (value) => assert.ok(!JSON.stringify(value).includes(marker)); +let storageFails = false; const collection = { createIndex: async () => {}, dropIndex: async () => {}, - insertOne: async (row) => { stored.push(row); return { insertedId: "opaque-id" }; }, + insertOne: async (row) => { if (storageFails) throw failure; stored.push(row); return { insertedId: "opaque-id" }; }, findOne: async () => ({ code: "ac25abcde" }), }; const database = { db: { collection: () => collection }, disconnect: async () => {} }; let pushThrows = false; let smtpCalls = 0; const backend = await load("../system/backend/mail.mjs", { + "./mail-events.mjs": eventMocks, "./authorization.mjs": { handleFor: async () => marker, userIDFromHandleOrEmail: async () => "recipient", @@ -76,6 +90,12 @@ await backend.deliver({ from: "sender", to: "recipient", text: marker }, databas await backend.deliverFromOutside({ to: "recipient", fromEmail: marker, text: marker }, database); await backend.sendOutside({ from: "sender", toEmail: "test@example.invalid", text: marker, subject: marker }, database); assert.equal(smtpCalls, 2, "SMTP fallback remains functional"); +storageFails = true; +const checkpoint = logs.length; +await assert.rejects(backend.sendOutside({ from: "sender", toEmail: canaryAddress(), text: marker }, database)); +assert.deepEqual(logs.slice(checkpoint).filter((entry) => entry?.event).map((entry) => entry.event), ["started", "relay_accepted", "failed"], "relay acceptance remains distinguishable from sent-list storage failure"); +storageFails = false; +function canaryAddress() { return `${marker}@example.invalid`; } privateFree(logs); privateFree(notes); assert.equal(privacy.mailErrorCode({ code: "ETIMEDOUT", message: marker }), "ETIMEDOUT"); @@ -86,6 +106,7 @@ assert.equal(privacy.mailErrorCode(failure), "UNKNOWN"); for (const api of ["mail", "tell"]) { for (const stage of ["operation", "connect", "disconnect", "success"]) { const handler = await load(`../system/netlify/functions/${api}.mjs`, { + "../../backend/mail-events.mjs": eventMocks, "../../backend/authorization.mjs": { authorize: async () => ({ sub: "sender" }) }, "../../backend/database.mjs": { connect: async () => { if (stage === "connect") throw failure; @@ -113,6 +134,7 @@ privateFree(logs); // Exercise inbound callbacks without opening a socket or parsing real mail. let parsed = { from: { value: [{ address: `${marker}@example.invalid` }] }, headers: new Map() }; const inbound = await load("../lith/mail-inbound.mjs", { + "../system/backend/mail-events.mjs": eventMocks, "smtp-server": { SMTPServer: class { constructor(options) { this.options = options; } } }, mailparser: { simpleParser: async () => parsed }, "node:net": { BlockList: class {} }, @@ -128,7 +150,7 @@ for (const stage of ["filed", "failure", "stamp", "lookup", "relay", "rate"]) { secret: stage === "stamp" ? "stamp" : null, lookup: async () => { throw failure; }, file: async () => { if (stage === "failure") throw failure; return { toHandle: marker }; }, - limiter: { take: () => ({ ok: stage !== "rate", why: "Too many letters", code: 451 }) }, + limiter: { takeMany: () => ({ ok: stage !== "rate", gates: [{ quiet: false }], why: "Too many letters", code: 451 }) }, }); let result; const cb = (error) => { result = error; }; diff --git a/system/backend/MAIL.md b/system/backend/MAIL.md index f815ee4978..5aea58300d 100644 --- a/system/backend/MAIL.md +++ b/system/backend/MAIL.md @@ -33,3 +33,65 @@ node --experimental-vm-modules spec/mail-media-spec.mjs The media spec exercises real MIME generation and loopback SMTP with an isolated in-memory mailbox. Deploying this change requires both lith's web process and `lith-mail.service` to reload the changed modules. + +## Delivery inspection + +Mail emits structured `mail.event` JSON into the `lith` / `lith-mail` journals +and mirrors it asynchronously into MongoDB `mail-events`, retained for 30 days. +The mirror has a bounded queue; database failures emit `telemetry_unavailable` +in the journal and never delay or fail a letter. A sudden process exit may lose +queued database events; use the journal to investigate gaps. Failures before an +API database connection exists are journal-only. + +Each API response includes `X-Mail-Trace`. That UUID follows the send through +storage, SMTP relay, and push. An inbound SMTP transaction has its own trace, +shared by its recipient copies. Stored copies include their opaque letter ID +in telemetry. Bodies, subjects, addresses, handles, IPs, filenames, raw MIME +message IDs, and provider response text are excluded. Nothing is sent to PostHog. + +From the repository on lith (or locally with a configured `system/.env`): + +```sh +node --env-file=system/.env system/backend/mail-events-cli.mjs --since 60 +node --env-file=system/.env system/backend/mail-events-cli.mjs --since 1440 --failures +node --env-file=system/.env system/backend/mail-events-cli.mjs --trace --json +node --env-file=system/.env system/backend/mail-events-cli.mjs --letter +journalctl -u lith -u lith-mail --since '1 hour ago' -o cat | rg '"kind":"mail.event"' +``` + +The CLI summary covers the full selected window; the chronological detail is +limited to the latest 100 events (`--limit` up to 1000). Counts are **events**, +not unique letters: a rejected message can also have a `smtp_response` status +and a disconnect event. Successful unread-count polls are omitted. + +- `stored`: inserted in the recipient mailbox (or the sender's SMTP sent list). +- `accepted` / SMTP 250: inbound copies filed or recognized as duplicates. +- `relay_accepted`: the external SMTP relay accepted it; this does **not** prove + delivery to the destination inbox. A later storage failure shares its trace. +- `duplicate`: an incoming retry was already filed; no second copy or push. +- `rejected` / 5xx: fixed reason such as unknown recipient, DMARC, route stamp, + attachment/wire limits, or sender-to-recipient hourly limit. +- `deferred` / 451: temporary lookup/storage failure or global rate limit; + the sending relay should retry. +- `smtp_response`: a 4xx/5xx status sent by the SMTP library, including SIZE and + protocol rejections that precede application callbacks. No response text. +- `push`: attempted/succeeded/failed/pruned counts. `no_devices` means the + letter was stored but the recipient has no matching registered push device. + `push_quiet` means the inbound notification quota was reached; mail is kept. +- `ready`: inbound TLS, open-relay testing mode, and route-secret presence. + `relays_refreshed` / `relays_failed`: hourly Google relay-list refresh health. + +Limits remain enforced. A multi-recipient DATA transaction checks all quotas +before storing any copy; it cannot return 250 while skipping a throttled box. +Recipient state resets at every MAIL command, including reused connections. + +This sees AC and SMTP handoff outcomes from deployment forward. It cannot see +Google Workspace quarantine/routing decisions before a connection reaches AC, +or a remote provider's spam placement after relay acceptance. For those, use +Workspace Email Log Search and the sender's bounce notice. + +Additional regression coverage: + +```sh +node spec/mail-events-spec.mjs +``` diff --git a/system/backend/mail-events-cli.mjs b/system/backend/mail-events-cli.mjs new file mode 100644 index 0000000000..6b18612a53 --- /dev/null +++ b/system/backend/mail-events-cli.mjs @@ -0,0 +1,70 @@ +#!/usr/bin/env node +// Run on lith: node --env-file=system/.env system/backend/mail-events-cli.mjs +import { pathToFileURL } from "node:url"; + +export function options(args) { + const out = { since: 1440, limit: 100, failures: false, json: false }; + for (let i = 0; i < args.length; i++) { + const key = args[i].replace(/^--/, ""); + if (["json", "failures", "help"].includes(key)) out[key] = true; + else if (["since", "limit", "trace", "letter"].includes(key)) out[key] = args[++i]; + else throw new Error("Unknown option; use --help"); + } + out.since = Number(out.since); out.limit = Number(out.limit); + if (!Number.isFinite(out.since) || out.since <= 0 || out.since > 43200) throw new Error("--since must be 1–43200 minutes"); + if (!Number.isInteger(out.limit) || out.limit < 1 || out.limit > 1000) throw new Error("--limit must be 1–1000"); + if (out.trace !== undefined && !/^[a-f\d]{8}(?:-[a-f\d]{4}){3}-[a-f\d]{12}$/i.test(out.trace)) throw new Error("Invalid --trace UUID"); + if (out.letter !== undefined && !/^[a-f\d]{24}$/i.test(out.letter)) throw new Error("Invalid --letter ID"); + return out; +} + +export function queryFor(opts, now = Date.now()) { + return { + when: { $gte: new Date(now - opts.since * 60000) }, + ...(opts.trace ? { trace: opts.trace } : {}), + ...(opts.letter ? { letterId: opts.letter } : {}), + ...(opts.failures ? { $or: [ + { event: { $in: ["failed", "rejected", "deferred", "push_failed", "relay_fallback", "relays_failed", "smtp_error"] } }, + { status: { $gte: 400 } }, { failed: { $gt: 0 } }, + ] } : {}), + }; +} + +export async function inspectMailEvents(database, opts) { + const collection = database.db.collection("mail-events"); + const query = queryFor(opts); + const [summary, rows] = await Promise.all([ + collection.aggregate([ + { $match: query }, + { $group: { _id: { transport: "$transport", event: "$event", reason: "$reason", status: "$status" }, count: { $sum: 1 }, latest: { $max: "$when" } } }, + { $sort: { count: -1 } }, + ], { maxTimeMS: 5000 }).toArray(), + collection.find(query, { projection: { _id: 0 }, maxTimeMS: 5000 }).sort({ when: -1 }).limit(opts.limit).toArray(), + ]); + return { sinceMinutes: opts.since, retainedDays: 30, summary, events: rows.reverse() }; +} + +async function main() { + const opts = options(process.argv.slice(2)); + if (opts.help) { + console.log("mail-events-cli [--since minutes (default 1440)] [--limit 1–1000] [--failures] [--trace UUID] [--letter ID] [--json]\nSummary counts cover the full selected window; event rows are limited. Times are UTC. Retention: 30 days."); + return; + } + const { connect, closePool } = await import("./database.mjs"); + // Keep machine-readable stdout free of database startup chatter. + const saved = console.log; + console.log = (...args) => console.error(...args); + try { + const database = await connect(); + const result = await inspectMailEvents(database, opts); + if (opts.json) saved(JSON.stringify(result)); + else { + saved(`Last ${opts.since} minutes; ${result.events.length} latest events shown (30-day retention)`); + for (const row of result.summary) saved(`${row.count} ${Object.values(row._id).filter((v) => v != null).join(" / ")}`); + for (const { when, transport, event, trace, ...fields } of result.events) saved(`${when.toISOString()} ${transport || "-"} ${event} ${trace || "-"} ${JSON.stringify(fields)}`); + } + } finally { await closePool(); console.log = saved; } +} +if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) { + main().catch(() => { console.error("Mail event inspection failed; check options and database access."); process.exitCode = 1; }); +} diff --git a/system/backend/mail-events.mjs b/system/backend/mail-events.mjs new file mode 100644 index 0000000000..cba92500b8 --- /dev/null +++ b/system/backend/mail-events.mjs @@ -0,0 +1,89 @@ +// Operational mail metadata only. Never pass a letter or provider response here. +import { randomUUID } from "node:crypto"; +import { mailErrorCode } from "../../shared/mail-privacy.mjs"; + +export const mailTrace = () => randomUUID(); +const events = new Set([ + "started", "stored", "duplicate", "push", "push_failed", "push_quiet", + "relay_accepted", "relay_fallback", "failed", "request", "connected", + "recipient", "rejected", "deferred", "accepted", "disconnected", + "smtp_response", "ready", "relays_refreshed", "relays_failed", "smtp_error", +]); +const reasons = new Set([ + "relay", "recipient", "lookup", "wire_size", "attachments", "route", + "dmarc", "rate_pair", "rate_global", "storage", "smtp", "no_devices", + "push_limit", "unauthorized", "invalid", "not_found", "method", "request", +]); +const counts = ["recipients", "duplicates", "attachments", "bytes", "attempted", "succeeded", "failed", "pruned", "durationMs", "status"]; +const uuid = /^[a-f\d]{8}(?:-[a-f\d]{4}){3}-[a-f\d]{12}$/i; +export function mailEvent(fields = {}) { + if (!events.has(fields.event)) return null; + const row = { when: new Date(), event: fields.event }; + if (["internal", "smtp-in", "smtp-out", "api", "tell"].includes(fields.transport)) row.transport = fields.transport; + if (uuid.test(fields.trace || "")) row.trace = fields.trace; + if (/^[a-f\d]{24}$/i.test(String(fields.letterId || ""))) row.letterId = String(fields.letterId); + if (["inbox", "count", "download", "read", "send"].includes(fields.action)) row.action = fields.action; + if (reasons.has(fields.reason)) row.reason = fields.reason; + if (fields.error) row.error = mailErrorCode({ code: fields.error }); + for (const key of counts) if (Number.isSafeInteger(fields[key]) && fields[key] >= 0) row[key] = fields[key]; + for (const key of ["verified", "tls", "routeSecret", "open"]) if (typeof fields[key] === "boolean") row[key] = fields[key]; + return row; +} + +const indexed = new WeakMap(); +const pending = new Set(); +// Log first so a Mongo outage still leaves an inspectable journal event. A +// bounded asynchronous mirror cannot hold up SMTP or turn delivery into failure. +export function recordMailEvent(database, fields, log = console.log) { + const row = mailEvent(fields); + if (!row) return; + const writeLog = (value) => { try { log(JSON.stringify({ kind: "mail.event", ...value })); } catch {} }; + writeLog(row); + if (!database?.db) return; + if (pending.size >= 128) { + writeLog({ event: "telemetry_unavailable", reason: "queue_full" }); + return; + } + const db = database.db; + const task = Promise.resolve().then(async () => { + const collection = db.collection("mail-events"); + const client = db.client || db; + let databases = indexed.get(client); + if (!databases) { databases = new Map(); indexed.set(client, databases); } + const name = db.databaseName || ""; + let setup = databases.get(name); + if (!setup) { + setup = (async () => { + await collection.createIndex({ when: 1 }, { expireAfterSeconds: 30 * 86400, maxTimeMS: 1500 }); + await collection.createIndex({ trace: 1, when: 1 }, { maxTimeMS: 1500 }); + await collection.createIndex({ letterId: 1, when: 1 }, { maxTimeMS: 1500 }); + })(); + databases.set(name, setup); + setup.catch(() => databases.delete(name)); + } + await setup; + await collection.insertOne(row, { maxTimeMS: 1500 }); + }).catch((error) => writeLog({ event: "telemetry_unavailable", error: mailErrorCode(error) })); + pending.add(task); + task.finally(() => pending.delete(task)); +} + +// For bounded service shutdown / one-shot scripts; normal requests never wait. +export async function flushMailEvents() { + let timer; + await Promise.race([ + Promise.allSettled([...pending]), + new Promise((resolve) => { timer = setTimeout(resolve, 2000); }), + ]); + clearTimeout(timer); +} + +// Covers setup/storage errors as well as successful delivery. The same trace +// follows API → delivery → push, or SMTP transaction → recipient copies. +export async function observeMail(transport, options, database, operation) { + const trace = uuid.test(options.trace || "") ? options.trace : mailTrace(); + const event = (event, fields = {}) => recordMailEvent(database, { ...fields, event, transport, trace }); + event("started"); + try { return await operation({ ...options, trace }, event); } + catch (error) { event("failed", { error: mailErrorCode(error) }); throw error; } +} diff --git a/system/backend/mail.mjs b/system/backend/mail.mjs index faf1b37d42..f4fd27ee24 100644 --- a/system/backend/mail.mjs +++ b/system/backend/mail.mjs @@ -6,9 +6,9 @@ // permahandle (`ac25namuc`, permanent) and their @handle (follows the rename). // The permahandle is canonical; the @handle is an alias over it. +import { observeMail } from "./mail-events.mjs"; import { handleFor, userIDFromHandleOrEmail } from "./authorization.mjs"; import { filter } from "./filter.mjs"; -import { shell } from "./shell.mjs"; import { sendToUser } from "../../shared/push.mjs"; import { letterNotification, mailErrorCode, quietMailPush } from "../../shared/mail-privacy.mjs"; import { resolveMailMedia, outsideMediaBody } from "./mail-media.mjs"; @@ -80,10 +80,18 @@ export function clean(text, max = MAX_TEXT_LENGTH) { return filter((text || "").trim()).slice(0, max); } +export const deliver = (options, database) => observeMail("internal", options, database, + (options, event) => deliverInternal(options, database, event)); +export const deliverFromOutside = (options, database) => observeMail("smtp-in", options, database, + (options, event) => receiveOutside(options, database, event)); +export const sendOutside = (options, database) => observeMail("smtp-out", options, database, + (options, event) => relayOutside(options, database, event)); + // Put one message in a mailbox and buzz whatever devices the reader carries. -export async function deliver( +async function deliverInternal( { from, to, text, subject, device, verb = "mailed" }, database, + event, ) { const tells = await mailbox(database); const [fromHandle, toHandle] = await Promise.all([ @@ -105,6 +113,7 @@ export async function deliver( read: false, }); + event("stored", { letterId: insertedId }); let push = { attempted: 0, succeeded: 0, failed: 0, pruned: 0 }; try { push = await sendToUser( @@ -114,10 +123,10 @@ export async function deliver( { device }, quietMailPush, ); - if (push.failed) shell.log("mail.push.failed", push.failed); + event("push", { letterId: insertedId, ...push, ...(push.attempted === 0 ? { reason: "no_devices" } : {}) }); } catch (err) { // A silent phone shouldn't eat the letter — it's already in the mailbox. - shell.log("mail.push.error", mailErrorCode(err)); + event("push_failed", { letterId: insertedId, error: mailErrorCode(err) }); } return { id: insertedId, fromHandle, toHandle, when, push }; @@ -132,9 +141,10 @@ export async function deliver( // enough this hour. `messageId` keeps a relay retry from filing the same // letter twice — scoped to the box, so nobody can pre-empt another's letter. let dedupeIndexed = false; -export async function deliverFromOutside( +async function receiveOutside( { to, fromEmail, fromName, subject, text, messageId, auth = null, quiet = false, attachments = [] }, database, + event, ) { const tells = await mailbox(database); if (!dedupeIndexed) { @@ -170,12 +180,16 @@ export async function deliverFromOutside( read: false, })); } catch (err) { - if (err?.code === 11000) return { duplicate: true, toHandle }; // relay retried + if (err?.code === 11000) { event("duplicate"); return { duplicate: true, toHandle }; } // relay retried throw err; } + event("stored", { letterId: insertedId }); let push = { attempted: 0, succeeded: 0, failed: 0, pruned: 0 }; - if (quiet) return { id: insertedId, fromHandle, toHandle, when, push, quiet }; + if (quiet) { + event("push_quiet", { letterId: insertedId, reason: "push_limit" }); + return { id: insertedId, fromHandle, toHandle, when, push, quiet }; + } try { push = await sendToUser( database.db, @@ -184,9 +198,9 @@ export async function deliverFromOutside( {}, quietMailPush, ); - if (push.failed) shell.log("mail.push.failed", push.failed); + event("push", { letterId: insertedId, ...push, ...(push.attempted === 0 ? { reason: "no_devices" } : {}) }); } catch (err) { - shell.log("mail.push.error", mailErrorCode(err)); + event("push_failed", { letterId: insertedId, error: mailErrorCode(err) }); } return { id: insertedId, fromHandle, toHandle, when, push }; @@ -202,7 +216,7 @@ export async function deliverFromOutside( // the door — the plus-tag carries the permahandle, which outlives a rename. // It leaves through Google's SMTP relay with the mail@ credentials; the relay // lets any address in the domain sign, which plain smtp.gmail.com would not. -export async function sendOutside({ from, toEmail, subject, text }, database) { +async function relayOutside({ from, toEmail, subject, text }, database, event) { const nodemailer = (await import("nodemailer")).default; const [handle, user] = await Promise.all([ handleFor(from), @@ -237,12 +251,13 @@ export async function sendOutside({ from, toEmail, subject, text }, database) { .createTransport({ ...common, host: process.env.AMAIL_SMTP_SERVER || "smtp-relay.gmail.com" }) .sendMail(letter); } catch (err) { - shell.log("mail.relay.fallback", mailErrorCode(err)); + event("relay_fallback", { error: mailErrorCode(err) }); info = await nodemailer .createTransport({ ...common, host: process.env.SMTP_SERVER || "smtp.gmail.com" }) .sendMail(letter); } + event("relay_accepted"); const tells = await mailbox(database); const when = new Date(); const { insertedId } = await tells.insertOne({ @@ -258,5 +273,6 @@ export async function sendOutside({ from, toEmail, subject, text }, database) { when, read: true, }); + event("stored", { letterId: insertedId }); return { id: insertedId, fromHandle, toHandle: toEmail, when }; } diff --git a/system/netlify/functions/mail.mjs b/system/netlify/functions/mail.mjs index ad6963e9fc..2694387c4e 100644 --- a/system/netlify/functions/mail.mjs +++ b/system/netlify/functions/mail.mjs @@ -7,6 +7,7 @@ // POST /api/mail { action: "read" } → mark all read (or one, with `id`) // headers: Authorization: Bearer +import { mailTrace, recordMailEvent } from "../../backend/mail-events.mjs"; import { authorize } from "../../backend/authorization.mjs"; import { connect } from "../../backend/database.mjs"; import { respond as httpRespond } from "../../backend/http.mjs"; @@ -28,16 +29,25 @@ const respond = (status, body, headers = {}) => httpRespond(status, body, { "Cac const NO_FILES = { projection: { "attachments.data": 0 } }; export async function handler(event) { - try { - return await handleMail(event); - } catch (err) { - // Includes authorization, connection, and cleanup failures outside the query. - console.error("mail.request.error", mailErrorCode(err)); - return respond(500, { message: "Could not complete mail request" }); + const context = { trace: mailTrace(), action: (event.httpMethod === "GET" ? (event.queryStringParameters?.attachment !== undefined ? "download" : event.queryStringParameters?.count !== undefined ? "count" : "inbox") : "send") }; + const started = Date.now(); + let response; + try { response = await handleMail(event, context); } + catch (err) { + recordMailEvent(context.database, { event: "failed", transport: "api", trace: context.trace, error: mailErrorCode(err) }); + response = respond(500, { message: "Could not complete mail request" }); + } + const status = response.statusCode; + // Count polls are frequent; keep failures, writes, inbox and file reads. + if (!(event.httpMethod === "GET" && event.queryStringParameters?.count !== undefined && status === 200) && event.httpMethod !== "OPTIONS") { + recordMailEvent(context.database, { event: "request", transport: "api", trace: context.trace, + action: context.action, status, durationMs: Date.now() - started, + reason: status === 401 ? "unauthorized" : status === 400 ? "invalid" : status === 404 ? "not_found" : status === 405 ? "method" : status >= 500 ? "request" : undefined }); } + return { ...response, headers: { ...response.headers, "X-Mail-Trace": context.trace } }; } -async function handleMail(event) { +async function handleMail(event, context) { if (event.httpMethod === "OPTIONS") return respond(200, {}); if (event.httpMethod !== "GET" && event.httpMethod !== "POST") { return respond(405, { message: "Method Not Allowed" }); @@ -47,6 +57,7 @@ async function handleMail(event) { if (!user?.sub) return respond(401, { message: "unauthorized" }); const database = await connect(); + context.database = database; try { const tells = await mailbox(database); @@ -137,6 +148,7 @@ async function handleMail(event) { } if (body.action === "read") { + context.action = "read"; const where = { to: user.sub, read: { $ne: true } }; if (body.id) where._id = new ObjectId(body.id); const result = await tells.updateMany(where, { @@ -159,7 +171,7 @@ async function handleMail(event) { return respond(404, { message: "Recipient not found" }); } const sent = await sendOutside( - { from: user.sub, toEmail, subject, text }, + { trace: context.trace, from: user.sub, toEmail, subject, text }, database, ); return respond(200, { @@ -171,7 +183,7 @@ async function handleMail(event) { } const sentMail = await deliver( - { from: user.sub, to, text, subject, device: body.device }, + { trace: context.trace, from: user.sub, to, text, subject, device: body.device }, database, ); @@ -182,7 +194,7 @@ async function handleMail(event) { push: sentMail.push, }); } catch (err) { - console.error("mail.request.error", mailErrorCode(err)); + recordMailEvent(database, { event: "failed", transport: "api", trace: context.trace, error: mailErrorCode(err) }); return respond(500, { message: "Could not complete mail request" }); } finally { await database.disconnect(); diff --git a/system/netlify/functions/tell.mjs b/system/netlify/functions/tell.mjs index afd5720f73..4f5cf0cec6 100644 --- a/system/netlify/functions/tell.mjs +++ b/system/netlify/functions/tell.mjs @@ -7,6 +7,7 @@ // body: { to: "@handle" | "ac25namuc", text: "message", device?: "id-or-label" } // headers: Authorization: Bearer +import { mailTrace, recordMailEvent } from "../../backend/mail-events.mjs"; import { authorize } from "../../backend/authorization.mjs"; import { connect } from "../../backend/database.mjs"; import { respond } from "../../backend/http.mjs"; @@ -14,15 +15,22 @@ import { clean, deliver, subFromAddress } from "../../backend/mail.mjs"; import { mailErrorCode } from "../../../shared/mail-privacy.mjs"; export async function handler(event) { - try { - return await handleTell(event); - } catch (err) { - console.error("mail.tell.error", mailErrorCode(err)); - return respond(500, { message: "Could not send letter" }); + const context = { trace: mailTrace(), action: "send" }; + const started = Date.now(); + let response; + try { response = await handleTell(event, context); } + catch (err) { + recordMailEvent(context.database, { event: "failed", transport: "tell", trace: context.trace, error: mailErrorCode(err) }); + response = respond(500, { message: "Could not complete mail request" }); } + const status = response.statusCode; + recordMailEvent(context.database, { event: "request", transport: "tell", trace: context.trace, + action: context.action, status, durationMs: Date.now() - started, + reason: status === 401 ? "unauthorized" : status === 400 ? "invalid" : status === 404 ? "not_found" : status === 405 ? "method" : status >= 500 ? "request" : undefined }); + return { ...response, headers: { ...response.headers, "X-Mail-Trace": context.trace } }; } -async function handleTell(event) { +async function handleTell(event, context) { if (event.httpMethod !== "POST") { return respond(405, { message: "Method Not Allowed" }); } @@ -42,12 +50,13 @@ async function handleTell(event) { if (!sender?.sub) return respond(401, { message: "Unauthorized" }); const database = await connect(); + context.database = database; try { const to = await subFromAddress(body.to, database); if (!to) return respond(404, { message: "Recipient not found" }); const told = await deliver( - { from: sender.sub, to, text, device: body.device, verb: "told" }, + { trace: context.trace, from: sender.sub, to, text, device: body.device, verb: "told" }, database, ); @@ -58,7 +67,7 @@ async function handleTell(event) { push: told.push, }); } catch (err) { - console.error("mail.tell.error", mailErrorCode(err)); + recordMailEvent(database, { event: "failed", transport: "tell", trace: context.trace, error: mailErrorCode(err) }); return respond(500, { message: "Could not send letter" }); } finally { await database.disconnect();