"""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 ..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, '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,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') 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,expires_at=now()+interval '30 days' WHERE id=%s", (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}