Spaces:
Sleeping
Sleeping
| import { NextFunction, Request, Router, Response } from 'express'; | |
| import multer from 'multer'; | |
| import { authenticateUser, AuthenticatedRequest } from '../middleware/auth'; | |
| import { deleteInputFile, uploadInputFileFromPath } from '../../db/storage'; | |
| import { | |
| createJobWithTicket, | |
| getJobBySubmissionRequestId, | |
| getUserProfile, | |
| } from '../../db/tickets'; | |
| import { logger } from '../../utils/logger'; | |
| import { createHash, randomUUID } from 'crypto'; | |
| import { createReadStream, promises as fsPromises } from 'fs'; | |
| import * as os from 'os'; | |
| const router = Router(); | |
| /** | |
| * Keep large uploads out of the worker heap. The route removes every temporary | |
| * file after validation/storage, including early returns and failed requests. | |
| */ | |
| const upload = multer({ | |
| storage: multer.diskStorage({ | |
| destination: os.tmpdir(), | |
| filename: (_req, file, callback) => { | |
| callback(null, `relv-upload-${randomUUID()}${getFileExtension(file.originalname)}`); | |
| }, | |
| }), | |
| limits: { fileSize: 100 * 1024 * 1024 }, // User uploads are capped at 100 MB | |
| }); | |
| function receiveUpload( | |
| req: Request, | |
| res: Response, | |
| next: NextFunction, | |
| ): void { | |
| upload.single('file')(req, res, (error: unknown) => { | |
| if (error instanceof multer.MulterError && error.code === 'LIMIT_FILE_SIZE') { | |
| res.status(413).json({ error: 'File exceeds the 100MB upload limit' }); | |
| return; | |
| } | |
| if (error) { | |
| next(error); | |
| return; | |
| } | |
| next(); | |
| }); | |
| } | |
| /** Allowed file extensions for Turnitin submissions */ | |
| const ALLOWED_EXTENSIONS = new Set([ | |
| '.docx', | |
| '.xlsx', | |
| '.pptx', | |
| '.ps', | |
| '.pdf', | |
| '.html', | |
| '.rtf', | |
| '.odt', | |
| '.hwp', | |
| '.txt', | |
| ]); | |
| const ALLOWED_MODES = new Set(['upload', 'resubmit']); | |
| const UUID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i; | |
| const DEFAULT_FILTERS: Record<string, unknown> = { | |
| excludeBibliography: false, | |
| excludeQuotes: false, | |
| excludeCitations: false, | |
| excludeSmallMatches: true, | |
| smallMatchMode: 'words', | |
| smallMatchThreshold: 8, | |
| }; | |
| function getFileExtension(filename: string): string { | |
| const lastDot = filename.lastIndexOf('.'); | |
| return lastDot >= 0 ? filename.slice(lastDot).toLowerCase() : ''; | |
| } | |
| async function sha256File(filePath: string): Promise<string> { | |
| const hash = createHash('sha256'); | |
| for await (const chunk of createReadStream(filePath)) { | |
| hash.update(chunk as Buffer); | |
| } | |
| return hash.digest('hex'); | |
| } | |
| function existingJobMatchesRequest( | |
| existingJob: { | |
| assignment_target_id: string; | |
| mode: string; | |
| input_file_name: string; | |
| input_file_size: number | null; | |
| input_file_sha256: string | null; | |
| }, | |
| request: { | |
| assignmentTargetId: string; | |
| mode: string; | |
| inputFileName: string; | |
| inputFileSize: number; | |
| inputFileSha256: string; | |
| }, | |
| ): boolean { | |
| return ( | |
| existingJob.assignment_target_id === request.assignmentTargetId && | |
| existingJob.mode === request.mode && | |
| existingJob.input_file_name === request.inputFileName && | |
| existingJob.input_file_size === request.inputFileSize && | |
| existingJob.input_file_sha256 === request.inputFileSha256 | |
| ); | |
| } | |
| /** | |
| * POST /api/submit | |
| * Accepts a file upload + metadata, validates, stores the file, | |
| * creates a job (deducting a ticket), and returns the job ID. | |
| */ | |
| router.post( | |
| '/api/submit', | |
| authenticateUser, | |
| receiveUpload, | |
| async (req, res: Response): Promise<void> => { | |
| const authReq = req as AuthenticatedRequest; | |
| const temporaryUploadPath = authReq.file?.path; | |
| try { | |
| // 1. Validate required fields | |
| const { | |
| assignment_target_id, | |
| mode, | |
| filters: filtersRaw, | |
| submission_request_id: bodySubmissionRequestId, | |
| } = authReq.body; | |
| if (!assignment_target_id || !mode) { | |
| res.status(400).json({ | |
| error: 'Missing required fields: assignment_target_id, mode', | |
| }); | |
| return; | |
| } | |
| if (!ALLOWED_MODES.has(String(mode))) { | |
| res.status(400).json({ | |
| error: 'Invalid mode. Allowed: upload, resubmit', | |
| }); | |
| return; | |
| } | |
| const suppliedRequestId = String( | |
| bodySubmissionRequestId || authReq.get('Idempotency-Key') || '', | |
| ).trim(); | |
| const submissionRequestId = suppliedRequestId || randomUUID(); | |
| if (!UUID_PATTERN.test(submissionRequestId)) { | |
| res.status(400).json({ error: 'Invalid submission request ID' }); | |
| return; | |
| } | |
| // Parse filters (may come as stringified JSON) | |
| let filters: Record<string, unknown>; | |
| try { | |
| const parsedFilters = | |
| typeof filtersRaw === 'string' | |
| ? JSON.parse(filtersRaw) | |
| : filtersRaw && typeof filtersRaw === 'object' | |
| ? filtersRaw | |
| : {}; | |
| filters = { ...DEFAULT_FILTERS, ...parsedFilters }; | |
| } catch { | |
| res.status(400).json({ error: 'Invalid filters JSON' }); | |
| return; | |
| } | |
| // 2. Validate file exists | |
| if (!authReq.file) { | |
| res.status(400).json({ error: 'File is required' }); | |
| return; | |
| } | |
| const uploadedFile = authReq.file; | |
| // 3. Validate file type | |
| const ext = getFileExtension(uploadedFile.originalname); | |
| if (!ALLOWED_EXTENSIONS.has(ext)) { | |
| res.status(400).json({ | |
| error: `Unsupported file type: ${ext}. Allowed: ${[...ALLOWED_EXTENSIONS].join(', ')}`, | |
| }); | |
| return; | |
| } | |
| // 4. Validate filters. Legacy Turnitin supports small-match words, | |
| // percent, or off; modern Turnitin uses word threshold only. | |
| const smallMatchMode = | |
| filters.smallMatchMode === 'percent' || | |
| filters.smallMatchMode === 'off' || | |
| filters.smallMatchMode === 'words' | |
| ? filters.smallMatchMode | |
| : 'words'; | |
| filters.smallMatchMode = smallMatchMode; | |
| if (filters.excludeSmallMatches === true && smallMatchMode !== 'off') { | |
| let threshold = Number(filters.smallMatchThreshold) || 8; | |
| const maxThreshold = smallMatchMode === 'percent' ? 100 : 40; | |
| threshold = Math.max(1, Math.min(maxThreshold, threshold)); | |
| filters.smallMatchThreshold = threshold; | |
| } else { | |
| filters.excludeSmallMatches = false; | |
| filters.smallMatchMode = 'off'; | |
| filters.smallMatchThreshold = null; | |
| } | |
| // 5. Fingerprint the payload before any state-changing operation. | |
| const inputFileSha256 = await sha256File(uploadedFile.path); | |
| const existingJob = await getJobBySubmissionRequestId( | |
| authReq.userId, | |
| submissionRequestId, | |
| ); | |
| if (existingJob) { | |
| if (!existingJobMatchesRequest(existingJob, { | |
| assignmentTargetId: String(assignment_target_id), | |
| mode: String(mode), | |
| inputFileName: uploadedFile.originalname, | |
| inputFileSize: uploadedFile.size, | |
| inputFileSha256, | |
| })) { | |
| res.status(409).json({ | |
| error: 'Submission request ID was already used for a different file or configuration', | |
| }); | |
| return; | |
| } | |
| const profile = await getUserProfile(authReq.userId); | |
| logger.info('Idempotent submit replay returned existing job', { | |
| jobId: existingJob.id, | |
| userId: authReq.userId, | |
| submissionRequestId, | |
| }); | |
| res.status(200).json({ | |
| jobId: existingJob.id, | |
| ticketBalance: profile?.ticket_balance ?? 0, | |
| idempotentReplay: true, | |
| }); | |
| return; | |
| } | |
| // A deterministic request/hash path makes concurrent retries upload the | |
| // same bytes to the same object before the database invariant resolves them. | |
| const storagePath = await uploadInputFileFromPath( | |
| authReq.userId, | |
| `${submissionRequestId}/${inputFileSha256}`, | |
| uploadedFile.originalname, | |
| uploadedFile.path, | |
| true, | |
| ); | |
| // 6. Atomically create one job/ticket ledger entry, or return the job | |
| // already created by another Space handling this exact request. | |
| const creation = await (async () => { | |
| try { | |
| return await createJobWithTicket({ | |
| userId: authReq.userId, | |
| assignmentTargetId: assignment_target_id, | |
| mode, | |
| filters, | |
| inputFileName: uploadedFile.originalname, | |
| inputStoragePath: storagePath, | |
| inputFileSize: uploadedFile.size, | |
| inputFileSha256, | |
| submissionRequestId, | |
| }); | |
| } catch (creationError) { | |
| // The RPC response may fail after the transaction commits. Confirm | |
| // database state before deleting the staged object. | |
| let recoveredJob; | |
| try { | |
| recoveredJob = await getJobBySubmissionRequestId( | |
| authReq.userId, | |
| submissionRequestId, | |
| ); | |
| } catch (recoveryError) { | |
| logger.warn('Could not verify job creation after RPC failure; preserving input object', { | |
| userId: authReq.userId, | |
| submissionRequestId, | |
| error: | |
| recoveryError instanceof Error | |
| ? recoveryError.message | |
| : String(recoveryError), | |
| }); | |
| throw creationError; | |
| } | |
| if (recoveredJob) { | |
| if (!existingJobMatchesRequest(recoveredJob, { | |
| assignmentTargetId: String(assignment_target_id), | |
| mode: String(mode), | |
| inputFileName: uploadedFile.originalname, | |
| inputFileSize: uploadedFile.size, | |
| inputFileSha256, | |
| })) { | |
| await deleteInputFile(storagePath).catch(() => {}); | |
| throw new Error('Idempotency key recovered a different job payload'); | |
| } | |
| logger.warn('Recovered committed job after job-creation response failure', { | |
| jobId: recoveredJob.id, | |
| userId: authReq.userId, | |
| submissionRequestId, | |
| }); | |
| return { jobId: recoveredJob.id, created: false }; | |
| } | |
| await deleteInputFile(storagePath).catch((cleanupError: unknown) => { | |
| logger.warn('Failed to roll back staged input after job creation failure', { | |
| storagePath, | |
| error: | |
| cleanupError instanceof Error | |
| ? cleanupError.message | |
| : String(cleanupError), | |
| }); | |
| }); | |
| throw creationError; | |
| } | |
| })(); | |
| // Fetch updated ticket balance | |
| const profile = await getUserProfile(authReq.userId); | |
| const ticketBalance = profile?.ticket_balance ?? 0; | |
| logger.info(creation.created ? 'Job submitted successfully' : 'Idempotent submit race resolved', { | |
| jobId: creation.jobId, | |
| userId: authReq.userId, | |
| fileName: uploadedFile.originalname, | |
| submissionRequestId, | |
| }); | |
| // 7. Return success | |
| res.status(creation.created ? 201 : 200).json({ | |
| jobId: creation.jobId, | |
| ticketBalance, | |
| idempotentReplay: !creation.created, | |
| }); | |
| } catch (err: unknown) { | |
| const message = err instanceof Error ? err.message : String(err); | |
| // Handle specific error cases from the RPC | |
| if (message.includes('insufficient') || message.includes('ticket')) { | |
| res.status(402).json({ error: 'Insufficient ticket balance' }); | |
| return; | |
| } | |
| if (message.includes('Idempotency key')) { | |
| res.status(409).json({ error: 'Submission request ID conflict' }); | |
| return; | |
| } | |
| logger.error('Submit endpoint error', { | |
| userId: authReq.userId, | |
| error: message, | |
| }); | |
| res.status(500).json({ error: 'Internal server error' }); | |
| } finally { | |
| if (temporaryUploadPath) { | |
| await fsPromises.unlink(temporaryUploadPath).catch((error: unknown) => { | |
| logger.warn('Failed to remove temporary upload file', { | |
| temporaryUploadPath, | |
| error: error instanceof Error ? error.message : String(error), | |
| }); | |
| }); | |
| } | |
| } | |
| }, | |
| ); | |
| export default router; | |