local/auth.py used os.environ without importing os, so new_session raised NameError. Every first visit to /api/session, every registration and every login returned 500, which left Site checkout, cart recovery and the customer portal unusable since the R2 stack change. Read COOKIE_SECURE once as a module constant and share it with local/app.py instead of resolving the same variable in two places. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
355 lines
19 KiB
Python
355 lines
19 KiB
Python
import hashlib
|
|
import json
|
|
import math
|
|
import os
|
|
import secrets
|
|
from contextlib import asynccontextmanager
|
|
from datetime import datetime, timedelta, timezone
|
|
from uuid import UUID, uuid4
|
|
|
|
from botocore.exceptions import ClientError
|
|
from fastapi import Depends, FastAPI, HTTPException, Request, Response
|
|
from fastapi.responses import JSONResponse
|
|
from starlette.middleware.trustedhost import TrustedHostMiddleware
|
|
from psycopg.types.json import Jsonb
|
|
|
|
from . import db
|
|
from .adapters import FakeFreight, FakePayment, LocalS3Storage, require_runtime
|
|
from .models import Freight, Move, Pay, QuoteRequest, Review, UploadStart, OperatorLogin
|
|
from .pricing import price
|
|
from .auth import COOKIE_SECURE, owner, session_row, new_session, operator, throttle, audit, rate_limit
|
|
from .scanning import require_clean
|
|
|
|
require_runtime()
|
|
storage = LocalS3Storage()
|
|
payment = FakePayment()
|
|
freight = FakeFreight()
|
|
ENVIRONMENT = os.environ.get('APP_ENV', 'local')
|
|
PUBLIC_ORIGIN = os.environ.get('PUBLIC_ORIGIN', 'http://localhost')
|
|
ALLOWED_HOSTS = [host for host in os.environ.get('ALLOWED_HOSTS', 'localhost,127.0.0.1').split(',') if host]
|
|
ALLOWED_ORIGINS = [origin for origin in os.environ.get('ALLOWED_ORIGINS', PUBLIC_ORIGIN).split(',') if origin]
|
|
PART_BYTES = int(os.environ.get('UPLOAD_PART_BYTES', '8388608'))
|
|
if not 5242880 <= PART_BYTES <= 67108864:
|
|
raise RuntimeError('UPLOAD_PART_BYTES must be between 5 and 64 MiB')
|
|
STATES = {'rec': 'Arte recebida', 'tra': 'Arte tratada', 'fil': 'Fila de impressão',
|
|
'imp': 'Imprimindo', 'cor': 'Correção', 'fin': 'Finalizado'}
|
|
TRANSITIONS = {'rec': ['tra','cor'], 'tra': ['fil','cor'], 'fil': ['imp','cor'],
|
|
'imp': ['fin','cor'], 'cor': ['rec','tra'], 'fin': []}
|
|
|
|
@asynccontextmanager
|
|
async def lifespan(app):
|
|
with db.connect() as c:
|
|
c.execute('SELECT 1 FROM dtf_local.operator_sessions LIMIT 1')
|
|
storage.health()
|
|
yield
|
|
|
|
app = FastAPI(title='DTF Portal/API', lifespan=lifespan, docs_url=None, redoc_url=None)
|
|
app.add_middleware(TrustedHostMiddleware, allowed_hosts=ALLOWED_HOSTS)
|
|
|
|
@app.post('/api/operator/login')
|
|
def operator_login(body: OperatorLogin, request: Request, response: Response):
|
|
configured_email = os.environ.get('OPERATOR_EMAIL', '').strip().lower()
|
|
if not configured_email:
|
|
raise HTTPException(503, 'Kanban operator email is not configured')
|
|
email = body.email
|
|
throttle('operator:'+email, request)
|
|
valid_user = secrets.compare_digest(email.encode(), configured_email.encode())
|
|
valid_password = secrets.compare_digest(body.password.encode(), os.environ['OPERATOR_PASSWORD'].encode())
|
|
if not (valid_user and valid_password):
|
|
audit('operator_login_failed')
|
|
raise HTTPException(401, 'Invalid operator login')
|
|
token = secrets.token_urlsafe(32)
|
|
with db.connect() as c:
|
|
previous = hashlib.sha256(request.cookies.get('dtf_operator','').encode()).hexdigest()
|
|
c.execute('DELETE FROM dtf_local.operator_sessions WHERE token_hash=%s', (previous,))
|
|
c.execute('INSERT INTO dtf_local.operator_sessions(token_hash,username) VALUES(%s,%s)',
|
|
(hashlib.sha256(token.encode()).hexdigest(), email))
|
|
response.set_cookie('dtf_operator', token, httponly=True, secure=COOKIE_SECURE,
|
|
samesite='strict', path='/api/operator', max_age=28800)
|
|
audit('operator_login_success', operator=email)
|
|
return {'ok': True}
|
|
|
|
@app.post('/api/operator/logout')
|
|
def operator_logout(request: Request, response: Response):
|
|
with db.connect() as c:
|
|
digest = hashlib.sha256(request.cookies.get('dtf_operator','').encode()).hexdigest()
|
|
c.execute('DELETE FROM dtf_local.operator_sessions WHERE token_hash=%s', (digest,))
|
|
response.delete_cookie('dtf_operator', path='/api/operator', httponly=True,
|
|
secure=COOKIE_SECURE, samesite='strict')
|
|
audit('operator_logout')
|
|
return {'ok': True}
|
|
|
|
@app.middleware('http')
|
|
async def safe_headers(request, call_next):
|
|
if request.method not in ('GET','HEAD','OPTIONS'):
|
|
origin = request.headers.get('origin')
|
|
if request.headers.get('sec-fetch-site') == 'cross-site' or (origin and origin not in ALLOWED_ORIGINS):
|
|
audit('cross_origin_rejected')
|
|
return JSONResponse({'detail':'Cross-origin request rejected'}, status_code=403)
|
|
response = await call_next(request)
|
|
if response.status_code in (401,403,429) or response.status_code>=500:
|
|
audit('http_security_event', method=request.method, status=response.status_code)
|
|
response.headers['Cache-Control'] = 'no-store'
|
|
response.headers['X-Content-Type-Options'] = 'nosniff'
|
|
response.headers['Referrer-Policy'] = 'no-referrer'
|
|
return response
|
|
|
|
@app.get('/health')
|
|
@app.get('/api/health')
|
|
def health():
|
|
try:
|
|
with db.connect() as c:
|
|
c.execute('SELECT 1')
|
|
storage.health()
|
|
except Exception:
|
|
raise HTTPException(503, 'Database or storage unavailable')
|
|
return {'status': 'ok', 'environment': ENVIRONMENT,
|
|
'storage': 'minio' if ENVIRONMENT == 'local' else 'r2', 'integrations': 'fake'}
|
|
|
|
@app.get('/api/session')
|
|
def session(request: Request, response: Response):
|
|
try:
|
|
session_id = owner(request)
|
|
except HTTPException:
|
|
rate_limit('guest-sessions', ENVIRONMENT, 120, 900)
|
|
with db.connect() as c:
|
|
session_id = new_session(c, response)
|
|
return {'environment': ENVIRONMENT, 'cart_scope': str(session_id), 'part_bytes': PART_BYTES,
|
|
'max_upload_bytes': int(os.environ.get('MAX_UPLOAD_BYTES', '5368709120'))}
|
|
|
|
@app.post('/api/freight')
|
|
def quote_freight(body: Freight):
|
|
try:
|
|
return freight.quote(body.service, body.postal_code)
|
|
except ValueError as exc:
|
|
raise HTTPException(422, str(exc))
|
|
|
|
def upload_row(c, upload_id, session_id, lock=False):
|
|
row = c.execute('SELECT * FROM dtf_local.uploads WHERE id=%s AND owner=%s' +
|
|
(' FOR UPDATE' if lock else ''), (upload_id, session_id)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, 'Upload not found')
|
|
days = 30 if row['complete'] else 1
|
|
if row['purged_at'] or row['expires_at'] <= datetime.now(timezone.utc) or row['created_at'] < datetime.now(timezone.utc) - timedelta(days=days):
|
|
raise HTTPException(410, 'Upload expired; select the file again')
|
|
return row
|
|
|
|
@app.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}
|
|
|
|
@app.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]}
|
|
|
|
@app.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)}
|
|
|
|
@app.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}
|
|
|
|
@app.post('/api/quotes')
|
|
def create_quote(body: QuoteRequest, session_id=Depends(owner)):
|
|
draft = body.model_dump(mode='json', exclude={'request_key'})
|
|
digest = hashlib.sha256(json.dumps(draft, sort_keys=True).encode()).hexdigest()
|
|
try:
|
|
freight.quote(body.freight.service, body.freight.postal_code)
|
|
except ValueError as exc:
|
|
raise HTTPException(422, str(exc))
|
|
with db.connect() as c:
|
|
for item in body.items:
|
|
for uid in item.uploads:
|
|
row = upload_row(c, uid, session_id)
|
|
if not row['complete']:
|
|
raise HTTPException(409, 'Complete every upload before requesting a quote')
|
|
require_clean(row)
|
|
uid = uuid4()
|
|
c.execute('INSERT INTO dtf_local.quotes(id,owner,request_key,request_hash,draft) VALUES(%s,%s,%s,%s,%s) ON CONFLICT(owner,request_key) DO NOTHING',
|
|
(uid, session_id, body.request_key, digest, Jsonb(draft)))
|
|
row = c.execute('SELECT * FROM dtf_local.quotes WHERE owner=%s AND request_key=%s', (session_id, body.request_key)).fetchone()
|
|
if row['request_hash'] != digest:
|
|
raise HTTPException(409, 'Request key already used for a different cart')
|
|
return {'id': row['id'], 'status': 'pending_review'}
|
|
|
|
def quote_view(c, row):
|
|
order = c.execute('SELECT id,number,state FROM dtf_local.orders WHERE quote_id=%s', (row['id'],)).fetchone()
|
|
expired = row['approved_at'] and row['approved_at'] < datetime.now(timezone.utc)-timedelta(hours=24)
|
|
return {'id': row['id'], 'draft': row['draft'], 'approved': row['approved'],
|
|
'status': 'paid' if order else 'expired' if expired else 'approved' if row['approved'] else 'pending_review',
|
|
'order': order}
|
|
|
|
@app.get('/api/quotes/{uid}')
|
|
def get_quote(uid: UUID, session_id=Depends(owner)):
|
|
with db.connect() as c:
|
|
row = c.execute('SELECT * FROM dtf_local.quotes WHERE id=%s AND owner=%s', (uid,session_id)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, 'Quote not found')
|
|
return quote_view(c, row)
|
|
|
|
def enqueue(c, event_key, provider, payload):
|
|
c.execute('INSERT INTO dtf_local.outbox(event_key,provider,payload) VALUES(%s,%s,%s) ON CONFLICT(event_key) DO NOTHING',
|
|
(event_key, provider, Jsonb(payload)))
|
|
|
|
@app.post('/api/orders/dev-paid')
|
|
def dev_paid(body: Pay, session_id=Depends(owner)):
|
|
if ENVIRONMENT != 'local':
|
|
raise HTTPException(503, 'Checkout is not configured yet')
|
|
with db.connect() as c:
|
|
row = c.execute('SELECT * FROM dtf_local.quotes WHERE id=%s AND owner=%s FOR UPDATE', (body.quote_id,session_id)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, 'Quote not found')
|
|
existing = c.execute('SELECT * FROM dtf_local.orders WHERE quote_id=%s', (body.quote_id,)).fetchone()
|
|
if existing:
|
|
return existing
|
|
if not row['approved']:
|
|
raise HTTPException(409, 'An operator must verify length and grade first')
|
|
if row['approved_at'] < datetime.now(timezone.utc)-timedelta(hours=24):
|
|
raise HTTPException(409, 'Quote expired; request a new quote')
|
|
approved = row['approved']
|
|
for item in approved['items']:
|
|
for upload_id in item['uploads']:
|
|
require_clean(upload_row(c, UUID(upload_id), session_id))
|
|
paid = payment.pay(str(body.quote_id), approved['total_cents'])
|
|
result = c.execute('INSERT INTO dtf_local.orders(id,quote_id,owner,snapshot,payment) VALUES(%s,%s,%s,%s,%s) RETURNING *',
|
|
(uuid4(),body.quote_id,session_id,Jsonb(approved),Jsonb(paid))).fetchone()
|
|
for provider in ('tiny','whatsapp'):
|
|
enqueue(c, f"{result['id']}:paid:{provider}", provider,
|
|
{'order_id': str(result['id']), 'number': result['number'], 'event': 'payment_approved', 'order': approved})
|
|
return result
|
|
|
|
@app.get('/api/operator/board')
|
|
def board(user=Depends(operator)):
|
|
with db.connect() as c:
|
|
orders = c.execute('SELECT * FROM dtf_local.orders ORDER BY created_at').fetchall()
|
|
quotes = c.execute('SELECT q.* FROM dtf_local.quotes q LEFT JOIN dtf_local.orders o ON o.quote_id=q.id WHERE o.id IS NULL ORDER BY q.created_at').fetchall()
|
|
return {'states': STATES, 'transitions': TRANSITIONS, 'orders': orders,
|
|
'quotes': [quote_view(c, q) for q in quotes],
|
|
'events': c.execute('SELECT * FROM dtf_local.outbox ORDER BY id DESC LIMIT 100').fetchall()}
|
|
|
|
@app.post('/api/operator/quotes/{uid}/approve')
|
|
def approve(uid: UUID, body: Review, user=Depends(operator)):
|
|
with db.connect() as c:
|
|
row = c.execute('SELECT * FROM dtf_local.quotes WHERE id=%s FOR UPDATE', (uid,)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, 'Quote not found')
|
|
if row['approved']:
|
|
raise HTTPException(409, 'Approved quotes are immutable; request a new quote')
|
|
draft = row['draft']
|
|
if len(body.items) != len(draft['items']):
|
|
raise HTTPException(422, 'Review must cover every item')
|
|
items = []
|
|
for item, original in zip(body.items, draft['items']):
|
|
if item.mode != original['mode'] or list(map(str,item.uploads)) != original['uploads']:
|
|
raise HTTPException(422, 'Product mode and attached files cannot change during review')
|
|
for upload_id in item.uploads:
|
|
require_clean(upload_row(c, upload_id, row['owner']))
|
|
items.append({**price(item.mode, str(item.metres), item.grade), 'uploads': original['uploads']})
|
|
quoted_freight = freight.quote(**draft['freight'])
|
|
approved = {'customer': draft['customer'], 'items': items, 'freight': quoted_freight,
|
|
'total_cents': sum(i['total_cents'] for i in items)+quoted_freight['total_cents']}
|
|
c.execute('UPDATE dtf_local.quotes SET approved=%s, reviewed_by=%s, approved_at=now() WHERE id=%s', (Jsonb(approved),user,uid))
|
|
return approved
|
|
|
|
@app.post('/api/operator/orders/{uid}/move')
|
|
def move(uid: UUID, body: Move, user=Depends(operator)):
|
|
with db.connect() as c:
|
|
row = c.execute('SELECT * FROM dtf_local.orders WHERE id=%s FOR UPDATE', (uid,)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, 'Order not found')
|
|
if body.version != row['version']:
|
|
raise HTTPException(409, 'Order changed; refresh the board')
|
|
if body.state == row['state']:
|
|
return row
|
|
if body.state not in TRANSITIONS[row['state']]:
|
|
raise HTTPException(409, 'Move is not allowed from this state')
|
|
if body.state == 'cor' and not body.reason.strip():
|
|
raise HTTPException(422, 'Correction requires a reason')
|
|
if body.state in ('fil','imp'):
|
|
coverage = c.execute('SELECT DISTINCT f.item_index FROM dtf_local.order_files f JOIN dtf_local.uploads u ON u.id=f.upload_id WHERE f.order_id=%s AND f.kind=\'final\' AND f.active AND u.expires_at>now() AND u.purged_at IS NULL AND u.scan_state=\'clean\'', (uid,)).fetchall()
|
|
if {r['item_index'] for r in coverage} != set(range(len(row['snapshot']['items']))):
|
|
raise HTTPException(409, 'Approve a complete final-file set for every item before queueing')
|
|
if body.state == 'cor':
|
|
c.execute("UPDATE dtf_local.order_files SET active=false WHERE order_id=%s AND kind='final'", (uid,))
|
|
c.execute('INSERT INTO dtf_local.movements(order_id,from_state,to_state,operator,reason) VALUES(%s,%s,%s,%s,%s)',
|
|
(uid,row['state'],body.state,user,body.reason))
|
|
changed = c.execute('UPDATE dtf_local.orders SET state=%s, version=version+1, updated_at=now() WHERE id=%s RETURNING *', (body.state,uid)).fetchone()
|
|
events = {'imp':'production_started','cor':'correction_needed','fin':'ready'}
|
|
if body.state in events:
|
|
for provider in ('tiny','whatsapp'):
|
|
enqueue(c, f'{uid}:{changed["version"]}:{provider}', provider,
|
|
{'order_id':str(uid), 'number':row['number'], 'event':events[body.state], 'reason':body.reason,
|
|
'customer_path': f'/portal.html?order={uid}'})
|
|
return changed
|
|
|
|
@app.get('/api/operator/orders/{uid}/history')
|
|
def history(uid: UUID, user=Depends(operator)):
|
|
with db.connect() as c:
|
|
return c.execute('SELECT * FROM dtf_local.movements WHERE order_id=%s ORDER BY id', (uid,)).fetchall()
|
|
|
|
@app.get('/api/operator/uploads/{uid}/download')
|
|
def download(uid: UUID, user=Depends(operator)):
|
|
with db.connect() as c:
|
|
row = c.execute('SELECT * FROM dtf_local.uploads WHERE id=%s AND complete', (uid,)).fetchone()
|
|
if not row:
|
|
raise HTTPException(404, 'Completed upload not found')
|
|
if row['expires_at'] <= datetime.now(timezone.utc):
|
|
raise HTTPException(410, 'Artwork retention expired')
|
|
require_clean(row)
|
|
return {'name':row['name'], 'url':storage.download(row['object_key'],row['name']), 'expires_in':300}
|
|
|
|
from .customer import install_routes
|
|
install_routes(app, operator, storage, begin_upload, upload_status, part_url, complete_upload, upload_row, STATES)
|