relv-space6 / src /db /storage.ts
relv-dev's picture
Preserve Supabase bucket routing for legacy storage paths
8bde444 verified
Raw
History Blame Contribute Delete
11.6 kB
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<string> {
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<string> {
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<void> {
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<void> {
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<void> {
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<string> {
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<string | null> {
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<string> {
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<string, string> = {
'.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';
}