relv-space4 / src /server /routes /submit.ts
relv-dev's picture
Add R2 dual-storage adapter and resilient legacy large-document wait
16e3957 verified
Raw
History Blame Contribute Delete
12.1 kB
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;