"""Customer uploads: reservation, signed parts, completion. The operator artwork routes reuse these functions directly with a different identity, which is why they are plain functions with a session_id argument. """ import math import os from uuid import UUID, uuid4 from botocore.exceptions import ClientError from fastapi import APIRouter, Depends, HTTPException, Request from starlette.concurrency import run_in_threadpool from ..core import db from ..core.auth import audit, owner, rate_limit from ..core.limits import upload_limit_bytes from ..core.models import UploadStart from ..payments import UNPAID_HOLD from ..runtime import PART_BYTES, storage, upload_row router = APIRouter() @router.post('/api/uploads') def begin_upload(body: UploadStart, session_id=Depends(owner)): if body.size > upload_limit_bytes(): raise HTTPException(413, 'File exceeds the malware scan limit; select a smaller file') uid = uuid4() key = f'originals/{uid}' rate_limit('upload-start', str(session_id), 60, 900) with db.connect() as c: # Serialize reservations across API processes, including changing guest IDs. c.execute('SELECT pg_advisory_xact_lock(804208)') usage = c.execute('''SELECT COALESCE(sum(size),0) AS total, COALESCE(sum(size) FILTER(WHERE owner=%s),0) AS owned, count(*) FILTER(WHERE owner=%s AND NOT complete) AS pending FROM dtf_local.uploads WHERE purged_at IS NULL''', (session_id,session_id)).fetchone() if (usage['total']+body.size > int(os.environ.get('STORAGE_QUOTA_BYTES','53687091200')) or usage['owned']+body.size > int(os.environ.get('OWNER_UPLOAD_QUOTA_BYTES','10737418240')) or usage['pending'] >= int(os.environ.get('MAX_PENDING_UPLOADS','10'))): audit('upload_quota_rejected') raise HTTPException(429, 'Limite de armazenamento ou de envios pendentes atingido. Tente de novo mais tarde.') multipart = storage.begin(key) c.execute('''INSERT INTO dtf_local.uploads(id,owner,name,size,object_key,multipart_id,expires_at) VALUES(%s,%s,%s,%s,%s,%s,now()+interval '1 hour')''', (uid, session_id, body.name, body.size, key, multipart)) return {'id': uid, 'part_bytes': PART_BYTES} @router.get('/api/uploads/{uid}') def upload_status(uid: UUID, session_id=Depends(owner)): with db.connect() as c: row = upload_row(c, uid, session_id) parts = [] if row['complete'] else storage.parts(row['object_key'], row['multipart_id']) return {'id': uid, 'complete': row['complete'], 'part_bytes': PART_BYTES, 'scan_state':row['scan_state'], 'scan_reason':row['scan_reason'], 'parts': [p['PartNumber'] for p in parts]} @router.post('/api/uploads/{uid}/parts/{part}') def part_url(uid: UUID, part: int, session_id=Depends(owner)): with db.connect() as c: row = upload_row(c, uid, session_id) if row['complete'] or not 1 <= part <= math.ceil(row['size'] / PART_BYTES): raise HTTPException(409, 'Invalid part or completed upload') # The reservation lease is an hour; a multi-GB upload on a slow line # takes longer, so each part it asks for keeps it alive. c.execute("UPDATE dtf_local.uploads SET expires_at=GREATEST(expires_at,now()+interval '1 hour') WHERE id=%s", (uid,)) size = min(PART_BYTES, row['size']-(part-1)*PART_BYTES) return {'url': storage.part_url(row['object_key'], row['multipart_id'], part, size)} @router.post('/api/uploads/{uid}/complete') def complete_upload(uid: UUID, session_id=Depends(owner)): with db.connect() as c: row = upload_row(c, uid, session_id, lock=True) if row['complete']: return {'id': uid, 'complete': True} try: existing_size = storage.size(row['object_key']) except ClientError as exc: if exc.response['ResponseMetadata']['HTTPStatusCode'] != 404: raise existing_size = None if existing_size is None: parts = storage.parts(row['object_key'], row['multipart_id']) expected = math.ceil(row['size'] / PART_BYTES) if [p['PartNumber'] for p in parts] != list(range(1, expected+1)) or any( p['Size'] != min(PART_BYTES, row['size'] - i*PART_BYTES) for i,p in enumerate(parts)): raise HTTPException(409, 'Parts are missing or their sizes do not match') storage.complete(row['object_key'], row['multipart_id'], parts) existing_size = storage.size(row['object_key']) if existing_size != row['size']: raise HTTPException(409, 'Stored size differs from declared size') # Held while the cart is unpaid: an abandoned cart's files go after # UNPAID_HOLD; a paid order keeps them for its 30 days (app/payments.py). c.execute('UPDATE dtf_local.uploads SET complete=true,expires_at=now()+%s::interval WHERE id=%s', (UNPAID_HOLD, uid)) return {'id': uid, 'complete': True} # The browser's small picture of the artwork, for the Kanban. Only an image, and # only a small one: the type is read from the bytes, never from the header. THUMBNAIL_BYTES = 300_000 def thumbnail_type(data): if data[:4] == b'RIFF' and data[8:12] == b'WEBP': return 'image/webp' if data[:3] == b'\xff\xd8\xff': return 'image/jpeg' if data[:8] == b'\x89PNG\r\n\x1a\n': return 'image/png' return None def store_thumbnail(uid, session_id, mime, data): rate_limit('upload-thumbnail', str(session_id), 120, 900) with db.connect() as c: upload_row(c, uid, session_id) c.execute('''INSERT INTO dtf_local.upload_thumbnails(upload_id,mime,data) VALUES(%s,%s,%s) ON CONFLICT(upload_id) DO UPDATE SET mime=EXCLUDED.mime, data=EXCLUDED.data, created_at=now()''', (uid, mime, data)) @router.put('/api/uploads/{uid}/thumbnail') async def put_thumbnail(uid: UUID, request: Request, session_id=Depends(owner)): if int(request.headers.get('content-length') or 0) > THUMBNAIL_BYTES: raise HTTPException(413, 'Thumbnail too large') data = await request.body() if len(data) > THUMBNAIL_BYTES: raise HTTPException(413, 'Thumbnail too large') mime = thumbnail_type(data) if not mime: raise HTTPException(415, 'Thumbnail must be a WebP, JPEG or PNG image') await run_in_threadpool(store_thumbnail, uid, session_id, mime, data) return {'id': uid} @router.delete('/api/uploads/{uid}') def cancel_upload(uid: UUID, session_id=Depends(owner)): with db.connect() as c: row=upload_row(c,uid,session_id,lock=True) if row['complete']: raise HTTPException(409, 'Completed upload cannot be cancelled') storage.discard(row['object_key'],row['multipart_id'],False) c.execute('UPDATE dtf_local.uploads SET purged_at=now() WHERE id=%s',(uid,)) c.execute('DELETE FROM dtf_local.upload_thumbnails WHERE upload_id=%s',(uid,)) audit('upload_cancelled', upload=str(uid)) return {'id':uid,'cancelled':True}