"""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 .. import db from ..auth import audit, owner, rate_limit from ..models import UploadStart 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 > int(os.environ.get('MAX_UPLOAD_BYTES', '5368709120')): raise HTTPException(413, 'File exceeds the upload limit') 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, 'Local storage quota or pending upload limit reached') multipart = storage.begin(key) c.execute('INSERT INTO dtf_local.uploads(id,owner,name,size,object_key,multipart_id) VALUES(%s,%s,%s,%s,%s,%s)', (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') 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') c.execute('UPDATE dtf_local.uploads SET complete=true WHERE id=%s', (uid,)) return {'id': uid, 'complete': True}