| 1 | 'use strict';
|
|---|
| 2 |
|
|---|
| 3 | const packageData = require('../../package.json');
|
|---|
| 4 | const shared = require('../shared');
|
|---|
| 5 |
|
|---|
| 6 | /**
|
|---|
| 7 | * Generates a Transport object for streaming
|
|---|
| 8 | *
|
|---|
| 9 | * Possible options can be the following:
|
|---|
| 10 | *
|
|---|
| 11 | * * **buffer** if true, then returns the message as a Buffer object instead of a stream
|
|---|
| 12 | * * **newline** either 'windows' or 'unix'
|
|---|
| 13 | *
|
|---|
| 14 | * @constructor
|
|---|
| 15 | * @param {Object} optional config parameter
|
|---|
| 16 | */
|
|---|
| 17 | class StreamTransport {
|
|---|
| 18 | constructor(options) {
|
|---|
| 19 | options = options || {};
|
|---|
| 20 |
|
|---|
| 21 | this.options = options || {};
|
|---|
| 22 |
|
|---|
| 23 | this.name = 'StreamTransport';
|
|---|
| 24 | this.version = packageData.version;
|
|---|
| 25 |
|
|---|
| 26 | this.logger = shared.getLogger(this.options, {
|
|---|
| 27 | component: this.options.component || 'stream-transport'
|
|---|
| 28 | });
|
|---|
| 29 |
|
|---|
| 30 | this.winbreak = ['win', 'windows', 'dos', '\r\n'].includes((options.newline || '').toString().toLowerCase());
|
|---|
| 31 | }
|
|---|
| 32 |
|
|---|
| 33 | /**
|
|---|
| 34 | * Compiles a mailcomposer message and forwards it to handler that sends it
|
|---|
| 35 | *
|
|---|
| 36 | * @param {Object} emailMessage MailComposer object
|
|---|
| 37 | * @param {Function} callback Callback function to run when the sending is completed
|
|---|
| 38 | */
|
|---|
| 39 | send(mail, done) {
|
|---|
| 40 | // We probably need this in the output
|
|---|
| 41 | mail.message.keepBcc = true;
|
|---|
| 42 |
|
|---|
| 43 | let envelope = mail.data.envelope || mail.message.getEnvelope();
|
|---|
| 44 | let messageId = mail.message.messageId();
|
|---|
| 45 |
|
|---|
| 46 | let recipients = [].concat(envelope.to || []);
|
|---|
| 47 | if (recipients.length > 3) {
|
|---|
| 48 | recipients.push('...and ' + recipients.splice(2).length + ' more');
|
|---|
| 49 | }
|
|---|
| 50 | this.logger.info(
|
|---|
| 51 | {
|
|---|
| 52 | tnx: 'send',
|
|---|
| 53 | messageId
|
|---|
| 54 | },
|
|---|
| 55 | 'Sending message %s to <%s> using %s line breaks',
|
|---|
| 56 | messageId,
|
|---|
| 57 | recipients.join(', '),
|
|---|
| 58 | this.winbreak ? '<CR><LF>' : '<LF>'
|
|---|
| 59 | );
|
|---|
| 60 |
|
|---|
| 61 | setImmediate(() => {
|
|---|
| 62 | let stream;
|
|---|
| 63 |
|
|---|
| 64 | try {
|
|---|
| 65 | stream = mail.message.createReadStream();
|
|---|
| 66 | } catch (E) {
|
|---|
| 67 | this.logger.error(
|
|---|
| 68 | {
|
|---|
| 69 | err: E,
|
|---|
| 70 | tnx: 'send',
|
|---|
| 71 | messageId
|
|---|
| 72 | },
|
|---|
| 73 | 'Creating send stream failed for %s. %s',
|
|---|
| 74 | messageId,
|
|---|
| 75 | E.message
|
|---|
| 76 | );
|
|---|
| 77 | return done(E);
|
|---|
| 78 | }
|
|---|
| 79 |
|
|---|
| 80 | if (!this.options.buffer) {
|
|---|
| 81 | stream.once('error', err => {
|
|---|
| 82 | this.logger.error(
|
|---|
| 83 | {
|
|---|
| 84 | err,
|
|---|
| 85 | tnx: 'send',
|
|---|
| 86 | messageId
|
|---|
| 87 | },
|
|---|
| 88 | 'Failed creating message for %s. %s',
|
|---|
| 89 | messageId,
|
|---|
| 90 | err.message
|
|---|
| 91 | );
|
|---|
| 92 | });
|
|---|
| 93 | return done(null, {
|
|---|
| 94 | envelope: mail.data.envelope || mail.message.getEnvelope(),
|
|---|
| 95 | messageId,
|
|---|
| 96 | message: stream
|
|---|
| 97 | });
|
|---|
| 98 | }
|
|---|
| 99 |
|
|---|
| 100 | let chunks = [];
|
|---|
| 101 | let chunklen = 0;
|
|---|
| 102 | stream.on('readable', () => {
|
|---|
| 103 | let chunk;
|
|---|
| 104 | while ((chunk = stream.read()) !== null) {
|
|---|
| 105 | chunks.push(chunk);
|
|---|
| 106 | chunklen += chunk.length;
|
|---|
| 107 | }
|
|---|
| 108 | });
|
|---|
| 109 |
|
|---|
| 110 | stream.once('error', err => {
|
|---|
| 111 | this.logger.error(
|
|---|
| 112 | {
|
|---|
| 113 | err,
|
|---|
| 114 | tnx: 'send',
|
|---|
| 115 | messageId
|
|---|
| 116 | },
|
|---|
| 117 | 'Failed creating message for %s. %s',
|
|---|
| 118 | messageId,
|
|---|
| 119 | err.message
|
|---|
| 120 | );
|
|---|
| 121 | return done(err);
|
|---|
| 122 | });
|
|---|
| 123 |
|
|---|
| 124 | stream.on('end', () =>
|
|---|
| 125 | done(null, {
|
|---|
| 126 | envelope: mail.data.envelope || mail.message.getEnvelope(),
|
|---|
| 127 | messageId,
|
|---|
| 128 | message: Buffer.concat(chunks, chunklen)
|
|---|
| 129 | })
|
|---|
| 130 | );
|
|---|
| 131 | });
|
|---|
| 132 | }
|
|---|
| 133 | }
|
|---|
| 134 |
|
|---|
| 135 | module.exports = StreamTransport;
|
|---|