feat: implement storage quota management for uploads

- Added storage quota enforcement for audio and image uploads in the respective routes.
- Introduced reservation system to manage concurrent uploads and prevent quota overages.
- Enhanced comment creation to account for audio and image attachment sizes against user quotas.
- Created new UploadReservation model to track in-flight upload reservations.
- Backfilled existing video assets with size information from R2.
- Added progress component for UI feedback during uploads.
- Updated API responses to include reservation IDs for better quota management.
- Adjusted error handling to return appropriate storage limit exceeded messages.
This commit is contained in:
Yusuf İpek
2026-04-15 19:53:43 +03:00
parent acf31b3d6b
commit 873945464d
21 changed files with 753 additions and 85 deletions
+3
View File
@@ -189,6 +189,9 @@ export async function PATCH(request: NextRequest, { params }: RouteParams) {
if (annotationData === null) {
updateData.annotationData = null;
} else {
if (!Array.isArray(annotationData)) {
return apiErrors.badRequest('annotationData must be an array of valid stroke objects');
}
const validStrokes = validateAnnotationStrokes(annotationData);
if (validStrokes === null) {
return apiErrors.badRequest('annotationData must be an array of valid stroke objects');
@@ -8,13 +8,14 @@ import { cleanupBunnyStreamVideos } from '@/lib/bunny-stream-cleanup';
import { createBunnyUploadToken, verifyBunnyUploadToken } from '@/lib/bunny-upload-token';
import { isBunnyUploadsFeatureEnabled } from '@/lib/feature-flags';
import { logError } from '@/lib/logger';
import { enforceStorageQuota } from '@/lib/storage-quota';
type RouteParams = { params: Promise<{ projectId: string }> };
async function getProjectWithEditAccess(projectId: string, userId: string) {
const project = await db.project.findUnique({
where: { id: projectId },
select: { id: true, name: true, ownerId: true, workspaceId: true, visibility: true },
select: { id: true, name: true, ownerId: true, workspaceId: true, visibility: true, workspace: { select: { ownerId: true } } },
});
if (!project) return null;
@@ -56,6 +57,9 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
return apiErrors.badRequest('Direct uploads are disabled by this host');
}
const quotaError = await enforceStorageQuota(project.workspace.ownerId, BigInt(0));
if (quotaError) return quotaError;
const apiKey = process.env.BUNNY_STREAM_API_KEY;
const libraryId = process.env.BUNNY_STREAM_LIBRARY_ID || process.env.NEXT_PUBLIC_BUNNY_STREAM_LIBRARY_ID;
+40
View File
@@ -0,0 +1,40 @@
import { auth } from '@/lib/auth';
import { apiErrors, successResponse, withCacheControl } from '@/lib/api-response';
import { getUserStorageInfo } from '@/lib/storage-quota';
import { hasBillingAccess } from '@/lib/billing';
import { db } from '@/lib/db';
// GET /api/settings/storage
export async function GET() {
const session = await auth();
if (!session?.user?.id) {
return apiErrors.unauthorized();
}
// Only users with active billing (or on a self-hosted instance where billing
// is disabled) should be able to enumerate their storage breakdown.
const user = await db.user.findUnique({
where: { id: session.user.id },
select: {
subscriptionStatus: true,
trialEndsAt: true,
stripeCurrentPeriodEnd: true,
billingAccessEndedAt: true,
},
});
if (!user || !hasBillingAccess(user)) {
return apiErrors.forbidden();
}
const info = await getUserStorageInfo(session.user.id);
const response = successResponse({
usedBytes: info.usedBytes.toString(),
limitBytes: info.limitBytes.toString(),
percentage: info.percentage,
});
// Cache for 60s — stale data is acceptable for a usage meter
return withCacheControl(response, 'private, max-age=60');
}
+32 -11
View File
@@ -13,6 +13,7 @@ import {
enforceGuestUploadQuota,
verifyGuestUploadToken,
} from '@/lib/guest-upload-token';
import { reserveStorageQuota, releaseStorageReservation } from '@/lib/storage-quota';
import { logError } from '@/lib/logger';
const MAX_FILE_SIZE = 10 * 1024 * 1024; // 10MB
@@ -118,7 +119,11 @@ export async function POST(request: NextRequest) {
const safeVideoId = videoId.trim();
const video = await db.video.findUnique({
where: { id: safeVideoId },
include: { project: true },
include: {
project: {
include: { workspace: { select: { ownerId: true } } },
},
},
});
if (!video) {
return apiErrors.notFound('Video');
@@ -170,11 +175,20 @@ export async function POST(request: NextRequest) {
return apiErrors.badRequest('File too large. Maximum size is 10MB.');
}
// Enforce per-user storage quota before uploading.
// All paths use the advisory-locked reservation so concurrent uploads always
// see each other's in-flight sizes, eliminating the TOCTOU race.
const workspaceOwnerId = video.project.workspace.ownerId;
const reserveResult = await reserveStorageQuota(workspaceOwnerId, BigInt(file.size));
if ('error' in reserveResult) return reserveResult.error;
const reservationId = reserveResult.reservationId;
// Normalize content type: strip codec params, then resolve aliases
const rawContentType = file.type || 'audio/webm';
const strippedType = rawContentType.split(';')[0].trim().toLowerCase();
const contentType = MIME_ALIASES[strippedType] ?? strippedType;
if (!ALLOWED_TYPES.has(contentType)) {
await releaseStorageReservation(reservationId);
return apiErrors.badRequest(`Unsupported audio format: ${rawContentType}`);
}
@@ -191,26 +205,33 @@ export async function POST(request: NextRequest) {
// Validate file content against magic bytes — rejects HTML/scripts masquerading as audio
if (isHtmlContent(buffer)) {
await releaseStorageReservation(reservationId);
return apiErrors.badRequest('File content does not match an audio format');
}
if (!hasValidAudioMagicBytes(buffer.slice(0, 16), contentType)) {
await releaseStorageReservation(reservationId);
return apiErrors.badRequest('File content does not match the declared audio format');
}
// Upload to R2
await r2Client.send(
new PutObjectCommand({
Bucket: R2_BUCKET_NAME,
Key: key,
Body: buffer,
ContentType: contentType,
})
);
try {
// Upload to R2
await r2Client.send(
new PutObjectCommand({
Bucket: R2_BUCKET_NAME,
Key: key,
Body: buffer,
ContentType: contentType,
})
);
} catch (uploadError) {
await releaseStorageReservation(reservationId);
throw uploadError;
}
// Return the URL through our proxy endpoint
const voiceUrl = `/api/upload/audio/${filename}`;
const response = successResponse({ url: voiceUrl }, 201);
const response = successResponse({ url: voiceUrl, reservationId }, 201);
return withCacheControl(response, 'private, no-store');
} catch (error) {
logError('Error uploading audio:', error);
+32 -12
View File
@@ -20,6 +20,7 @@ import {
verifyGuestUploadToken,
} from '@/lib/guest-upload-token';
import { logError } from '@/lib/logger';
import { reserveStorageQuota, releaseStorageReservation } from '@/lib/storage-quota';
const MAX_FILE_SIZE = 10 * 1024 * 1024; // 10MB
const MAX_MULTIPART_BODY_SIZE = MAX_FILE_SIZE + (512 * 1024); // file + multipart overhead
@@ -64,7 +65,11 @@ export async function POST(request: NextRequest) {
const safeVideoId = videoId.trim();
const video = await db.video.findUnique({
where: { id: safeVideoId },
include: { project: true },
include: {
project: {
include: { workspace: { select: { ownerId: true } } },
},
},
});
if (!video) {
return apiErrors.notFound('Video');
@@ -116,9 +121,18 @@ export async function POST(request: NextRequest) {
return apiErrors.badRequest('File too large. Maximum size is 10MB.');
}
// Enforce per-user storage quota before uploading.
// All paths use the advisory-locked reservation so concurrent uploads always
// see each other's in-flight sizes, eliminating the TOCTOU race.
const workspaceOwnerId = video.project.workspace.ownerId;
const reserveResult = await reserveStorageQuota(workspaceOwnerId, BigInt(file.size));
if ('error' in reserveResult) return reserveResult.error;
const reservationId = reserveResult.reservationId;
// Check content type
const normalizedMime = normalizeImageMime(file.type);
if (normalizedMime && !isAllowedImageType(normalizedMime)) {
await releaseStorageReservation(reservationId);
return apiErrors.badRequest(`Unsupported image format: ${file.type}`);
}
@@ -127,6 +141,7 @@ export async function POST(request: NextRequest) {
const buffer = Buffer.from(arrayBuffer);
const detectedMime = detectImageMime(buffer);
if (!detectedMime) {
await releaseStorageReservation(reservationId);
return apiErrors.badRequest('Uploaded file content does not match an allowed image type');
}
@@ -135,21 +150,26 @@ export async function POST(request: NextRequest) {
const filename = `${randomUUID()}.${ext}`;
const key = `images/${filename}`;
// Upload to R2
await r2Client.send(
new PutObjectCommand({
Bucket: R2_BUCKET_NAME,
Key: key,
Body: buffer,
ContentType: detectedMime,
})
);
try {
// Upload to R2
await r2Client.send(
new PutObjectCommand({
Bucket: R2_BUCKET_NAME,
Key: key,
Body: buffer,
ContentType: detectedMime,
})
);
} catch (uploadError) {
await releaseStorageReservation(reservationId);
throw uploadError;
}
// Return the URL through our proxy endpoint
const imageUrl = `/api/upload/image/${filename}`;
const response = successResponse({ url: imageUrl }, 201);
return withCacheControl(response, 'public, max-age=31536000, immutable');
const response = successResponse({ url: imageUrl, reservationId }, 201);
return withCacheControl(response, 'private, no-store');
} catch (error) {
logError('Error uploading image:', error);
return apiErrors.internalError('Failed to upload image');
+68 -11
View File
@@ -9,18 +9,21 @@ import { getShareSessionFromRequest } from '@/lib/share-session';
import { HeadObjectCommand } from '@aws-sdk/client-s3';
import { r2Client, R2_BUCKET_NAME } from '@/lib/r2';
import { ensureGuestIdentityFromRequest, getGuestIdentityFromRequest, setGuestIdentityCookie } from '@/lib/guest-identity';
import { extractImageFileNameFromProxyUrl, sanitizeAssetDisplayName } from '@/lib/video-assets';
import { extractImageFileNameFromProxyUrl, extractAudioFileNameFromProxyUrl, sanitizeAssetDisplayName } from '@/lib/video-assets';
import { validateAnnotationStrokes } from '@/lib/validation';
import { logError } from '@/lib/logger';
import { reserveStorageQuota, releaseStorageReservation } from '@/lib/storage-quota';
type RouteParams = { params: Promise<{ versionId: string }> };
const SAFE_IMAGE_PATH = /^\/api\/upload\/image\/[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\.[a-z0-9]+$/i;
const SAFE_AUDIO_PATH = /^\/api\/upload\/audio\/[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\.[a-z0-9]+$/i;
const UNATTACHED_UPLOAD_TTL_MS = 15 * 60 * 1000;
async function isFreshAttachment(url: string, kind: 'audio' | 'image'): Promise<boolean> {
type AttachmentCheck = { isFresh: boolean; sizeBytes: bigint };
async function isFreshAttachment(url: string, kind: 'audio' | 'image'): Promise<AttachmentCheck> {
const prefix = kind === 'audio' ? '/api/upload/audio/' : '/api/upload/image/';
if (!url.startsWith(prefix)) return false;
if (!url.startsWith(prefix)) return { isFresh: false, sizeBytes: BigInt(0) };
const filename = url.slice(prefix.length);
const key = kind === 'audio' ? `voice/${filename}` : `images/${filename}`;
@@ -32,10 +35,11 @@ async function isFreshAttachment(url: string, kind: 'audio' | 'image'): Promise<
Key: key,
})
);
if (!head.LastModified) return false;
return Date.now() - head.LastModified.getTime() <= UNATTACHED_UPLOAD_TTL_MS;
if (!head.LastModified) return { isFresh: false, sizeBytes: BigInt(0) };
const isFresh = Date.now() - head.LastModified.getTime() <= UNATTACHED_UPLOAD_TTL_MS;
return { isFresh, sizeBytes: BigInt(head.ContentLength ?? 0) };
} catch {
return false;
return { isFresh: false, sizeBytes: BigInt(0) };
}
}
@@ -188,6 +192,7 @@ export async function GET(request: NextRequest, { params }: RouteParams) {
// POST /api/versions/[versionId]/comments
export async function POST(request: NextRequest, { params }: RouteParams) {
let attachmentReservationId: string | null = null;
try {
const limited = await rateLimit(request, 'comment');
if (limited) return limited;
@@ -275,8 +280,12 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
// Validate annotation data structure to prevent prototype pollution and stored XSS.
// Reject anything that is not a well-formed array of AnnotationStroke objects.
// The HTTP body is already JSON-parsed by Next.js; double-encoded strings are rejected.
let serializedAnnotationData: string | null = null;
if (annotationData !== undefined && annotationData !== null) {
if (!Array.isArray(annotationData)) {
return apiErrors.badRequest('annotationData must be an array of valid stroke objects');
}
const validStrokes = validateAnnotationStrokes(annotationData);
if (validStrokes === null) {
return apiErrors.badRequest('annotationData must be an array of valid stroke objects');
@@ -317,21 +326,45 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
if (voiceUrl && !SAFE_AUDIO_PATH.test(voiceUrl)) {
return apiErrors.badRequest('Voice URL must reference an uploaded audio file');
}
if (voiceUrl && !(await isFreshAttachment(voiceUrl, 'audio'))) {
return apiErrors.badRequest('Voice upload expired. Please upload again.');
let voiceSizeBytes = BigInt(0);
if (voiceUrl) {
const voiceCheck = await isFreshAttachment(voiceUrl, 'audio');
if (!voiceCheck.isFresh) {
return apiErrors.badRequest('Voice upload expired. Please upload again.');
}
voiceSizeBytes = voiceCheck.sizeBytes;
}
if (imageUrl && !SAFE_IMAGE_PATH.test(imageUrl)) {
return apiErrors.badRequest('Image URL must reference an uploaded image file');
}
if (imageUrl && !(await isFreshAttachment(imageUrl, 'image'))) {
return apiErrors.badRequest('Image upload expired. Please upload again.');
let imageSizeBytes = BigInt(0);
if (imageUrl) {
const imageCheck = await isFreshAttachment(imageUrl, 'image');
if (!imageCheck.isFresh) {
return apiErrors.badRequest('Image upload expired. Please upload again.');
}
imageSizeBytes = imageCheck.sizeBytes;
}
const guestIdentity = isGuest ? ensureGuestIdentityFromRequest(request) : null;
// Use a transaction to create both comment and asset (if image is attached)
// Enforce per-workspace storage quota for any R2 attachments on this comment.
// Uses the advisory-locked reservation path so concurrent comment submissions
// see each other's in-flight sizes, eliminating the TOCTOU race.
const totalAttachmentBytes = voiceSizeBytes + imageSizeBytes;
if (totalAttachmentBytes > BigInt(0)) {
const reserveResult = await reserveStorageQuota(project.workspace.ownerId, totalAttachmentBytes);
if ('error' in reserveResult) return reserveResult.error;
attachmentReservationId = reserveResult.reservationId;
}
// Use a transaction to create both the comment and any asset rows atomically.
// Consume the reservation inside the transaction so quota is never double-counted.
const result = await db.$transaction(async (tx) => {
if (attachmentReservationId) {
await tx.uploadReservation.deleteMany({ where: { id: attachmentReservationId, billedUserId: project.workspace.ownerId } });
}
const comment = await tx.comment.create({
data: {
content: content?.trim() || null,
@@ -375,6 +408,29 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
displayName,
sourceUrl: imageUrl,
thumbnailUrl: imageUrl,
sizeBytes: imageSizeBytes,
uploadedByUserId: session?.user?.id || null,
uploadedByGuestIdentityId: isGuest ? guestIdentity?.identityId ?? null : null,
uploadedByGuestName: isGuest ? safeGuestName : null,
billedUserId: project.workspace.ownerId,
},
});
}
// If a voice recording was attached, also track it in the assets pane
if (voiceUrl) {
const fileName = extractAudioFileNameFromProxyUrl(voiceUrl);
const displayName = sanitizeAssetDisplayName(null, fileName || 'Voice Comment');
const safeGuestName = sanitizeAssetDisplayName(guestName, 'Guest').slice(0, 80);
await tx.videoAsset.create({
data: {
videoId: version.video.id,
kind: 'AUDIO',
provider: 'R2_AUDIO',
displayName,
sourceUrl: voiceUrl,
sizeBytes: voiceSizeBytes,
uploadedByUserId: session?.user?.id || null,
uploadedByGuestIdentityId: isGuest ? guestIdentity?.identityId ?? null : null,
uploadedByGuestName: isGuest ? safeGuestName : null,
@@ -450,6 +506,7 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
}
return withCacheControl(response, 'private, no-store');
} catch (error) {
await releaseStorageReservation(attachmentReservationId);
logError('Error creating comment:', error);
return apiErrors.internalError('Failed to create comment');
}
@@ -14,6 +14,7 @@ import { isBunnyUploadsFeatureEnabled } from '@/lib/feature-flags';
import { getShareSessionFromRequest } from '@/lib/share-session';
import { getVideoAssetAccessContext, SAFE_BUNNY_VIDEO_ID } from '@/lib/video-assets';
import { logError } from '@/lib/logger';
import { enforceStorageQuota } from '@/lib/storage-quota';
type RouteParams = { params: Promise<{ videoId: string }> };
@@ -36,6 +37,10 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
return apiErrors.badRequest('Direct uploads are disabled by this host');
}
const billedUserId = context.video.project.workspace.ownerId;
const quotaError = await enforceStorageQuota(billedUserId, BigInt(0));
if (quotaError) return quotaError;
const shareSession = getShareSessionFromRequest(request, context.video.id);
if (!context.viewerUserId) {
const quotaError = await enforceGuestUploadQuota(request, context.video.id, 'bunny', shareSession?.token ?? null);
+146 -43
View File
@@ -25,6 +25,13 @@ import {
sanitizeAssetDisplayName,
} from '@/lib/video-assets';
import { logError } from '@/lib/logger';
import { enforceStorageQuota, reserveStorageQuota, releaseStorageReservation, PLAN_STORAGE_LIMIT_BYTES } from '@/lib/storage-quota';
import { getCachedUserBunnyStorage } from '@/lib/admin-stats';
import { isStripeFeatureEnabled } from '@/lib/feature-flags';
// Sentinel thrown inside a Prisma transaction when a fake reservationId is
// supplied and the fallback quota check finds the limit would be exceeded.
class QuotaExceededInTxError extends Error {}
type RouteParams = { params: Promise<{ videoId: string }> };
@@ -147,35 +154,39 @@ async function fetchYouTubeTitle(videoId: string): Promise<string | null> {
return title;
}
async function isFreshImageAttachment(url: string): Promise<boolean> {
type AttachmentCheck = { isFresh: boolean; sizeBytes: bigint };
async function isFreshImageAttachment(url: string): Promise<AttachmentCheck> {
const key = extractImageKeyFromProxyUrl(url);
if (!key) return false;
if (!key) return { isFresh: false, sizeBytes: BigInt(0) };
try {
const head = await r2Client.send(new HeadObjectCommand({
Bucket: R2_BUCKET_NAME,
Key: key,
}));
if (!head.LastModified) return false;
return Date.now() - head.LastModified.getTime() <= UNATTACHED_UPLOAD_TTL_MS;
if (!head.LastModified) return { isFresh: false, sizeBytes: BigInt(0) };
const isFresh = Date.now() - head.LastModified.getTime() <= UNATTACHED_UPLOAD_TTL_MS;
return { isFresh, sizeBytes: BigInt(head.ContentLength ?? 0) };
} catch {
return false;
return { isFresh: false, sizeBytes: BigInt(0) };
}
}
async function isFreshAudioAttachment(url: string): Promise<boolean> {
async function isFreshAudioAttachment(url: string): Promise<AttachmentCheck> {
const key = extractAudioKeyFromProxyUrl(url);
if (!key) return false;
if (!key) return { isFresh: false, sizeBytes: BigInt(0) };
try {
const head = await r2Client.send(new HeadObjectCommand({
Bucket: R2_BUCKET_NAME,
Key: key,
}));
if (!head.LastModified) return false;
return Date.now() - head.LastModified.getTime() <= UNATTACHED_UPLOAD_TTL_MS;
if (!head.LastModified) return { isFresh: false, sizeBytes: BigInt(0) };
const isFresh = Date.now() - head.LastModified.getTime() <= UNATTACHED_UPLOAD_TTL_MS;
return { isFresh, sizeBytes: BigInt(head.ContentLength ?? 0) };
} catch {
return false;
return { isFresh: false, sizeBytes: BigInt(0) };
}
}
@@ -268,6 +279,7 @@ export async function GET(request: NextRequest, { params }: RouteParams) {
// POST /api/videos/[videoId]/assets
export async function POST(request: NextRequest, { params }: RouteParams) {
let reservationId: string | null = null;
try {
const limited = await rateLimit(request, 'asset-create');
if (limited) return limited;
@@ -293,20 +305,38 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
const guestIdentity = isGuest ? ensureGuestIdentityFromRequest(request) : null;
const requestedDisplayName = typeof body?.displayName === 'string' ? body.displayName : null;
// Optional reservation ID created by the upload route for atomic quota accounting
reservationId = typeof body?.reservationId === 'string' ? body.reservationId.trim() : null;
let displayName = '';
let sourceUrl = '';
let providerVideoId: string | null = null;
let thumbnailUrl: string | null = null;
let kind: 'IMAGE' | 'VIDEO' | 'AUDIO' = 'IMAGE';
let assetSizeBytes = BigInt(0);
const billedUserId = context.video.project.workspace.ownerId;
if (provider === VideoAssetProvider.R2_IMAGE) {
sourceUrl = typeof body?.sourceUrl === 'string' ? body.sourceUrl.trim() : '';
if (!SAFE_IMAGE_PROXY_PATH.test(sourceUrl)) {
return apiErrors.badRequest('Image URL must reference an uploaded image file');
}
if (!(await isFreshImageAttachment(sourceUrl))) {
const imageCheck = await isFreshImageAttachment(sourceUrl);
if (!imageCheck.isFresh) {
return apiErrors.badRequest('Image upload expired. Please upload again.');
}
assetSizeBytes = imageCheck.sizeBytes;
// Always use the advisory-locked reservation path so concurrent uploads
// see each other's in-flight sizes, eliminating the TOCTOU race. When
// the client already supplied a reservationId (new upload flow) the
// existing reservation is consumed in the transaction below. For the
// backward-compat path (no reservationId) we create one here.
if (!reservationId) {
const reserveResult = await reserveStorageQuota(billedUserId, assetSizeBytes);
if ('error' in reserveResult) return reserveResult.error;
reservationId = reserveResult.reservationId;
}
const fileName = extractImageFileNameFromProxyUrl(sourceUrl);
displayName = sanitizeAssetDisplayName(requestedDisplayName, fileName || 'Image');
@@ -319,9 +349,18 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
if (!SAFE_AUDIO_PROXY_PATH.test(sourceUrl)) {
return apiErrors.badRequest('Audio URL must reference an uploaded audio file');
}
if (!(await isFreshAudioAttachment(sourceUrl))) {
const audioCheck = await isFreshAudioAttachment(sourceUrl);
if (!audioCheck.isFresh) {
return apiErrors.badRequest('Audio upload expired. Please upload again.');
}
assetSizeBytes = audioCheck.sizeBytes;
// Same reservation logic as R2_IMAGE above
if (!reservationId) {
const reserveResult = await reserveStorageQuota(billedUserId, assetSizeBytes);
if ('error' in reserveResult) return reserveResult.error;
reservationId = reserveResult.reservationId;
}
const fileName = extractAudioFileNameFromProxyUrl(sourceUrl);
displayName = sanitizeAssetDisplayName(requestedDisplayName, fileName || 'Voice Recording');
@@ -405,40 +444,100 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
}
}
kind = 'VIDEO';
const quotaError = await enforceStorageQuota(billedUserId, BigInt(0));
if (quotaError) return quotaError;
}
const created = await db.videoAsset.create({
data: {
videoId: context.video.id,
kind,
provider,
displayName,
sourceUrl,
providerVideoId,
thumbnailUrl,
uploadedByUserId: context.viewerUserId,
uploadedByGuestIdentityId: context.viewerUserId ? null : guestIdentity?.identityId ?? null,
uploadedByGuestName: context.viewerUserId
? null
: sanitizeAssetDisplayName(typeof body?.guestName === 'string' ? body.guestName : null, 'Guest'),
billedUserId: context.video.project.workspace.ownerId,
},
select: {
id: true,
videoId: true,
kind: true,
provider: true,
displayName: true,
sourceUrl: true,
providerVideoId: true,
thumbnailUrl: true,
uploadedByGuestName: true,
createdAt: true,
updatedAt: true,
uploadedByUser: {
select: { id: true, name: true, image: true },
// Pre-fetch Bunny storage BEFORE entering the transaction to avoid making an
// HTTP call while holding a DB connection open (connection-pool exhaustion
// risk under adversarial load). Mirrors the discipline in reserveStorageQuota.
// Only needed for R2 providers where the invalid-reservation fallback quota
// check requires Bunny usage data.
const preFetchedBunnyData =
provider === VideoAssetProvider.R2_IMAGE || provider === VideoAssetProvider.R2_AUDIO
? await getCachedUserBunnyStorage()
: null;
// Create the VideoAsset and atomically consume the upload reservation (if any)
// so the spot is never double-counted.
const created = await db.$transaction(async (tx) => {
if (reservationId) {
// Acquire the per-user advisory lock unconditionally so both the happy path
// (valid reservation) and the fallback path (fake/expired reservation ID) are
// serialised — eliminating the TOCTOU race in the deleted.count === 0 branch.
await tx.$executeRaw`
SELECT pg_advisory_xact_lock(
('x' || left(md5(${billedUserId}), 16))::bit(64)::bigint
)
`;
// Validate the reservation by checking it actually exists and belongs to the
// billed user. A client-supplied fake ID would delete 0 rows — in that case
// we fall back to a standard (non-locked) quota check so the bypass attempt
// is caught rather than silently allowed.
const deleted = await tx.uploadReservation.deleteMany({
where: { id: reservationId, billedUserId, expiresAt: { gt: new Date() } },
});
if (deleted.count === 0) {
// Reservation didn't exist — enforce quota the normal way inside the tx.
// We read inside the same transaction so the check is at least consistent
// with the asset insert that follows.
const [r2Row] = await tx.$queryRaw<[{ total: bigint }]>`
SELECT COALESCE(SUM(size_bytes), 0)::bigint AS total
FROM video_assets
WHERE "billedUserId" = ${billedUserId}
AND provider IN ('R2_IMAGE', 'R2_AUDIO')
`;
const [resRow] = await tx.$queryRaw<[{ total: bigint }]>`
SELECT COALESCE(SUM("sizeBytes"), 0)::bigint AS total
FROM upload_reservations
WHERE "billedUserId" = ${billedUserId}
AND "expiresAt" > NOW()
`;
const bunnyData = preFetchedBunnyData ?? {};
const totalUsed =
(r2Row?.total ?? BigInt(0)) +
(resRow?.total ?? BigInt(0)) +
BigInt(bunnyData[billedUserId] ?? 0);
if (isStripeFeatureEnabled() && totalUsed + assetSizeBytes >= PLAN_STORAGE_LIMIT_BYTES) {
throw new QuotaExceededInTxError();
}
}
}
return tx.videoAsset.create({
data: {
videoId: context.video.id,
kind,
provider,
displayName,
sourceUrl,
providerVideoId,
thumbnailUrl,
sizeBytes: assetSizeBytes,
uploadedByUserId: context.viewerUserId,
uploadedByGuestIdentityId: context.viewerUserId ? null : guestIdentity?.identityId ?? null,
uploadedByGuestName: context.viewerUserId
? null
: sanitizeAssetDisplayName(typeof body?.guestName === 'string' ? body.guestName : null, 'Guest'),
billedUserId,
},
},
select: {
id: true,
videoId: true,
kind: true,
provider: true,
displayName: true,
sourceUrl: true,
providerVideoId: true,
thumbnailUrl: true,
uploadedByGuestName: true,
createdAt: true,
updatedAt: true,
uploadedByUser: {
select: { id: true, name: true, image: true },
},
},
});
});
const response = successResponse(shapeAssetForViewer(
@@ -451,6 +550,10 @@ export async function POST(request: NextRequest, { params }: RouteParams) {
}
return withCacheControl(response, 'private, no-store');
} catch (error) {
if (error instanceof QuotaExceededInTxError) {
return apiErrors.storageExceeded() as NextResponse;
}
await releaseStorageReservation(reservationId);
logError('Error creating video asset:', error);
return apiErrors.internalError('Failed to create asset');
}