node_modules: update (#338)

Co-authored-by: dawidd6 <9713907+dawidd6@users.noreply.github.com>
This commit is contained in:
Dawid Dziurlaanddawidd6 authored and GitHub committed 2026-10-03 10:23:35 +02:00
1 parent e1936242ed
commit a9c5eb6c35
39 files changed
+771 -360

No files matched your search

+3 -3
View File
@@ -94,9 +94,9 @@
} }
}, },
"node_modules/nodemailer": { "node_modules/nodemailer": {
"version": "10.0.11", "version": "10.0.13",
"resolved": "https://registry.npmjs.org/nodemailer/-/nodemailer-10.0.11.tgz", "resolved": "https://registry.npmjs.org/nodemailer/-/nodemailer-10.0.13.tgz",
"integrity": "sha512-/c4P7U7aGpiNu2Rl08q525q2zECkUpSaf2zQPgIStCcAjHAkkQQKwD4sTyi3l91mGcLQrKzL1syXPHymo+Ogtw==", "integrity": "sha512-SzG86OlvcW/NNhUFC6uROMwRTL4n7MswfQqC/T8mhkmnY1YVa23zUEMYi4ijSeXSl9GLz9ZeTJDUatEDuY5FeQ==",
"license": "MIT-0", "license": "MIT-0",
"engines": { "engines": {
"node": ">=20.0.0" "node": ">=20.0.0"
+15
View File
@@ -1,5 +1,20 @@
# CHANGELOG # CHANGELOG
## [10.0.13](https://github.com/nodemailer/nodemailer/compare/v10.0.12...v10.0.13) (2026-09-30)
### Bug Fixes
* read the advertised SASL methods without backtracking regexes ([b5a896f](https://github.com/nodemailer/nodemailer/commit/b5a896fcb02df1628329d7cd665ca08490b14b56))
* strip comments inside an angle-addr before it becomes the address ([a502247](https://github.com/nodemailer/nodemailer/commit/a502247555fd81773eac337486781f3b6008e16a))
## [10.0.12](https://github.com/nodemailer/nodemailer/compare/v10.0.11...v10.0.12) (2026-09-28)
### Bug Fixes
* settle every send on a connection error, back off pool requeues, turn a bare CR into CRLF, bound fetch, honour requireTLS ([63ccd66](https://github.com/nodemailer/nodemailer/commit/63ccd66c894bbaac18c82084ee385bf6488226cf))
## [10.0.11](https://github.com/nodemailer/nodemailer/compare/v10.0.10...v10.0.11) (2026-09-27) ## [10.0.11](https://github.com/nodemailer/nodemailer/compare/v10.0.10...v10.0.11) (2026-09-27)
+97
View File
@@ -242,6 +242,81 @@ function _recoverAddrSpec(data) {
.filter(part => part) .filter(part => part)
.join(' '); .join(' ');
} }
/**
* Takes the RFC 5322 comments out of the contents of an angle-addr.
*
* The tokenizer tracks a single open operator, so once '<' is open a '(' is plain text and
* a comment inside the brackets reached the address verbatim: 'Name <user@example.com(x)evil.com>'
* put 'user@example.com(x)evil.com' into the envelope, a value that is no mailbox and that a
* receiver stripping the comment reads as 'user@example.comevil.com' (GHSA-g73g-hqqh-jr95).
* A comment is folding whitespace, so it is read as a space that splits two atoms and as
* nothing next to an '@', the same rule the token walk applies outside the brackets. Quoted
* strings and domain literals are copied through untouched, a '(' in there is not a comment.
*
* @param address Contents of the angle brackets
* @return The address with the comments removed, and the text of the comments
*/
function _stripAddressComments(address) {
const comments = [];
let result = '';
let comment = '';
let depth = 0;
let closer = '';
// carried along rather than read back off the growing result, which would flatten it on
// every comment (see lastChars in _handleAddress)
let lastChar = '';
for (let i = 0, len = address.length; i < len; i++) {
const chr = address.charAt(i);
if (depth) {
if (chr === '\\' && i < len - 1) {
comment += address.charAt(++i);
}
else if (chr === '(') {
depth++;
comment += chr;
}
else if (chr === ')' && !--depth) {
comments.push(comment.trim());
comment = '';
if (lastChar !== '@' && address.charAt(i + 1) !== '@') {
result += ' ';
lastChar = ' ';
}
}
else {
comment += chr;
}
continue;
}
if (closer) {
if (chr === '\\' && closer === '"' && i < len - 1) {
result += chr + address.charAt(++i);
lastChar = address.charAt(i);
continue;
}
if (chr === closer) {
closer = '';
}
}
else if (chr === '"') {
closer = '"';
}
else if (chr === '[') {
closer = ']';
}
else if (chr === '(') {
depth = 1;
continue;
}
result += chr;
lastChar = chr;
}
if (depth) {
// an unterminated comment runs to the end of the address
comments.push(comment.trim());
}
return { address: result.trim(), comments: comments.filter(text => text) };
}
/** /**
* Converts tokens for a single address into an address object * Converts tokens for a single address into an address object
* *
@@ -369,6 +444,24 @@ function _handleAddress(tokens, depth) {
}); });
} }
else { else {
// Comments come out of the angle-addr before anything asks whether one was found, so
// that brackets holding nothing but a comment read the same as empty ones
const addressComments = [];
const addressParts = [];
for (const part of data.address) {
if (part.indexOf('(') < 0) {
addressParts.push(part);
continue;
}
const stripped = _stripAddressComments(part);
for (const comment of stripped.comments) {
addressComments.push(comment);
}
if (stripped.address) {
addressParts.push(stripped.address);
}
}
data.address = addressParts;
// If no address was found, try to detect one from regular text // If no address was found, try to detect one from regular text
if (!data.address.length && data.text.length) { if (!data.address.length && data.text.length) {
for (let i = data.text.length - 1; i >= 0; i--) { for (let i = data.text.length - 1; i >= 0; i--) {
@@ -435,6 +528,10 @@ function _handleAddress(tokens, depth) {
data.text = ''; data.text = '';
} }
_recoverAddrSpec(data); _recoverAddrSpec(data);
// a comment names the mailbox only when nothing else does, as one outside the brackets
if (!data.text && addressComments.length) {
data.text = addressComments.join(' ');
}
const address = { const address = {
address: data.address || data.text || '', address: data.address || data.text || '',
name: data.text || data.address || '' name: data.text || data.address || ''
+18 -1
View File
@@ -90,13 +90,22 @@ class DKIMSigner {
return; return;
} }
const key = this.keys[keyPos++]; const key = this.keys[keyPos++];
const dkimField = (0, sign_js_1.default)(this.headers, this.hashAlgo, this.bodyHash, { let dkimField;
try {
dkimField = (0, sign_js_1.default)(this.headers, this.hashAlgo, this.bodyHash, {
domainName: key.domainName, domainName: key.domainName,
keySelector: key.keySelector, keySelector: key.keySelector,
privateKey: key.privateKey, privateKey: key.privateKey,
headerFieldNames: this.options.headerFieldNames, headerFieldNames: this.options.headerFieldNames,
skipFields: this.options.skipFields skipFields: this.options.skipFields
}); });
}
catch (err) {
this.hasErrored = true;
this.cleanup();
this.output.emit('error', err);
return;
}
if (dkimField) { if (dkimField) {
this.output.write(Buffer.from(dkimField + '\r\n')); this.output.write(Buffer.from(dkimField + '\r\n'));
} }
@@ -196,7 +205,15 @@ class DKIM {
} }
const signer = new DKIMSigner(options, this.keys, inputStream, output); const signer = new DKIMSigner(options, this.keys, inputStream, output);
setImmediate(() => { setImmediate(() => {
try {
signer.signStream(); signer.signStream();
}
catch (_E) {
// the body hash is created here, an unknown hashAlgo throws inside this timer
// where nothing else could catch it
output.emit('error', sign_js_1.default.unsupportedHashAlgoError(signer.hashAlgo));
return;
}
if (writeValue) { if (writeValue) {
setImmediate(() => { setImmediate(() => {
inputStream.end(writeValue); inputStream.end(writeValue);
+2
View File
@@ -1,4 +1,5 @@
import crypto from 'node:crypto'; import crypto from 'node:crypto';
import type { NodemailerError } from '../errors.js';
import type { MessageParserHeaderLine } from './message-parser.js'; import type { MessageParserHeaderLine } from './message-parser.js';
/** /**
* Private key accepted by crypto.Sign#sign: a PEM string, a Buffer, a KeyObject or an * Private key accepted by crypto.Sign#sign: a PEM string, a Buffer, a KeyObject or an
@@ -45,5 +46,6 @@ export interface DKIMRelaxedHeaders {
declare function sign(headers: MessageParserHeaderLine[], hashAlgo: string, bodyHash: string, options?: DKIMSignOptions): string | false; declare function sign(headers: MessageParserHeaderLine[], hashAlgo: string, bodyHash: string, options?: DKIMSignOptions): string | false;
declare namespace sign { declare namespace sign {
var relaxedHeaders: (headers: MessageParserHeaderLine[], fieldNames?: string, skipFields?: string) => DKIMRelaxedHeaders; var relaxedHeaders: (headers: MessageParserHeaderLine[], fieldNames?: string, skipFields?: string) => DKIMRelaxedHeaders;
var unsupportedHashAlgoError: (hashAlgo: string) => NodemailerError;
} }
export default sign; export default sign;
+17 -1
View File
@@ -39,6 +39,15 @@ Object.defineProperty(exports, "__esModule", { value: true });
const punycode = __importStar(require("../punycode/index.js")); const punycode = __importStar(require("../punycode/index.js"));
const mimeFuncs = __importStar(require("../mime-funcs/index.js")); const mimeFuncs = __importStar(require("../mime-funcs/index.js"));
const node_crypto_1 = __importDefault(require("node:crypto")); const node_crypto_1 = __importDefault(require("node:crypto"));
const errors = __importStar(require("../errors.js"));
/**
* Error for a hashAlgo value the crypto module does not know
*/
function unsupportedHashAlgoError(hashAlgo) {
const err = new Error('Unsupported DKIM hash algorithm "' + hashAlgo + '"');
err.code = errors.ECONFIG;
return err;
}
/** /**
* Returns DKIM signature header line * Returns DKIM signature header line
* *
@@ -63,7 +72,13 @@ function sign(headers, hashAlgo, bodyHash, options) {
const canonicalizedHeaderData = relaxedHeaders(headers, fieldNames, options.skipFields); const canonicalizedHeaderData = relaxedHeaders(headers, fieldNames, options.skipFields);
const dkimHeader = generateDKIMHeader(options.domainName, options.keySelector, canonicalizedHeaderData.fieldNames, hashAlgo, bodyHash); const dkimHeader = generateDKIMHeader(options.domainName, options.keySelector, canonicalizedHeaderData.fieldNames, hashAlgo, bodyHash);
canonicalizedHeaderData.headers += 'dkim-signature:' + relaxedHeaderLine(dkimHeader); canonicalizedHeaderData.headers += 'dkim-signature:' + relaxedHeaderLine(dkimHeader);
const signer = node_crypto_1.default.createSign(('rsa-' + hashAlgo).toUpperCase()); let signer;
try {
signer = node_crypto_1.default.createSign(('rsa-' + hashAlgo).toUpperCase());
}
catch (_E) {
throw unsupportedHashAlgoError(hashAlgo);
}
// the header lines are 'binary' strings, so this reproduces the original header bytes // the header lines are 'binary' strings, so this reproduces the original header bytes
signer.update(canonicalizedHeaderData.headers, 'latin1'); signer.update(canonicalizedHeaderData.headers, 'latin1');
let signature; let signature;
@@ -76,6 +91,7 @@ function sign(headers, hashAlgo, bodyHash, options) {
return dkimHeader + signature.replace(/(^.{73}|.{75}(?!\r?\n|\r))/g, '$&\r\n ').trim(); return dkimHeader + signature.replace(/(^.{73}|.{75}(?!\r?\n|\r))/g, '$&\r\n ').trim();
} }
sign.relaxedHeaders = relaxedHeaders; sign.relaxedHeaders = relaxedHeaders;
sign.unsupportedHashAlgoError = unsupportedHashAlgoError;
exports.default = sign; exports.default = sign;
function generateDKIMHeader(domainName, keySelector, fieldNames, hashAlgo, bodyHash) { function generateDKIMHeader(domainName, keySelector, fieldNames, hashAlgo, bodyHash) {
// the caller supplied tag values are interpolated straight into the tag list, and none of // the caller supplied tag values are interpolated straight into the tag list, and none of
+4 -1
View File
@@ -25,8 +25,10 @@ export interface FetchOptions {
tls?: { tls?: {
[key: string]: any; [key: string]: any;
} | undefined; } | undefined;
/** Request timeout in milliseconds */ /** Socket inactivity timeout in milliseconds, defaults to 60000, 0 disables it */
timeout?: number | undefined; timeout?: number | undefined;
/** Maximum size of the (decoded) response body in bytes, defaults to 64 MB, Infinity disables the limit */
maxBytes?: number | undefined;
/** Maximum number of redirects to follow (default 5) */ /** Maximum number of redirects to follow (default 5) */
maxRedirects?: number | undefined; maxRedirects?: number | undefined;
/** Resolve responses with a status code of 300 or above instead of emitting an error */ /** Resolve responses with a status code of 300 or above instead of emitting an error */
@@ -47,6 +49,7 @@ export interface FetchResponse extends PassThrough {
declare function nmfetch(url: string, options?: FetchOptions): FetchResponse; declare function nmfetch(url: string, options?: FetchOptions): FetchResponse;
declare namespace nmfetch { declare namespace nmfetch {
var Cookies: typeof import("./cookies.js").default; var Cookies: typeof import("./cookies.js").default;
var DEFAULT_TIMEOUT: number;
} }
type CookiesJar = Cookies; type CookiesJar = Cookies;
/** /**
+35 -60
View File
@@ -47,6 +47,10 @@ const node_net_1 = __importDefault(require("node:net"));
const errors = __importStar(require("../errors.js")); const errors = __importStar(require("../errors.js"));
const objects_js_1 = require("../shared/objects.js"); const objects_js_1 = require("../shared/objects.js");
const MAX_REDIRECTS = 5; const MAX_REDIRECTS = 5;
// a stalled server would otherwise hold the request, and whatever waits on it, open forever
const DEFAULT_TIMEOUT = 60 * 1000;
// the body is usually buffered in memory by the caller, so an unbounded download is refused
const DEFAULT_MAX_BYTES = 64 * 1024 * 1024;
// Only genuine TLS settings are taken from options.tls. That object reaches us straight // Only genuine TLS settings are taken from options.tls. That object reaches us straight
// from a user supplied attachment (content.tls), so keys like host, port, path, socketPath // from a user supplied attachment (content.tls), so keys like host, port, path, socketPath
// or lookup would otherwise repoint the request at a destination that never went through // or lookup would otherwise repoint the request at a destination that never went through
@@ -252,28 +256,22 @@ function nmfetch(url, options) {
}); });
return fetchRes; return fetchRes;
} }
if (options.timeout) { // reports the first failure of this request on fetchRes and releases the request
req.setTimeout(options.timeout, () => { const fail = (err, sourceUrl = url) => {
if (finished) { if (finished) {
return; return;
} }
finished = true; finished = true;
err.code = errors.EFETCH;
err.sourceUrl = sourceUrl;
fetchRes.emit('error', err);
req.abort(); req.abort();
const err = new Error('Request Timeout'); };
err.code = errors.EFETCH; const timeout = typeof options.timeout === 'number' && options.timeout >= 0 ? options.timeout : DEFAULT_TIMEOUT;
err.sourceUrl = url; if (timeout) {
fetchRes.emit('error', err); req.setTimeout(timeout, () => fail(new Error('Request Timeout')));
});
} }
req.on('error', (err) => { req.on('error', (err) => fail(err));
if (finished) {
return;
}
finished = true;
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
});
req.on('response', res => { req.on('response', res => {
let inflate; let inflate;
if (finished) { if (finished) {
@@ -294,13 +292,7 @@ function nmfetch(url, options) {
// redirect // redirect
options.redirects++; options.redirects++;
if (options.redirects > options.maxRedirects) { if (options.redirects > options.maxRedirects) {
finished = true; return fail(new Error('Maximum redirect count exceeded'));
const err = new Error('Maximum redirect count exceeded');
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
req.abort();
return;
} }
// redirect does not include POST body // redirect does not include POST body
options.method = 'GET'; options.method = 'GET';
@@ -320,13 +312,7 @@ function nmfetch(url, options) {
// call: that call gets its own `finished` flag and no handle on this // call: that call gets its own `finished` flag and no handle on this
// request, so this one would stay open and could emit a second error on // request, so this one would stay open and could emit a second error on
// the shared fetchRes once it times out. Callers listen with req.once(). // the shared fetchRes once it times out. Callers listen with req.once().
finished = true; return fail(new Error('Unsupported protocol for URL ' + redirectUrl), redirectUrl);
const err = new Error('Unsupported protocol for URL ' + redirectUrl);
err.code = errors.EFETCH;
err.sourceUrl = redirectUrl;
fetchRes.emit('error', err);
req.abort();
return;
} }
// Do not forward credentials when the redirect leaves the original // Do not forward credentials when the redirect leaves the original
// security context: a different host, or a downgrade from https to // security context: a different host, or a downgrade from https to
@@ -343,41 +329,33 @@ function nmfetch(url, options) {
} }
}); });
} }
// this request is done with, release it so its socket and timeout do not
// outlive it and report a late error on the shared fetchRes
finished = true;
res.resume();
req.abort();
return nmfetch(redirectUrl, options); return nmfetch(redirectUrl, options);
} }
fetchRes.statusCode = res.statusCode; fetchRes.statusCode = res.statusCode;
fetchRes.headers = res.headers; fetchRes.headers = res.headers;
if (res.statusCode >= 300 && !options.allowErrorResponse) { if (res.statusCode >= 300 && !options.allowErrorResponse) {
finished = true; return fail(new Error('Invalid status code ' + res.statusCode));
const err = new Error('Invalid status code ' + res.statusCode); }
err.code = errors.EFETCH; res.on('error', (err) => fail(err));
err.sourceUrl = url; const maxBytes = typeof options.maxBytes === 'number' && options.maxBytes > 0 ? options.maxBytes : DEFAULT_MAX_BYTES;
fetchRes.emit('error', err); const source = inflate || res;
req.abort(); let received = 0;
source.on('data', (chunk) => {
received += chunk.length;
if (received <= maxBytes || finished) {
return; return;
} }
res.on('error', (err) => { source.unpipe(fetchRes);
if (finished) { fail(new Error('Response size exceeds the allowed ' + maxBytes + ' bytes'));
return;
}
finished = true;
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
req.abort();
}); });
if (inflate) { if (inflate) {
res.pipe(inflate).pipe(fetchRes); res.pipe(inflate).pipe(fetchRes);
inflate.on('error', (err) => { inflate.on('error', (err) => fail(err));
if (finished) {
return;
}
finished = true;
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
req.abort();
});
} }
else { else {
res.pipe(fetchRes); res.pipe(fetchRes);
@@ -392,11 +370,7 @@ function nmfetch(url, options) {
req.write(body); req.write(body);
} }
catch (err) { catch (err) {
finished = true; return fail(err);
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
return;
} }
} }
req.end(); req.end();
@@ -404,6 +378,7 @@ function nmfetch(url, options) {
return fetchRes; return fetchRes;
} }
nmfetch.Cookies = cookies_js_1.default; nmfetch.Cookies = cookies_js_1.default;
nmfetch.DEFAULT_TIMEOUT = DEFAULT_TIMEOUT;
exports.default = nmfetch; exports.default = nmfetch;
module.exports = exports.default; module.exports = exports.default;
Object.defineProperty(module.exports, 'default', { value: exports.default, enumerable: false, writable: true, configurable: true }); Object.defineProperty(module.exports, 'default', { value: exports.default, enumerable: false, writable: true, configurable: true });
+1 -1
View File
@@ -1,3 +1,3 @@
export declare const name = "nodemailer"; export declare const name = "nodemailer";
export declare const version = "10.0.11"; export declare const version = "10.0.13";
export declare const homepage = "https://nodemailer.com/"; export declare const homepage = "https://nodemailer.com/";
+1 -1
View File
@@ -3,5 +3,5 @@
Object.defineProperty(exports, "__esModule", { value: true }); Object.defineProperty(exports, "__esModule", { value: true });
exports.homepage = exports.version = exports.name = void 0; exports.homepage = exports.version = exports.name = void 0;
exports.name = 'nodemailer'; exports.name = 'nodemailer';
exports.version = '10.0.11'; exports.version = '10.0.13';
exports.homepage = 'https://nodemailer.com/'; exports.homepage = 'https://nodemailer.com/';
+2 -1
View File
@@ -1,7 +1,8 @@
import { Transform, type TransformOptions } from 'node:stream'; import { Transform, type TransformOptions } from 'node:stream';
/** /**
* Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF> * Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF>
* Also makes sure that only <CR><LF> sequences are used for linebreaks * Also makes sure that only <CR><LF> sequences are used for linebreaks, bare CR and bare LF
* are both turned into <CR><LF>
* *
* @param options Stream options * @param options Stream options
*/ */
+28 -20
View File
@@ -1,9 +1,15 @@
"use strict"; "use strict";
Object.defineProperty(exports, "__esModule", { value: true }); Object.defineProperty(exports, "__esModule", { value: true });
const node_stream_1 = require("node:stream"); const node_stream_1 = require("node:stream");
// bytes inserted into the output, shared as they are only ever copied by Buffer.concat
const INSERT_LF = Buffer.from('\n');
const INSERT_LF_DOT = Buffer.from('\n.');
const INSERT_CR = Buffer.from('\r');
const INSERT_DOT = Buffer.from('.');
/** /**
* Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF> * Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF>
* Also makes sure that only <CR><LF> sequences are used for linebreaks * Also makes sure that only <CR><LF> sequences are used for linebreaks, bare CR and bare LF
* are both turned into <CR><LF>
* *
* @param options Stream options * @param options Stream options
*/ */
@@ -32,33 +38,35 @@ class DataStream extends node_stream_1.Transform {
} }
this.inByteCount += chunk.length; this.inByteCount += chunk.length;
for (i = 0, len = chunk.length; i < len; i++) { for (i = 0, len = chunk.length; i < len; i++) {
if (chunk[i] === 0x2e) { const byte = chunk[i];
// . const prev = i ? chunk[i - 1] : this.lastByte;
if ((i && chunk[i - 1] === 0x0a) || (!i && (!this.lastByte || this.lastByte === 0x0a))) { let insert = false;
buf = chunk.slice(lastPos, i + 1); if (prev === 0x0d && byte !== 0x0a) {
chunks.push(buf); // a bare CR becomes CRLF. A receiver that treats a lone CR as a line end would
chunks.push(Buffer.from('.')); // otherwise see "\r.\r" as the end of the data (SMTP smuggling), so a dot
chunklen += buf.length + 1; // following it is stuffed like at the start of any other line
lastPos = i + 1; insert = byte === 0x2e ? INSERT_LF_DOT : INSERT_LF;
} }
else if (byte === 0x0a && prev !== 0x0d) {
// a bare LF becomes CRLF
insert = INSERT_CR;
} }
else if (chunk[i] === 0x0a) { else if (byte === 0x2e && (prev === 0x0a || prev === false)) {
// \n // a dot at the start of a line
if ((i && chunk[i - 1] !== 0x0d) || (!i && this.lastByte !== 0x0d)) { insert = INSERT_DOT;
}
if (insert) {
if (i > lastPos) { if (i > lastPos) {
buf = chunk.slice(lastPos, i); buf = chunk.slice(lastPos, i);
chunks.push(buf); chunks.push(buf);
chunklen += buf.length + 2; chunklen += buf.length;
} }
else { chunks.push(insert);
chunklen += 2; chunklen += insert.length;
} lastPos = i;
chunks.push(Buffer.from('\r\n'));
lastPos = i + 1;
} }
} }
} if (chunks.length) {
if (chunklen) {
// add last piece // add last piece
if (lastPos < chunk.length) { if (lastPos < chunk.length) {
buf = chunk.slice(lastPos); buf = chunk.slice(lastPos);
+5 -5
View File
@@ -26,11 +26,11 @@ export interface SMTPConnectionOptions {
secured?: boolean | undefined; secured?: boolean | undefined;
/** Server name for SNI, defaults to host when that is not an IP address */ /** Server name for SNI, defaults to host when that is not an IP address */
servername?: string | undefined; servername?: string | undefined;
/** Ignore STARTTLS even when the server advertises it */ /** Ignore STARTTLS even when the server advertises it, has no effect when requireTLS is set */
ignoreTLS?: boolean | undefined; ignoreTLS?: boolean | undefined;
/** Force STARTTLS, fail when the server does not support it */ /** Force STARTTLS, fail when the server does not support it. Takes precedence over ignoreTLS and opportunisticTLS */
requireTLS?: boolean | undefined; requireTLS?: boolean | undefined;
/** Continue unencrypted when the STARTTLS upgrade fails */ /** Continue unencrypted when the STARTTLS upgrade fails, has no effect when requireTLS is set */
opportunisticTLS?: boolean | undefined; opportunisticTLS?: boolean | undefined;
/** Name of the client server, sent with EHLO/HELO, CRLF is stripped */ /** Name of the client server, sent with EHLO/HELO, CRLF is stripped */
name?: string | undefined; name?: string | undefined;
@@ -288,8 +288,8 @@ export type SMTPConnectionResponseAction = (str: string) => void;
* * **port** - is the port to connect to (defaults to 587 or 465) * * **port** - is the port to connect to (defaults to 587 or 465)
* * **host** - is the hostname or IP address to connect to (defaults to 'localhost') * * **host** - is the hostname or IP address to connect to (defaults to 'localhost')
* * **secure** - use SSL * * **secure** - use SSL
* * **ignoreTLS** - ignore server support for STARTTLS * * **ignoreTLS** - ignore server support for STARTTLS (has no effect when requireTLS is set)
* * **requireTLS** - forces the client to use STARTTLS * * **requireTLS** - forces the client to use STARTTLS, takes precedence over ignoreTLS and opportunisticTLS
* * **name** - the name of the client server * * **name** - the name of the client server
* * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener) * * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener)
* * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds) * * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds)
+75 -21
View File
@@ -104,8 +104,8 @@ function isPartialLine(line) {
* * **port** - is the port to connect to (defaults to 587 or 465) * * **port** - is the port to connect to (defaults to 587 or 465)
* * **host** - is the hostname or IP address to connect to (defaults to 'localhost') * * **host** - is the hostname or IP address to connect to (defaults to 'localhost')
* * **secure** - use SSL * * **secure** - use SSL
* * **ignoreTLS** - ignore server support for STARTTLS * * **ignoreTLS** - ignore server support for STARTTLS (has no effect when requireTLS is set)
* * **requireTLS** - forces the client to use STARTTLS * * **requireTLS** - forces the client to use STARTTLS, takes precedence over ignoreTLS and opportunisticTLS
* * **name** - the name of the client server * * **name** - the name of the client server
* * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener) * * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener)
* * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds) * * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds)
@@ -130,6 +130,11 @@ class SMTPConnection extends node_events_1.EventEmitter {
this.id = node_crypto_1.default.randomBytes(8).toString('base64').replace(/\W/g, ''); this.id = node_crypto_1.default.randomBytes(8).toString('base64').replace(/\W/g, '');
this.stage = 'init'; this.stage = 'init';
this.options = options || {}; this.options = options || {};
if (this.options.requireTLS && (this.options.ignoreTLS || this.options.opportunisticTLS)) {
// requireTLS wins, a contradictory configuration must not quietly fall back to plaintext.
// Copied so the caller's (possibly shared) options object is left as it was
this.options = Object.assign({}, this.options, { ignoreTLS: false, opportunisticTLS: false });
}
this.secureConnection = !!this.options.secure; this.secureConnection = !!this.options.secure;
this.alreadySecured = !!this.options.secured; this.alreadySecured = !!this.options.secured;
this.port = Number(this.options.port) || (this.secureConnection ? 465 : 587); this.port = Number(this.options.port) || (this.secureConnection ? 465 : 587);
@@ -173,6 +178,8 @@ class SMTPConnection extends node_events_1.EventEmitter {
this._destroyed = false; this._destroyed = false;
this._closing = false; this._closing = false;
this._currentDataStream = false; this._currentDataStream = false;
this._pendingSend = false;
this._connectCallback = false;
this._onSocketData = chunk => this._onData(chunk); this._onSocketData = chunk => this._onData(chunk);
this._onSocketError = error => this._onError(error, 'ESOCKET', false, 'CONN'); this._onSocketError = error => this._onError(error, 'ESOCKET', false, 'CONN');
this._onSocketClose = () => this._onClose(); this._onSocketClose = () => this._onClose();
@@ -187,7 +194,9 @@ class SMTPConnection extends node_events_1.EventEmitter {
*/ */
connect(connectCallback) { connect(connectCallback) {
if (typeof connectCallback === 'function') { if (typeof connectCallback === 'function') {
this._connectCallback = connectCallback;
this.once('connect', () => { this.once('connect', () => {
this._connectCallback = false;
this.logger.debug({ this.logger.debug({
tnx: 'smtp' tnx: 'smtp'
}, 'SMTP handshake finished'); }, 'SMTP handshake finished');
@@ -419,6 +428,16 @@ class SMTPConnection extends node_events_1.EventEmitter {
} }
this._currentDataStream = false; this._currentDataStream = false;
} }
// Detach from the message stream as well. The listener is swapped for a no-op rather than
// removed, a stream destroyed with an error later on would otherwise throw it as unhandled
if (this._pendingSend) {
const { stream, onStreamError } = this._pendingSend;
if (stream) {
stream.removeListener('error', onStreamError);
stream.on('error', TEARDOWN_NOOP);
}
this._pendingSend = false;
}
if (socket && !socket.destroyed) { if (socket && !socket.destroyed) {
try { try {
// Clear socket timeout to prevent timer leaks // Clear socket timeout to prevent timer leaks
@@ -606,6 +625,9 @@ class SMTPConnection extends node_events_1.EventEmitter {
return; return;
} }
returned = true; returned = true;
if (this._pendingSend && this._pendingSend.callback === callback) {
this._pendingSend = false;
}
done(err, info); done(err, info);
}; };
if (!message) { if (!message) {
@@ -622,9 +644,16 @@ class SMTPConnection extends node_events_1.EventEmitter {
}); });
return; return;
} }
const pendingSend = {
callback,
stream: false,
onStreamError: err => callback(this._formatError(err, 'ESTREAM', false, 'API'))
};
if (typeof message.on === 'function') { if (typeof message.on === 'function') {
message.on('error', err => callback(this._formatError(err, 'ESTREAM', false, 'API'))); pendingSend.stream = message;
pendingSend.stream.on('error', pendingSend.onStreamError);
} }
this._pendingSend = pendingSend;
const startTime = Date.now(); const startTime = Date.now();
this._setEnvelope(envelope, (err, info) => { this._setEnvelope(envelope, (err, info) => {
if (err) { if (err) {
@@ -817,8 +846,14 @@ class SMTPConnection extends node_events_1.EventEmitter {
else { else {
this.logger.error(data, err.message); this.logger.error(data, err.message);
} }
// close() forgets the send in flight, it is completed with this same error afterwards so
// a late message stream error has nothing left to report
const pendingSend = this._pendingSend;
this.emit('error', err); this.emit('error', err);
this.close(); this.close();
if (pendingSend) {
pendingSend.callback(err);
}
} }
/** @internal */ /** @internal */
_formatError(message, type, response, command) { _formatError(message, type, response, command) {
@@ -864,14 +899,25 @@ class SMTPConnection extends node_events_1.EventEmitter {
this.logger.info({ this.logger.info({
tnx: 'network' tnx: 'network'
}, 'Connection closed'); }, 'Connection closed');
// the unterminated remainder is only reported as a reply (and so gives the error a responseCode)
// when it starts like a complete failure reply, not for a fragment such as "55" or a 250
const failureResponse = typeof serverResponse === 'string' && /^[45]\d{2}[ -]/.test(serverResponse) ? serverResponse : false;
if (this.upgrading && !this._destroyed) { if (this.upgrading && !this._destroyed) {
return this._onError(new Error('Connection closed unexpectedly'), 'ETLS', serverResponse, 'CONN'); return this._onError(new Error('Connection closed unexpectedly'), 'ETLS', failureResponse, 'CONN');
} }
else if (![this._actionGreeting, this.close].includes(this._responseActions[0]) && !this._destroyed) { if (!failureResponse && this._responseActions[0] === this._actionGreeting && this._connectCallback && !this._destroyed) {
return this._onError(new Error('Connection closed unexpectedly'), 'ECONNECTION', serverResponse, 'CONN'); // A silent close before the greeting is handed to the connect() callback rather than
// emitted as 'error', callers that never saw an error for it must not start throwing one
const connectCallback = this._connectCallback;
this._connectCallback = false;
const err = this._formatError(new Error('Connection closed unexpectedly'), 'ECONNECTION', false, 'CONN');
this.logger.warn({ tnx: 'network' }, err.message);
connectCallback(err);
this.close();
return;
} }
else if (/^[45]\d{2}\b/.test(serverResponse)) { if (failureResponse || (this._responseActions[0] !== this.close && !this._destroyed)) {
return this._onError(new Error('Connection closed unexpectedly'), 'ECONNECTION', serverResponse, 'CONN'); return this._onError(new Error('Connection closed unexpectedly'), 'ECONNECTION', failureResponse, 'CONN');
} }
this._destroy(); this._destroy();
} }
@@ -1348,21 +1394,25 @@ class SMTPConnection extends node_events_1.EventEmitter {
if (/[ -]AUTH\b/i.test(str)) { if (/[ -]AUTH\b/i.test(str)) {
this.allowsAuth = true; this.allowsAuth = true;
} }
// Detect if the server supports PLAIN auth // Detect the advertised SASL mechanisms. The list is split into whole tokens rather
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)PLAIN/i.test(str)) { // than searched for each name with a pattern: the patterns this replaced let two
this._supportedAuth.push('PLAIN'); // whitespace runs overlap and backtracked quadratically over an AUTH line padded with
// spaces, so a hostile server could stall the event loop from its EHLO reply
// (GHSA-4ffr-jq9g-5ffx).
const authMechanisms = new Set();
for (const line of this._ehloLines) {
const authMatch = /^AUTH[\s=](.*)/i.exec(line);
if (authMatch) {
for (const mechanism of authMatch[1].split(/[\s=]+/)) {
authMechanisms.add(mechanism.toUpperCase());
} }
// Detect if the server supports LOGIN auth
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)LOGIN/i.test(str)) {
this._supportedAuth.push('LOGIN');
} }
// Detect if the server supports CRAM-MD5 auth
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)CRAM-MD5/i.test(str)) {
this._supportedAuth.push('CRAM-MD5');
} }
// Detect if the server supports XOAUTH2 auth // listed in order of preference, the first one the credentials allow is used
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)XOAUTH2/i.test(str)) { for (const mechanism of ['PLAIN', 'LOGIN', 'CRAM-MD5', 'XOAUTH2']) {
this._supportedAuth.push('XOAUTH2'); if (authMechanisms.has(mechanism)) {
this._supportedAuth.push(mechanism);
}
} }
// Detect if the server supports SIZE extensions (and the max allowed size) // Detect if the server supports SIZE extensions (and the max allowed size)
if ((match = str.match(/[ -]SIZE(?:[ \t]+(\d+))?/im))) { if ((match = str.match(/[ -]SIZE(?:[ \t]+(\d+))?/im))) {
@@ -1622,7 +1672,11 @@ class SMTPConnection extends node_events_1.EventEmitter {
this._sendCommand('DATA'); this._sendCommand('DATA');
} }
else { else {
err = this._formatError("Can't send mail - all recipients were rejected", 'EENVELOPE', str, 'RCPT TO'); // report a temporary rejection when there is one, taking the last reply would mark the
// whole message as permanently failed although some recipients were only deferred
const deferred = envelope.rejectedErrors.find(rejectedErr => rejectedErr.responseCode && rejectedErr.responseCode < 500);
const reply = deferred?.response ?? str;
err = this._formatError("Can't send mail - all recipients were rejected", 'EENVELOPE', reply, 'RCPT TO');
err.rejected = envelope.rejected; err.rejected = envelope.rejected;
err.rejectedErrors = envelope.rejectedErrors; err.rejectedErrors = envelope.rejectedErrors;
return callback(err); return callback(err);
+2 -1
View File
@@ -17,7 +17,7 @@ export interface SMTPPoolOptions extends SMTPTransportOptions {
rateLimit?: number | undefined; rateLimit?: number | undefined;
/** Time window for rateLimit in milliseconds, defaults to 1000 */ /** Time window for rateLimit in milliseconds, defaults to 1000 */
rateDelta?: number | undefined; rateDelta?: number | undefined;
/** How many times a message is requeued when its connection closes while sending, unlimited when not set or negative */ /** How many times a message is requeued when its connection closes while sending, defaults to 5, a negative value means unlimited */
maxRequeues?: number | undefined; maxRequeues?: number | undefined;
} }
/** /**
@@ -26,6 +26,7 @@ export interface SMTPPoolOptions extends SMTPTransportOptions {
export type SMTPPoolResolvedOptions = SMTPPoolOptions & { export type SMTPPoolResolvedOptions = SMTPPoolOptions & {
maxConnections: number; maxConnections: number;
maxMessages: number; maxMessages: number;
maxRequeues: number;
}; };
/** /**
* Result of a message sent through the pool, same as for the SMTP transport * Result of a message sent through the pool, same as for the SMTP transport
+30 -4
View File
@@ -43,6 +43,10 @@ const index_js_2 = __importDefault(require("../well-known/index.js"));
const shared = __importStar(require("../shared/index.js")); const shared = __importStar(require("../shared/index.js"));
const errors = __importStar(require("../errors.js")); const errors = __importStar(require("../errors.js"));
const packageData = __importStar(require("../package-info.js")); const packageData = __importStar(require("../package-info.js"));
/** First delay before a requeued message is retried, doubled on every further requeue */
const REQUEUE_BASE_DELAY = 50;
/** Upper bound for the requeue delay */
const REQUEUE_MAX_DELAY = 2000;
/** /**
* Creates a SMTP pool transport object for Nodemailer * Creates a SMTP pool transport object for Nodemailer
* *
@@ -74,6 +78,8 @@ class SMTPPool extends node_events_1.EventEmitter {
); );
this.options.maxConnections = this.options.maxConnections || 5; this.options.maxConnections = this.options.maxConnections || 5;
this.options.maxMessages = this.options.maxMessages || 100; this.options.maxMessages = this.options.maxMessages || 100;
// a default bound, a server that closes every connection before the greeting would otherwise be retried forever
this.options.maxRequeues = typeof this.options.maxRequeues === 'number' ? this.options.maxRequeues : 5;
this.logger = shared.getLogger(this.options, { this.logger = shared.getLogger(this.options, {
component: this.options.component || 'smtp-pool' component: this.options.component || 'smtp-pool'
}); });
@@ -117,6 +123,10 @@ class SMTPPool extends node_events_1.EventEmitter {
*/ */
send(mail, callback) { send(mail, callback) {
if (this._closed) { if (this._closed) {
// Mail.sendMail ignores the return value, so without a callback its promise would never settle
const err = new Error('Connection pool was closed');
err.code = errors.ECONNECTION;
setImmediate(() => callback(err));
return false; return false;
} }
this._queue.push({ this._queue.push({
@@ -330,15 +340,22 @@ class SMTPPool extends node_events_1.EventEmitter {
// Note that we must wait a bit.. because the callback of the 'error' handler might be called // Note that we must wait a bit.. because the callback of the 'error' handler might be called
// in the next event loop // in the next event loop
setTimeout(() => { setTimeout(() => {
let delay = 0;
if (connection.queueEntry) { if (connection.queueEntry) {
if (this._shouldRequeuOnConnectionClose(connection.queueEntry)) { if (this._shouldRequeuOnConnectionClose(connection.queueEntry)) {
this._requeueEntryOnConnectionClose(connection); delay = this._requeueEntryOnConnectionClose(connection);
} }
else { else {
this._failDeliveryOnConnectionClose(connection); this._failDeliveryOnConnectionClose(connection);
} }
} }
if (delay) {
// back off, a server that keeps dropping connections is not hammered
setTimeout(() => this._continueProcessing(), delay);
}
else {
this._continueProcessing(); this._continueProcessing();
}
}, 50); }, 50);
} }
else { else {
@@ -353,7 +370,7 @@ class SMTPPool extends node_events_1.EventEmitter {
} }
/** @internal */ /** @internal */
_shouldRequeuOnConnectionClose(queueEntry) { _shouldRequeuOnConnectionClose(queueEntry) {
if (this.options.maxRequeues === undefined || this.options.maxRequeues < 0) { if (this.options.maxRequeues < 0) {
return true; return true;
} }
return queueEntry.requeueAttempts < this.options.maxRequeues; return queueEntry.requeueAttempts < this.options.maxRequeues;
@@ -362,7 +379,9 @@ class SMTPPool extends node_events_1.EventEmitter {
_failDeliveryOnConnectionClose(connection) { _failDeliveryOnConnectionClose(connection) {
if (connection.queueEntry && connection.queueEntry.callback) { if (connection.queueEntry && connection.queueEntry.callback) {
try { try {
connection.queueEntry.callback(new Error('Reached maximum number of retries after connection was closed')); const err = new Error('Reached maximum number of retries after connection was closed');
err.code = errors.ECONNECTION;
connection.queueEntry.callback(err);
} }
catch (E) { catch (E) {
this.logger.error({ this.logger.error({
@@ -377,6 +396,7 @@ class SMTPPool extends node_events_1.EventEmitter {
} }
/** @internal */ /** @internal */
_requeueEntryOnConnectionClose(connection) { _requeueEntryOnConnectionClose(connection) {
const delay = Math.min(REQUEUE_BASE_DELAY * 2 ** connection.queueEntry.requeueAttempts, REQUEUE_MAX_DELAY);
connection.queueEntry.requeueAttempts += 1; connection.queueEntry.requeueAttempts += 1;
this.logger.debug({ this.logger.debug({
tnx: 'pool', tnx: 'pool',
@@ -386,6 +406,7 @@ class SMTPPool extends node_events_1.EventEmitter {
}, 'Re-queued message <%s> for #%s. Attempt: #%s', connection.queueEntry.messageId, connection.id, connection.queueEntry.requeueAttempts); }, 'Re-queued message <%s> for #%s. Attempt: #%s', connection.queueEntry.messageId, connection.id, connection.queueEntry.requeueAttempts);
this._queue.unshift(connection.queueEntry); this._queue.unshift(connection.queueEntry);
connection.queueEntry = false; connection.queueEntry = false;
return delay;
} }
/** /**
* Continue to process message if the pool hasn't closed * Continue to process message if the pool hasn't closed
@@ -506,10 +527,15 @@ class SMTPPool extends node_events_1.EventEmitter {
connection.quit(); connection.quit();
return done(null, true); return done(null, true);
}; };
connection.connect(() => { connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
if (err) {
returned = true;
connection.close();
return done(err);
}
if (auth && (connection.allowsAuth || options.forceAuth)) { if (auth && (connection.allowsAuth || options.forceAuth)) {
connection.login(auth, err => { connection.login(auth, err => {
if (returned) { if (returned) {
+26 -28
View File
@@ -66,7 +66,7 @@ class PoolResource extends node_events_1.EventEmitter {
method: 'XOAUTH2' method: 'XOAUTH2'
}; };
oauth2.on('token', (token) => this.pool.mailer.emit('token', token)); oauth2.on('token', (token) => this.pool.mailer.emit('token', token));
oauth2.on('error', err => this.emit('error', err)); oauth2.on('error', err => this._fail(err));
break; break;
} }
default: default:
@@ -89,6 +89,19 @@ class PoolResource extends node_events_1.EventEmitter {
this._connected = false; this._connected = false;
this.messages = 0; this.messages = 0;
this.available = true; this.available = true;
this._failed = false;
}
/**
* Emits 'error' for the first failure only. A dead resource can report the same failure more
* than once (the connection error, then the send callback), the pool handles it once
* @internal
*/
_fail(err) {
if (this._failed) {
return;
}
this._failed = true;
this.emit('error', err);
} }
/** /**
* Initiates a connection to the SMTP server * Initiates a connection to the SMTP server
@@ -101,7 +114,7 @@ class PoolResource extends node_events_1.EventEmitter {
// nothing was connected, so no 'close' event is coming that would free the // nothing was connected, so no 'close' event is coming that would free the
// slot this resource holds in the pool, report the failure the way a failed // slot this resource holds in the pool, report the failure the way a failed
// login does // login does
this.emit('error', err); this._fail(err);
return callback(err); return callback(err);
} }
let returned = false; let returned = false;
@@ -118,8 +131,8 @@ class PoolResource extends node_events_1.EventEmitter {
options = Object.assign((0, index_js_2.assign)(false, options), socketOptions); options = Object.assign((0, index_js_2.assign)(false, options), socketOptions);
} }
this.connection = new index_js_1.default(options); this.connection = new index_js_1.default(options);
this.connection.once('error', err => { this.connection.on('error', err => {
this.emit('error', err); this._fail(err);
if (returned) { if (returned) {
return; return;
} }
@@ -128,31 +141,16 @@ class PoolResource extends node_events_1.EventEmitter {
}); });
this.connection.once('end', () => { this.connection.once('end', () => {
this.close(); this.close();
if (returned) {
return;
}
returned = true; returned = true;
const timer = setTimeout(() => { });
this.connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
// still have not returned, this means we have an unexpected connection close if (err) {
const err = new Error('Unexpected socket close'); // a close before the greeting, the 'end' that follows closes this resource
if (this.connection && this.connection.upgrading) { // and the pool requeues or fails the entry, bounded by maxRequeues
// starttls connection errors returned = true;
err.code = errors.ETLS;
}
callback(err);
}, 1000);
try {
timer.unref();
}
catch (_E) {
// Ignore. Happens on envs with non-node timer implementation
}
});
this.connection.connect(() => {
if (returned) {
return; return;
} }
if (this.auth && (this.connection.allowsAuth || options.forceAuth)) { if (this.auth && (this.connection.allowsAuth || options.forceAuth)) {
@@ -163,7 +161,7 @@ class PoolResource extends node_events_1.EventEmitter {
returned = true; returned = true;
if (err) { if (err) {
this.connection.close(); this.connection.close();
this.emit('error', err); this._fail(err);
return callback(err); return callback(err);
} }
this._connected = true; this._connected = true;
@@ -215,7 +213,7 @@ class PoolResource extends node_events_1.EventEmitter {
this.messages++; this.messages++;
if (err) { if (err) {
this.connection.close(); this.connection.close();
this.emit('error', err); this._fail(err);
return callback(err); return callback(err);
} }
info.envelope = { info.envelope = {
@@ -228,7 +226,7 @@ class PoolResource extends node_events_1.EventEmitter {
const err = new Error('Resource exhausted'); const err = new Error('Resource exhausted');
err.code = errors.EMAXLIMIT; err.code = errors.EMAXLIMIT;
this.connection.close(); this.connection.close();
this.emit('error', err); this._fail(err);
} }
else { else {
this.pool._checkRateLimit(() => { this.pool._checkRateLimit(() => {
+17 -27
View File
@@ -177,31 +177,6 @@ class SMTPTransport extends node_events_1.EventEmitter {
connection.close(); connection.close();
return callback(err); return callback(err);
}); });
connection.once('end', () => {
if (returned) {
return;
}
const timer = setTimeout(() => {
if (returned) {
return;
}
returned = true;
cleanupPerCallAuth();
// still have not returned, this means we have an unexpected connection close
const err = new Error('Unexpected socket close');
if (connection && connection.upgrading) {
// starttls connection errors
err.code = errors.ETLS;
}
callback(err);
}, 1000);
try {
timer.unref();
}
catch (_E) {
// Ignore. Happens on envs with non-node timer implementation
}
});
const sendMessage = () => { const sendMessage = () => {
const envelope = mail.message.getEnvelope(); const envelope = mail.message.getEnvelope();
const messageId = mail.message.messageId(); const messageId = mail.message.messageId();
@@ -221,6 +196,10 @@ class SMTPTransport extends node_events_1.EventEmitter {
messageId messageId
}, 'Sending message %s to <%s>', messageId, recipients.join(', ')); }, 'Sending message %s to <%s>', messageId, recipients.join(', '));
connection.send(envelope, mail.message.createReadStream(), (err, info) => { connection.send(envelope, mail.message.createReadStream(), (err, info) => {
if (returned) {
// the connection error handler has already reported this send
return;
}
returned = true; returned = true;
cleanupPerCallAuth(); cleanupPerCallAuth();
connection.close(); connection.close();
@@ -247,10 +226,16 @@ class SMTPTransport extends node_events_1.EventEmitter {
} }
}); });
}; };
connection.connect(() => { connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
if (err) {
// the server closed the connection before the greeting
returned = true;
connection.close();
return callback(err);
}
perCallAuth = this.getAuth(mail.data.auth); perCallAuth = this.getAuth(mail.data.auth);
if (perCallAuth && (connection.allowsAuth || options.forceAuth)) { if (perCallAuth && (connection.allowsAuth || options.forceAuth)) {
connection.login(perCallAuth, err => { connection.login(perCallAuth, err => {
@@ -332,10 +317,15 @@ class SMTPTransport extends node_events_1.EventEmitter {
connection.quit(); connection.quit();
return done(null, true); return done(null, true);
}; };
connection.connect(() => { connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
if (err) {
returned = true;
connection.close();
return done(err);
}
perCallAuth = this.getAuth({}); perCallAuth = this.getAuth({});
if (perCallAuth && (connection.allowsAuth || options.forceAuth)) { if (perCallAuth && (connection.allowsAuth || options.forceAuth)) {
connection.login(perCallAuth, err => { connection.login(perCallAuth, err => {
+2
View File
@@ -61,6 +61,8 @@ export interface XOAuth2Options {
component?: string | undefined; component?: string | undefined;
/** Extra headers for the token request */ /** Extra headers for the token request */
customHeaders?: OutgoingHttpHeaders | undefined; customHeaders?: OutgoingHttpHeaders | undefined;
/** Timeout for the token request in milliseconds, defaults to 60000, 0 disables it */
requestTimeout?: number | undefined;
/** Extra form fields for the token request */ /** Extra form fields for the token request */
customParams?: { customParams?: {
[key: string]: any; [key: string]: any;
+11 -3
View File
@@ -74,12 +74,14 @@ class XOAuth2 extends node_stream_1.Stream {
constructor(options, logger) { constructor(options, logger) {
super(); super();
this.options = options || {}; this.options = options || {};
this.configError = false;
if (options && options.serviceClient) { if (options && options.serviceClient) {
if (!options.privateKey || !options.user) { if (!options.privateKey || !options.user) {
// reported through getToken() rather than an 'error' event: the transports forward
// that event up to the Mail object, where nothing may be listening
const err = new Error('Options "privateKey" and "user" are required for service account!'); const err = new Error('Options "privateKey" and "user" are required for service account!');
err.code = errors.EOAUTH2; err.code = errors.EOAUTH2;
setImmediate(() => this.emit('error', err)); this.configError = err;
return;
} }
const serviceRequestTimeout = Math.min(Math.max(Number(this.options.serviceRequestTimeout) || 0, 0), 3600); const serviceRequestTimeout = Math.min(Math.max(Number(this.options.serviceRequestTimeout) || 0, 0), 3600);
this.options.serviceRequestTimeout = serviceRequestTimeout || 5 * 60; this.options.serviceRequestTimeout = serviceRequestTimeout || 5 * 60;
@@ -113,6 +115,9 @@ class XOAuth2 extends node_stream_1.Stream {
getToken(renew, callback) { getToken(renew, callback) {
// the error paths hand over the error alone // the error paths hand over the error alone
const done = callback; const done = callback;
if (this.configError) {
return done(this.configError);
}
if (!renew && this.accessToken && (!this.expires || this.expires > Date.now())) { if (!renew && this.accessToken && (!this.expires || this.expires > Date.now())) {
this.logger.debug({ this.logger.debug({
tnx: 'OAUTH2', tnx: 'OAUTH2',
@@ -349,7 +354,10 @@ class XOAuth2 extends node_stream_1.Stream {
method: 'post', method: 'post',
headers: params.customHeaders, headers: params.customHeaders,
body: payload, body: payload,
allowErrorResponse: true allowErrorResponse: true,
// unset falls back to the fetch default, a stalled token endpoint would otherwise keep
// `renewing` set and queue every later request
timeout: params.requestTimeout
}; };
// OAuth2 token endpoints are credential-bearing. src/fetch already // OAuth2 token endpoints are credential-bearing. src/fetch already
// validates certs by default; pin rejectUnauthorized:true here so the // validates certs by default; pin rejectUnauthorized:true here so the
+97
View File
@@ -240,6 +240,81 @@ function _recoverAddrSpec(data) {
.filter(part => part) .filter(part => part)
.join(' '); .join(' ');
} }
/**
* Takes the RFC 5322 comments out of the contents of an angle-addr.
*
* The tokenizer tracks a single open operator, so once '<' is open a '(' is plain text and
* a comment inside the brackets reached the address verbatim: 'Name <user@example.com(x)evil.com>'
* put 'user@example.com(x)evil.com' into the envelope, a value that is no mailbox and that a
* receiver stripping the comment reads as 'user@example.comevil.com' (GHSA-g73g-hqqh-jr95).
* A comment is folding whitespace, so it is read as a space that splits two atoms and as
* nothing next to an '@', the same rule the token walk applies outside the brackets. Quoted
* strings and domain literals are copied through untouched, a '(' in there is not a comment.
*
* @param address Contents of the angle brackets
* @return The address with the comments removed, and the text of the comments
*/
function _stripAddressComments(address) {
const comments = [];
let result = '';
let comment = '';
let depth = 0;
let closer = '';
// carried along rather than read back off the growing result, which would flatten it on
// every comment (see lastChars in _handleAddress)
let lastChar = '';
for (let i = 0, len = address.length; i < len; i++) {
const chr = address.charAt(i);
if (depth) {
if (chr === '\\' && i < len - 1) {
comment += address.charAt(++i);
}
else if (chr === '(') {
depth++;
comment += chr;
}
else if (chr === ')' && !--depth) {
comments.push(comment.trim());
comment = '';
if (lastChar !== '@' && address.charAt(i + 1) !== '@') {
result += ' ';
lastChar = ' ';
}
}
else {
comment += chr;
}
continue;
}
if (closer) {
if (chr === '\\' && closer === '"' && i < len - 1) {
result += chr + address.charAt(++i);
lastChar = address.charAt(i);
continue;
}
if (chr === closer) {
closer = '';
}
}
else if (chr === '"') {
closer = '"';
}
else if (chr === '[') {
closer = ']';
}
else if (chr === '(') {
depth = 1;
continue;
}
result += chr;
lastChar = chr;
}
if (depth) {
// an unterminated comment runs to the end of the address
comments.push(comment.trim());
}
return { address: result.trim(), comments: comments.filter(text => text) };
}
/** /**
* Converts tokens for a single address into an address object * Converts tokens for a single address into an address object
* *
@@ -367,6 +442,24 @@ function _handleAddress(tokens, depth) {
}); });
} }
else { else {
// Comments come out of the angle-addr before anything asks whether one was found, so
// that brackets holding nothing but a comment read the same as empty ones
const addressComments = [];
const addressParts = [];
for (const part of data.address) {
if (part.indexOf('(') < 0) {
addressParts.push(part);
continue;
}
const stripped = _stripAddressComments(part);
for (const comment of stripped.comments) {
addressComments.push(comment);
}
if (stripped.address) {
addressParts.push(stripped.address);
}
}
data.address = addressParts;
// If no address was found, try to detect one from regular text // If no address was found, try to detect one from regular text
if (!data.address.length && data.text.length) { if (!data.address.length && data.text.length) {
for (let i = data.text.length - 1; i >= 0; i--) { for (let i = data.text.length - 1; i >= 0; i--) {
@@ -433,6 +526,10 @@ function _handleAddress(tokens, depth) {
data.text = ''; data.text = '';
} }
_recoverAddrSpec(data); _recoverAddrSpec(data);
// a comment names the mailbox only when nothing else does, as one outside the brackets
if (!data.text && addressComments.length) {
data.text = addressComments.join(' ');
}
const address = { const address = {
address: data.address || data.text || '', address: data.address || data.text || '',
name: data.text || data.address || '' name: data.text || data.address || ''
+18 -1
View File
@@ -85,13 +85,22 @@ class DKIMSigner {
return; return;
} }
const key = this.keys[keyPos++]; const key = this.keys[keyPos++];
const dkimField = sign(this.headers, this.hashAlgo, this.bodyHash, { let dkimField;
try {
dkimField = sign(this.headers, this.hashAlgo, this.bodyHash, {
domainName: key.domainName, domainName: key.domainName,
keySelector: key.keySelector, keySelector: key.keySelector,
privateKey: key.privateKey, privateKey: key.privateKey,
headerFieldNames: this.options.headerFieldNames, headerFieldNames: this.options.headerFieldNames,
skipFields: this.options.skipFields skipFields: this.options.skipFields
}); });
}
catch (err) {
this.hasErrored = true;
this.cleanup();
this.output.emit('error', err);
return;
}
if (dkimField) { if (dkimField) {
this.output.write(Buffer.from(dkimField + '\r\n')); this.output.write(Buffer.from(dkimField + '\r\n'));
} }
@@ -191,7 +200,15 @@ class DKIM {
} }
const signer = new DKIMSigner(options, this.keys, inputStream, output); const signer = new DKIMSigner(options, this.keys, inputStream, output);
setImmediate(() => { setImmediate(() => {
try {
signer.signStream(); signer.signStream();
}
catch (_E) {
// the body hash is created here, an unknown hashAlgo throws inside this timer
// where nothing else could catch it
output.emit('error', sign.unsupportedHashAlgoError(signer.hashAlgo));
return;
}
if (writeValue) { if (writeValue) {
setImmediate(() => { setImmediate(() => {
inputStream.end(writeValue); inputStream.end(writeValue);
+2
View File
@@ -1,4 +1,5 @@
import crypto from 'node:crypto'; import crypto from 'node:crypto';
import type { NodemailerError } from '../errors.js';
import type { MessageParserHeaderLine } from './message-parser.js'; import type { MessageParserHeaderLine } from './message-parser.js';
/** /**
* Private key accepted by crypto.Sign#sign: a PEM string, a Buffer, a KeyObject or an * Private key accepted by crypto.Sign#sign: a PEM string, a Buffer, a KeyObject or an
@@ -45,5 +46,6 @@ export interface DKIMRelaxedHeaders {
declare function sign(headers: MessageParserHeaderLine[], hashAlgo: string, bodyHash: string, options?: DKIMSignOptions): string | false; declare function sign(headers: MessageParserHeaderLine[], hashAlgo: string, bodyHash: string, options?: DKIMSignOptions): string | false;
declare namespace sign { declare namespace sign {
var relaxedHeaders: (headers: MessageParserHeaderLine[], fieldNames?: string, skipFields?: string) => DKIMRelaxedHeaders; var relaxedHeaders: (headers: MessageParserHeaderLine[], fieldNames?: string, skipFields?: string) => DKIMRelaxedHeaders;
var unsupportedHashAlgoError: (hashAlgo: string) => NodemailerError;
} }
export default sign; export default sign;
+17 -1
View File
@@ -1,6 +1,15 @@
import * as punycode from '../punycode/index.js'; import * as punycode from '../punycode/index.js';
import * as mimeFuncs from '../mime-funcs/index.js'; import * as mimeFuncs from '../mime-funcs/index.js';
import crypto from 'node:crypto'; import crypto from 'node:crypto';
import * as errors from '../errors.js';
/**
* Error for a hashAlgo value the crypto module does not know
*/
function unsupportedHashAlgoError(hashAlgo) {
const err = new Error('Unsupported DKIM hash algorithm "' + hashAlgo + '"');
err.code = errors.ECONFIG;
return err;
}
/** /**
* Returns DKIM signature header line * Returns DKIM signature header line
* *
@@ -25,7 +34,13 @@ function sign(headers, hashAlgo, bodyHash, options) {
const canonicalizedHeaderData = relaxedHeaders(headers, fieldNames, options.skipFields); const canonicalizedHeaderData = relaxedHeaders(headers, fieldNames, options.skipFields);
const dkimHeader = generateDKIMHeader(options.domainName, options.keySelector, canonicalizedHeaderData.fieldNames, hashAlgo, bodyHash); const dkimHeader = generateDKIMHeader(options.domainName, options.keySelector, canonicalizedHeaderData.fieldNames, hashAlgo, bodyHash);
canonicalizedHeaderData.headers += 'dkim-signature:' + relaxedHeaderLine(dkimHeader); canonicalizedHeaderData.headers += 'dkim-signature:' + relaxedHeaderLine(dkimHeader);
const signer = crypto.createSign(('rsa-' + hashAlgo).toUpperCase()); let signer;
try {
signer = crypto.createSign(('rsa-' + hashAlgo).toUpperCase());
}
catch (_E) {
throw unsupportedHashAlgoError(hashAlgo);
}
// the header lines are 'binary' strings, so this reproduces the original header bytes // the header lines are 'binary' strings, so this reproduces the original header bytes
signer.update(canonicalizedHeaderData.headers, 'latin1'); signer.update(canonicalizedHeaderData.headers, 'latin1');
let signature; let signature;
@@ -38,6 +53,7 @@ function sign(headers, hashAlgo, bodyHash, options) {
return dkimHeader + signature.replace(/(^.{73}|.{75}(?!\r?\n|\r))/g, '$&\r\n ').trim(); return dkimHeader + signature.replace(/(^.{73}|.{75}(?!\r?\n|\r))/g, '$&\r\n ').trim();
} }
sign.relaxedHeaders = relaxedHeaders; sign.relaxedHeaders = relaxedHeaders;
sign.unsupportedHashAlgoError = unsupportedHashAlgoError;
export default sign; export default sign;
function generateDKIMHeader(domainName, keySelector, fieldNames, hashAlgo, bodyHash) { function generateDKIMHeader(domainName, keySelector, fieldNames, hashAlgo, bodyHash) {
// the caller supplied tag values are interpolated straight into the tag list, and none of // the caller supplied tag values are interpolated straight into the tag list, and none of
+4 -1
View File
@@ -25,8 +25,10 @@ export interface FetchOptions {
tls?: { tls?: {
[key: string]: any; [key: string]: any;
} | undefined; } | undefined;
/** Request timeout in milliseconds */ /** Socket inactivity timeout in milliseconds, defaults to 60000, 0 disables it */
timeout?: number | undefined; timeout?: number | undefined;
/** Maximum size of the (decoded) response body in bytes, defaults to 64 MB, Infinity disables the limit */
maxBytes?: number | undefined;
/** Maximum number of redirects to follow (default 5) */ /** Maximum number of redirects to follow (default 5) */
maxRedirects?: number | undefined; maxRedirects?: number | undefined;
/** Resolve responses with a status code of 300 or above instead of emitting an error */ /** Resolve responses with a status code of 300 or above instead of emitting an error */
@@ -47,6 +49,7 @@ export interface FetchResponse extends PassThrough {
declare function nmfetch(url: string, options?: FetchOptions): FetchResponse; declare function nmfetch(url: string, options?: FetchOptions): FetchResponse;
declare namespace nmfetch { declare namespace nmfetch {
var Cookies: typeof import("./cookies.js").default; var Cookies: typeof import("./cookies.js").default;
var DEFAULT_TIMEOUT: number;
} }
type CookiesJar = Cookies; type CookiesJar = Cookies;
/** /**
+35 -60
View File
@@ -9,6 +9,10 @@ import net from 'node:net';
import * as errors from '../errors.js'; import * as errors from '../errors.js';
import { isProtoKey } from '../shared/objects.js'; import { isProtoKey } from '../shared/objects.js';
const MAX_REDIRECTS = 5; const MAX_REDIRECTS = 5;
// a stalled server would otherwise hold the request, and whatever waits on it, open forever
const DEFAULT_TIMEOUT = 60 * 1000;
// the body is usually buffered in memory by the caller, so an unbounded download is refused
const DEFAULT_MAX_BYTES = 64 * 1024 * 1024;
// Only genuine TLS settings are taken from options.tls. That object reaches us straight // Only genuine TLS settings are taken from options.tls. That object reaches us straight
// from a user supplied attachment (content.tls), so keys like host, port, path, socketPath // from a user supplied attachment (content.tls), so keys like host, port, path, socketPath
// or lookup would otherwise repoint the request at a destination that never went through // or lookup would otherwise repoint the request at a destination that never went through
@@ -214,28 +218,22 @@ function nmfetch(url, options) {
}); });
return fetchRes; return fetchRes;
} }
if (options.timeout) { // reports the first failure of this request on fetchRes and releases the request
req.setTimeout(options.timeout, () => { const fail = (err, sourceUrl = url) => {
if (finished) { if (finished) {
return; return;
} }
finished = true; finished = true;
err.code = errors.EFETCH;
err.sourceUrl = sourceUrl;
fetchRes.emit('error', err);
req.abort(); req.abort();
const err = new Error('Request Timeout'); };
err.code = errors.EFETCH; const timeout = typeof options.timeout === 'number' && options.timeout >= 0 ? options.timeout : DEFAULT_TIMEOUT;
err.sourceUrl = url; if (timeout) {
fetchRes.emit('error', err); req.setTimeout(timeout, () => fail(new Error('Request Timeout')));
});
} }
req.on('error', (err) => { req.on('error', (err) => fail(err));
if (finished) {
return;
}
finished = true;
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
});
req.on('response', res => { req.on('response', res => {
let inflate; let inflate;
if (finished) { if (finished) {
@@ -256,13 +254,7 @@ function nmfetch(url, options) {
// redirect // redirect
options.redirects++; options.redirects++;
if (options.redirects > options.maxRedirects) { if (options.redirects > options.maxRedirects) {
finished = true; return fail(new Error('Maximum redirect count exceeded'));
const err = new Error('Maximum redirect count exceeded');
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
req.abort();
return;
} }
// redirect does not include POST body // redirect does not include POST body
options.method = 'GET'; options.method = 'GET';
@@ -282,13 +274,7 @@ function nmfetch(url, options) {
// call: that call gets its own `finished` flag and no handle on this // call: that call gets its own `finished` flag and no handle on this
// request, so this one would stay open and could emit a second error on // request, so this one would stay open and could emit a second error on
// the shared fetchRes once it times out. Callers listen with req.once(). // the shared fetchRes once it times out. Callers listen with req.once().
finished = true; return fail(new Error('Unsupported protocol for URL ' + redirectUrl), redirectUrl);
const err = new Error('Unsupported protocol for URL ' + redirectUrl);
err.code = errors.EFETCH;
err.sourceUrl = redirectUrl;
fetchRes.emit('error', err);
req.abort();
return;
} }
// Do not forward credentials when the redirect leaves the original // Do not forward credentials when the redirect leaves the original
// security context: a different host, or a downgrade from https to // security context: a different host, or a downgrade from https to
@@ -305,41 +291,33 @@ function nmfetch(url, options) {
} }
}); });
} }
// this request is done with, release it so its socket and timeout do not
// outlive it and report a late error on the shared fetchRes
finished = true;
res.resume();
req.abort();
return nmfetch(redirectUrl, options); return nmfetch(redirectUrl, options);
} }
fetchRes.statusCode = res.statusCode; fetchRes.statusCode = res.statusCode;
fetchRes.headers = res.headers; fetchRes.headers = res.headers;
if (res.statusCode >= 300 && !options.allowErrorResponse) { if (res.statusCode >= 300 && !options.allowErrorResponse) {
finished = true; return fail(new Error('Invalid status code ' + res.statusCode));
const err = new Error('Invalid status code ' + res.statusCode); }
err.code = errors.EFETCH; res.on('error', (err) => fail(err));
err.sourceUrl = url; const maxBytes = typeof options.maxBytes === 'number' && options.maxBytes > 0 ? options.maxBytes : DEFAULT_MAX_BYTES;
fetchRes.emit('error', err); const source = inflate || res;
req.abort(); let received = 0;
source.on('data', (chunk) => {
received += chunk.length;
if (received <= maxBytes || finished) {
return; return;
} }
res.on('error', (err) => { source.unpipe(fetchRes);
if (finished) { fail(new Error('Response size exceeds the allowed ' + maxBytes + ' bytes'));
return;
}
finished = true;
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
req.abort();
}); });
if (inflate) { if (inflate) {
res.pipe(inflate).pipe(fetchRes); res.pipe(inflate).pipe(fetchRes);
inflate.on('error', (err) => { inflate.on('error', (err) => fail(err));
if (finished) {
return;
}
finished = true;
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
req.abort();
});
} }
else { else {
res.pipe(fetchRes); res.pipe(fetchRes);
@@ -354,11 +332,7 @@ function nmfetch(url, options) {
req.write(body); req.write(body);
} }
catch (err) { catch (err) {
finished = true; return fail(err);
err.code = errors.EFETCH;
err.sourceUrl = url;
fetchRes.emit('error', err);
return;
} }
} }
req.end(); req.end();
@@ -366,4 +340,5 @@ function nmfetch(url, options) {
return fetchRes; return fetchRes;
} }
nmfetch.Cookies = Cookies; nmfetch.Cookies = Cookies;
nmfetch.DEFAULT_TIMEOUT = DEFAULT_TIMEOUT;
export default nmfetch; export default nmfetch;
+1 -1
View File
@@ -1,3 +1,3 @@
export declare const name = "nodemailer"; export declare const name = "nodemailer";
export declare const version = "10.0.11"; export declare const version = "10.0.13";
export declare const homepage = "https://nodemailer.com/"; export declare const homepage = "https://nodemailer.com/";
+1 -1
View File
@@ -1,4 +1,4 @@
// Generated by scripts/build.js from package.json. Do not edit by hand. // Generated by scripts/build.js from package.json. Do not edit by hand.
export const name = 'nodemailer'; export const name = 'nodemailer';
export const version = '10.0.11'; export const version = '10.0.13';
export const homepage = 'https://nodemailer.com/'; export const homepage = 'https://nodemailer.com/';
+2 -1
View File
@@ -1,7 +1,8 @@
import { Transform, type TransformOptions } from 'node:stream'; import { Transform, type TransformOptions } from 'node:stream';
/** /**
* Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF> * Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF>
* Also makes sure that only <CR><LF> sequences are used for linebreaks * Also makes sure that only <CR><LF> sequences are used for linebreaks, bare CR and bare LF
* are both turned into <CR><LF>
* *
* @param options Stream options * @param options Stream options
*/ */
+28 -20
View File
@@ -1,7 +1,13 @@
import { Transform } from 'node:stream'; import { Transform } from 'node:stream';
// bytes inserted into the output, shared as they are only ever copied by Buffer.concat
const INSERT_LF = Buffer.from('\n');
const INSERT_LF_DOT = Buffer.from('\n.');
const INSERT_CR = Buffer.from('\r');
const INSERT_DOT = Buffer.from('.');
/** /**
* Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF> * Escapes dots in the beginning of lines. Ends the stream with <CR><LF>.<CR><LF>
* Also makes sure that only <CR><LF> sequences are used for linebreaks * Also makes sure that only <CR><LF> sequences are used for linebreaks, bare CR and bare LF
* are both turned into <CR><LF>
* *
* @param options Stream options * @param options Stream options
*/ */
@@ -30,33 +36,35 @@ export default class DataStream extends Transform {
} }
this.inByteCount += chunk.length; this.inByteCount += chunk.length;
for (i = 0, len = chunk.length; i < len; i++) { for (i = 0, len = chunk.length; i < len; i++) {
if (chunk[i] === 0x2e) { const byte = chunk[i];
// . const prev = i ? chunk[i - 1] : this.lastByte;
if ((i && chunk[i - 1] === 0x0a) || (!i && (!this.lastByte || this.lastByte === 0x0a))) { let insert = false;
buf = chunk.slice(lastPos, i + 1); if (prev === 0x0d && byte !== 0x0a) {
chunks.push(buf); // a bare CR becomes CRLF. A receiver that treats a lone CR as a line end would
chunks.push(Buffer.from('.')); // otherwise see "\r.\r" as the end of the data (SMTP smuggling), so a dot
chunklen += buf.length + 1; // following it is stuffed like at the start of any other line
lastPos = i + 1; insert = byte === 0x2e ? INSERT_LF_DOT : INSERT_LF;
} }
else if (byte === 0x0a && prev !== 0x0d) {
// a bare LF becomes CRLF
insert = INSERT_CR;
} }
else if (chunk[i] === 0x0a) { else if (byte === 0x2e && (prev === 0x0a || prev === false)) {
// \n // a dot at the start of a line
if ((i && chunk[i - 1] !== 0x0d) || (!i && this.lastByte !== 0x0d)) { insert = INSERT_DOT;
}
if (insert) {
if (i > lastPos) { if (i > lastPos) {
buf = chunk.slice(lastPos, i); buf = chunk.slice(lastPos, i);
chunks.push(buf); chunks.push(buf);
chunklen += buf.length + 2; chunklen += buf.length;
} }
else { chunks.push(insert);
chunklen += 2; chunklen += insert.length;
} lastPos = i;
chunks.push(Buffer.from('\r\n'));
lastPos = i + 1;
} }
} }
} if (chunks.length) {
if (chunklen) {
// add last piece // add last piece
if (lastPos < chunk.length) { if (lastPos < chunk.length) {
buf = chunk.slice(lastPos); buf = chunk.slice(lastPos);
+5 -5
View File
@@ -26,11 +26,11 @@ export interface SMTPConnectionOptions {
secured?: boolean | undefined; secured?: boolean | undefined;
/** Server name for SNI, defaults to host when that is not an IP address */ /** Server name for SNI, defaults to host when that is not an IP address */
servername?: string | undefined; servername?: string | undefined;
/** Ignore STARTTLS even when the server advertises it */ /** Ignore STARTTLS even when the server advertises it, has no effect when requireTLS is set */
ignoreTLS?: boolean | undefined; ignoreTLS?: boolean | undefined;
/** Force STARTTLS, fail when the server does not support it */ /** Force STARTTLS, fail when the server does not support it. Takes precedence over ignoreTLS and opportunisticTLS */
requireTLS?: boolean | undefined; requireTLS?: boolean | undefined;
/** Continue unencrypted when the STARTTLS upgrade fails */ /** Continue unencrypted when the STARTTLS upgrade fails, has no effect when requireTLS is set */
opportunisticTLS?: boolean | undefined; opportunisticTLS?: boolean | undefined;
/** Name of the client server, sent with EHLO/HELO, CRLF is stripped */ /** Name of the client server, sent with EHLO/HELO, CRLF is stripped */
name?: string | undefined; name?: string | undefined;
@@ -288,8 +288,8 @@ export type SMTPConnectionResponseAction = (str: string) => void;
* * **port** - is the port to connect to (defaults to 587 or 465) * * **port** - is the port to connect to (defaults to 587 or 465)
* * **host** - is the hostname or IP address to connect to (defaults to 'localhost') * * **host** - is the hostname or IP address to connect to (defaults to 'localhost')
* * **secure** - use SSL * * **secure** - use SSL
* * **ignoreTLS** - ignore server support for STARTTLS * * **ignoreTLS** - ignore server support for STARTTLS (has no effect when requireTLS is set)
* * **requireTLS** - forces the client to use STARTTLS * * **requireTLS** - forces the client to use STARTTLS, takes precedence over ignoreTLS and opportunisticTLS
* * **name** - the name of the client server * * **name** - the name of the client server
* * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener) * * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener)
* * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds) * * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds)
+75 -21
View File
@@ -66,8 +66,8 @@ function isPartialLine(line) {
* * **port** - is the port to connect to (defaults to 587 or 465) * * **port** - is the port to connect to (defaults to 587 or 465)
* * **host** - is the hostname or IP address to connect to (defaults to 'localhost') * * **host** - is the hostname or IP address to connect to (defaults to 'localhost')
* * **secure** - use SSL * * **secure** - use SSL
* * **ignoreTLS** - ignore server support for STARTTLS * * **ignoreTLS** - ignore server support for STARTTLS (has no effect when requireTLS is set)
* * **requireTLS** - forces the client to use STARTTLS * * **requireTLS** - forces the client to use STARTTLS, takes precedence over ignoreTLS and opportunisticTLS
* * **name** - the name of the client server * * **name** - the name of the client server
* * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener) * * **localAddress** - outbound address to bind to (see: http://nodejs.org/api/net.html#net_net_connect_options_connectionlistener)
* * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds) * * **greetingTimeout** - Time to wait in ms until greeting message is received from the server (defaults to 30 seconds)
@@ -92,6 +92,11 @@ class SMTPConnection extends EventEmitter {
this.id = crypto.randomBytes(8).toString('base64').replace(/\W/g, ''); this.id = crypto.randomBytes(8).toString('base64').replace(/\W/g, '');
this.stage = 'init'; this.stage = 'init';
this.options = options || {}; this.options = options || {};
if (this.options.requireTLS && (this.options.ignoreTLS || this.options.opportunisticTLS)) {
// requireTLS wins, a contradictory configuration must not quietly fall back to plaintext.
// Copied so the caller's (possibly shared) options object is left as it was
this.options = Object.assign({}, this.options, { ignoreTLS: false, opportunisticTLS: false });
}
this.secureConnection = !!this.options.secure; this.secureConnection = !!this.options.secure;
this.alreadySecured = !!this.options.secured; this.alreadySecured = !!this.options.secured;
this.port = Number(this.options.port) || (this.secureConnection ? 465 : 587); this.port = Number(this.options.port) || (this.secureConnection ? 465 : 587);
@@ -135,6 +140,8 @@ class SMTPConnection extends EventEmitter {
this._destroyed = false; this._destroyed = false;
this._closing = false; this._closing = false;
this._currentDataStream = false; this._currentDataStream = false;
this._pendingSend = false;
this._connectCallback = false;
this._onSocketData = chunk => this._onData(chunk); this._onSocketData = chunk => this._onData(chunk);
this._onSocketError = error => this._onError(error, 'ESOCKET', false, 'CONN'); this._onSocketError = error => this._onError(error, 'ESOCKET', false, 'CONN');
this._onSocketClose = () => this._onClose(); this._onSocketClose = () => this._onClose();
@@ -149,7 +156,9 @@ class SMTPConnection extends EventEmitter {
*/ */
connect(connectCallback) { connect(connectCallback) {
if (typeof connectCallback === 'function') { if (typeof connectCallback === 'function') {
this._connectCallback = connectCallback;
this.once('connect', () => { this.once('connect', () => {
this._connectCallback = false;
this.logger.debug({ this.logger.debug({
tnx: 'smtp' tnx: 'smtp'
}, 'SMTP handshake finished'); }, 'SMTP handshake finished');
@@ -381,6 +390,16 @@ class SMTPConnection extends EventEmitter {
} }
this._currentDataStream = false; this._currentDataStream = false;
} }
// Detach from the message stream as well. The listener is swapped for a no-op rather than
// removed, a stream destroyed with an error later on would otherwise throw it as unhandled
if (this._pendingSend) {
const { stream, onStreamError } = this._pendingSend;
if (stream) {
stream.removeListener('error', onStreamError);
stream.on('error', TEARDOWN_NOOP);
}
this._pendingSend = false;
}
if (socket && !socket.destroyed) { if (socket && !socket.destroyed) {
try { try {
// Clear socket timeout to prevent timer leaks // Clear socket timeout to prevent timer leaks
@@ -568,6 +587,9 @@ class SMTPConnection extends EventEmitter {
return; return;
} }
returned = true; returned = true;
if (this._pendingSend && this._pendingSend.callback === callback) {
this._pendingSend = false;
}
done(err, info); done(err, info);
}; };
if (!message) { if (!message) {
@@ -584,9 +606,16 @@ class SMTPConnection extends EventEmitter {
}); });
return; return;
} }
const pendingSend = {
callback,
stream: false,
onStreamError: err => callback(this._formatError(err, 'ESTREAM', false, 'API'))
};
if (typeof message.on === 'function') { if (typeof message.on === 'function') {
message.on('error', err => callback(this._formatError(err, 'ESTREAM', false, 'API'))); pendingSend.stream = message;
pendingSend.stream.on('error', pendingSend.onStreamError);
} }
this._pendingSend = pendingSend;
const startTime = Date.now(); const startTime = Date.now();
this._setEnvelope(envelope, (err, info) => { this._setEnvelope(envelope, (err, info) => {
if (err) { if (err) {
@@ -779,8 +808,14 @@ class SMTPConnection extends EventEmitter {
else { else {
this.logger.error(data, err.message); this.logger.error(data, err.message);
} }
// close() forgets the send in flight, it is completed with this same error afterwards so
// a late message stream error has nothing left to report
const pendingSend = this._pendingSend;
this.emit('error', err); this.emit('error', err);
this.close(); this.close();
if (pendingSend) {
pendingSend.callback(err);
}
} }
/** @internal */ /** @internal */
_formatError(message, type, response, command) { _formatError(message, type, response, command) {
@@ -826,14 +861,25 @@ class SMTPConnection extends EventEmitter {
this.logger.info({ this.logger.info({
tnx: 'network' tnx: 'network'
}, 'Connection closed'); }, 'Connection closed');
// the unterminated remainder is only reported as a reply (and so gives the error a responseCode)
// when it starts like a complete failure reply, not for a fragment such as "55" or a 250
const failureResponse = typeof serverResponse === 'string' && /^[45]\d{2}[ -]/.test(serverResponse) ? serverResponse : false;
if (this.upgrading && !this._destroyed) { if (this.upgrading && !this._destroyed) {
return this._onError(new Error('Connection closed unexpectedly'), 'ETLS', serverResponse, 'CONN'); return this._onError(new Error('Connection closed unexpectedly'), 'ETLS', failureResponse, 'CONN');
} }
else if (![this._actionGreeting, this.close].includes(this._responseActions[0]) && !this._destroyed) { if (!failureResponse && this._responseActions[0] === this._actionGreeting && this._connectCallback && !this._destroyed) {
return this._onError(new Error('Connection closed unexpectedly'), 'ECONNECTION', serverResponse, 'CONN'); // A silent close before the greeting is handed to the connect() callback rather than
// emitted as 'error', callers that never saw an error for it must not start throwing one
const connectCallback = this._connectCallback;
this._connectCallback = false;
const err = this._formatError(new Error('Connection closed unexpectedly'), 'ECONNECTION', false, 'CONN');
this.logger.warn({ tnx: 'network' }, err.message);
connectCallback(err);
this.close();
return;
} }
else if (/^[45]\d{2}\b/.test(serverResponse)) { if (failureResponse || (this._responseActions[0] !== this.close && !this._destroyed)) {
return this._onError(new Error('Connection closed unexpectedly'), 'ECONNECTION', serverResponse, 'CONN'); return this._onError(new Error('Connection closed unexpectedly'), 'ECONNECTION', failureResponse, 'CONN');
} }
this._destroy(); this._destroy();
} }
@@ -1310,21 +1356,25 @@ class SMTPConnection extends EventEmitter {
if (/[ -]AUTH\b/i.test(str)) { if (/[ -]AUTH\b/i.test(str)) {
this.allowsAuth = true; this.allowsAuth = true;
} }
// Detect if the server supports PLAIN auth // Detect the advertised SASL mechanisms. The list is split into whole tokens rather
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)PLAIN/i.test(str)) { // than searched for each name with a pattern: the patterns this replaced let two
this._supportedAuth.push('PLAIN'); // whitespace runs overlap and backtracked quadratically over an AUTH line padded with
// spaces, so a hostile server could stall the event loop from its EHLO reply
// (GHSA-4ffr-jq9g-5ffx).
const authMechanisms = new Set();
for (const line of this._ehloLines) {
const authMatch = /^AUTH[\s=](.*)/i.exec(line);
if (authMatch) {
for (const mechanism of authMatch[1].split(/[\s=]+/)) {
authMechanisms.add(mechanism.toUpperCase());
} }
// Detect if the server supports LOGIN auth
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)LOGIN/i.test(str)) {
this._supportedAuth.push('LOGIN');
} }
// Detect if the server supports CRAM-MD5 auth
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)CRAM-MD5/i.test(str)) {
this._supportedAuth.push('CRAM-MD5');
} }
// Detect if the server supports XOAUTH2 auth // listed in order of preference, the first one the credentials allow is used
if (/[ -]AUTH(?:(\s+|=)[^\n]*\s+|\s+|=)XOAUTH2/i.test(str)) { for (const mechanism of ['PLAIN', 'LOGIN', 'CRAM-MD5', 'XOAUTH2']) {
this._supportedAuth.push('XOAUTH2'); if (authMechanisms.has(mechanism)) {
this._supportedAuth.push(mechanism);
}
} }
// Detect if the server supports SIZE extensions (and the max allowed size) // Detect if the server supports SIZE extensions (and the max allowed size)
if ((match = str.match(/[ -]SIZE(?:[ \t]+(\d+))?/im))) { if ((match = str.match(/[ -]SIZE(?:[ \t]+(\d+))?/im))) {
@@ -1584,7 +1634,11 @@ class SMTPConnection extends EventEmitter {
this._sendCommand('DATA'); this._sendCommand('DATA');
} }
else { else {
err = this._formatError("Can't send mail - all recipients were rejected", 'EENVELOPE', str, 'RCPT TO'); // report a temporary rejection when there is one, taking the last reply would mark the
// whole message as permanently failed although some recipients were only deferred
const deferred = envelope.rejectedErrors.find(rejectedErr => rejectedErr.responseCode && rejectedErr.responseCode < 500);
const reply = deferred?.response ?? str;
err = this._formatError("Can't send mail - all recipients were rejected", 'EENVELOPE', reply, 'RCPT TO');
err.rejected = envelope.rejected; err.rejected = envelope.rejected;
err.rejectedErrors = envelope.rejectedErrors; err.rejectedErrors = envelope.rejectedErrors;
return callback(err); return callback(err);
+2 -1
View File
@@ -17,7 +17,7 @@ export interface SMTPPoolOptions extends SMTPTransportOptions {
rateLimit?: number | undefined; rateLimit?: number | undefined;
/** Time window for rateLimit in milliseconds, defaults to 1000 */ /** Time window for rateLimit in milliseconds, defaults to 1000 */
rateDelta?: number | undefined; rateDelta?: number | undefined;
/** How many times a message is requeued when its connection closes while sending, unlimited when not set or negative */ /** How many times a message is requeued when its connection closes while sending, defaults to 5, a negative value means unlimited */
maxRequeues?: number | undefined; maxRequeues?: number | undefined;
} }
/** /**
@@ -26,6 +26,7 @@ export interface SMTPPoolOptions extends SMTPTransportOptions {
export type SMTPPoolResolvedOptions = SMTPPoolOptions & { export type SMTPPoolResolvedOptions = SMTPPoolOptions & {
maxConnections: number; maxConnections: number;
maxMessages: number; maxMessages: number;
maxRequeues: number;
}; };
/** /**
* Result of a message sent through the pool, same as for the SMTP transport * Result of a message sent through the pool, same as for the SMTP transport
+30 -4
View File
@@ -5,6 +5,10 @@ import wellKnown from '../well-known/index.js';
import * as shared from '../shared/index.js'; import * as shared from '../shared/index.js';
import * as errors from '../errors.js'; import * as errors from '../errors.js';
import * as packageData from '../package-info.js'; import * as packageData from '../package-info.js';
/** First delay before a requeued message is retried, doubled on every further requeue */
const REQUEUE_BASE_DELAY = 50;
/** Upper bound for the requeue delay */
const REQUEUE_MAX_DELAY = 2000;
/** /**
* Creates a SMTP pool transport object for Nodemailer * Creates a SMTP pool transport object for Nodemailer
* *
@@ -36,6 +40,8 @@ class SMTPPool extends EventEmitter {
); );
this.options.maxConnections = this.options.maxConnections || 5; this.options.maxConnections = this.options.maxConnections || 5;
this.options.maxMessages = this.options.maxMessages || 100; this.options.maxMessages = this.options.maxMessages || 100;
// a default bound, a server that closes every connection before the greeting would otherwise be retried forever
this.options.maxRequeues = typeof this.options.maxRequeues === 'number' ? this.options.maxRequeues : 5;
this.logger = shared.getLogger(this.options, { this.logger = shared.getLogger(this.options, {
component: this.options.component || 'smtp-pool' component: this.options.component || 'smtp-pool'
}); });
@@ -79,6 +85,10 @@ class SMTPPool extends EventEmitter {
*/ */
send(mail, callback) { send(mail, callback) {
if (this._closed) { if (this._closed) {
// Mail.sendMail ignores the return value, so without a callback its promise would never settle
const err = new Error('Connection pool was closed');
err.code = errors.ECONNECTION;
setImmediate(() => callback(err));
return false; return false;
} }
this._queue.push({ this._queue.push({
@@ -292,15 +302,22 @@ class SMTPPool extends EventEmitter {
// Note that we must wait a bit.. because the callback of the 'error' handler might be called // Note that we must wait a bit.. because the callback of the 'error' handler might be called
// in the next event loop // in the next event loop
setTimeout(() => { setTimeout(() => {
let delay = 0;
if (connection.queueEntry) { if (connection.queueEntry) {
if (this._shouldRequeuOnConnectionClose(connection.queueEntry)) { if (this._shouldRequeuOnConnectionClose(connection.queueEntry)) {
this._requeueEntryOnConnectionClose(connection); delay = this._requeueEntryOnConnectionClose(connection);
} }
else { else {
this._failDeliveryOnConnectionClose(connection); this._failDeliveryOnConnectionClose(connection);
} }
} }
if (delay) {
// back off, a server that keeps dropping connections is not hammered
setTimeout(() => this._continueProcessing(), delay);
}
else {
this._continueProcessing(); this._continueProcessing();
}
}, 50); }, 50);
} }
else { else {
@@ -315,7 +332,7 @@ class SMTPPool extends EventEmitter {
} }
/** @internal */ /** @internal */
_shouldRequeuOnConnectionClose(queueEntry) { _shouldRequeuOnConnectionClose(queueEntry) {
if (this.options.maxRequeues === undefined || this.options.maxRequeues < 0) { if (this.options.maxRequeues < 0) {
return true; return true;
} }
return queueEntry.requeueAttempts < this.options.maxRequeues; return queueEntry.requeueAttempts < this.options.maxRequeues;
@@ -324,7 +341,9 @@ class SMTPPool extends EventEmitter {
_failDeliveryOnConnectionClose(connection) { _failDeliveryOnConnectionClose(connection) {
if (connection.queueEntry && connection.queueEntry.callback) { if (connection.queueEntry && connection.queueEntry.callback) {
try { try {
connection.queueEntry.callback(new Error('Reached maximum number of retries after connection was closed')); const err = new Error('Reached maximum number of retries after connection was closed');
err.code = errors.ECONNECTION;
connection.queueEntry.callback(err);
} }
catch (E) { catch (E) {
this.logger.error({ this.logger.error({
@@ -339,6 +358,7 @@ class SMTPPool extends EventEmitter {
} }
/** @internal */ /** @internal */
_requeueEntryOnConnectionClose(connection) { _requeueEntryOnConnectionClose(connection) {
const delay = Math.min(REQUEUE_BASE_DELAY * 2 ** connection.queueEntry.requeueAttempts, REQUEUE_MAX_DELAY);
connection.queueEntry.requeueAttempts += 1; connection.queueEntry.requeueAttempts += 1;
this.logger.debug({ this.logger.debug({
tnx: 'pool', tnx: 'pool',
@@ -348,6 +368,7 @@ class SMTPPool extends EventEmitter {
}, 'Re-queued message <%s> for #%s. Attempt: #%s', connection.queueEntry.messageId, connection.id, connection.queueEntry.requeueAttempts); }, 'Re-queued message <%s> for #%s. Attempt: #%s', connection.queueEntry.messageId, connection.id, connection.queueEntry.requeueAttempts);
this._queue.unshift(connection.queueEntry); this._queue.unshift(connection.queueEntry);
connection.queueEntry = false; connection.queueEntry = false;
return delay;
} }
/** /**
* Continue to process message if the pool hasn't closed * Continue to process message if the pool hasn't closed
@@ -468,10 +489,15 @@ class SMTPPool extends EventEmitter {
connection.quit(); connection.quit();
return done(null, true); return done(null, true);
}; };
connection.connect(() => { connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
if (err) {
returned = true;
connection.close();
return done(err);
}
if (auth && (connection.allowsAuth || options.forceAuth)) { if (auth && (connection.allowsAuth || options.forceAuth)) {
connection.login(auth, err => { connection.login(auth, err => {
if (returned) { if (returned) {
+26 -28
View File
@@ -28,7 +28,7 @@ export default class PoolResource extends EventEmitter {
method: 'XOAUTH2' method: 'XOAUTH2'
}; };
oauth2.on('token', (token) => this.pool.mailer.emit('token', token)); oauth2.on('token', (token) => this.pool.mailer.emit('token', token));
oauth2.on('error', err => this.emit('error', err)); oauth2.on('error', err => this._fail(err));
break; break;
} }
default: default:
@@ -51,6 +51,19 @@ export default class PoolResource extends EventEmitter {
this._connected = false; this._connected = false;
this.messages = 0; this.messages = 0;
this.available = true; this.available = true;
this._failed = false;
}
/**
* Emits 'error' for the first failure only. A dead resource can report the same failure more
* than once (the connection error, then the send callback), the pool handles it once
* @internal
*/
_fail(err) {
if (this._failed) {
return;
}
this._failed = true;
this.emit('error', err);
} }
/** /**
* Initiates a connection to the SMTP server * Initiates a connection to the SMTP server
@@ -63,7 +76,7 @@ export default class PoolResource extends EventEmitter {
// nothing was connected, so no 'close' event is coming that would free the // nothing was connected, so no 'close' event is coming that would free the
// slot this resource holds in the pool, report the failure the way a failed // slot this resource holds in the pool, report the failure the way a failed
// login does // login does
this.emit('error', err); this._fail(err);
return callback(err); return callback(err);
} }
let returned = false; let returned = false;
@@ -80,8 +93,8 @@ export default class PoolResource extends EventEmitter {
options = Object.assign(assign(false, options), socketOptions); options = Object.assign(assign(false, options), socketOptions);
} }
this.connection = new SMTPConnection(options); this.connection = new SMTPConnection(options);
this.connection.once('error', err => { this.connection.on('error', err => {
this.emit('error', err); this._fail(err);
if (returned) { if (returned) {
return; return;
} }
@@ -90,31 +103,16 @@ export default class PoolResource extends EventEmitter {
}); });
this.connection.once('end', () => { this.connection.once('end', () => {
this.close(); this.close();
if (returned) {
return;
}
returned = true; returned = true;
const timer = setTimeout(() => { });
this.connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
// still have not returned, this means we have an unexpected connection close if (err) {
const err = new Error('Unexpected socket close'); // a close before the greeting, the 'end' that follows closes this resource
if (this.connection && this.connection.upgrading) { // and the pool requeues or fails the entry, bounded by maxRequeues
// starttls connection errors returned = true;
err.code = errors.ETLS;
}
callback(err);
}, 1000);
try {
timer.unref();
}
catch (_E) {
// Ignore. Happens on envs with non-node timer implementation
}
});
this.connection.connect(() => {
if (returned) {
return; return;
} }
if (this.auth && (this.connection.allowsAuth || options.forceAuth)) { if (this.auth && (this.connection.allowsAuth || options.forceAuth)) {
@@ -125,7 +123,7 @@ export default class PoolResource extends EventEmitter {
returned = true; returned = true;
if (err) { if (err) {
this.connection.close(); this.connection.close();
this.emit('error', err); this._fail(err);
return callback(err); return callback(err);
} }
this._connected = true; this._connected = true;
@@ -177,7 +175,7 @@ export default class PoolResource extends EventEmitter {
this.messages++; this.messages++;
if (err) { if (err) {
this.connection.close(); this.connection.close();
this.emit('error', err); this._fail(err);
return callback(err); return callback(err);
} }
info.envelope = { info.envelope = {
@@ -190,7 +188,7 @@ export default class PoolResource extends EventEmitter {
const err = new Error('Resource exhausted'); const err = new Error('Resource exhausted');
err.code = errors.EMAXLIMIT; err.code = errors.EMAXLIMIT;
this.connection.close(); this.connection.close();
this.emit('error', err); this._fail(err);
} }
else { else {
this.pool._checkRateLimit(() => { this.pool._checkRateLimit(() => {
+17 -27
View File
@@ -139,31 +139,6 @@ class SMTPTransport extends EventEmitter {
connection.close(); connection.close();
return callback(err); return callback(err);
}); });
connection.once('end', () => {
if (returned) {
return;
}
const timer = setTimeout(() => {
if (returned) {
return;
}
returned = true;
cleanupPerCallAuth();
// still have not returned, this means we have an unexpected connection close
const err = new Error('Unexpected socket close');
if (connection && connection.upgrading) {
// starttls connection errors
err.code = errors.ETLS;
}
callback(err);
}, 1000);
try {
timer.unref();
}
catch (_E) {
// Ignore. Happens on envs with non-node timer implementation
}
});
const sendMessage = () => { const sendMessage = () => {
const envelope = mail.message.getEnvelope(); const envelope = mail.message.getEnvelope();
const messageId = mail.message.messageId(); const messageId = mail.message.messageId();
@@ -183,6 +158,10 @@ class SMTPTransport extends EventEmitter {
messageId messageId
}, 'Sending message %s to <%s>', messageId, recipients.join(', ')); }, 'Sending message %s to <%s>', messageId, recipients.join(', '));
connection.send(envelope, mail.message.createReadStream(), (err, info) => { connection.send(envelope, mail.message.createReadStream(), (err, info) => {
if (returned) {
// the connection error handler has already reported this send
return;
}
returned = true; returned = true;
cleanupPerCallAuth(); cleanupPerCallAuth();
connection.close(); connection.close();
@@ -209,10 +188,16 @@ class SMTPTransport extends EventEmitter {
} }
}); });
}; };
connection.connect(() => { connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
if (err) {
// the server closed the connection before the greeting
returned = true;
connection.close();
return callback(err);
}
perCallAuth = this.getAuth(mail.data.auth); perCallAuth = this.getAuth(mail.data.auth);
if (perCallAuth && (connection.allowsAuth || options.forceAuth)) { if (perCallAuth && (connection.allowsAuth || options.forceAuth)) {
connection.login(perCallAuth, err => { connection.login(perCallAuth, err => {
@@ -294,10 +279,15 @@ class SMTPTransport extends EventEmitter {
connection.quit(); connection.quit();
return done(null, true); return done(null, true);
}; };
connection.connect(() => { connection.connect(err => {
if (returned) { if (returned) {
return; return;
} }
if (err) {
returned = true;
connection.close();
return done(err);
}
perCallAuth = this.getAuth({}); perCallAuth = this.getAuth({});
if (perCallAuth && (connection.allowsAuth || options.forceAuth)) { if (perCallAuth && (connection.allowsAuth || options.forceAuth)) {
connection.login(perCallAuth, err => { connection.login(perCallAuth, err => {
+2
View File
@@ -61,6 +61,8 @@ export interface XOAuth2Options {
component?: string | undefined; component?: string | undefined;
/** Extra headers for the token request */ /** Extra headers for the token request */
customHeaders?: OutgoingHttpHeaders | undefined; customHeaders?: OutgoingHttpHeaders | undefined;
/** Timeout for the token request in milliseconds, defaults to 60000, 0 disables it */
requestTimeout?: number | undefined;
/** Extra form fields for the token request */ /** Extra form fields for the token request */
customParams?: { customParams?: {
[key: string]: any; [key: string]: any;
+11 -3
View File
@@ -36,12 +36,14 @@ class XOAuth2 extends Stream {
constructor(options, logger) { constructor(options, logger) {
super(); super();
this.options = options || {}; this.options = options || {};
this.configError = false;
if (options && options.serviceClient) { if (options && options.serviceClient) {
if (!options.privateKey || !options.user) { if (!options.privateKey || !options.user) {
// reported through getToken() rather than an 'error' event: the transports forward
// that event up to the Mail object, where nothing may be listening
const err = new Error('Options "privateKey" and "user" are required for service account!'); const err = new Error('Options "privateKey" and "user" are required for service account!');
err.code = errors.EOAUTH2; err.code = errors.EOAUTH2;
setImmediate(() => this.emit('error', err)); this.configError = err;
return;
} }
const serviceRequestTimeout = Math.min(Math.max(Number(this.options.serviceRequestTimeout) || 0, 0), 3600); const serviceRequestTimeout = Math.min(Math.max(Number(this.options.serviceRequestTimeout) || 0, 0), 3600);
this.options.serviceRequestTimeout = serviceRequestTimeout || 5 * 60; this.options.serviceRequestTimeout = serviceRequestTimeout || 5 * 60;
@@ -75,6 +77,9 @@ class XOAuth2 extends Stream {
getToken(renew, callback) { getToken(renew, callback) {
// the error paths hand over the error alone // the error paths hand over the error alone
const done = callback; const done = callback;
if (this.configError) {
return done(this.configError);
}
if (!renew && this.accessToken && (!this.expires || this.expires > Date.now())) { if (!renew && this.accessToken && (!this.expires || this.expires > Date.now())) {
this.logger.debug({ this.logger.debug({
tnx: 'OAUTH2', tnx: 'OAUTH2',
@@ -311,7 +316,10 @@ class XOAuth2 extends Stream {
method: 'post', method: 'post',
headers: params.customHeaders, headers: params.customHeaders,
body: payload, body: payload,
allowErrorResponse: true allowErrorResponse: true,
// unset falls back to the fetch default, a stalled token endpoint would otherwise keep
// `renewing` set and queue every later request
timeout: params.requestTimeout
}; };
// OAuth2 token endpoints are credential-bearing. src/fetch already // OAuth2 token endpoints are credential-bearing. src/fetch already
// validates certs by default; pin rejectUnauthorized:true here so the // validates certs by default; pin rejectUnauthorized:true here so the
+7 -7
View File
@@ -1,6 +1,6 @@
{ {
"name": "nodemailer", "name": "nodemailer",
"version": "10.0.11", "version": "10.0.13",
"description": "Easy as cake e-mail sending from your Node.js applications", "description": "Easy as cake e-mail sending from your Node.js applications",
"type": "module", "type": "module",
"main": "./dist/cjs/nodemailer.js", "main": "./dist/cjs/nodemailer.js",
@@ -151,24 +151,24 @@
}, },
"homepage": "https://nodemailer.com/", "homepage": "https://nodemailer.com/",
"devDependencies": { "devDependencies": {
"@aws-sdk/client-sesv2": "3.1141.0", "@aws-sdk/client-sesv2": "3.1143.0",
"@types/node": "20.19.43", "@types/node": "20.19.43",
"bunyan": "1.8.15", "bunyan": "1.8.15",
"c8": "12.0.0", "c8": "12.0.0",
"eslint": "10.11.0", "eslint": "10.11.0",
"eslint-config-prettier": "10.1.8", "eslint-config-prettier": "10.1.8",
"globals": "17.12.0", "globals": "17.12.0",
"libbase64": "1.3.0", "libbase64": "1.3.1",
"libmime": "5.4.4", "libmime": "5.4.6",
"libqp": "2.1.1", "libqp": "2.1.2",
"mailauth": "7.1.0", "mailauth": "7.1.0",
"prettier": "3.9.9", "prettier": "3.9.9",
"proxy": "1.0.2", "proxy": "1.0.2",
"proxy-test-server": "1.0.0", "proxy-test-server": "1.0.0",
"smtp-server": "3.19.13", "smtp-server": "3.19.15",
"tsx": "4.23.15", "tsx": "4.23.15",
"typescript": "6.0.3", "typescript": "6.0.3",
"typescript-eslint": "8.70.1" "typescript-eslint": "8.71.0"
}, },
"engines": { "engines": {
"node": ">=20.0.0" "node": ">=20.0.0"