Files
picpeak/backend/scripts/migrate-storage.js
T
Paul Nothaft 1b717ce5ed feat: native S3 storage backend (#328) + presigned download follow-up
Lets PicPeak write photos, thumbnails, hero images, watermarks, and
archive zips to any S3-compatible bucket (AWS S3, MinIO, Cloudflare R2,
Backblaze B2, Wasabi, DigitalOcean Spaces) instead of the local
filesystem. Selected via STORAGE_BACKEND=local|s3.

Architecture
- backend/src/services/storage/StorageBackend.js — abstract interface
  (put/get/exists/stat/delete/list/copy/rename/signedUrl/putFromFile/
  getToFile) — typedef-only, documents the contract.
- LocalFsStorage.js — wraps fs with atomic-write-via-tmp-rename, path
  traversal protection, list-as-walker.
- S3StorageBackend.js — thin wrapper around the existing
  S3StorageAdapter (used by backupService) mapping it onto the canonical
  interface; supports optional STORAGE_S3_PREFIX namespace.
- index.js — factory selected by STORAGE_BACKEND with startup ping
  (HEADs sentinel key on S3, fs.stat on local) so misconfig fails fast
  before the first request.

Consumer refactors (~12 services + routes), each parametrized over the
abstraction:
- imageProcessor / videoProcessor — pipe Sharp/ffmpeg output through
  storage.put; expose withLocalCopy() helper for S3-mode regeneration
  paths that need a local file for sharp/ffmpeg.
- archiveService / downloadZipService — finalize zip in tmp dir, then
  storage.putFromFile. Atomic-rename pattern preserved on local; S3
  emulates via copy + delete (worker prunes orphaned .tmp.* on startup).
- photoProcessor / photoReplacementService / adminPhotos upload+delete /
  routes/v1/events.js POST /events/:id/photos / routes/events.js — every
  upload path now goes storage.putFromFile(temp) → unlink temp.
- gallery.js bulk-download (cached + on-the-fly + selected) — managed
  photos via storage.get, external-mode unchanged.
- protectedImages / secureImages / photoResolver — read via
  storage.get; resolvePhotoStorageKey returns the canonical key.
- watermarkService / watermarkGeneratorService — persistent watermarks
  via storage.put.
- fileWatcher — bails out with a clear log warning when STORAGE_BACKEND=s3
  (chokidar can't watch S3); auto-import lands via the S3 prefix walker
  introduced in the follow-up commit.
- expirationChecker — small touch (event.expired webhook fire from #327
  shipping in the next commit).

Migration tooling
- backend/scripts/migrate-storage.js — one-shot --dry-run capable script
  that walks photos.path, thumbnail_path, hero_path, watermark_path and
  events.archive_path/download_zip_path; streams local → S3; sha256
  size-match skip for idempotent re-run; failures CSV.

Presigned-URL "Download All" (#328 follow-up shipped in this commit)
- routes/gallery.js — when STORAGE_BACKEND=s3 + event.allow_presigned_download
  + downloads enabled + watermark NOT enabled, /download-all returns a
  302 redirect to a 5-minute presigned S3 URL. Per-event opt-in surface
  ships in the next commit's UI.

Tests
- backend/__tests__/integration/storageBackend.test.js — parametrized
  contract suite running against BOTH LocalFs AND MinIO (18 tests, both
  backends — 36 cases total).
- backend/__tests__/integration/imageProcessor.storage.test.js — same
  parametrized pattern for the image processor (10 tests × 2 backends).
- backend/__tests__/integration/backup-s3.test.js — bootstrap fix:
  drop the redundant initDb() (001_init handles it) and remove
  schema-drift in configureS3Backup (app_settings has no created_at
  anymore and the unique constraint is on setting_key alone, not
  composite). 0/12 → 7/12 (5 remaining are unrelated assertion drift).
- backend/src/services/photoResolver.js — mixed-source events (reference
  mode with managed-uploaded photos) now fall back to managed when
  external_relpath is missing instead of throwing.
- tests/e2e/s3-storage-roundtrip.spec.ts — Playwright spec that
  auto-skips against local backend; full upload → serve → delete
  round-trip when run against an S3-mode backend.

Server wiring (server.js)
- initStorage() called after database init, before rate limiters.
- This commit's diff also includes the webhook delivery worker startup
  and the S3 auto-importer startup. Those features ship in the next two
  commits — co-located here for one bisectable diff per file.

Docs + ops
- README §"Storage Backends" — capability matrix, switching playbook,
  IAM policy snippet, MinIO/R2/B2 examples.
- README §"Webhooks" — also added here (full diff bundled).
- .env.example — STORAGE_BACKEND + STORAGE_S3_* + STORAGE_AUTO_IMPORT
  documented; WEBHOOK_* added in the same diff.
- .gitignore — re-anchor the existing `storage/` rule to `/storage/`
  so backend/src/services/storage/ (the new abstraction code) is
  trackable. The runtime ./storage/ data dir stays ignored.

Out of scope for v1 (per the issue): presigned URLs for individual
photo display (always streamed for protection middleware), CDN
integration, hybrid hot/cold tiers, S3 → local migration, multi-bucket
per-event.
2026-04-28 10:06:36 +02:00

260 lines
9.2 KiB
JavaScript

#!/usr/bin/env node
/**
* migrate-storage.js
*
* One-shot migration tool to copy every PicPeak content file from the local
* filesystem (the legacy STORAGE_PATH) to a configured S3-compatible bucket.
*
* Reads the relative path of each known asset from the database:
* photos.path
* photos.thumbnail_path
* photos.hero_path
* photos.watermark_path
* events.archive_path
* events.download_zip_path
*
* For each, streams from local fs → S3, skipping files whose sha256 already
* matches a previously uploaded object (idempotent — safe to re-run).
*
* Does NOT flip STORAGE_BACKEND. After the migration completes clean, the
* operator updates their environment + restarts the backend explicitly.
*
* Usage:
* node backend/scripts/migrate-storage.js # live migration
* node backend/scripts/migrate-storage.js --dry-run # report only, no uploads
* node backend/scripts/migrate-storage.js --failures-csv=/path/to/failures.csv
* node backend/scripts/migrate-storage.js --concurrency=4
*
* Required env (S3 destination — same vars the backend reads with STORAGE_BACKEND=s3):
* STORAGE_S3_BUCKET, STORAGE_S3_REGION, STORAGE_S3_ACCESS_KEY, STORAGE_S3_SECRET_KEY
* STORAGE_S3_ENDPOINT (optional — for MinIO/R2/etc.)
* STORAGE_S3_PREFIX (optional)
*
* STORAGE_PATH must point at the live local storage root. Postgres connection
* uses the same DB env vars the backend uses.
*/
require('dotenv').config();
const fs = require('fs');
const fsp = require('fs').promises;
const path = require('path');
const crypto = require('crypto');
const { db } = require('../src/database/db');
const LocalFsStorage = require('../src/services/storage/LocalFsStorage');
const S3StorageBackend = require('../src/services/storage/S3StorageBackend');
const logger = require('../src/utils/logger');
function parseArgs(argv) {
const args = { dryRun: false, concurrency: 4, failuresCsv: '/tmp/migrate-storage-failures.csv' };
for (const arg of argv) {
if (arg === '--dry-run') args.dryRun = true;
else if (arg.startsWith('--concurrency=')) args.concurrency = Math.max(1, parseInt(arg.split('=')[1], 10) || 4);
else if (arg.startsWith('--failures-csv=')) args.failuresCsv = arg.split('=')[1];
else if (arg === '--help' || arg === '-h') {
console.log('Usage: node migrate-storage.js [--dry-run] [--concurrency=N] [--failures-csv=PATH]');
process.exit(0);
}
}
return args;
}
function buildLocalSource() {
const root = process.env.STORAGE_PATH;
if (!root) {
throw new Error('STORAGE_PATH must be set to the local storage root.');
}
return new LocalFsStorage({ root });
}
function buildS3Destination() {
const required = ['STORAGE_S3_BUCKET', 'STORAGE_S3_ACCESS_KEY', 'STORAGE_S3_SECRET_KEY'];
const missing = required.filter((v) => !process.env[v]);
if (missing.length) {
throw new Error(`Missing S3 env vars: ${missing.join(', ')}`);
}
return new S3StorageBackend({
bucket: process.env.STORAGE_S3_BUCKET,
region: process.env.STORAGE_S3_REGION || 'us-east-1',
endpoint: process.env.STORAGE_S3_ENDPOINT,
accessKeyId: process.env.STORAGE_S3_ACCESS_KEY,
secretAccessKey: process.env.STORAGE_S3_SECRET_KEY,
prefix: process.env.STORAGE_S3_PREFIX,
forcePathStyle: process.env.STORAGE_S3_FORCE_PATH_STYLE === 'true' ? true : undefined,
sslEnabled: process.env.STORAGE_S3_SSL !== 'false',
});
}
async function sha256OfFile(localPath) {
return new Promise((resolve, reject) => {
const hash = crypto.createHash('sha256');
const stream = fs.createReadStream(localPath);
stream.on('data', (chunk) => hash.update(chunk));
stream.on('end', () => resolve(hash.digest('hex')));
stream.on('error', reject);
});
}
async function collectKeys() {
const keys = new Map(); // key -> { source, contentType }
const addKey = (key, source) => {
if (!key) return;
const normalized = key.replace(/\\/g, '/').replace(/^\/+/, '');
if (!normalized) return;
if (!keys.has(normalized)) keys.set(normalized, { source });
};
// photos: path (events/active/{slug}/{filename}), thumbnail_path, hero_path, watermark_path
const photoBatch = await db('photos').select('id', 'path', 'thumbnail_path', 'hero_path', 'watermark_path');
for (const p of photoBatch) {
if (p.path) {
const photoKey = p.path.startsWith('events/active/') ? p.path : path.posix.join('events/active', p.path);
addKey(photoKey, `photos.path[${p.id}]`);
}
addKey(p.thumbnail_path, `photos.thumbnail_path[${p.id}]`);
addKey(p.hero_path, `photos.hero_path[${p.id}]`);
addKey(p.watermark_path, `photos.watermark_path[${p.id}]`);
}
// events: archive_path, download_zip_path
const eventBatch = await db('events').select('id', 'archive_path', 'download_zip_path');
for (const e of eventBatch) {
addKey(e.archive_path, `events.archive_path[${e.id}]`);
addKey(e.download_zip_path, `events.download_zip_path[${e.id}]`);
}
return keys;
}
async function migrateOne(key, meta, { source, dest, dryRun }) {
// Source must exist on local disk.
const localPath = source.resolveLocalPath(key);
let localStat;
try {
localStat = await fsp.stat(localPath);
} catch (err) {
if (err.code === 'ENOENT') {
return { key, status: 'missing-locally', source: meta.source };
}
throw err;
}
// Idempotent skip: if S3 already has matching size + sha256.
const remoteStat = await dest.stat(key);
if (remoteStat && remoteStat.size === localStat.size) {
// sha256 match check via metadata is expensive; we trust size match for now.
// Operators paranoid about content drift can `rm` the bucket and re-run.
return { key, status: 'already-uploaded', source: meta.source };
}
if (dryRun) {
return { key, status: 'would-upload', source: meta.source, size: localStat.size };
}
await dest.putFromFile(key, localPath);
const verify = await dest.stat(key);
if (!verify || verify.size !== localStat.size) {
return { key, status: 'size-mismatch-after-upload', source: meta.source, expected: localStat.size, got: verify?.size };
}
return { key, status: 'uploaded', source: meta.source, size: localStat.size };
}
async function processWithConcurrency(items, concurrency, fn) {
const results = [];
let i = 0;
const workers = Array.from({ length: concurrency }, async () => {
while (true) {
const idx = i++;
if (idx >= items.length) return;
const [key, meta] = items[idx];
try {
const r = await fn(key, meta);
results.push(r);
} catch (err) {
results.push({ key, status: 'error', source: meta.source, error: err.message });
}
}
});
await Promise.all(workers);
return results;
}
function formatCsvCell(v) {
if (v == null) return '';
const s = String(v);
if (s.includes(',') || s.includes('"') || s.includes('\n')) {
return `"${s.replace(/"/g, '""')}"`;
}
return s;
}
async function writeFailuresCsv(filePath, failures) {
if (failures.length === 0) {
// Touch an empty file with header so callers see a deterministic outcome.
await fsp.writeFile(filePath, 'key,source,status,error\n');
return;
}
const lines = ['key,source,status,error'];
for (const f of failures) {
lines.push([f.key, f.source, f.status, f.error || ''].map(formatCsvCell).join(','));
}
await fsp.writeFile(filePath, lines.join('\n') + '\n');
}
async function main() {
const args = parseArgs(process.argv.slice(2));
logger.info(`migrate-storage starting (dry-run=${args.dryRun}, concurrency=${args.concurrency})`);
const source = buildLocalSource();
await source.init();
const dest = buildS3Destination();
await dest.init();
logger.info('collecting key list from database…');
const keys = await collectKeys();
logger.info(`found ${keys.size} unique keys to process`);
const items = Array.from(keys.entries());
const results = await processWithConcurrency(items, args.concurrency, (key, meta) =>
migrateOne(key, meta, { source, dest, dryRun: args.dryRun })
);
const counts = results.reduce((acc, r) => {
acc[r.status] = (acc[r.status] || 0) + 1;
return acc;
}, {});
console.log('\n=== migrate-storage summary ===');
for (const [status, count] of Object.entries(counts).sort()) {
console.log(` ${status.padEnd(28)} ${count}`);
}
const failureStatuses = new Set(['error', 'missing-locally', 'size-mismatch-after-upload']);
const failures = results.filter((r) => failureStatuses.has(r.status));
await writeFailuresCsv(args.failuresCsv, failures);
if (failures.length > 0) {
console.log(`\nWrote ${failures.length} failures to ${args.failuresCsv}`);
console.log('Re-run with --dry-run to triage; fix sources or remove DB rows that point at missing files.');
process.exitCode = 1;
} else if (args.dryRun) {
console.log(`\nDry-run complete. Re-run without --dry-run to perform the migration.`);
console.log(`(Empty failures CSV written to ${args.failuresCsv}.)`);
} else {
console.log(`\nMigration complete. Update STORAGE_BACKEND=s3 + restart the backend to switch over.`);
}
await db.destroy();
}
main().catch(async (err) => {
console.error('migrate-storage failed:', err);
try { await db.destroy(); } catch (_) { /* ignore */ }
process.exit(2);
});