const fs = require('fs'); const path = require('path'); const { pipeline } = require('stream/promises'); const { S3Client, PutObjectCommand, GetObjectCommand, HeadObjectCommand, DeleteObjectCommand, CopyObjectCommand, ListObjectsV2Command, HeadBucketCommand, CreateBucketCommand, } = require('@aws-sdk/client-s3'); const MIME_TYPES = { '.jpg': 'image/jpeg', '.jpeg': 'image/jpeg', '.jfif': 'image/jpeg', '.png': 'image/png', '.gif': 'image/gif', '.webp': 'image/webp', '.bmp': 'image/bmp', '.avif': 'image/avif', '.ico': 'image/x-icon', '.heic': 'image/heic', '.heif': 'image/heif', '.pdf': 'application/pdf', '.txt': 'text/plain; charset=utf-8', '.md': 'text/markdown; charset=utf-8', '.html': 'text/html; charset=utf-8', '.htm': 'text/html; charset=utf-8', '.zip': 'application/zip', '.rar': 'application/vnd.rar', '.7z': 'application/x-7z-compressed', '.doc': 'application/msword', '.docx': 'application/vnd.openxmlformats-officedocument.wordprocessingml.document', '.mp4': 'video/mp4', '.m4v': 'video/mp4', '.mov': 'video/quicktime', '.webm': 'video/webm', '.ogv': 'video/ogg', '.mpeg': 'video/mpeg', '.mpg': 'video/mpeg', '.mkv': 'video/x-matroska', '.avi': 'video/x-msvideo', '.3gp': 'video/3gpp', '.m3u8': 'application/vnd.apple.mpegurl', '.mp3': 'audio/mpeg', '.m4a': 'audio/mp4', '.oga': 'audio/ogg', '.wav': 'audio/wav', '.flac': 'audio/flac', }; const SAFE_SEGMENT = /^[\w,.()-]+$/; const NON_OBJECT_PREFIXES = ['.thumbs/', '.cache/']; function mimeFor(name) { const m = String(name || '').toLowerCase().match(/\.[a-z0-9]+$/); return (m && MIME_TYPES[m[0]]) || 'application/octet-stream'; } function normalizeKey(input) { if (typeof input !== 'string' || !input) return null; let key = input.replace(/\\/g, '/').replace(/^\/+/, ''); if (key.startsWith('uploads/')) key = key.slice('uploads/'.length); if (!key || key.length > 255 || key.includes('..')) return null; const segments = key.split('/'); if (segments.some(seg => !SAFE_SEGMENT.test(seg))) return null; return key; } function isObjectKey(key) { return !NON_OBJECT_PREFIXES.some(p => key.startsWith(p)); } function createStorage(options = {}) { const driver = options.driver || process.env.STORAGE_DRIVER || 'local'; const dir = path.resolve(options.dir || path.join(__dirname, 'uploads')); const cacheDir = path.resolve(options.cacheDir || path.join(dir, '.cache')); const localFallback = options.localFallback !== undefined ? !!options.localFallback : process.env.STORAGE_LOCAL_FALLBACK !== '0'; const keepLocal = options.keepLocal !== undefined ? !!options.keepLocal : process.env.STORAGE_KEEP_LOCAL === '1'; const bucket = options.bucket || process.env.S3_BUCKET || ''; const prefix = String(options.prefix !== undefined ? options.prefix : process.env.S3_PREFIX || '').replace(/^\/+|\/+$/g, ''); const cacheMaxAgeMs = Math.max(0, parseInt(options.cacheMaxAgeHours || process.env.STORAGE_CACHE_MAX_AGE_HOURS || '168', 10) || 0) * 3600 * 1000; const remote = driver === 's3'; const client = remote ? new S3Client({ region: options.region || process.env.S3_REGION || 'us-east-1', endpoint: options.endpoint || process.env.S3_ENDPOINT || undefined, forcePathStyle: options.forcePathStyle !== undefined ? !!options.forcePathStyle : process.env.S3_FORCE_PATH_STYLE !== '0', credentials: { accessKeyId: options.accessKey || process.env.S3_ACCESS_KEY || '', secretAccessKey: options.secretKey || process.env.S3_SECRET_KEY || '', }, }) : null; function isRemote() { return remote; } function objectKey(key) { const k = normalizeKey(key); if (!k) return null; return prefix ? `${prefix}/${k}` : k; } function localFile(key) { const k = normalizeKey(key); if (!k) return null; const abs = path.resolve(dir, ...k.split('/')); if (abs !== dir && !abs.startsWith(dir + path.sep)) return null; return abs; } function cacheFile(key) { const k = normalizeKey(key); if (!k) return null; return path.join(cacheDir, ...k.split('/')); } function keyFromPath(p) { return normalizeKey(p); } function localExists(key) { const fp = localFile(key); if (!fp) return false; try { return fs.statSync(fp).isFile(); } catch { return false; } } function unlinkLocalOnly(key) { for (const fp of [localFile(key), cacheFile(key)]) { if (!fp) continue; try { if (fs.existsSync(fp)) fs.unlinkSync(fp); } catch {} } } function isMissingError(err) { if (!err) return false; const code = err.name || err.Code || err.code; const status = err.$metadata && err.$metadata.httpStatusCode; return code === 'NoSuchKey' || code === 'NotFound' || code === 'ENOENT' || status === 404; } async function put(key, body, opts = {}) { const objKey = objectKey(key); if (!objKey) return false; if (!remote) { const fp = localFile(key); if (!fp) return false; fs.mkdirSync(path.dirname(fp), { recursive: true }); if (Buffer.isBuffer(body) || typeof body === 'string') fs.writeFileSync(fp, body); else await pipeline(body, fs.createWriteStream(fp)); return true; } await client.send(new PutObjectCommand({ Bucket: bucket, Key: objKey, Body: body, ContentType: opts.contentType || mimeFor(key), ContentLength: opts.contentLength, })); return true; } async function putFile(key, localPath, opts = {}) { const objKey = objectKey(key); if (!objKey) return false; let size = null; try { size = fs.statSync(localPath).size; } catch { return false; } if (!remote) { const fp = localFile(key); if (!fp) return false; fs.mkdirSync(path.dirname(fp), { recursive: true }); if (path.resolve(localPath) !== fp) fs.copyFileSync(localPath, fp); return true; } await client.send(new PutObjectCommand({ Bucket: bucket, Key: objKey, Body: fs.createReadStream(localPath), ContentLength: size, ContentType: opts.contentType || mimeFor(key), })); return true; } async function head(key) { const objKey = objectKey(key); if (!objKey) return null; if (!remote) { const fp = localFile(key); if (!fp) return null; try { const st = fs.statSync(fp); if (!st.isFile()) return null; return { key: normalizeKey(key), size: st.size, contentType: mimeFor(key), etag: null }; } catch { return null; } } try { const out = await client.send(new HeadObjectCommand({ Bucket: bucket, Key: objKey })); return { key: normalizeKey(key), size: out.ContentLength || 0, contentType: out.ContentType || mimeFor(key), etag: out.ETag ? out.ETag.replace(/"/g, '') : null, }; } catch (err) { if (isMissingError(err)) return null; throw err; } } async function exists(key) { if (localFallback && localExists(key)) return true; return !!(await head(key)); } async function sizeOf(key) { const info = await head(key); return info ? info.size : 0; } async function getStream(key) { const objKey = objectKey(key); if (!objKey) return null; if (!remote) { const fp = localFile(key); if (!fp || !fs.existsSync(fp)) return null; return { stream: fs.createReadStream(fp), contentLength: fs.statSync(fp).size, contentType: mimeFor(key) }; } try { const out = await client.send(new GetObjectCommand({ Bucket: bucket, Key: objKey })); return { stream: out.Body, contentLength: out.ContentLength || 0, contentType: out.ContentType || mimeFor(key) }; } catch (err) { if (isMissingError(err)) return null; throw err; } } async function getBuffer(key) { const fp = await localize(key); if (fp) return fs.readFileSync(fp); return null; } async function getRange(key, start, end) { const k = normalizeKey(key); if (!k || !Number.isInteger(start) || start < 0) return null; const last = Number.isInteger(end) && end >= start ? end : null; if (localFallback && localExists(k)) { const fp = localFile(k); const size = fs.statSync(fp).size; const to = last === null ? size - 1 : Math.min(last, size - 1); if (start > to) return null; return { stream: fs.createReadStream(fp, { start, end: to }), contentLength: to - start + 1, contentType: mimeFor(k) }; } const objKey = objectKey(k); if (!objKey) return null; try { const out = await client.send(new GetObjectCommand({ Bucket: bucket, Key: objKey, Range: `bytes=${start}-${last === null ? '' : last}`, })); return { stream: out.Body, contentLength: out.ContentLength || 0, contentType: out.ContentType || mimeFor(k) }; } catch (err) { if (isMissingError(err)) return null; if (err && (err.name === 'InvalidRange' || (err.$metadata && err.$metadata.httpStatusCode === 416))) return null; throw err; } } async function del(key) { const k = normalizeKey(key); if (!k) return false; unlinkLocalOnly(k); if (!remote) return true; const objKey = objectKey(k); try { await client.send(new DeleteObjectCommand({ Bucket: bucket, Key: objKey })); } catch (err) { if (!isMissingError(err)) throw err; } return true; } async function copyObject(srcKey, dstKey) { const src = normalizeKey(srcKey); const dst = normalizeKey(dstKey); if (!src || !dst) return false; if (!remote) { const from = localFile(src); const to = localFile(dst); if (!from || !to || !fs.existsSync(from)) return false; fs.mkdirSync(path.dirname(to), { recursive: true }); fs.copyFileSync(from, to); return true; } try { await client.send(new CopyObjectCommand({ Bucket: bucket, Key: objectKey(dst), CopySource: `${bucket}/${objectKey(src)}`, ContentType: mimeFor(dst), MetadataDirective: 'REPLACE', })); return true; } catch (err) { if (isMissingError(err)) return false; throw err; } } async function listAll(listPrefix = '') { const out = []; if (!remote) { const base = listPrefix ? localFile(listPrefix) : dir; if (!base || !fs.existsSync(base)) return out; const walk = (abs, rel) => { for (const name of fs.readdirSync(abs)) { const childAbs = path.join(abs, name); const childRel = rel ? `${rel}/${name}` : name; let st; try { st = fs.statSync(childAbs); } catch { continue; } if (st.isDirectory()) walk(childAbs, childRel); else out.push({ key: childRel, size: st.size }); } }; walk(base, listPrefix ? normalizeKey(listPrefix) || '' : ''); return out; } let token = null; const fullPrefix = prefix ? `${prefix}/${listPrefix}` : listPrefix; do { const res = await client.send(new ListObjectsV2Command({ Bucket: bucket, Prefix: fullPrefix || undefined, ContinuationToken: token || undefined, })); for (const item of res.Contents || []) { const key = prefix ? String(item.Key).slice(prefix.length + 1) : String(item.Key); out.push({ key, size: item.Size || 0 }); } token = res.IsTruncated ? res.NextContinuationToken : null; } while (token); return out; } async function localize(key, opts = {}) { const k = normalizeKey(key); if (!k) return null; if (localExists(k)) return localFile(k); if (!remote) return null; const cf = cacheFile(k); if (cf && fs.existsSync(cf)) { const info = await head(k); if (info && info.size === fs.statSync(cf).size) return cf; try { fs.unlinkSync(cf); } catch {} } const source = await getStream(k); if (!source) return null; fs.mkdirSync(path.dirname(cf), { recursive: true }); const tmp = `${cf}.${process.pid}.${Date.now()}.tmp`; await pipeline(source.stream, fs.createWriteStream(tmp)); if (opts.useCache === false) return tmp; fs.renameSync(tmp, cf); return cf; } async function persist(key, localPath) { const k = normalizeKey(key); if (!k) return false; const src = path.resolve(localPath); if (!fs.existsSync(src)) return false; if (!remote) return true; const info = await head(k); if (info && info.size === fs.statSync(src).size) { if (!keepLocal) unlinkLocalOnly(k); return true; } await putFile(k, src); if (!keepLocal) { try { fs.unlinkSync(src); } catch {} const cf = cacheFile(k); if (cf && fs.existsSync(cf)) { try { fs.unlinkSync(cf); } catch {} } } return true; } async function streamTo(res, key, opts = {}) { const k = normalizeKey(key); if (!k) return false; if (localFallback && localExists(k)) { const fp = localFile(k); if (opts.cacheControl) res.setHeader('Cache-Control', opts.cacheControl); if (opts.download) { res.download(fp, opts.name || path.basename(fp), opts.callback || (() => {})); return true; } res.setHeader('Content-Type', opts.contentType || mimeFor(k)); res.sendFile(fp); return true; } const source = await getStream(k); if (!source) return false; if (opts.cacheControl) res.setHeader('Cache-Control', opts.cacheControl); res.setHeader('Content-Type', opts.contentType || source.contentType || mimeFor(k)); if (source.contentLength) res.setHeader('Content-Length', String(source.contentLength)); if (opts.download) { const name = opts.name || path.basename(k); const ascii = name.replace(/[^\x20-\x7E]/g, '_').replace(/"/g, ''); res.setHeader('Content-Disposition', `attachment; filename="${ascii}"; filename*=UTF-8''${encodeURIComponent(name)}`); } await pipeline(source.stream, res); return true; } async function streamRangeTo(res, key, start, end, opts = {}) { const k = normalizeKey(key); if (!k) return false; const total = await sizeOf(k); if (!total) return false; const source = await getRange(k, start, end); if (!source || !source.contentLength) return false; const lastByte = start + source.contentLength - 1; res.status(206); res.setHeader('Accept-Ranges', 'bytes'); res.setHeader('Content-Range', `bytes ${start}-${lastByte}/${total}`); res.setHeader('Content-Type', opts.contentType || source.contentType || mimeFor(k)); res.setHeader('Content-Length', String(source.contentLength)); if (opts.cacheControl) res.setHeader('Cache-Control', opts.cacheControl); await pipeline(source.stream, res); return true; } async function downloadAll(destDir, opts = {}) { fs.mkdirSync(destDir, { recursive: true }); const root = path.resolve(destDir); const objects = await listAll(''); let count = 0; for (const obj of objects) { if (!isObjectKey(obj.key)) continue; const k = normalizeKey(obj.key); if (!k) continue; const target = path.resolve(root, ...k.split('/')); if (!target.startsWith(root + path.sep)) continue; const source = await getStream(k); if (!source) continue; fs.mkdirSync(path.dirname(target), { recursive: true }); await pipeline(source.stream, fs.createWriteStream(target)); count++; } if (opts.pruneThumbs !== false) { const thumbs = path.join(root, '.thumbs'); if (fs.existsSync(thumbs)) fs.rmSync(thumbs, { recursive: true, force: true }); } return count; } async function uploadTree(srcDir) { const root = path.resolve(srcDir); if (!fs.existsSync(root)) return 0; let count = 0; const walk = async (abs, rel) => { for (const name of fs.readdirSync(abs)) { const childAbs = path.join(abs, name); const childRel = rel ? `${rel}/${name}` : name; let st; try { st = fs.statSync(childAbs); } catch { continue; } if (st.isDirectory()) { if (name === '.thumbs' || name === '.cache') continue; await walk(childAbs, childRel); continue; } const k = normalizeKey(childRel); if (!k) continue; await putFile(k, childAbs); count++; } }; await walk(root, ''); return count; } async function ensureBucket() { if (!remote || !bucket) return false; try { await client.send(new HeadBucketCommand({ Bucket: bucket })); return true; } catch (err) { const status = err && err.$metadata && err.$metadata.httpStatusCode; const name = err && (err.name || err.Code); if (status && status !== 404 && name !== 'NotFound' && name !== 'NoSuchBucket') throw err; } try { await client.send(new CreateBucketCommand({ Bucket: bucket })); return true; } catch (err) { const name = err && (err.name || err.Code); if (name === 'BucketAlreadyOwnedByYou' || name === 'BucketAlreadyExists') return true; throw err; } } async function usage() { const objects = (await listAll('')).filter(o => isObjectKey(o.key)); return { count: objects.length, size_bytes: objects.reduce((sum, o) => sum + (o.size || 0), 0), bucket: remote ? bucket : null, prefix: prefix || null, }; } function pruneCache(now = Date.now()) { if (!cacheMaxAgeMs || !fs.existsSync(cacheDir)) return 0; let removed = 0; const walk = (abs) => { for (const name of fs.readdirSync(abs)) { const childAbs = path.join(abs, name); let st; try { st = fs.statSync(childAbs); } catch { continue; } if (st.isDirectory()) { walk(childAbs); continue; } if (now - st.mtimeMs > cacheMaxAgeMs) { try { fs.unlinkSync(childAbs); removed++; } catch {} } } }; walk(cacheDir); return removed; } return { bucket, cacheDir, isRemote, keyFromPath, localFile, unlinkLocalOnly, normalizeKey, put, putFile, head, exists, sizeOf, getStream, getBuffer, getRange, del, copyObject, listAll, localize, persist, streamTo, streamRangeTo, downloadAll, uploadTree, ensureBucket, usage, pruneCache, }; } module.exports = { createStorage, mimeFor, normalizeKey };