mirror of
https://github.com/dawidd6/action-send-mail.git
synced 2026-09-17 09:06:48 +07:00
node_modules: update (#322)
Co-authored-by: dawidd6 <9713907+dawidd6@users.noreply.github.com>
This commit is contained in:
+125
@@ -0,0 +1,125 @@
|
||||
import { EventEmitter } from 'node:events';
|
||||
import * as shared from '../shared/index.js';
|
||||
import type { SMTPTransportOptions, SMTPTransportGetSocketCallback, SMTPTransportSendCallback, SMTPSentMessageInfo } from '../smtp-transport/index.js';
|
||||
import type MailMessage from '../mailer/mail-message.js';
|
||||
import type { default as Mail, SendMailOptions, VerifyCallback } from '../mailer/index.js';
|
||||
/**
|
||||
* Options for the pooled SMTP transport, the SMTP transport options plus the pool settings
|
||||
*/
|
||||
export interface SMTPPoolOptions extends SMTPTransportOptions {
|
||||
/** Set to true to get this pooled transport from createTransport */
|
||||
pool?: boolean | undefined;
|
||||
/** Maximum number of open connections, defaults to 5 */
|
||||
maxConnections?: number | undefined;
|
||||
/** Number of messages a connection sends before it is closed and replaced, defaults to 100 */
|
||||
maxMessages?: number | undefined;
|
||||
/** Maximum number of messages to send in rateDelta milliseconds, unlimited when not set */
|
||||
rateLimit?: number | undefined;
|
||||
/** Time window for rateLimit in milliseconds, defaults to 1000 */
|
||||
rateDelta?: number | undefined;
|
||||
/** How many times a message is requeued when its connection closes while sending, unlimited when not set or negative */
|
||||
maxRequeues?: number | undefined;
|
||||
}
|
||||
/**
|
||||
* The pool options once the constructor has applied the defaults
|
||||
*/
|
||||
export type SMTPPoolResolvedOptions = SMTPPoolOptions & {
|
||||
maxConnections: number;
|
||||
maxMessages: number;
|
||||
};
|
||||
/**
|
||||
* Result of a message sent through the pool, same as for the SMTP transport
|
||||
*/
|
||||
export type SMTPPoolSentMessageInfo = SMTPSentMessageInfo;
|
||||
/**
|
||||
* Callback for send()
|
||||
*/
|
||||
export type SMTPPoolSendCallback = SMTPTransportSendCallback;
|
||||
/**
|
||||
* A message waiting in the pool queue
|
||||
*/
|
||||
export interface SMTPPoolQueueEntry {
|
||||
/** The message to send */
|
||||
mail: MailMessage;
|
||||
/** How many times the entry was put back on the queue after its connection closed */
|
||||
requeueAttempts: number;
|
||||
/** Callback to run once the message is sent or failed */
|
||||
callback: SMTPPoolSendCallback;
|
||||
/** Message-ID value without the angle brackets, set when the entry is assigned to a connection */
|
||||
messageId?: string | undefined;
|
||||
}
|
||||
/**
|
||||
* Rate limiter state of the pool
|
||||
*/
|
||||
export interface SMTPPoolRateLimit {
|
||||
/** Messages assigned within the current window */
|
||||
counter: number;
|
||||
/** Timer that clears the current window */
|
||||
timeout: NodeJS.Timeout | null;
|
||||
/** Availability callbacks waiting for the window to clear */
|
||||
waiting: Array<() => void>;
|
||||
/** Start of the current window as a timestamp, false when no window is open */
|
||||
checkpoint: number | false;
|
||||
/** Window length in milliseconds */
|
||||
delta: number;
|
||||
/** Maximum number of messages per window, 0 for no limit */
|
||||
limit: number;
|
||||
}
|
||||
/**
|
||||
* Creates a SMTP pool transport object for Nodemailer
|
||||
*
|
||||
* @constructor
|
||||
* @param options SMTP Connection options
|
||||
*/
|
||||
declare class SMTPPool extends EventEmitter {
|
||||
options: SMTPPoolResolvedOptions;
|
||||
logger: shared.Logger;
|
||||
name: string;
|
||||
version: string;
|
||||
idling: boolean;
|
||||
/**
|
||||
* The Mail instance using this transport, assigned by Mail
|
||||
*/
|
||||
mailer?: Mail<SMTPPoolSentMessageInfo> | undefined;
|
||||
constructor(options?: SMTPPoolOptions | string);
|
||||
/**
|
||||
* Placeholder function for creating proxy sockets. This method immediatelly returns
|
||||
* without a socket
|
||||
*
|
||||
* @param options Connection options
|
||||
* @param callback Callback function to run with the socket keys
|
||||
*/
|
||||
getSocket(options: SMTPPoolOptions, callback: SMTPTransportGetSocketCallback): void;
|
||||
/**
|
||||
* Queues an e-mail to be sent using the selected settings
|
||||
*
|
||||
* @param mail Mail object
|
||||
* @param callback Callback function
|
||||
*/
|
||||
send(mail: MailMessage, callback: SMTPPoolSendCallback): boolean;
|
||||
/**
|
||||
* Closes all connections in the pool. If there is a message being sent, the connection
|
||||
* is closed later
|
||||
*/
|
||||
close(): void;
|
||||
/**
|
||||
* Returns true if there are free slots in the queue
|
||||
*/
|
||||
isIdle(): boolean;
|
||||
/**
|
||||
* Verifies SMTP configuration
|
||||
*
|
||||
* @param callback Callback function
|
||||
*/
|
||||
verify(): Promise<true>;
|
||||
verify(callback: VerifyCallback): void;
|
||||
}
|
||||
/**
|
||||
* Type aliases in the layout of @types/nodemailer, so `SMTPPool.Options` style references keep working
|
||||
*/
|
||||
declare namespace SMTPPool {
|
||||
type Options = SMTPPoolOptions;
|
||||
type MailOptions = SendMailOptions;
|
||||
type SentMessageInfo = SMTPPoolSentMessageInfo;
|
||||
}
|
||||
export default SMTPPool;
|
||||
+541
@@ -0,0 +1,541 @@
|
||||
"use strict";
|
||||
var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
var desc = Object.getOwnPropertyDescriptor(m, k);
|
||||
if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) {
|
||||
desc = { enumerable: true, get: function() { return m[k]; } };
|
||||
}
|
||||
Object.defineProperty(o, k2, desc);
|
||||
}) : (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
o[k2] = m[k];
|
||||
}));
|
||||
var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) {
|
||||
Object.defineProperty(o, "default", { enumerable: true, value: v });
|
||||
}) : function(o, v) {
|
||||
o["default"] = v;
|
||||
});
|
||||
var __importStar = (this && this.__importStar) || (function () {
|
||||
var ownKeys = function(o) {
|
||||
ownKeys = Object.getOwnPropertyNames || function (o) {
|
||||
var ar = [];
|
||||
for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k;
|
||||
return ar;
|
||||
};
|
||||
return ownKeys(o);
|
||||
};
|
||||
return function (mod) {
|
||||
if (mod && mod.__esModule) return mod;
|
||||
var result = {};
|
||||
if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]);
|
||||
__setModuleDefault(result, mod);
|
||||
return result;
|
||||
};
|
||||
})();
|
||||
var __importDefault = (this && this.__importDefault) || function (mod) {
|
||||
return (mod && mod.__esModule) ? mod : { "default": mod };
|
||||
};
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
const node_events_1 = require("node:events");
|
||||
const pool_resource_js_1 = __importDefault(require("./pool-resource.js"));
|
||||
const index_js_1 = __importDefault(require("../smtp-connection/index.js"));
|
||||
const index_js_2 = __importDefault(require("../well-known/index.js"));
|
||||
const shared = __importStar(require("../shared/index.js"));
|
||||
const errors = __importStar(require("../errors.js"));
|
||||
const packageData = __importStar(require("../package-info.js"));
|
||||
/**
|
||||
* Creates a SMTP pool transport object for Nodemailer
|
||||
*
|
||||
* @constructor
|
||||
* @param options SMTP Connection options
|
||||
*/
|
||||
class SMTPPool extends node_events_1.EventEmitter {
|
||||
constructor(options) {
|
||||
super();
|
||||
options = options || {};
|
||||
if (typeof options === 'string') {
|
||||
options = {
|
||||
url: options
|
||||
};
|
||||
}
|
||||
let urlData;
|
||||
let service = options.service;
|
||||
if (typeof options.getSocket === 'function') {
|
||||
this.getSocket = options.getSocket;
|
||||
}
|
||||
if (options.url) {
|
||||
urlData = shared.parseConnectionUrl(options.url);
|
||||
service = service || urlData.service;
|
||||
}
|
||||
this.options = shared.assign(false, // create new object
|
||||
options, // regular options
|
||||
urlData, // url options
|
||||
(service && (0, index_js_2.default)(service)) // wellknown options
|
||||
);
|
||||
this.options.maxConnections = this.options.maxConnections || 5;
|
||||
this.options.maxMessages = this.options.maxMessages || 100;
|
||||
this.logger = shared.getLogger(this.options, {
|
||||
component: this.options.component || 'smtp-pool'
|
||||
});
|
||||
this.name = 'SMTP (pool)';
|
||||
this.version = packageData.version + '[client:' + packageData.version + ']';
|
||||
this._rateLimit = {
|
||||
counter: 0,
|
||||
timeout: null,
|
||||
waiting: [],
|
||||
checkpoint: false,
|
||||
delta: Number(this.options.rateDelta) || 1000,
|
||||
limit: Number(this.options.rateLimit) || 0
|
||||
};
|
||||
this._closed = false;
|
||||
this._queue = [];
|
||||
this._connections = [];
|
||||
this._connectionCounter = 0;
|
||||
this.idling = true;
|
||||
setImmediate(() => {
|
||||
if (this.idling) {
|
||||
this.emit('idle');
|
||||
}
|
||||
});
|
||||
}
|
||||
/**
|
||||
* Placeholder function for creating proxy sockets. This method immediatelly returns
|
||||
* without a socket
|
||||
*
|
||||
* @param options Connection options
|
||||
* @param callback Callback function to run with the socket keys
|
||||
*/
|
||||
getSocket(options, callback) {
|
||||
// return immediatelly
|
||||
setImmediate(() => callback(null, false));
|
||||
}
|
||||
/**
|
||||
* Queues an e-mail to be sent using the selected settings
|
||||
*
|
||||
* @param mail Mail object
|
||||
* @param callback Callback function
|
||||
*/
|
||||
send(mail, callback) {
|
||||
if (this._closed) {
|
||||
return false;
|
||||
}
|
||||
this._queue.push({
|
||||
mail,
|
||||
requeueAttempts: 0,
|
||||
callback
|
||||
});
|
||||
if (this.idling && this._queue.length >= this.options.maxConnections) {
|
||||
this.idling = false;
|
||||
}
|
||||
setImmediate(() => this._processMessages());
|
||||
return true;
|
||||
}
|
||||
/**
|
||||
* Closes all connections in the pool. If there is a message being sent, the connection
|
||||
* is closed later
|
||||
*/
|
||||
close() {
|
||||
let connection;
|
||||
const len = this._connections.length;
|
||||
this._closed = true;
|
||||
// release connections gated by the rate limiter so they become
|
||||
// available and are torn down below instead of leaking with a
|
||||
// cleared timer that never fires
|
||||
this._clearRateLimit();
|
||||
if (!len && !this._queue.length) {
|
||||
return;
|
||||
}
|
||||
// remove all available connections
|
||||
for (let i = len - 1; i >= 0; i--) {
|
||||
if (this._connections[i] && this._connections[i].available) {
|
||||
connection = this._connections[i];
|
||||
connection.close();
|
||||
this.logger.info({
|
||||
tnx: 'connection',
|
||||
cid: connection.id,
|
||||
action: 'removed'
|
||||
}, 'Connection #%s removed', connection.id);
|
||||
}
|
||||
}
|
||||
if (len && !this._connections.length) {
|
||||
this.logger.debug({
|
||||
tnx: 'connection'
|
||||
}, 'All connections removed');
|
||||
}
|
||||
if (!this._queue.length) {
|
||||
return;
|
||||
}
|
||||
// make sure that entire queue would be cleaned
|
||||
const invokeCallbacks = () => {
|
||||
if (!this._queue.length) {
|
||||
this.logger.debug({
|
||||
tnx: 'connection'
|
||||
}, 'Pending queue entries cleared');
|
||||
return;
|
||||
}
|
||||
const entry = this._queue.shift();
|
||||
if (entry && typeof entry.callback === 'function') {
|
||||
try {
|
||||
entry.callback(new Error('Connection pool was closed'));
|
||||
}
|
||||
catch (E) {
|
||||
// the queue is drained without a connection, so there is no cid to log
|
||||
this.logger.error({
|
||||
err: E,
|
||||
tnx: 'callback'
|
||||
}, 'Callback error: %s', E.message);
|
||||
}
|
||||
}
|
||||
setImmediate(invokeCallbacks);
|
||||
};
|
||||
setImmediate(invokeCallbacks);
|
||||
}
|
||||
/**
|
||||
* Check the queue and available connections. If there is a message to be sent and there is
|
||||
* an available connection, then use this connection to send the mail
|
||||
* @internal
|
||||
*/
|
||||
_processMessages() {
|
||||
// do nothing if already closed
|
||||
if (this._closed) {
|
||||
return;
|
||||
}
|
||||
// do nothing if queue is empty
|
||||
if (!this._queue.length) {
|
||||
if (!this.idling) {
|
||||
// no pending jobs
|
||||
this.idling = true;
|
||||
this.emit('idle');
|
||||
}
|
||||
return;
|
||||
}
|
||||
// find first available connection
|
||||
let connection = this._connections.find(c => c.available);
|
||||
if (!connection && this._connections.length < this.options.maxConnections) {
|
||||
connection = this._createConnection();
|
||||
}
|
||||
if (!connection) {
|
||||
// no more free connection slots available
|
||||
this.idling = false;
|
||||
return;
|
||||
}
|
||||
// check if there is free space in the processing queue
|
||||
if (!this.idling && this._queue.length < this.options.maxConnections) {
|
||||
this.idling = true;
|
||||
this.emit('idle');
|
||||
}
|
||||
const entry = (connection.queueEntry = this._queue.shift());
|
||||
entry.messageId = (connection.queueEntry.mail.message.getHeader('message-id') || '').replace(/[<>\s]/g, '');
|
||||
connection.available = false;
|
||||
this.logger.debug({
|
||||
tnx: 'pool',
|
||||
cid: connection.id,
|
||||
messageId: entry.messageId,
|
||||
action: 'assign'
|
||||
}, 'Assigned message <%s> to #%s (%s)', entry.messageId, connection.id, connection.messages + 1);
|
||||
if (this._rateLimit.limit) {
|
||||
this._rateLimit.counter++;
|
||||
if (!this._rateLimit.checkpoint) {
|
||||
this._rateLimit.checkpoint = Date.now();
|
||||
}
|
||||
}
|
||||
connection.send(entry.mail, (err, info) => {
|
||||
// only process callback if current handler is not changed
|
||||
if (entry === connection.queueEntry) {
|
||||
try {
|
||||
entry.callback(err, info);
|
||||
}
|
||||
catch (E) {
|
||||
this.logger.error({
|
||||
err: E,
|
||||
tnx: 'callback',
|
||||
cid: connection.id
|
||||
}, 'Callback error for #%s: %s', connection.id, E.message);
|
||||
}
|
||||
connection.queueEntry = false;
|
||||
}
|
||||
});
|
||||
}
|
||||
/**
|
||||
* Creates a new pool resource
|
||||
* @internal
|
||||
*/
|
||||
_createConnection() {
|
||||
const connection = new pool_resource_js_1.default(this);
|
||||
connection.id = ++this._connectionCounter;
|
||||
this.logger.info({
|
||||
tnx: 'pool',
|
||||
cid: connection.id,
|
||||
action: 'conection'
|
||||
}, 'Created new pool resource #%s', connection.id);
|
||||
// resource comes available
|
||||
connection.on('available', () => {
|
||||
this.logger.debug({
|
||||
tnx: 'connection',
|
||||
cid: connection.id,
|
||||
action: 'available'
|
||||
}, 'Connection #%s became available', connection.id);
|
||||
if (this._closed) {
|
||||
// if already closed run close() that will remove this connections from connections list
|
||||
this.close();
|
||||
}
|
||||
else {
|
||||
// check if there's anything else to send
|
||||
this._processMessages();
|
||||
}
|
||||
});
|
||||
// resource is terminated with an error
|
||||
connection.once('error', (err) => {
|
||||
if (err.code !== errors.EMAXLIMIT) {
|
||||
this.logger.warn({
|
||||
err,
|
||||
tnx: 'pool',
|
||||
cid: connection.id
|
||||
}, 'Pool Error for #%s: %s', connection.id, err.message);
|
||||
}
|
||||
else {
|
||||
this.logger.debug({
|
||||
tnx: 'pool',
|
||||
cid: connection.id,
|
||||
action: 'maxlimit'
|
||||
}, 'Max messages limit exchausted for #%s', connection.id);
|
||||
}
|
||||
if (connection.queueEntry) {
|
||||
try {
|
||||
connection.queueEntry.callback(err);
|
||||
}
|
||||
catch (E) {
|
||||
this.logger.error({
|
||||
err: E,
|
||||
tnx: 'callback',
|
||||
cid: connection.id
|
||||
}, 'Callback error for #%s: %s', connection.id, E.message);
|
||||
}
|
||||
connection.queueEntry = false;
|
||||
}
|
||||
// remove the erroneus connection from connections list
|
||||
this._removeConnection(connection);
|
||||
this._continueProcessing();
|
||||
});
|
||||
connection.once('close', () => {
|
||||
this.logger.info({
|
||||
tnx: 'connection',
|
||||
cid: connection.id,
|
||||
action: 'closed'
|
||||
}, 'Connection #%s was closed', connection.id);
|
||||
this._removeConnection(connection);
|
||||
if (connection.queueEntry) {
|
||||
// If the connection closed when sending, add the message to the queue again
|
||||
// if max number of requeues is not reached yet
|
||||
// Note that we must wait a bit.. because the callback of the 'error' handler might be called
|
||||
// in the next event loop
|
||||
setTimeout(() => {
|
||||
if (connection.queueEntry) {
|
||||
if (this._shouldRequeuOnConnectionClose(connection.queueEntry)) {
|
||||
this._requeueEntryOnConnectionClose(connection);
|
||||
}
|
||||
else {
|
||||
this._failDeliveryOnConnectionClose(connection);
|
||||
}
|
||||
}
|
||||
this._continueProcessing();
|
||||
}, 50);
|
||||
}
|
||||
else {
|
||||
if (!this._closed && this.idling && !this._connections.length) {
|
||||
this.emit('clear');
|
||||
}
|
||||
this._continueProcessing();
|
||||
}
|
||||
});
|
||||
this._connections.push(connection);
|
||||
return connection;
|
||||
}
|
||||
/** @internal */
|
||||
_shouldRequeuOnConnectionClose(queueEntry) {
|
||||
if (this.options.maxRequeues === undefined || this.options.maxRequeues < 0) {
|
||||
return true;
|
||||
}
|
||||
return queueEntry.requeueAttempts < this.options.maxRequeues;
|
||||
}
|
||||
/** @internal */
|
||||
_failDeliveryOnConnectionClose(connection) {
|
||||
if (connection.queueEntry && connection.queueEntry.callback) {
|
||||
try {
|
||||
connection.queueEntry.callback(new Error('Reached maximum number of retries after connection was closed'));
|
||||
}
|
||||
catch (E) {
|
||||
this.logger.error({
|
||||
err: E,
|
||||
tnx: 'callback',
|
||||
messageId: connection.queueEntry.messageId,
|
||||
cid: connection.id
|
||||
}, 'Callback error for #%s: %s', connection.id, E.message);
|
||||
}
|
||||
connection.queueEntry = false;
|
||||
}
|
||||
}
|
||||
/** @internal */
|
||||
_requeueEntryOnConnectionClose(connection) {
|
||||
connection.queueEntry.requeueAttempts += 1;
|
||||
this.logger.debug({
|
||||
tnx: 'pool',
|
||||
cid: connection.id,
|
||||
messageId: connection.queueEntry.messageId,
|
||||
action: 'requeue'
|
||||
}, 'Re-queued message <%s> for #%s. Attempt: #%s', connection.queueEntry.messageId, connection.id, connection.queueEntry.requeueAttempts);
|
||||
this._queue.unshift(connection.queueEntry);
|
||||
connection.queueEntry = false;
|
||||
}
|
||||
/**
|
||||
* Continue to process message if the pool hasn't closed
|
||||
* @internal
|
||||
*/
|
||||
_continueProcessing() {
|
||||
if (this._closed) {
|
||||
this.close();
|
||||
}
|
||||
else {
|
||||
setTimeout(() => this._processMessages(), 100);
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Remove resource from pool
|
||||
*
|
||||
* @param connection The PoolResource to remove
|
||||
* @internal
|
||||
*/
|
||||
_removeConnection(connection) {
|
||||
const index = this._connections.indexOf(connection);
|
||||
if (index !== -1) {
|
||||
this._connections.splice(index, 1);
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Checks if connections have hit current rate limit and if so, queues the availability callback
|
||||
*
|
||||
* @param callback Callback function to run once rate limiter has been cleared
|
||||
* @internal
|
||||
*/
|
||||
_checkRateLimit(callback) {
|
||||
if (!this._rateLimit.limit) {
|
||||
return callback();
|
||||
}
|
||||
const now = Date.now();
|
||||
if (this._rateLimit.counter < this._rateLimit.limit) {
|
||||
return callback();
|
||||
}
|
||||
this._rateLimit.waiting.push(callback);
|
||||
if (this._rateLimit.checkpoint <= now - this._rateLimit.delta) {
|
||||
return this._clearRateLimit();
|
||||
}
|
||||
if (!this._rateLimit.timeout) {
|
||||
this._rateLimit.timeout = setTimeout(() => this._clearRateLimit(), this._rateLimit.delta - (now - this._rateLimit.checkpoint));
|
||||
this._rateLimit.checkpoint = now;
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Clears current rate limit limitation and runs paused callback
|
||||
* @internal
|
||||
*/
|
||||
_clearRateLimit() {
|
||||
clearTimeout(this._rateLimit.timeout);
|
||||
this._rateLimit.timeout = null;
|
||||
this._rateLimit.counter = 0;
|
||||
this._rateLimit.checkpoint = false;
|
||||
// resume all paused connections
|
||||
while (this._rateLimit.waiting.length) {
|
||||
const cb = this._rateLimit.waiting.shift();
|
||||
setImmediate(cb);
|
||||
}
|
||||
}
|
||||
/**
|
||||
* Returns true if there are free slots in the queue
|
||||
*/
|
||||
isIdle() {
|
||||
return this.idling;
|
||||
}
|
||||
verify(callback) {
|
||||
let promise;
|
||||
if (!callback) {
|
||||
promise = new Promise((resolve, reject) => {
|
||||
callback = shared.callbackPromise(resolve, reject);
|
||||
});
|
||||
}
|
||||
const auth = new pool_resource_js_1.default(this).auth;
|
||||
this.getSocket(this.options, (err, socketOptions) => {
|
||||
if (err) {
|
||||
return callback(err);
|
||||
}
|
||||
let options = this.options;
|
||||
if (socketOptions && socketOptions.connection) {
|
||||
this.logger.info({
|
||||
tnx: 'proxy',
|
||||
remoteAddress: socketOptions.connection.remoteAddress,
|
||||
remotePort: socketOptions.connection.remotePort,
|
||||
destHost: options.host || '',
|
||||
destPort: options.port || '',
|
||||
action: 'connected'
|
||||
}, 'Using proxied socket from %s:%s to %s:%s', socketOptions.connection.remoteAddress, socketOptions.connection.remotePort, options.host || '', options.port || '');
|
||||
options = Object.assign(shared.assign(false, options), socketOptions);
|
||||
}
|
||||
const connection = new index_js_1.default(options);
|
||||
let returned = false;
|
||||
connection.once('error', err => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
returned = true;
|
||||
connection.close();
|
||||
return callback(err);
|
||||
});
|
||||
connection.once('end', () => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
returned = true;
|
||||
return callback(new Error('Connection closed'));
|
||||
});
|
||||
const finalize = () => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
returned = true;
|
||||
connection.quit();
|
||||
return callback(null, true);
|
||||
};
|
||||
connection.connect(() => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
if (auth && (connection.allowsAuth || options.forceAuth)) {
|
||||
connection.login(auth, err => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
if (err) {
|
||||
returned = true;
|
||||
connection.close();
|
||||
return callback(err);
|
||||
}
|
||||
finalize();
|
||||
});
|
||||
}
|
||||
else if (!auth && connection.allowsAuth && options.forceAuth) {
|
||||
const err = new Error('Authentication info was not provided');
|
||||
err.code = errors.ENOAUTH;
|
||||
returned = true;
|
||||
connection.close();
|
||||
return callback(err);
|
||||
}
|
||||
else {
|
||||
finalize();
|
||||
}
|
||||
});
|
||||
});
|
||||
return promise;
|
||||
}
|
||||
}
|
||||
exports.default = SMTPPool;
|
||||
module.exports = exports.default;
|
||||
Object.defineProperty(module.exports, 'default', { value: exports.default, enumerable: false, writable: true, configurable: true });
|
||||
+62
@@ -0,0 +1,62 @@
|
||||
import SMTPConnection from '../smtp-connection/index.js';
|
||||
import { type Logger } from '../shared/index.js';
|
||||
import { EventEmitter } from 'node:events';
|
||||
import type { SMTPTransportAuth, SMTPTransportSendCallback } from '../smtp-transport/index.js';
|
||||
import type MailMessage from '../mailer/mail-message.js';
|
||||
import type SMTPPool from './index.js';
|
||||
import type { SMTPPoolResolvedOptions, SMTPPoolQueueEntry } from './index.js';
|
||||
/**
|
||||
* Callback for connect(), the result is true once the connection is ready for messages
|
||||
*/
|
||||
export type PoolResourceConnectCallback = (err: Error | null, connected?: true) => void;
|
||||
/**
|
||||
* Callback for send()
|
||||
*/
|
||||
export type PoolResourceSendCallback = SMTPTransportSendCallback;
|
||||
/**
|
||||
* Creates an element for the pool
|
||||
*
|
||||
* @constructor
|
||||
* @param pool SMTPPool instance
|
||||
*/
|
||||
export default class PoolResource extends EventEmitter {
|
||||
pool: SMTPPool;
|
||||
options: SMTPPoolResolvedOptions;
|
||||
logger: Logger;
|
||||
/**
|
||||
* Authentication data for the connection, set when the pool options include auth
|
||||
*/
|
||||
auth?: SMTPTransportAuth | undefined;
|
||||
messages: number;
|
||||
available: boolean;
|
||||
/**
|
||||
* The SMTP connection, set by connect()
|
||||
*/
|
||||
connection: SMTPConnection;
|
||||
/**
|
||||
* Resource id, assigned by the pool
|
||||
*/
|
||||
id: number;
|
||||
/**
|
||||
* The queue entry being sent, assigned by the pool. False once it has been handled
|
||||
*/
|
||||
queueEntry?: SMTPPoolQueueEntry | false | undefined;
|
||||
constructor(pool: SMTPPool);
|
||||
/**
|
||||
* Initiates a connection to the SMTP server
|
||||
*
|
||||
* @param callback Callback function to run once the connection is established or failed
|
||||
*/
|
||||
connect(callback: PoolResourceConnectCallback): void;
|
||||
/**
|
||||
* Sends an e-mail to be sent using the selected settings
|
||||
*
|
||||
* @param mail Mail object
|
||||
* @param callback Callback function
|
||||
*/
|
||||
send(mail: MailMessage, callback: PoolResourceSendCallback): void;
|
||||
/**
|
||||
* Closes the connection
|
||||
*/
|
||||
close(): void;
|
||||
}
|
||||
+261
@@ -0,0 +1,261 @@
|
||||
"use strict";
|
||||
var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
var desc = Object.getOwnPropertyDescriptor(m, k);
|
||||
if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) {
|
||||
desc = { enumerable: true, get: function() { return m[k]; } };
|
||||
}
|
||||
Object.defineProperty(o, k2, desc);
|
||||
}) : (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
o[k2] = m[k];
|
||||
}));
|
||||
var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) {
|
||||
Object.defineProperty(o, "default", { enumerable: true, value: v });
|
||||
}) : function(o, v) {
|
||||
o["default"] = v;
|
||||
});
|
||||
var __importStar = (this && this.__importStar) || (function () {
|
||||
var ownKeys = function(o) {
|
||||
ownKeys = Object.getOwnPropertyNames || function (o) {
|
||||
var ar = [];
|
||||
for (var k in o) if (Object.prototype.hasOwnProperty.call(o, k)) ar[ar.length] = k;
|
||||
return ar;
|
||||
};
|
||||
return ownKeys(o);
|
||||
};
|
||||
return function (mod) {
|
||||
if (mod && mod.__esModule) return mod;
|
||||
var result = {};
|
||||
if (mod != null) for (var k = ownKeys(mod), i = 0; i < k.length; i++) if (k[i] !== "default") __createBinding(result, mod, k[i]);
|
||||
__setModuleDefault(result, mod);
|
||||
return result;
|
||||
};
|
||||
})();
|
||||
var __importDefault = (this && this.__importDefault) || function (mod) {
|
||||
return (mod && mod.__esModule) ? mod : { "default": mod };
|
||||
};
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
const index_js_1 = __importDefault(require("../smtp-connection/index.js"));
|
||||
const index_js_2 = require("../shared/index.js");
|
||||
const index_js_3 = __importDefault(require("../xoauth2/index.js"));
|
||||
const errors = __importStar(require("../errors.js"));
|
||||
const node_events_1 = require("node:events");
|
||||
/**
|
||||
* Creates an element for the pool
|
||||
*
|
||||
* @constructor
|
||||
* @param pool SMTPPool instance
|
||||
*/
|
||||
class PoolResource extends node_events_1.EventEmitter {
|
||||
constructor(pool) {
|
||||
super();
|
||||
this.pool = pool;
|
||||
this.options = pool.options;
|
||||
this.logger = this.pool.logger;
|
||||
if (this.options.auth) {
|
||||
switch ((this.options.auth.type || '').toString().toUpperCase()) {
|
||||
case 'OAUTH2': {
|
||||
const oauth2 = new index_js_3.default(this.options.auth, this.logger);
|
||||
oauth2.provisionCallback =
|
||||
(this.pool.mailer && this.pool.mailer.get('oauth2_provision_cb')) || oauth2.provisionCallback;
|
||||
this.auth = {
|
||||
type: 'OAUTH2',
|
||||
user: this.options.auth.user,
|
||||
oauth2,
|
||||
method: 'XOAUTH2'
|
||||
};
|
||||
oauth2.on('token', (token) => this.pool.mailer.emit('token', token));
|
||||
oauth2.on('error', err => this.emit('error', err));
|
||||
break;
|
||||
}
|
||||
default:
|
||||
if (!this.options.auth.user && !this.options.auth.pass) {
|
||||
break;
|
||||
}
|
||||
this.auth = {
|
||||
type: (this.options.auth.type || '').toString().toUpperCase() || 'LOGIN',
|
||||
user: this.options.auth.user,
|
||||
credentials: {
|
||||
user: this.options.auth.user || '',
|
||||
pass: this.options.auth.pass,
|
||||
options: this.options.auth.options
|
||||
},
|
||||
method: (this.options.auth.method || '').trim().toUpperCase() || this.options.authMethod || false
|
||||
};
|
||||
}
|
||||
}
|
||||
this._connection = false;
|
||||
this._connected = false;
|
||||
this.messages = 0;
|
||||
this.available = true;
|
||||
}
|
||||
/**
|
||||
* Initiates a connection to the SMTP server
|
||||
*
|
||||
* @param callback Callback function to run once the connection is established or failed
|
||||
*/
|
||||
connect(callback) {
|
||||
this.pool.getSocket(this.options, (err, socketOptions) => {
|
||||
if (err) {
|
||||
// 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
|
||||
// login does
|
||||
this.emit('error', err);
|
||||
return callback(err);
|
||||
}
|
||||
let returned = false;
|
||||
let options = this.options;
|
||||
if (socketOptions && socketOptions.connection) {
|
||||
this.logger.info({
|
||||
tnx: 'proxy',
|
||||
remoteAddress: socketOptions.connection.remoteAddress,
|
||||
remotePort: socketOptions.connection.remotePort,
|
||||
destHost: options.host || '',
|
||||
destPort: options.port || '',
|
||||
action: 'connected'
|
||||
}, 'Using proxied socket from %s:%s to %s:%s', socketOptions.connection.remoteAddress, socketOptions.connection.remotePort, options.host || '', options.port || '');
|
||||
options = Object.assign((0, index_js_2.assign)(false, options), socketOptions);
|
||||
}
|
||||
this.connection = new index_js_1.default(options);
|
||||
this.connection.once('error', err => {
|
||||
this.emit('error', err);
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
returned = true;
|
||||
return callback(err);
|
||||
});
|
||||
this.connection.once('end', () => {
|
||||
this.close();
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
returned = true;
|
||||
const timer = setTimeout(() => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
// still have not returned, this means we have an unexpected connection close
|
||||
const err = new Error('Unexpected socket close');
|
||||
if (this.connection &&
|
||||
this.connection._socket &&
|
||||
this.connection._socket.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
|
||||
}
|
||||
});
|
||||
this.connection.connect(() => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
if (this.auth && (this.connection.allowsAuth || options.forceAuth)) {
|
||||
this.connection.login(this.auth, err => {
|
||||
if (returned) {
|
||||
return;
|
||||
}
|
||||
returned = true;
|
||||
if (err) {
|
||||
this.connection.close();
|
||||
this.emit('error', err);
|
||||
return callback(err);
|
||||
}
|
||||
this._connected = true;
|
||||
callback(null, true);
|
||||
});
|
||||
}
|
||||
else {
|
||||
returned = true;
|
||||
this._connected = true;
|
||||
return callback(null, true);
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
/**
|
||||
* Sends an e-mail to be sent using the selected settings
|
||||
*
|
||||
* @param mail Mail object
|
||||
* @param callback Callback function
|
||||
*/
|
||||
send(mail, callback) {
|
||||
if (!this._connected) {
|
||||
return this.connect(err => {
|
||||
if (err) {
|
||||
return callback(err);
|
||||
}
|
||||
return this.send(mail, callback);
|
||||
});
|
||||
}
|
||||
const envelope = mail.message.getEnvelope();
|
||||
const messageId = mail.message.messageId();
|
||||
const recipients = [].concat(envelope.to || []);
|
||||
if (recipients.length > 3) {
|
||||
recipients.push('...and ' + recipients.splice(2).length + ' more');
|
||||
}
|
||||
this.logger.info({
|
||||
tnx: 'send',
|
||||
messageId,
|
||||
cid: this.id
|
||||
}, 'Sending message %s using #%s to <%s>', messageId, this.id, recipients.join(', '));
|
||||
if (mail.data.dsn) {
|
||||
envelope.dsn = mail.data.dsn;
|
||||
}
|
||||
// RFC 8689: Pass requireTLSExtensionEnabled to envelope for MAIL FROM parameter
|
||||
if (mail.data.requireTLSExtensionEnabled) {
|
||||
envelope.requireTLSExtensionEnabled = mail.data.requireTLSExtensionEnabled;
|
||||
}
|
||||
this.connection.send(envelope, mail.message.createReadStream(), (err, info) => {
|
||||
this.messages++;
|
||||
if (err) {
|
||||
this.connection.close();
|
||||
this.emit('error', err);
|
||||
return callback(err);
|
||||
}
|
||||
info.envelope = {
|
||||
from: envelope.from,
|
||||
to: envelope.to
|
||||
};
|
||||
info.messageId = messageId;
|
||||
setImmediate(() => {
|
||||
if (this.messages >= this.options.maxMessages) {
|
||||
const err = new Error('Resource exhausted');
|
||||
err.code = errors.EMAXLIMIT;
|
||||
this.connection.close();
|
||||
this.emit('error', err);
|
||||
}
|
||||
else {
|
||||
this.pool._checkRateLimit(() => {
|
||||
this.available = true;
|
||||
this.emit('available');
|
||||
});
|
||||
}
|
||||
});
|
||||
callback(null, info);
|
||||
});
|
||||
}
|
||||
/**
|
||||
* Closes the connection
|
||||
*/
|
||||
close() {
|
||||
this._connected = false;
|
||||
if (this.auth && this.auth.oauth2) {
|
||||
this.auth.oauth2.removeAllListeners();
|
||||
}
|
||||
if (this.connection) {
|
||||
this.connection.close();
|
||||
}
|
||||
this.emit('close');
|
||||
}
|
||||
}
|
||||
exports.default = PoolResource;
|
||||
module.exports = exports.default;
|
||||
Object.defineProperty(module.exports, 'default', { value: exports.default, enumerable: false, writable: true, configurable: true });
|
||||
Reference in New Issue
Block a user