diff --git a/config.js b/config.js index fce4aba..5748f5d 100644 --- a/config.js +++ b/config.js @@ -3,17 +3,38 @@ const AWS = require('aws-sdk'); dotenv.config(); +function integerFromEnv(name, fallback) { + const value = process.env[name]; + if (value === undefined || value === '') return fallback; + return Number(value); +} + module.exports = { - WEBHOOK_URL: process.env.WEBHOOK_URL || 'https://enkhprqr4n2t.x.pipedream.net/', - PORT: process.env.PORT || 25, - MAX_FILE_SIZE: process.env.MAX_FILE_SIZE || 5 * 1024 * 1024, + WEBHOOK_URL: process.env.WEBHOOK_URL, + PORT: integerFromEnv('PORT', 25), + MAX_FILE_SIZE: integerFromEnv('MAX_FILE_SIZE', 5 * 1024 * 1024), + MAX_MESSAGE_SIZE: integerFromEnv('MAX_MESSAGE_SIZE', 25 * 1024 * 1024), + MAX_SMTP_CLIENTS: integerFromEnv('MAX_SMTP_CLIENTS', 50), + SMTP_SOCKET_TIMEOUT: integerFromEnv('SMTP_SOCKET_TIMEOUT', 5 * 60 * 1000), + SMTP_CLOSE_TIMEOUT: integerFromEnv('SMTP_CLOSE_TIMEOUT', 10 * 1000), + SMTP_ALLOWED_IPS: (process.env.SMTP_ALLOWED_IPS || '') + .split(',') + .map(ip => ip.trim()) + .filter(Boolean), BUCKET_NAME: process.env.S3_BUCKET_NAME, SMTP_SECURE: process.env.SMTP_SECURE === 'true', - WEBHOOK_CONCURRENCY: process.env.WEBHOOK_CONCURRENCY || 5, + TLS_KEY_PATH: process.env.TLS_KEY_PATH, + TLS_CERT_PATH: process.env.TLS_CERT_PATH, + WEBHOOK_CONCURRENCY: integerFromEnv('WEBHOOK_CONCURRENCY', 5), + WEBHOOK_QUEUE_RETRIES: integerFromEnv('WEBHOOK_QUEUE_RETRIES', 10), + WEBHOOK_QUEUE_RETRY_DELAY: integerFromEnv('WEBHOOK_QUEUE_RETRY_DELAY', 60 * 1000), + WEBHOOK_QUEUE_TIMEOUT: integerFromEnv('WEBHOOK_QUEUE_TIMEOUT', 35 * 1000), + SPOOL_DIR: process.env.SPOOL_DIR || '/var/lib/smtpwebhook/spool', + LOG_DIR: process.env.LOG_DIR || '/var/lib/smtpwebhook/logs', s3: new AWS.S3({ region: process.env.AWS_REGION, accessKeyId: process.env.AWS_ACCESS_KEY_ID, secretAccessKey: process.env.AWS_SECRET_ACCESS_KEY }) -}; \ No newline at end of file +}; diff --git a/package-lock.json b/package-lock.json index e44c3c5..b24b718 100644 --- a/package-lock.json +++ b/package-lock.json @@ -9,9 +9,12 @@ "@aws-sdk/lib-storage": "^3.x.x", "aws-sdk": "^2.x", "axios": "^1.7.7", + "better-queue": "^3.8.12", "dotenv": "^16.4.5", "mailparser": "^3.7.1", - "smtp-server": "^3.13.5" + "smtp-server": "^3.13.5", + "winston": "^3.x.x", + "winston-daily-rotate-file": "^4.x.x" } }, "node_modules/@aws-crypto/crc32": { @@ -886,6 +889,24 @@ "node": ">=16.0.0" } }, + "node_modules/@colors/colors": { + "version": "1.6.0", + "resolved": "https://registry.npmjs.org/@colors/colors/-/colors-1.6.0.tgz", + "integrity": "sha512-Ir+AOibqzrIsL6ajt3Rz3LskB7OiMVHqltZmspbW/TJuTVuyOMirVqAkjfY6JISiLHgyNqicAC8AyHHGzNd/dA==", + "engines": { + "node": ">=0.1.90" + } + }, + "node_modules/@dabh/diagnostics": { + "version": "2.0.3", + "resolved": "https://registry.npmjs.org/@dabh/diagnostics/-/diagnostics-2.0.3.tgz", + "integrity": "sha512-hrlQOIi7hAfzsMqlGSFyVucrx38O+j6wiGOf//H2ecvIEqYN4ADBSS2iLMh5UFyDunCNniUIPk/q3riFv45xRA==", + "dependencies": { + "colorspace": "1.1.x", + "enabled": "2.0.x", + "kuler": "^2.0.0" + } + }, "node_modules/@selderee/plugin-htmlparser2": { "version": "0.11.0", "resolved": "https://registry.npmjs.org/@selderee/plugin-htmlparser2/-/plugin-htmlparser2-0.11.0.tgz", @@ -1538,6 +1559,16 @@ "node": ">=16.0.0" } }, + "node_modules/@types/triple-beam": { + "version": "1.3.5", + "resolved": "https://registry.npmjs.org/@types/triple-beam/-/triple-beam-1.3.5.tgz", + "integrity": "sha512-6WaYesThRMCl19iryMYP7/x2OVgCtbIVflDGFpWnb9irXI3UjYE4AzmYuiUKY1AJstGijoY+MgUszMgRxIYTYw==" + }, + "node_modules/async": { + "version": "3.2.6", + "resolved": "https://registry.npmjs.org/async/-/async-3.2.6.tgz", + "integrity": "sha512-htCUDlxyyCLMgaM3xXg0C0LW2xqfuQ6p05pCEIsXuyQ+a1koYKTuBMzRNwmybfLgvJDMd0r1LTn4+E0Ti6C2AA==" + }, "node_modules/asynckit": { "version": "0.4.0", "resolved": "https://registry.npmjs.org/asynckit/-/asynckit-0.4.0.tgz", @@ -1605,9 +1636,10 @@ } }, "node_modules/axios": { - "version": "1.7.7", - "resolved": "https://registry.npmjs.org/axios/-/axios-1.7.7.tgz", - "integrity": "sha512-S4kL7XrjgBmvdGut0sN3yJxqYzrDOnivkBiN0OFs6hLiUam3UPvswUo0kqGyhqUZGEOytHyumEdXsAkgCOUf3Q==", + "version": "1.10.0", + "resolved": "https://registry.npmjs.org/axios/-/axios-1.10.0.tgz", + "integrity": "sha512-/1xYAC4MP/HEG+3duIhFr4ZQXR4sQXOIe+o6sdqzeykGLx6Upp/1p8MHqhINOvGeP7xyNHe7tsiJByc4SSVUxw==", + "license": "MIT", "dependencies": { "follow-redirects": "^1.15.6", "form-data": "^4.0.0", @@ -1641,6 +1673,21 @@ } ] }, + "node_modules/better-queue": { + "version": "3.8.12", + "resolved": "https://registry.npmjs.org/better-queue/-/better-queue-3.8.12.tgz", + "integrity": "sha512-D9KZ+Us+2AyaCz693/9AyjTg0s8hEmkiM/MB3i09cs4MdK1KgTSGJluXRYmOulR69oLZVo2XDFtqsExDt8oiLA==", + "dependencies": { + "better-queue-memory": "^1.0.1", + "node-eta": "^0.9.0", + "uuid": "^9.0.0" + } + }, + "node_modules/better-queue-memory": { + "version": "1.0.4", + "resolved": "https://registry.npmjs.org/better-queue-memory/-/better-queue-memory-1.0.4.tgz", + "integrity": "sha512-SWg5wFIShYffEmJpI6LgbL8/3Dqhku7xI1oEiy6FroP9DbcZlG0ZDjxvPdP9t7hTGW40IpIcC6zVoGT1oxjOuA==" + }, "node_modules/bowser": { "version": "2.11.0", "resolved": "https://registry.npmjs.org/bowser/-/bowser-2.11.0.tgz", @@ -1673,6 +1720,46 @@ "url": "https://github.com/sponsors/ljharb" } }, + "node_modules/color": { + "version": "3.2.1", + "resolved": "https://registry.npmjs.org/color/-/color-3.2.1.tgz", + "integrity": "sha512-aBl7dZI9ENN6fUGC7mWpMTPNHmWUSNan9tuWN6ahh5ZLNk9baLJOnSMlrQkHcrfFgz2/RigjUVAjdx36VcemKA==", + "dependencies": { + "color-convert": "^1.9.3", + "color-string": "^1.6.0" + } + }, + "node_modules/color-convert": { + "version": "1.9.3", + "resolved": "https://registry.npmjs.org/color-convert/-/color-convert-1.9.3.tgz", + "integrity": "sha512-QfAUtd+vFdAtFQcC8CCyYt1fYWxSqAiK2cSD6zDB8N3cpsEBAvRxp9zOGg6G/SHHJYAT88/az/IuDGALsNVbGg==", + "dependencies": { + "color-name": "1.1.3" + } + }, + "node_modules/color-name": { + "version": "1.1.3", + "resolved": "https://registry.npmjs.org/color-name/-/color-name-1.1.3.tgz", + "integrity": "sha512-72fSenhMw2HZMTVHeCA9KCmpEIbzWiQsjN+BHcBbS9vr1mtt+vJjPdksIBNUmKAW8TFUDPJK5SUU3QhE9NEXDw==" + }, + "node_modules/color-string": { + "version": "1.9.1", + "resolved": "https://registry.npmjs.org/color-string/-/color-string-1.9.1.tgz", + "integrity": "sha512-shrVawQFojnZv6xM40anx4CkoDP+fZsw/ZerEMsW/pyzsRbElpsL/DBVW7q3ExxwusdNXI3lXpuhEZkzs8p5Eg==", + "dependencies": { + "color-name": "^1.0.0", + "simple-swizzle": "^0.2.2" + } + }, + "node_modules/colorspace": { + "version": "1.1.4", + "resolved": "https://registry.npmjs.org/colorspace/-/colorspace-1.1.4.tgz", + "integrity": "sha512-BgvKJiuVu1igBUF2kEjRCZXol6wiiGbY5ipL/oVPwm0BL9sIpMIzM8IK7vwuxIIzOXMV3Ey5w+vxhm0rR/TN8w==", + "dependencies": { + "color": "^3.1.3", + "text-hex": "1.0.x" + } + }, "node_modules/combined-stream": { "version": "1.0.8", "resolved": "https://registry.npmjs.org/combined-stream/-/combined-stream-1.0.8.tgz", @@ -1778,6 +1865,11 @@ "url": "https://dotenvx.com" } }, + "node_modules/enabled": { + "version": "2.0.0", + "resolved": "https://registry.npmjs.org/enabled/-/enabled-2.0.0.tgz", + "integrity": "sha512-AKrN98kuwOzMIdAizXGI86UFBoo26CL21UM763y1h/GMSJ4/OHU9k2YlsmBpyScFo/wbLzWQJBMCW4+IO3/+OQ==" + }, "node_modules/encoding-japanese": { "version": "2.1.0", "resolved": "https://registry.npmjs.org/encoding-japanese/-/encoding-japanese-2.1.0.tgz", @@ -1845,6 +1937,24 @@ "fxparser": "src/cli/cli.js" } }, + "node_modules/fecha": { + "version": "4.2.3", + "resolved": "https://registry.npmjs.org/fecha/-/fecha-4.2.3.tgz", + "integrity": "sha512-OP2IUU6HeYKJi3i0z4A19kHMQoLVs4Hc+DPqqxI2h/DPZHTm/vjsfC6P0b4jCMy14XizLBqvndQ+UilD7707Jw==" + }, + "node_modules/file-stream-rotator": { + "version": "0.6.1", + "resolved": "https://registry.npmjs.org/file-stream-rotator/-/file-stream-rotator-0.6.1.tgz", + "integrity": "sha512-u+dBid4PvZw17PmDeRcNOtCP9CCK/9lRN2w+r1xIS7yOL9JFrIBKTvrYsxT4P0pGtThYTn++QS5ChHaUov3+zQ==", + "dependencies": { + "moment": "^2.29.1" + } + }, + "node_modules/fn.name": { + "version": "1.1.0", + "resolved": "https://registry.npmjs.org/fn.name/-/fn.name-1.1.0.tgz", + "integrity": "sha512-GRnmB5gPyJpAhTQdSZTSp9uaPSvl09KoYcMQtsB9rQoOmzs9dH6ffeccH+Z+cv6P68Hu5bC6JjRh4Ah/mHSNRw==" + }, "node_modules/follow-redirects": { "version": "1.15.9", "resolved": "https://registry.npmjs.org/follow-redirects/-/follow-redirects-1.15.9.tgz", @@ -2062,6 +2172,11 @@ "url": "https://github.com/sponsors/ljharb" } }, + "node_modules/is-arrayish": { + "version": "0.3.2", + "resolved": "https://registry.npmjs.org/is-arrayish/-/is-arrayish-0.3.2.tgz", + "integrity": "sha512-eVRqCvVlZbuw3GrM63ovNSNAeA1K16kaR/LRY/92w0zxQ5/1YzwblUX652i4Xs9RwAGjW9d9y6X88t8OaAJfWQ==" + }, "node_modules/is-callable": { "version": "1.2.7", "resolved": "https://registry.npmjs.org/is-callable/-/is-callable-1.2.7.tgz", @@ -2087,6 +2202,17 @@ "url": "https://github.com/sponsors/ljharb" } }, + "node_modules/is-stream": { + "version": "2.0.1", + "resolved": "https://registry.npmjs.org/is-stream/-/is-stream-2.0.1.tgz", + "integrity": "sha512-hFoiJiTl63nn+kstHGBtewWSKnQLpyb155KHheA1l39uvtO9nWIop1p3udqPcUd/xbF1VLMO4n7OI6p7RbngDg==", + "engines": { + "node": ">=8" + }, + "funding": { + "url": "https://github.com/sponsors/sindresorhus" + } + }, "node_modules/is-typed-array": { "version": "1.1.13", "resolved": "https://registry.npmjs.org/is-typed-array/-/is-typed-array-1.1.13.tgz", @@ -2114,6 +2240,11 @@ "node": ">= 0.6.0" } }, + "node_modules/kuler": { + "version": "2.0.0", + "resolved": "https://registry.npmjs.org/kuler/-/kuler-2.0.0.tgz", + "integrity": "sha512-Xq9nH7KlWZmXAtodXDDRE7vs6DU1gTU8zYDHDiWLSip45Egwq3plLHzPn27NgvzL2r1LMPC1vdqh98sQxtqj4A==" + }, "node_modules/leac": { "version": "0.6.0", "resolved": "https://registry.npmjs.org/leac/-/leac-0.6.0.tgz", @@ -2151,6 +2282,22 @@ "uc.micro": "^2.0.0" } }, + "node_modules/logform": { + "version": "2.7.0", + "resolved": "https://registry.npmjs.org/logform/-/logform-2.7.0.tgz", + "integrity": "sha512-TFYA4jnP7PVbmlBIfhlSe+WKxs9dklXMTEGcBCIvLhE/Tn3H6Gk1norupVW7m5Cnd4bLcr08AytbyV/xj7f/kQ==", + "dependencies": { + "@colors/colors": "1.6.0", + "@types/triple-beam": "^1.3.2", + "fecha": "^4.2.0", + "ms": "^2.1.1", + "safe-stable-stringify": "^2.3.1", + "triple-beam": "^1.3.0" + }, + "engines": { + "node": ">= 12.0.0" + } + }, "node_modules/mailparser": { "version": "3.7.1", "resolved": "https://registry.npmjs.org/mailparser/-/mailparser-3.7.1.tgz", @@ -2226,6 +2373,24 @@ "node": ">= 0.6" } }, + "node_modules/moment": { + "version": "2.30.1", + "resolved": "https://registry.npmjs.org/moment/-/moment-2.30.1.tgz", + "integrity": "sha512-uEmtNhbDOrWPFS+hdjFCBfy9f2YoyzRpwcl+DqpC6taX21FzsTLQVbMV/W7PzNSX6x/bhC1zA3c2UQ5NzH6how==", + "engines": { + "node": "*" + } + }, + "node_modules/ms": { + "version": "2.1.3", + "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", + "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==" + }, + "node_modules/node-eta": { + "version": "0.9.0", + "resolved": "https://registry.npmjs.org/node-eta/-/node-eta-0.9.0.tgz", + "integrity": "sha512-mTCTZk29tmX1OGfVkPt63H3c3VqXrI2Kvua98S7iUIB/Gbp0MNw05YtUomxQIxnnKMyRIIuY9izPcFixzhSBrA==" + }, "node_modules/nodemailer": { "version": "6.9.13", "resolved": "https://registry.npmjs.org/nodemailer/-/nodemailer-6.9.13.tgz", @@ -2234,6 +2399,22 @@ "node": ">=6.0.0" } }, + "node_modules/object-hash": { + "version": "2.2.0", + "resolved": "https://registry.npmjs.org/object-hash/-/object-hash-2.2.0.tgz", + "integrity": "sha512-gScRMn0bS5fH+IuwyIFgnh9zBdo4DV+6GhygmWM9HyNJSgS0hScp1f5vjtm7oIIOiT9trXrShAkLFSc2IqKNgw==", + "engines": { + "node": ">= 6" + } + }, + "node_modules/one-time": { + "version": "1.0.0", + "resolved": "https://registry.npmjs.org/one-time/-/one-time-1.0.0.tgz", + "integrity": "sha512-5DXOiRKwuSEcQ/l0kGCF6Q3jcADFv5tSmRaJck/OqkVFcOzutB134KRSfF0xDrL39MNnqxbHBbUUcjZIhTgb2g==", + "dependencies": { + "fn.name": "1.x.x" + } + }, "node_modules/parseley": { "version": "0.12.1", "resolved": "https://registry.npmjs.org/parseley/-/parseley-0.12.1.tgz", @@ -2321,6 +2502,14 @@ } ] }, + "node_modules/safe-stable-stringify": { + "version": "2.5.0", + "resolved": "https://registry.npmjs.org/safe-stable-stringify/-/safe-stable-stringify-2.5.0.tgz", + "integrity": "sha512-b3rppTKm9T+PsVCBEOUR46GWI7fdOs00VKZ1+9c1EWDaDMvjQc6tUwuFyIprgGgTcWoVHSKrU8H31ZHA2e0RHA==", + "engines": { + "node": ">=10" + } + }, "node_modules/safer-buffer": { "version": "2.1.2", "resolved": "https://registry.npmjs.org/safer-buffer/-/safer-buffer-2.1.2.tgz", @@ -2358,6 +2547,14 @@ "node": ">= 0.4" } }, + "node_modules/simple-swizzle": { + "version": "0.2.2", + "resolved": "https://registry.npmjs.org/simple-swizzle/-/simple-swizzle-0.2.2.tgz", + "integrity": "sha512-JA//kQgZtbuY83m+xT+tXJkmJncGMTFT+C+g2h2R9uxkYIrE2yy9sgmcLhCnw57/WSD+Eh3J97FPEDFnbXnDUg==", + "dependencies": { + "is-arrayish": "^0.3.1" + } + }, "node_modules/smtp-server": { "version": "3.13.5", "resolved": "https://registry.npmjs.org/smtp-server/-/smtp-server-3.13.5.tgz", @@ -2380,6 +2577,14 @@ "node": ">=6.0.0" } }, + "node_modules/stack-trace": { + "version": "0.0.10", + "resolved": "https://registry.npmjs.org/stack-trace/-/stack-trace-0.0.10.tgz", + "integrity": "sha512-KGzahc7puUKkzyMt+IqAep+TVNbKP+k2Lmwhub39m1AsTSkaDutx56aDCo+HLDzf/D26BIHTJWNiTG1KAJiQCg==", + "engines": { + "node": "*" + } + }, "node_modules/stream-browserify": { "version": "3.0.0", "resolved": "https://registry.npmjs.org/stream-browserify/-/stream-browserify-3.0.0.tgz", @@ -2402,6 +2607,11 @@ "resolved": "https://registry.npmjs.org/strnum/-/strnum-1.0.5.tgz", "integrity": "sha512-J8bbNyKKXl5qYcR36TIO8W3mVGVHrmmxsd5PAItGkmyzwJvybiw2IVq5nqd0i4LSNSkB/sx9VHllbfFdr9k1JA==" }, + "node_modules/text-hex": { + "version": "1.0.0", + "resolved": "https://registry.npmjs.org/text-hex/-/text-hex-1.0.0.tgz", + "integrity": "sha512-uuVGNWzgJ4yhRaNSiubPY7OjISw4sw4E5Uv0wbjp+OzcbmVU/rsT8ujgcXJhn9ypzsgr5vlzpPqP+MBBKcGvbg==" + }, "node_modules/tlds": { "version": "1.252.0", "resolved": "https://registry.npmjs.org/tlds/-/tlds-1.252.0.tgz", @@ -2410,6 +2620,14 @@ "tlds": "bin.js" } }, + "node_modules/triple-beam": { + "version": "1.4.1", + "resolved": "https://registry.npmjs.org/triple-beam/-/triple-beam-1.4.1.tgz", + "integrity": "sha512-aZbgViZrg1QNcG+LULa7nhZpJTZSLm/mXnHXnbAbjmN5aSa0y7V+wvv6+4WaBtpISJzThKy+PIPxc1Nq1EJ9mg==", + "engines": { + "node": ">= 14.0.0" + } + }, "node_modules/tslib": { "version": "2.7.0", "resolved": "https://registry.npmjs.org/tslib/-/tslib-2.7.0.tgz", @@ -2476,6 +2694,57 @@ "url": "https://github.com/sponsors/ljharb" } }, + "node_modules/winston": { + "version": "3.17.0", + "resolved": "https://registry.npmjs.org/winston/-/winston-3.17.0.tgz", + "integrity": "sha512-DLiFIXYC5fMPxaRg832S6F5mJYvePtmO5G9v9IgUFPhXm9/GkXarH/TUrBAVzhTCzAj9anE/+GjrgXp/54nOgw==", + "dependencies": { + "@colors/colors": "^1.6.0", + "@dabh/diagnostics": "^2.0.2", + "async": "^3.2.3", + "is-stream": "^2.0.0", + "logform": "^2.7.0", + "one-time": "^1.0.0", + "readable-stream": "^3.4.0", + "safe-stable-stringify": "^2.3.1", + "stack-trace": "0.0.x", + "triple-beam": "^1.3.0", + "winston-transport": "^4.9.0" + }, + "engines": { + "node": ">= 12.0.0" + } + }, + "node_modules/winston-daily-rotate-file": { + "version": "4.7.1", + "resolved": "https://registry.npmjs.org/winston-daily-rotate-file/-/winston-daily-rotate-file-4.7.1.tgz", + "integrity": "sha512-7LGPiYGBPNyGHLn9z33i96zx/bd71pjBn9tqQzO3I4Tayv94WPmBNwKC7CO1wPHdP9uvu+Md/1nr6VSH9h0iaA==", + "dependencies": { + "file-stream-rotator": "^0.6.1", + "object-hash": "^2.0.1", + "triple-beam": "^1.3.0", + "winston-transport": "^4.4.0" + }, + "engines": { + "node": ">=8" + }, + "peerDependencies": { + "winston": "^3" + } + }, + "node_modules/winston-transport": { + "version": "4.9.0", + "resolved": "https://registry.npmjs.org/winston-transport/-/winston-transport-4.9.0.tgz", + "integrity": "sha512-8drMJ4rkgaPo1Me4zD/3WLfI/zPdA9o2IipKODunnGDcuqbHwjsbB79ylv04LCGGzU0xQ6vTznOMpQGaLhhm6A==", + "dependencies": { + "logform": "^2.7.0", + "readable-stream": "^3.6.2", + "triple-beam": "^1.3.0" + }, + "engines": { + "node": ">= 12.0.0" + } + }, "node_modules/xml2js": { "version": "0.6.2", "resolved": "https://registry.npmjs.org/xml2js/-/xml2js-0.6.2.tgz", diff --git a/server.js b/server.js index f1555a0..cb10860 100644 --- a/server.js +++ b/server.js @@ -1,3 +1,6 @@ +const crypto = require('crypto'); +const fs = require('fs'); +const path = require('path'); const SMTPServer = require('smtp-server').SMTPServer; const config = require('./config'); const { parseEmail } = require('./services/emailParser'); @@ -14,9 +17,10 @@ const logger = winston.createLogger({ winston.format.json() ), transports: [ - new winston.transports.Console(), + // Supervisor captures stdout/stderr. A regular file avoids EPIPE loops + // when Supervisor rotates its child logs. new winston.transports.DailyRotateFile({ - filename: 'logs/application-%DATE%.log', + filename: path.join(config.LOG_DIR, 'application-%DATE%.log'), datePattern: 'YYYY-MM-DD', zippedArchive: true, maxSize: '20m', @@ -26,102 +30,343 @@ const logger = winston.createLogger({ }); function validateConfig() { - const requiredKeys = ['PORT', 'SMTP_SECURE', 'WEBHOOK_URL', 'WEBHOOK_CONCURRENCY']; - for (const key of requiredKeys) { - if (!(key in config)) { - throw new Error(`Missing required configuration: ${key}`); - } + if (!config.WEBHOOK_URL) { + throw new Error('Missing required configuration: WEBHOOK_URL'); + } + if (!/^https?:\/\//i.test(config.WEBHOOK_URL)) { + throw new Error('WEBHOOK_URL must use http:// or https://'); + } + if (!Number.isInteger(config.PORT) || config.PORT < 1 || config.PORT > 65535) { + throw new Error('PORT must be an integer between 1 and 65535'); + } + if (!Number.isInteger(config.MAX_FILE_SIZE) || config.MAX_FILE_SIZE < 0) { + throw new Error('MAX_FILE_SIZE must be a non-negative integer'); + } + if (!Number.isInteger(config.MAX_MESSAGE_SIZE) || config.MAX_MESSAGE_SIZE < 0) { + throw new Error('MAX_MESSAGE_SIZE must be a non-negative integer'); + } + if (!Number.isInteger(config.MAX_SMTP_CLIENTS) || config.MAX_SMTP_CLIENTS < 1) { + throw new Error('MAX_SMTP_CLIENTS must be a positive integer'); + } + if (!Number.isInteger(config.WEBHOOK_CONCURRENCY) || config.WEBHOOK_CONCURRENCY < 1) { + throw new Error('WEBHOOK_CONCURRENCY must be a positive integer'); + } + if (!Number.isInteger(config.WEBHOOK_QUEUE_RETRIES) || config.WEBHOOK_QUEUE_RETRIES < 0) { + throw new Error('WEBHOOK_QUEUE_RETRIES must be a non-negative integer'); + } + if (!Number.isInteger(config.WEBHOOK_QUEUE_RETRY_DELAY) || config.WEBHOOK_QUEUE_RETRY_DELAY < 1000) { + throw new Error('WEBHOOK_QUEUE_RETRY_DELAY must be at least 1000 milliseconds'); + } + if (!Number.isInteger(config.WEBHOOK_QUEUE_TIMEOUT) || config.WEBHOOK_QUEUE_TIMEOUT < 1000) { + throw new Error('WEBHOOK_QUEUE_TIMEOUT must be at least 1000 milliseconds'); + } + if (config.SMTP_SECURE && (!config.TLS_KEY_PATH || !config.TLS_CERT_PATH)) { + throw new Error('SMTP_SECURE=true requires TLS_KEY_PATH and TLS_CERT_PATH'); } } -const webhookQueue = new Queue(async function (parsed, cb) { - const maxRetries = 3; - let retries = 0; - - const attemptWebhook = async () => { - try { - await sendToWebhook(parsed); - logger.info('Successfully sent to webhook'); - cb(null); - } catch (error) { - logger.error('Webhook error:', { message: error.message, stack: error.stack }); - if (error.response) { - logger.error('Webhook response error:', { - status: error.response.status, - data: error.response.data - }); - } - if (retries < maxRetries) { - retries++; - logger.info(`Retrying webhook (attempt ${retries}/${maxRetries})`); - setTimeout(attemptWebhook, 1000 * retries); - } else { - cb(error); - } - } - }; - - attemptWebhook(); -}, { concurrent: config.WEBHOOK_CONCURRENCY || 5 }); - -const server = new SMTPServer({ - onData(stream, session, callback) { - parseEmail(stream) - .then(parsed => { - webhookQueue.push(parsed); - logger.info('Email added to queue', { queueSize: webhookQueue.getStats().total }); - callback(); - }) - .catch(error => { - logger.error('Parsing error:', { message: error.message, stack: error.stack }); - callback(new Error('Failed to parse email')); - }); - }, - onError(error) { - logger.error('SMTP server error:', { message: error.message, stack: error.stack }); - }, - disabledCommands: ['AUTH'], - secure: config.SMTP_SECURE -}); - -server.listen(config.PORT, '0.0.0.0', err => { - if (err) { - logger.error('Failed to start SMTP server:', { message: err.message, stack: err.stack }); - process.exit(1); - } - logger.info(`SMTP server listening on port ${config.PORT} on all interfaces`); -}); - -function gracefulShutdown(reason) { - logger.info(`Shutting down: ${reason}`); - server.close(() => { - logger.info('Server closed. Exiting process.'); - process.exit(0); - }); -} - -process.on('uncaughtException', (err) => { - logger.error('Uncaught exception:', { message: err.message, stack: err.stack }); - gracefulShutdown('Uncaught exception'); -}); - -process.on('unhandledRejection', (reason, promise) => { - logger.error('Unhandled Rejection:', { reason: reason, promise: promise }); - gracefulShutdown('Unhandled rejection'); -}); - -process.on('SIGTERM', () => { - gracefulShutdown('SIGTERM signal received'); -}); - -process.on('SIGINT', () => { - gracefulShutdown('SIGINT signal received'); -}); - -// Add configuration validation at startup try { validateConfig(); } catch (error) { logger.error('Configuration error:', { message: error.message }); process.exit(1); -} \ No newline at end of file +} + +const pendingDir = path.join(config.SPOOL_DIR, 'pending'); +const failedDir = path.join(config.SPOOL_DIR, 'failed'); +let shuttingDown = false; +let spoolScanTimer = null; +let activeJobs = 0; +const activeJobWaiters = new Set(); +const enqueuedJobIds = new Set(); + +try { + fs.mkdirSync(pendingDir, { recursive: true, mode: 0o700 }); + fs.mkdirSync(failedDir, { recursive: true, mode: 0o700 }); +} catch (error) { + logger.error('Failed to initialize spool directories:', { + message: error.message, + stack: error.stack + }); + process.exit(1); +} + +function notifyActiveJobWaiters() { + if (activeJobs !== 0) return; + for (const waiter of activeJobWaiters) waiter(); + activeJobWaiters.clear(); +} + +function waitForActiveJobs(timeout) { + if (activeJobs === 0) return Promise.resolve(); + + return new Promise(resolve => { + let timer; + const onDone = () => { + clearTimeout(timer); + activeJobWaiters.delete(onDone); + resolve(); + }; + + timer = setTimeout(onDone, timeout); + activeJobWaiters.add(onDone); + }); +} + +function delay(milliseconds) { + return new Promise(resolve => setTimeout(resolve, milliseconds)); +} + +function logWebhookError(error) { + const message = error instanceof Error ? error.message : String(error); + const stack = error instanceof Error ? error.stack : undefined; + logger.error('Webhook error:', { message, stack }); + + if (error && error.response) { + let data = error.response.data; + if (typeof data === 'string') data = data.slice(0, 500); + logger.error('Webhook response error:', { + status: error.response.status, + data + }); + } +} + +async function sendWithRetry(parsed) { + const maxRetries = 3; + let retries = 0; + + while (true) { + try { + await sendToWebhook(parsed); + return; + } catch (error) { + logWebhookError(error); + if (retries >= maxRetries) throw error; + + retries++; + logger.info(`Retrying webhook (attempt ${retries}/${maxRetries})`); + await delay(1000 * retries); + } + } +} + +async function persistEmail(parsed) { + const id = `${Date.now()}-${process.pid}-${crypto.randomUUID()}`; + const filePath = path.join(pendingDir, `${id}.json`); + const temporaryPath = `${filePath}.tmp`; + + let fileHandle; + try { + fileHandle = await fs.promises.open(temporaryPath, 'w', 0o600); + await fileHandle.writeFile(JSON.stringify(parsed), 'utf8'); + await fileHandle.sync(); + await fileHandle.close(); + fileHandle = undefined; + await fs.promises.rename(temporaryPath, filePath); + return { id, filePath }; + } catch (error) { + if (fileHandle) { + await fileHandle.close().catch(() => {}); + } + await fs.promises.unlink(temporaryPath).catch(() => {}); + throw error; + } +} + +async function quarantineFile(filePath) { + const target = path.join( + failedDir, + `${path.basename(filePath)}.${Date.now()}.failed` + ); + await fs.promises.rename(filePath, target); + return target; +} + +async function processWebhookJob(job, cb) { + activeJobs++; + + try { + let parsed; + try { + parsed = JSON.parse(await fs.promises.readFile(job.filePath, 'utf8')); + } catch (error) { + if (error.code === 'ENOENT') { + logger.warn('Spool job disappeared before processing', { jobId: job.id }); + cb(null); + return; + } + + try { + const failedPath = await quarantineFile(job.filePath); + logger.error('Invalid spool job moved to failed spool', { + jobId: job.id, + failedPath, + message: error.message + }); + cb(null); + } catch (quarantineError) { + cb(quarantineError); + } + return; + } + + await sendWithRetry(parsed); + await fs.promises.unlink(job.filePath); + logger.info('Successfully sent to webhook', { jobId: job.id }); + cb(null); + } catch (error) { + logWebhookError(error); + cb(error); + } finally { + activeJobs--; + notifyActiveJobWaiters(); + } +} + +const webhookQueue = new Queue(processWebhookJob, { + concurrent: config.WEBHOOK_CONCURRENCY, + maxRetries: config.WEBHOOK_QUEUE_RETRIES, + retryDelay: config.WEBHOOK_QUEUE_RETRY_DELAY, + maxTimeout: config.WEBHOOK_QUEUE_TIMEOUT +}); + +webhookQueue.on('error', error => { + logger.error('Webhook queue error:', { message: error.message, stack: error.stack }); +}); + +webhookQueue.on('task_failed', (jobId, error) => { + logger.error('Webhook queue task exhausted retries:', { + jobId, + message: error && error.message ? error.message : String(error) + }); +}); + +function enqueueSpoolFile(filePath) { + const id = path.basename(filePath, '.json'); + if (enqueuedJobIds.has(id)) return; + + enqueuedJobIds.add(id); + const ticket = webhookQueue.push({ id, filePath }); + ticket.once('finish', () => enqueuedJobIds.delete(id)); + ticket.once('failed', () => enqueuedJobIds.delete(id)); +} + +async function scanSpool() { + if (shuttingDown) return; + + try { + const entries = await fs.promises.readdir(pendingDir, { withFileTypes: true }); + for (const entry of entries) { + if (entry.isFile() && entry.name.endsWith('.json')) { + enqueueSpoolFile(path.join(pendingDir, entry.name)); + } + } + } catch (error) { + logger.error('Failed to scan spool:', { message: error.message, stack: error.stack }); + } +} + +const smtpOptions = { + maxClients: config.MAX_SMTP_CLIENTS, + socketTimeout: config.SMTP_SOCKET_TIMEOUT, + closeTimeout: config.SMTP_CLOSE_TIMEOUT, + size: config.MAX_MESSAGE_SIZE || undefined, + onConnect(session, callback) { + if ( + config.SMTP_ALLOWED_IPS.length > 0 && + !config.SMTP_ALLOWED_IPS.includes(session.remoteAddress) + ) { + const error = new Error('Connection not allowed'); + error.responseCode = 421; + callback(error); + return; + } + callback(); + }, + onData(stream, session, callback) { + parseEmail(stream) + .then(async parsed => { + const job = await persistEmail(parsed); + enqueueSpoolFile(job.filePath); + logger.info('Email persisted and added to queue', { + jobId: job.id, + queued: webhookQueue.length + }); + callback(); + }) + .catch(error => { + logger.error('Parsing or persistence error:', { + message: error.message, + stack: error.stack + }); + callback(new Error('Failed to persist email')); + }); + }, + disabledCommands: ['AUTH'], + secure: config.SMTP_SECURE +}; + +if (config.SMTP_SECURE) { + smtpOptions.key = fs.readFileSync(config.TLS_KEY_PATH); + smtpOptions.cert = fs.readFileSync(config.TLS_CERT_PATH); +} + +const server = new SMTPServer(smtpOptions); + +let serverReady = false; +server.on('error', error => { + logger.error('SMTP server error:', { message: error.message, stack: error.stack }); + if (!serverReady) { + // Let the file transport flush the startup error before exiting. Since the + // listener never came up, there are no active handles keeping the process + // alive; the non-zero exit code still lets Supervisor restart it. + process.exitCode = 1; + } +}); + +function gracefulShutdown(reason) { + if (shuttingDown) return; + shuttingDown = true; + logger.info(`Shutting down: ${reason}`); + clearInterval(spoolScanTimer); + webhookQueue.pause(); + + const cleanShutdown = reason === 'SIGTERM signal received' || reason === 'SIGINT signal received'; + const exitCode = cleanShutdown ? 0 : 1; + const shutdownDeadline = setTimeout(() => { + logger.error('Shutdown timeout reached; exiting with pending spool jobs preserved'); + process.exit(exitCode); + }, 15000); + + server.close(async () => { + await waitForActiveJobs(10000); + clearTimeout(shutdownDeadline); + logger.info('Server closed. Exiting process.'); + process.exit(exitCode); + }); +} + +process.on('uncaughtException', error => { + logger.error('Uncaught exception:', { message: error.message, stack: error.stack }); + gracefulShutdown('Uncaught exception'); +}); + +process.on('unhandledRejection', reason => { + const error = reason instanceof Error + ? { message: reason.message, stack: reason.stack } + : { message: String(reason) }; + logger.error('Unhandled Rejection:', error); + gracefulShutdown('Unhandled rejection'); +}); + +process.on('SIGTERM', () => gracefulShutdown('SIGTERM signal received')); +process.on('SIGINT', () => gracefulShutdown('SIGINT signal received')); + +server.listen(config.PORT, '0.0.0.0', () => { + serverReady = true; + logger.info(`SMTP server listening on port ${config.PORT} on all interfaces`); + void scanSpool(); + spoolScanTimer = setInterval(() => void scanSpool(), 60 * 1000); + spoolScanTimer.unref(); +});