// pages/api/chunk-upload-v2.ts // Chunked upload with server-side encryption. // // Raw chunks are uploaded, encrypted individually on arrival, and reassembled // into one encrypted file per uploaded file. See lib/server-encryption.ts for // the on-disk format. import { NextApiRequest, NextApiResponse } from 'next'; import { IncomingForm, Fields, Files } from 'formidable'; import fs from 'fs'; import { promises as fsPromises } from 'fs'; import path from 'path'; import { v4 as uuidv4 } from 'uuid'; import { UPLOAD_DIR } from '@/lib/config'; import { prisma } from '@/lib/prisma'; import bcrypt from 'bcrypt'; import { PLAN_CONFIG } from '@/lib/subscription'; import { sendMailjetEmail } from '@/lib/mailjet'; import { deriveKeyFromPassword, encryptChunkFrame, buildHeader } from '@/lib/server-encryption'; import { getServerSession } from 'next-auth/next'; import { authOptions } from '@/lib/auth'; import crypto from 'crypto'; export const config = { api: { bodyParser: false, }, }; interface UploadFile { name: string; totalChunks: number; receivedChunks: Set; receivedBytes: number; } interface UploadSession { uploadId: string; sender: string; recipient: string; password: string; message?: string; filenames?: string; fileCount: number; /** Keyed by fileIndex. */ files: Map; createdAt: Date; // One salt/key per transfer. Sharing a key across files is safe because every // chunk frame carries its own unique IV. salt: Buffer; encryptionKey: Buffer; maxBytes: number; totalReceivedBytes: number; } // NOTE: in-process state, matching the local-disk storage model. Running more // than one instance requires a shared store (Redis) plus shared object storage. const uploadSessions = new Map(); // Cleanup abandoned uploads every 5 minutes const SESSION_TIMEOUT = 24 * 60 * 60 * 1000; // 24 hours setInterval(() => { const now = Date.now(); for (const [uploadId, session] of uploadSessions.entries()) { if (now - session.createdAt.getTime() > SESSION_TIMEOUT) { uploadSessions.delete(uploadId); const chunkDir = path.join(UPLOAD_DIR, '.tmp', uploadId); fsPromises.rm(chunkDir, { recursive: true, force: true }).catch((err) => console.warn(`Failed to clean abandoned upload ${uploadId}:`, err) ); } } }, 5 * 60 * 1000); const firstValue = (value: string | string[] | undefined): string | undefined => Array.isArray(value) ? value[0] : value; const fileChunkDir = (uploadId: string, fileIndex: number) => path.join(UPLOAD_DIR, '.tmp', uploadId, `f${fileIndex}`); export default async function handler(req: NextApiRequest, res: NextApiResponse) { const authSession = await getServerSession(req, res, authOptions); if (!authSession?.user?.email) { return res.status(401).json({ success: false, message: 'Unauthorized' }); } const sessionEmail = authSession.user.email; if (req.method === 'POST') { return handleChunkUpload(req, res, sessionEmail); } else if (req.method === 'GET') { return handleStatusCheck(req, res, sessionEmail); } else if (req.method === 'PUT') { return handleChunkComplete(req, res, sessionEmail); } return res.status(405).json({ success: false, message: 'Method not allowed' }); } async function handleChunkUpload( req: NextApiRequest, res: NextApiResponse, sessionEmail: string ) { await fsPromises.mkdir(path.join(UPLOAD_DIR, '.tmp'), { recursive: true }).catch((err) => { console.error('Failed to create temp directory:', err); }); const form = new IncomingForm({ uploadDir: path.join(UPLOAD_DIR, '.tmp'), keepExtensions: true, }); return new Promise((resolve) => { form.parse(req, async (err: Error | null, fields: Fields, files: Files) => { if (err) { res.status(400).json({ success: false, message: 'Parse error: ' + err.message }); return resolve(); } try { const chunkFile = Array.isArray(files.chunk) ? files.chunk[0] : files.chunk; const uploadId = firstValue(fields.uploadId); const chunkIndex = parseInt(firstValue(fields.chunkIndex) ?? '', 10); const fileIndex = parseInt(firstValue(fields.fileIndex) ?? '0', 10); const fileChunks = parseInt(firstValue(fields.fileChunks) ?? '', 10); const fileCount = parseInt(firstValue(fields.fileCount) ?? '1', 10); const fileName = firstValue(fields.fileName); if ( !chunkFile || !uploadId || Number.isNaN(chunkIndex) || Number.isNaN(fileIndex) || Number.isNaN(fileChunks) || fileIndex < 0 || chunkIndex < 0 || chunkIndex >= fileChunks ) { res.status(400).json({ success: false, message: 'Missing or invalid chunk metadata' }); return resolve(); } let session = uploadSessions.get(uploadId); if (session && session.sender !== sessionEmail) { res.status(403).json({ success: false, message: 'Forbidden' }); return resolve(); } if (!session) { const recipient = firstValue(fields.email); const password = firstValue(fields.password); const message = firstValue(fields.message); const filenames = firstValue(fields.filenames); if (!recipient) { res.status(400).json({ success: false, message: 'Missing required fields' }); return resolve(); } // A password is mandatory: it is the key material. Without it every // file would be encrypted under a key derived from the empty string. if (!password) { res.status(400).json({ success: false, message: 'A password is required' }); return resolve(); } const user = await prisma.user.findUnique({ where: { email: sessionEmail } }); if (!user) { res.status(404).json({ success: false, message: 'User not found' }); return resolve(); } const plan = (user.plan || 'free') as 'free' | 'rookie' | 'pro'; const salt = crypto.randomBytes(16); session = { uploadId, sender: sessionEmail, recipient, password, message, filenames, fileCount: Number.isNaN(fileCount) ? 1 : fileCount, files: new Map(), createdAt: new Date(), salt, encryptionKey: await deriveKeyFromPassword(password, salt), maxBytes: PLAN_CONFIG[plan].maxFileSize, totalReceivedBytes: 0, }; uploadSessions.set(uploadId, session); } let entry = session.files.get(fileIndex); if (!entry) { entry = { name: fileName || `file-${fileIndex}`, totalChunks: fileChunks, receivedChunks: new Set(), receivedBytes: 0, }; session.files.set(fileIndex, entry); } const chunkData = await fsPromises.readFile(chunkFile.filepath); await fsPromises.unlink(chunkFile.filepath).catch(() => {}); // Enforce the plan cap against bytes actually received, not the // client-declared size. const isNewChunk = !entry.receivedChunks.has(chunkIndex); if ( isNewChunk && session.maxBytes !== Infinity && session.totalReceivedBytes + chunkData.length > session.maxBytes ) { res.status(413).json({ success: false, message: 'Transfer exceeds the maximum size for your plan', }); return resolve(); } const dir = fileChunkDir(uploadId, fileIndex); await fsPromises.mkdir(dir, { recursive: true }); await fsPromises.writeFile( path.join(dir, `chunk-${chunkIndex}.enc`), encryptChunkFrame(chunkData, session.encryptionKey) ); if (isNewChunk) { entry.receivedBytes += chunkData.length; session.totalReceivedBytes += chunkData.length; } entry.receivedChunks.add(chunkIndex); res.status(200).json({ success: true, uploadId, fileIndex, chunkIndex, receivedChunks: entry.receivedChunks.size, totalChunks: entry.totalChunks, }); return resolve(); } catch (error: unknown) { console.error('Chunk upload error:', error); res.status(500).json({ success: false, message: 'Upload error: ' + (error as Error).message, }); return resolve(); } }); }); } async function handleStatusCheck( req: NextApiRequest, res: NextApiResponse, sessionEmail: string ) { const { uploadId } = req.query; if (!uploadId || typeof uploadId !== 'string') { return res.status(400).json({ success: false, message: 'Missing uploadId' }); } const session = uploadSessions.get(uploadId); if (!session || session.sender !== sessionEmail) { return res.status(404).json({ success: false, message: 'Upload session not found' }); } const files = [...session.files.entries()].map(([fileIndex, entry]) => ({ fileIndex, name: entry.name, receivedChunks: entry.receivedChunks.size, totalChunks: entry.totalChunks, isComplete: entry.receivedChunks.size === entry.totalChunks, })); return res.status(200).json({ success: true, uploadId, fileCount: session.fileCount, files, receivedBytes: session.totalReceivedBytes, isComplete: session.files.size === session.fileCount && files.every((f) => f.isComplete), }); } async function handleChunkComplete( req: NextApiRequest, res: NextApiResponse, sessionEmail: string ) { const { uploadId } = req.query; if (!uploadId || typeof uploadId !== 'string') { return res.status(400).json({ success: false, message: 'Missing uploadId' }); } const session = uploadSessions.get(uploadId); if (!session || session.sender !== sessionEmail) { return res.status(404).json({ success: false, message: 'Upload session not found' }); } // Every file must be present and whole before anything is assembled. if (session.files.size !== session.fileCount) { return res.status(400).json({ success: false, message: `Missing files. Received ${session.files.size}/${session.fileCount}`, }); } for (const [fileIndex, entry] of session.files) { if (entry.receivedChunks.size !== entry.totalChunks) { return res.status(400).json({ success: false, message: `File ${fileIndex} incomplete: ${entry.receivedChunks.size}/${entry.totalChunks} chunks`, }); } } const writtenPaths: string[] = []; try { const user = await prisma.user.findUnique({ where: { email: session.sender } }); if (!user) { return res.status(404).json({ success: false, message: 'User not found' }); } const plan = (user.plan || 'free') as 'free' | 'rookie' | 'pro'; const limits = PLAN_CONFIG[plan]; const startOfMonth = new Date(); startOfMonth.setDate(1); startOfMonth.setHours(0, 0, 0, 0); const transfersThisMonth = await prisma.transfer.findMany({ where: { senderEmail: session.sender, createdAt: { gte: startOfMonth } }, }); // Sizes are plaintext bytes received, not on-disk size, which is inflated // by the format header and each frame's IV and auth tag. const totalPlaintext = session.totalReceivedBytes; const sentThisMonth = transfersThisMonth.reduce( (acc, t) => acc + BigInt(t.totalSize), BigInt(0) ); if ( (limits.maxTransfersPerMonth !== Infinity && transfersThisMonth.length >= limits.maxTransfersPerMonth) || (limits.maxTransferSizePerMonth !== Infinity && sentThisMonth + BigInt(totalPlaintext) > BigInt(limits.maxTransferSizePerMonth)) ) { return res.status(429).json({ success: false, message: 'Plan limits exceeded' }); } // Assemble each file into its own encrypted payload. const assembled: { name: string; path: string; size: number }[] = []; for (const [fileIndex, entry] of [...session.files.entries()].sort((a, b) => a[0] - b[0])) { const dir = fileChunkDir(uploadId, fileIndex); const finalPath = path.join(UPLOAD_DIR, uuidv4() + '.enc'); writtenPaths.push(finalPath); const writeStream = fs.createWriteStream(finalPath); // Honour backpressure so a multi-GB reassembly flushes to disk rather // than queueing in the stream's internal buffer. const write = (buf: Buffer): Promise => new Promise((resolve, reject) => { if (writeStream.write(buf)) return resolve(); writeStream.once('drain', resolve); writeStream.once('error', reject); }); await write(buildHeader(session.salt)); for (let i = 0; i < entry.totalChunks; i++) { const chunkPath = path.join(dir, `chunk-${i}.enc`); let chunkData: Buffer | null = null; for (let attempt = 0; attempt < 3; attempt++) { try { chunkData = await fsPromises.readFile(chunkPath); break; } catch (error: unknown) { const code = (error as NodeJS.ErrnoException)?.code; if (attempt < 2 && code === 'EACCES') { await new Promise((r) => setTimeout(r, Math.pow(2, attempt) * 100)); } else { throw error; } } } // A silently skipped chunk would corrupt the file, so fail loudly. if (!chunkData) { throw new Error(`Chunk ${i} of file ${fileIndex} missing during reassembly`); } await write(chunkData); } await new Promise((resolve, reject) => { writeStream.on('finish', resolve); writeStream.on('error', reject); writeStream.end(); }); assembled.push({ name: entry.name, path: finalPath, size: entry.receivedBytes }); } const hash = await bcrypt.hash(session.password, 10); const expiresAt = new Date(Date.now() + limits.maxExpiryMs); const transfer = await prisma.transfer.create({ data: { senderEmail: session.sender, recipientEmail: session.recipient, passwordHash: hash, downloadUrl: uuidv4(), expiresAt, filenames: session.filenames || JSON.stringify(assembled.map((f) => f.name)), message: session.message || null, totalFiles: assembled.length, totalSize: totalPlaintext, userId: user.id, files: { create: assembled.map((f) => ({ name: f.name, path: f.path, size: f.size })), }, }, }); const baseUrl = process.env.NEXT_PUBLIC_APP_URL || 'https://transfertribe.com'; const link = `${baseUrl}/download/${transfer.downloadUrl}`; const fileLine = assembled.length === 1 ? '1 file' : `${assembled.length} files`; await sendMailjetEmail({ to: session.recipient, subject: "You've received an encrypted file", text: ` ${session.sender} sent you ${fileLine} via TransferTribe. Link: ${link} ${session.message ? `Message:\n${session.message}\n\n` : ''}Note: You'll need the password they shared with you to decrypt it. `.trim(), }); uploadSessions.delete(uploadId); await fsPromises .rm(path.join(UPLOAD_DIR, '.tmp', uploadId), { recursive: true, force: true }) .catch((err) => console.warn('Failed to clean up temp directory:', err)); return res.status(200).json({ success: true, message: 'Upload complete', downloadUrl: transfer.downloadUrl, files: assembled.length, }); } catch (error: unknown) { console.error('Chunk completion error:', error); // Don't leave half-assembled payloads behind on failure. await Promise.all( writtenPaths.map((p) => fsPromises.unlink(p).catch(() => {})) ); return res.status(500).json({ success: false, message: 'Completion error: ' + (error as Error).message, }); } }