import * as fs from 'fs'; import * as path from 'path'; import { supabase } from './client'; import { config } from '../config'; import { logger } from '../utils/logger'; import { createR2SignedDownloadUrl, deleteR2Object, downloadR2File, isR2ObjectRef, putR2Buffer, putR2File, readR2Text, } from './r2'; /** * Upload a user's input file to the turnitin-inputs bucket. * Returns the storage path within the bucket. */ export async function uploadInputFile( userId: string, storageKey: string, fileName: string, fileBuffer: Buffer, upsert = false, ): Promise { const ext = path.extname(fileName); const storagePath = `${userId}/${storageKey}/input${ext}`; if (config.storageProvider === 'r2') { try { return await putR2Buffer( config.inputBucket, storagePath, fileBuffer, getMimeType(ext), ); } catch (error) { logger.error('Failed to upload input file to R2', { storagePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } const { error } = await supabase.storage .from(config.supabaseInputBucket) .upload(storagePath, fileBuffer, { contentType: getMimeType(ext), upsert, }); if (error) { logger.error('Failed to upload input file', { storagePath, error: error.message }); throw error; } return storagePath; } /** * Upload an input file from disk. R2 receives a stream so a 100 MB user upload * does not need a second full-size in-memory copy inside the worker. */ export async function uploadInputFileFromPath( userId: string, storageKey: string, fileName: string, localPath: string, upsert = false, ): Promise { const ext = path.extname(fileName); const storagePath = `${userId}/${storageKey}/input${ext}`; if (config.storageProvider === 'r2') { try { return await putR2File( config.inputBucket, storagePath, localPath, getMimeType(ext), ); } catch (error) { logger.error('Failed to upload input file to R2', { storagePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } // Supabase Storage's Node client expects a buffer. This fallback only serves // legacy objects/deployments; R2 is the configured destination for new jobs. const fileBuffer = await fs.promises.readFile(localPath); const { error } = await supabase.storage .from(config.supabaseInputBucket) .upload(storagePath, fileBuffer, { contentType: getMimeType(ext), upsert, }); if (error) { logger.error('Failed to upload input file', { storagePath, error: error.message }); throw error; } return storagePath; } /** * Download a file from Supabase Storage to a local path. */ export async function downloadInputFile(storagePath: string, localPath: string): Promise { if (isR2ObjectRef(storagePath)) { try { await downloadR2File(storagePath, localPath); return; } catch (error) { logger.error('Failed to download input file from R2', { storagePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } const { data, error } = await supabase.storage .from(config.supabaseInputBucket) .download(storagePath); if (error) { logger.error('Failed to download input file', { storagePath, error: error.message }); throw error; } const buffer = Buffer.from(await data.arrayBuffer()); const dir = path.dirname(localPath); if (!fs.existsSync(dir)) { fs.mkdirSync(dir, { recursive: true }); } fs.writeFileSync(localPath, buffer); } /** Delete a staged input when atomic job creation definitively did not happen. */ export async function deleteInputFile(storagePath: string): Promise { if (isR2ObjectRef(storagePath)) { await deleteR2Object(storagePath); return; } const { error } = await supabase.storage .from(config.supabaseInputBucket) .remove([storagePath]); if (error) { logger.error('Failed to delete staged input file', { storagePath, error: error.message, }); throw error; } } /** * Upload a generated report PDF to the turnitin-reports bucket. * Returns the storage path and expiry timestamp. */ export async function uploadReportPdf( userId: string, jobId: string, localPdfPath: string, ): Promise<{ storagePath: string; expiresAt: string }> { const storagePath = `${userId}/${jobId}/report.pdf`; if (config.storageProvider === 'r2') { try { const objectRef = await putR2File( config.reportBucket, storagePath, localPdfPath, 'application/pdf', ); return { storagePath: objectRef, expiresAt: createReportExpiry(), }; } catch (error) { logger.error('Failed to upload report PDF to R2', { storagePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } const fileBuffer = fs.readFileSync(localPdfPath); const { error } = await supabase.storage .from(config.supabaseReportBucket) .upload(storagePath, fileBuffer, { contentType: 'application/pdf', upsert: true, }); if (error) { logger.error('Failed to upload report PDF', { storagePath, error: error.message }); throw error; } return { storagePath, expiresAt: createReportExpiry() }; } /** * Upload a legacy Turnitin Digital Receipt PDF with the same retention policy * as the similarity report. */ export async function uploadReceiptPdf( userId: string, jobId: string, localPdfPath: string, ): Promise<{ storagePath: string; expiresAt: string }> { const storagePath = `${userId}/${jobId}/receipt.pdf`; if (config.storageProvider === 'r2') { try { const objectRef = await putR2File( config.reportBucket, storagePath, localPdfPath, 'application/pdf', ); return { storagePath: objectRef, expiresAt: createReportExpiry(), }; } catch (error) { logger.error('Failed to upload Digital Receipt PDF to R2', { storagePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } const fileBuffer = fs.readFileSync(localPdfPath); const { error } = await supabase.storage .from(config.supabaseReportBucket) .upload(storagePath, fileBuffer, { contentType: 'application/pdf', upsert: true, }); if (error) { logger.error('Failed to upload Digital Receipt PDF', { storagePath, error: error.message, }); throw error; } return { storagePath, expiresAt: createReportExpiry() }; } /** * Delete a report PDF from Supabase Storage. */ export async function deleteReportPdf(storagePath: string): Promise { if (isR2ObjectRef(storagePath)) { try { await deleteR2Object(storagePath); return; } catch (error) { logger.error('Failed to delete report PDF from R2', { storagePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } const { error } = await supabase.storage .from(config.supabaseReportBucket) .remove([storagePath]); if (error) { logger.error('Failed to delete report PDF', { storagePath, error: error.message }); throw error; } } /** * Upload a Playwright browser storage state to the sessions bucket. * Returns the storage path. */ export async function uploadStorageState( accountId: string, stateJson: string, ): Promise { const storagePath = `${accountId}/state.json`; if (config.storageProvider === 'r2') { try { return await putR2Buffer( config.sessionBucket, storagePath, Buffer.from(stateJson, 'utf-8'), 'application/json', ); } catch (error) { logger.error('Failed to upload storage state to R2', { accountId, error: error instanceof Error ? error.message : String(error), }); throw error; } } const { error } = await supabase.storage .from(config.supabaseSessionBucket) .upload(storagePath, Buffer.from(stateJson, 'utf-8'), { contentType: 'application/json', upsert: true, }); if (error) { logger.error('Failed to upload storage state', { accountId, error: error.message }); throw error; } return storagePath; } /** * Download a previously saved storage state. * Returns the JSON string, or null if not found. */ export async function downloadStorageState(storagePath: string): Promise { if (isR2ObjectRef(storagePath)) { try { return await readR2Text(storagePath); } catch (error) { logger.error('Failed to download storage state from R2', { storagePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } const { data, error } = await supabase.storage .from(config.supabaseSessionBucket) .download(storagePath); if (error) { // Not found is not fatal — the account may not have a saved session if (error.message?.includes('not found') || error.message?.includes('Object not found')) { return null; } logger.error('Failed to download storage state', { storagePath, error: error.message }); throw error; } return await data.text(); } /** * Create a time-limited signed URL for a file in any bucket. */ export async function createSignedUrl( bucket: string, filePath: string, expiresInSeconds: number, downloadFileName?: string, ): Promise { if (isR2ObjectRef(filePath)) { try { return await createR2SignedDownloadUrl( filePath, expiresInSeconds, downloadFileName, ); } catch (error) { logger.error('Failed to create R2 signed URL', { bucket, filePath, error: error instanceof Error ? error.message : String(error), }); throw error; } } const legacyBucket = bucket === config.reportBucket ? config.supabaseReportBucket : bucket === config.inputBucket ? config.supabaseInputBucket : bucket === config.sessionBucket ? config.supabaseSessionBucket : bucket; const { data, error } = await supabase.storage .from(legacyBucket) .createSignedUrl( filePath, expiresInSeconds, downloadFileName ? { download: downloadFileName } : undefined, ); if (error) { logger.error('Failed to create signed URL', { bucket, filePath, error: error.message }); throw error; } return data.signedUrl; } function createReportExpiry(): string { return new Date( Date.now() + config.reportRetentionHours * 60 * 60 * 1000, ).toISOString(); } /** Map file extensions to MIME types */ function getMimeType(ext: string): string { const mimeTypes: Record = { '.pdf': 'application/pdf', '.docx': 'application/vnd.openxmlformats-officedocument.wordprocessingml.document', '.xlsx': 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet', '.pptx': 'application/vnd.openxmlformats-officedocument.presentationml.presentation', '.ps': 'application/postscript', '.html': 'text/html', '.txt': 'text/plain', '.rtf': 'application/rtf', '.odt': 'application/vnd.oasis.opendocument.text', '.hwp': 'application/x-hwp', }; return mimeTypes[ext.toLowerCase()] || 'application/octet-stream'; }