"""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 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} @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,)) audit('upload_cancelled', upload=str(uid)) return {'id':uid,'cancelled':True}