feat: run DTF stack with Cloudflare R2
This commit is contained in:
@@ -6,18 +6,40 @@ import boto3
|
||||
from botocore.config import Config
|
||||
from botocore.exceptions import ClientError
|
||||
|
||||
def require_local():
|
||||
if os.environ.get('APP_ENV') != 'local':
|
||||
raise RuntimeError('This runtime only supports APP_ENV=local')
|
||||
def require_runtime():
|
||||
"""Validate the supported local and R2-backed deployment modes.
|
||||
|
||||
Real payment, freight, ERP, and WhatsApp providers are deliberately not
|
||||
enabled yet. A non-local deployment keeps those adapters fake and disables
|
||||
checkout until their audited implementations are added.
|
||||
"""
|
||||
environment = os.environ.get('APP_ENV', 'local')
|
||||
for name in ('PAYMENT', 'FREIGHT', 'TINY', 'WHATSAPP'):
|
||||
if os.environ.get(f'{name}_ADAPTER') != 'fake':
|
||||
raise RuntimeError(f'{name} must use the fake adapter')
|
||||
if os.environ.get('STORAGE_ADAPTER') != 's3-local':
|
||||
raise RuntimeError('Only local S3 storage is supported')
|
||||
raise RuntimeError(f'{name} must use the currently supported fake adapter')
|
||||
if environment == 'local':
|
||||
if os.environ.get('STORAGE_ADAPTER') != 's3-local':
|
||||
raise RuntimeError('Local runtime requires local S3 storage')
|
||||
for name in ('S3_ENDPOINT', 'S3_PUBLIC_ENDPOINT'):
|
||||
endpoint = urlparse(os.environ[name])
|
||||
if endpoint.scheme != 'http' or endpoint.hostname not in ('storage', 'localhost', '127.0.0.1'):
|
||||
raise RuntimeError(f'{name} must point to local MinIO')
|
||||
return
|
||||
if environment != 'production':
|
||||
raise RuntimeError('APP_ENV must be local or production')
|
||||
if os.environ.get('STORAGE_ADAPTER') != 's3-r2':
|
||||
raise RuntimeError('Production runtime requires R2 storage')
|
||||
for name in ('S3_ENDPOINT', 'S3_PUBLIC_ENDPOINT'):
|
||||
endpoint = urlparse(os.environ[name])
|
||||
if endpoint.scheme != 'http' or endpoint.hostname not in ('storage', 'localhost', '127.0.0.1'):
|
||||
raise RuntimeError(f'{name} must point to local MinIO')
|
||||
if endpoint.scheme != 'https' or not (endpoint.hostname or '').endswith('.r2.cloudflarestorage.com'):
|
||||
raise RuntimeError(f'{name} must be a Cloudflare R2 S3 API endpoint')
|
||||
for name in ('AWS_ACCESS_KEY_ID', 'AWS_SECRET_ACCESS_KEY'):
|
||||
if not os.environ.get(name):
|
||||
raise RuntimeError(f'{name} is required for R2')
|
||||
|
||||
|
||||
# Compatibility alias for local-only callers outside the active runtime.
|
||||
require_local = require_runtime
|
||||
|
||||
class PaymentAdapter(Protocol):
|
||||
def pay(self, quote_id: str, total_cents: int) -> dict: ...
|
||||
|
||||
34
local/app.py
34
local/app.py
@@ -14,16 +14,21 @@ from starlette.middleware.trustedhost import TrustedHostMiddleware
|
||||
from psycopg.types.json import Jsonb
|
||||
|
||||
from . import db
|
||||
from .adapters import FakeFreight, FakePayment, LocalS3Storage, require_local
|
||||
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 owner, session_row, new_session, operator, throttle, audit, rate_limit
|
||||
from .scanning import require_clean
|
||||
|
||||
require_local()
|
||||
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]
|
||||
COOKIE_SECURE = os.environ.get('COOKIE_SECURE', 'false').lower() == 'true'
|
||||
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')
|
||||
@@ -39,8 +44,8 @@ async def lifespan(app):
|
||||
storage.health()
|
||||
yield
|
||||
|
||||
app = FastAPI(title='DTF Local Portal/API', lifespan=lifespan, docs_url=None, redoc_url=None)
|
||||
app.add_middleware(TrustedHostMiddleware, allowed_hosts=['localhost', '127.0.0.1'])
|
||||
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):
|
||||
@@ -49,14 +54,15 @@ def operator_login(body: OperatorLogin, request: Request, response: Response):
|
||||
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 local operator login')
|
||||
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(), body.username))
|
||||
response.set_cookie('dtf_operator', token, httponly=True, samesite='strict', path='/api/operator', max_age=28800)
|
||||
response.set_cookie('dtf_operator', token, httponly=True, secure=COOKIE_SECURE,
|
||||
samesite='strict', path='/api/operator', max_age=28800)
|
||||
audit('operator_login_success', operator=body.username)
|
||||
return {'ok': True}
|
||||
|
||||
@@ -65,7 +71,8 @@ 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, samesite='strict')
|
||||
response.delete_cookie('dtf_operator', path='/api/operator', httponly=True,
|
||||
secure=COOKIE_SECURE, samesite='strict')
|
||||
audit('operator_logout')
|
||||
return {'ok': True}
|
||||
|
||||
@@ -73,7 +80,7 @@ def operator_logout(request: Request, response: Response):
|
||||
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 != 'http://'+request.headers.get('host','')):
|
||||
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)
|
||||
@@ -93,17 +100,18 @@ def health():
|
||||
storage.health()
|
||||
except Exception:
|
||||
raise HTTPException(503, 'Database or storage unavailable')
|
||||
return {'status': 'ok', 'environment': 'local', 'storage': 'minio', 'integrations': 'fake'}
|
||||
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', 'local-stack', 120, 900)
|
||||
rate_limit('guest-sessions', ENVIRONMENT, 120, 900)
|
||||
with db.connect() as c:
|
||||
session_id = new_session(c, response)
|
||||
return {'environment': 'local', 'cart_scope': str(session_id), 'part_bytes': PART_BYTES,
|
||||
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')
|
||||
@@ -126,7 +134,7 @@ def upload_row(c, upload_id, session_id, lock=False):
|
||||
@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 local upload limit')
|
||||
raise HTTPException(413, 'File exceeds the upload limit')
|
||||
uid = uuid4()
|
||||
key = f'originals/{uid}'
|
||||
rate_limit('upload-start', str(session_id), 60, 900)
|
||||
@@ -234,6 +242,8 @@ def enqueue(c, event_key, provider, 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:
|
||||
|
||||
@@ -42,7 +42,9 @@ def new_session(c, response, identity=None):
|
||||
sid = uuid4()
|
||||
identity = identity or uuid4()
|
||||
c.execute('INSERT INTO dtf_local.sessions(id,owner) VALUES(%s,%s)', (sid,identity))
|
||||
response.set_cookie('dtf_session', str(sid), httponly=True, samesite='strict', max_age=86400*7)
|
||||
response.set_cookie('dtf_session', str(sid), httponly=True,
|
||||
secure=os.environ.get('COOKIE_SECURE', 'false').lower() == 'true',
|
||||
samesite='strict', max_age=86400*7)
|
||||
return identity
|
||||
|
||||
def transfer_guest(c, previous, identity):
|
||||
|
||||
@@ -31,7 +31,7 @@ class ClamAV:
|
||||
|
||||
def scan(self, stream, size):
|
||||
if size > min(134217728, int(os.environ.get('SCAN_MAX_BYTES','134217728'))):
|
||||
return 'rejected', 'File exceeds the local malware scan limit'
|
||||
return 'rejected', 'File exceeds the malware scan limit'
|
||||
with socket.create_connection(('scanner',3310),timeout=10) as sock:
|
||||
sock.settimeout(150)
|
||||
sock.sendall(b'zINSTREAM\0')
|
||||
@@ -64,7 +64,7 @@ def scan_one(storage, scanner=None):
|
||||
try:state,reason=scanner.scan(stream,row['size'])
|
||||
finally:stream.close()
|
||||
except Exception:
|
||||
state,reason='error','Local malware scanner unavailable; file remains blocked'
|
||||
state,reason='error','Malware scanner unavailable; file remains blocked'
|
||||
c.execute("""UPDATE dtf_local.uploads SET scan_state=%s,scan_reason=%s,scanned_at=now(),
|
||||
scan_after=now()+interval '1 minute',
|
||||
expires_at=CASE WHEN %s IN ('rejected','error') THEN LEAST(expires_at,now()+interval '3 days') ELSE expires_at END
|
||||
|
||||
@@ -5,11 +5,11 @@ import threading
|
||||
import time
|
||||
from http.server import BaseHTTPRequestHandler, HTTPServer
|
||||
from psycopg.types.json import Jsonb
|
||||
from .adapters import FakeTiny, FakeWhatsApp, LocalS3Storage, require_local
|
||||
from .adapters import FakeTiny, FakeWhatsApp, LocalS3Storage, require_runtime
|
||||
from .db import connect
|
||||
from .scanning import ClamAV, scan_loop
|
||||
|
||||
require_local()
|
||||
require_runtime()
|
||||
adapters = {'tiny': FakeTiny(), 'whatsapp': FakeWhatsApp()}
|
||||
last_tick = 0.0
|
||||
last_cleanup = 0.0
|
||||
|
||||
Reference in New Issue
Block a user