import * as plugins from '../plugins.js'; import { logger } from '../logger.js'; import type { IBlobStorageManager } from '@push.rocks/smartmta'; const SMARTMTA_BLOB_NAMESPACES = [ '/email/queue/', '/email/attachments/', '/email/messages/', ] as const; const DEFAULT_OBJECT_PREFIX = 'dcrouter/smartmta/'; const MAX_STORAGE_KEY_BYTES = 1024; type TBlobBucket = Pick< plugins.smartbucket.Bucket, 'fastExists' | 'fastGet' | 'fastPut' | 'fastRemove' | 'listAllObjects' >; export interface ISmartMtaBlobStorageConfig { storageDescriptor: plugins.tsclass.storage.IStorageDescriptor; bucketName: string; objectPrefix?: string; createBucketIfMissing?: boolean; } function normalizeObjectPrefix(prefixArg = DEFAULT_OBJECT_PREFIX): string { const prefix = prefixArg.trim().replace(/^\/+|\/+$/g, ''); if (!prefix || prefix.includes('\\') || prefix.includes('\0')) { throw new Error('SmartMTA blob objectPrefix must be a non-empty object-key prefix'); } const segments = prefix.split('/'); if (segments.some((segment) => !segment || segment === '.' || segment === '..')) { throw new Error('SmartMTA blob objectPrefix contains an invalid path segment'); } return `${prefix}/`; } function normalizeBlobKey(keyArg: string, allowNamespaceRoot = false): string { if ( typeof keyArg !== 'string' || !keyArg.startsWith('/') || keyArg.includes('\\') || keyArg.includes('\0') || Buffer.byteLength(keyArg, 'utf8') > MAX_STORAGE_KEY_BYTES ) { throw new Error('Invalid SmartMTA blob storage key'); } const namespace = SMARTMTA_BLOB_NAMESPACES.find((candidate) => keyArg.startsWith(candidate)); if (!namespace) { throw new Error(`Unsupported SmartMTA blob storage namespace: ${keyArg}`); } const remainder = keyArg.slice(namespace.length); if (!allowNamespaceRoot && !remainder) { throw new Error('SmartMTA blob storage key must identify an object'); } if (remainder && remainder.split('/').some((segment) => !segment || segment === '.' || segment === '..')) { throw new Error('SmartMTA blob storage key contains an invalid path segment'); } return keyArg; } /** * SmartMTA blob persistence backed by a SmartBucket bucket. * * The adapter accepts only SmartMTA queue, attachment, and retained-message * namespaces. Every key remains a SmartBucket object key. */ export class SmartMtaBlobStorageManager implements IBlobStorageManager { private readonly objectPrefix: string; public constructor( private readonly bucketRef: TBlobBucket, objectPrefix = DEFAULT_OBJECT_PREFIX, ) { this.objectPrefix = normalizeObjectPrefix(objectPrefix); } public static async create( configArg: ISmartMtaBlobStorageConfig, ): Promise { const bucketName = configArg.bucketName?.trim(); if (!bucketName) { throw new Error('SmartMTA blob storage requires a bucketName'); } const smartBucket = new plugins.smartbucket.SmartBucket(configArg.storageDescriptor); let bucket: plugins.smartbucket.Bucket; if (configArg.createBucketIfMissing) { bucket = await smartBucket.bucketExists(bucketName) ? await smartBucket.getBucketByName(bucketName) : await smartBucket.createBucket(bucketName); } else { bucket = await smartBucket.getBucketByName(bucketName); } return new SmartMtaBlobStorageManager(bucket, configArg.objectPrefix); } public async get(keyArg: string): Promise { const objectKey = this.toObjectKey(normalizeBlobKey(keyArg)); if (!(await this.bucketRef.fastExists({ path: objectKey }))) { return null; } return Buffer.from(await this.bucketRef.fastGet({ path: objectKey })); } public async set(keyArg: string, valueArg: Buffer): Promise { if (!Buffer.isBuffer(valueArg)) { throw new TypeError('SmartMTA blob storage values must be Buffers'); } const objectKey = this.toObjectKey(normalizeBlobKey(keyArg)); await this.bucketRef.fastPut({ path: objectKey, contents: Buffer.from(valueArg), overwrite: true, }); } public async list(prefixArg: string): Promise { const prefix = normalizeBlobKey(prefixArg, true); const objectPrefix = this.toObjectKey(prefix); const keys: string[] = []; for await (const objectKey of this.bucketRef.listAllObjects(objectPrefix)) { // One stray or malformed object under the prefix must not abort the // whole listing (SmartMTA queue recovery lists at startup); skip it. if (!objectKey.startsWith(this.objectPrefix)) { logger.log('warn', `Skipping SmartBucket object outside the configured prefix during SmartMTA blob listing: ${objectKey}`); continue; } const key = `/${objectKey.slice(this.objectPrefix.length)}`; try { keys.push(normalizeBlobKey(key)); } catch (error: unknown) { logger.log('warn', `Skipping unnormalizable SmartBucket object key during SmartMTA blob listing: ${objectKey} (${(error as Error).message})`); } } return keys.sort(); } public async delete(keyArg: string): Promise { const objectKey = this.toObjectKey(normalizeBlobKey(keyArg)); await this.bucketRef.fastRemove({ path: objectKey }); } private toObjectKey(keyArg: string): string { return `${this.objectPrefix}${keyArg.slice(1)}`; } }