Files
dtf-system/app/api/operator.py
Cauê Faleiros 64c1260a9f
All checks were successful
Build and deploy / Validate source (push) Successful in 5s
Build and deploy / Integration suite on a real stack (push) Successful in 2m16s
Build and deploy / Secret scan and release gate (push) Successful in 5s
Build and deploy / Publish images (push) Successful in 45s
feat: report Mercado Pago's notifications in the account check
Without access to the Mercado Pago panel, "Verificar conta" now also says
how many payment notifications arrived and passed the signature in the last
24 hours, how many were refused by it, and the last one received: the PIX
payments created in test mode notify the webhook, so this shows whether
Mercado Pago reaches the server.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-29 14:35:47 -03:00

365 lines
21 KiB
Python

"""Kanban: sign-in, the board, commercial review and card movement."""
import hashlib
import os
import secrets
from datetime import datetime, timedelta, timezone
from typing import Literal
from uuid import UUID
from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response
from fastapi.responses import RedirectResponse
from ..core import db
from ..core.auth import (COOKIE_SECURE, DUMMY_PASSWORD_HASH, audit, client_ip, operator,
login_failed, password_matches, throttle)
from ..core.models import Move, OperatorLogin, Resolution, Review
from ..printjobs import queue as queue_print_files
from .. import quote_review
from .. import tiny
from ..runtime import (BACK, BOARD_FINISHED_LIMIT, BOARD_QUOTE_LIMIT, ENVIRONMENT, STATES, TRANSITIONS, payment,
enqueue, freight, quote_view, storage)
from ..scanning import require_clean
router = APIRouter()
@router.post('/api/operator/login')
def operator_login(body: OperatorLogin, request: Request, response: Response):
email = body.email
throttle('operator:'+email, request)
with db.connect() as c:
account = c.execute('SELECT * FROM dtf_local.operators WHERE email=%s', (email,)).fetchone()
if not c.execute('SELECT 1 FROM dtf_local.operators WHERE active LIMIT 1').fetchone():
raise HTTPException(503, 'Nenhuma conta de operador configurada.')
# Comparable password work whether or not the account exists or is active.
stored = account['password_hash'] if account else DUMMY_PASSWORD_HASH
matches = password_matches(body.password, stored)
if not account or not account['active'] or not matches:
login_failed('operator:'+email)
audit('operator_login_failed', ip=client_ip(request), operator=email)
raise HTTPException(401, 'E-mail ou senha inválidos.')
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))
c.execute('UPDATE dtf_local.operators SET last_login_at=now() WHERE id=%s', (account['id'],))
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}
@router.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}
def freight_status():
"""What delivery the Site offers, and the package rule it prices with."""
if freight.name != 'jadlog':
return {'provider': freight.name}
return {'provider': 'jadlog', 'base_kg': str(freight.base_kg), 'per_metre_kg': str(freight.per_metre_kg),
'production_days': freight.production_days}
@router.get('/api/operator/board')
def board(user=Depends(operator)):
with db.connect() as c:
# Everything still in progress, however old: an operator must never lose a
# card they can act on. Finished orders are terminal and only accumulate,
# so the board carries a recent window of them and reports the true total.
active = c.execute("SELECT * FROM dtf_local.orders WHERE state<>'fin' ORDER BY created_at").fetchall()
finished = c.execute("SELECT * FROM dtf_local.orders WHERE state='fin' ORDER BY created_at DESC LIMIT %s",
(BOARD_FINISHED_LIMIT,)).fetchall()
finished_total = c.execute("SELECT count(*) AS n FROM dtf_local.orders WHERE state='fin'").fetchone()['n']
pending = 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 AND q.approved IS NULL
ORDER BY q.created_at DESC,q.id DESC LIMIT %s''', (BOARD_QUOTE_LIMIT,)).fetchall()
approved = 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 AND q.approved IS NOT NULL
ORDER BY q.created_at DESC,q.id DESC LIMIT 20''').fetchall()
pending_total = c.execute('''SELECT count(*) AS n FROM dtf_local.quotes q
LEFT JOIN dtf_local.orders o ON o.quote_id=q.id
WHERE o.id IS NULL AND q.approved IS NULL''').fetchone()['n']
approved_total = c.execute('''SELECT count(*) AS n FROM dtf_local.quotes q
LEFT JOIN dtf_local.orders o ON o.quote_id=q.id
WHERE o.id IS NULL AND q.approved IS NOT NULL''').fetchone()['n']
orders = with_print_files(c, active + list(reversed(finished)))
# A paid notification that did not become an order is money received
# for nothing the factory will make. The board carries the count; the
# Pagamentos tab pages through them until someone records a resolution.
issues_total = c.execute('''SELECT count(*) AS n FROM dtf_local.payment_events
WHERE (outcome LIKE 'refused%' OR outcome LIKE 'attention%') AND resolved_at IS NULL''').fetchone()['n']
return {'states': STATES, 'transitions': TRANSITIONS, 'back': BACK,
'orders': orders, 'payment_issues_total': issues_total, 'tiny': tiny_status(), 'operator': user,
'providers': {'payment': payment.name, 'freight': freight_status()}, 'environment': ENVIRONMENT,
'finished_shown': len(finished), 'finished_total': finished_total,
'quotes': [quote_view(c, q) for q in pending + approved],
'pending_total': pending_total, 'approved_total': approved_total}
def with_print_files(c, orders):
generated = c.execute('''SELECT p.order_id,p.item_index,p.status,p.upload_id,p.detail,u.name
FROM dtf_local.print_files p LEFT JOIN dtf_local.uploads u ON u.id=p.upload_id
WHERE p.order_id=ANY(%s) ORDER BY p.item_index''', ([o['id'] for o in orders],)).fetchall()
for order in orders:
order['print_files'] = [row for row in generated if row['order_id'] == order['id']]
return orders
def cursor_pair(first, second):
if (first is None) != (second is None):
raise HTTPException(422, 'Both cursor fields are required')
return first is not None
@router.get('/api/operator/orders/finished')
def finished_page(before_created_at: datetime | None = None, before_id: UUID | None = None,
limit: int = Query(default=50, ge=1, le=100), user=Depends(operator)):
"""Older finished orders, newest first, beyond the board's recent window."""
paged = cursor_pair(before_created_at, before_id)
with db.connect() as c:
rows = c.execute('''SELECT * FROM dtf_local.orders WHERE state='fin'
''' + ('AND (created_at,id)<(%s,%s) ' if paged else '') + '''
ORDER BY created_at DESC,id DESC LIMIT %s''',
((before_created_at, before_id) if paged else ()) + (limit + 1,)).fetchall()
return {'orders': with_print_files(c, rows[:limit]), 'has_more': len(rows) > limit}
@router.get('/api/operator/payment-events')
def payment_events(state: Literal['open','resolved','all'] = 'open',
offset: int = Query(default=0, ge=0),
limit: int = Query(default=20, ge=1, le=100), user=Depends(operator)):
"""Payments that needed a person, newest first: open ones to act on, resolved ones as history."""
where = {'open': 'AND resolved_at IS NULL', 'resolved': 'AND resolved_at IS NOT NULL', 'all': ''}[state]
with db.connect() as c:
rows = c.execute('''SELECT id,provider,event_id,reference,status,amount_cents,received_at,outcome,
resolved_at,resolved_by,resolution
FROM dtf_local.payment_events WHERE (outcome LIKE 'refused%%' OR outcome LIKE 'attention%%') ''' + where + '''
ORDER BY received_at DESC,id DESC LIMIT %s OFFSET %s''', (limit, offset)).fetchall()
total = c.execute('''SELECT count(*) AS n FROM dtf_local.payment_events
WHERE (outcome LIKE 'refused%' OR outcome LIKE 'attention%') ''' + where).fetchone()['n']
return {'issues': rows, 'total': total}
@router.get('/api/operator/events')
def events(provider: Literal['tiny','whatsapp'] | None = None,
status: Literal['delivered','queued','failing'] | None = None,
event: Literal['payment_approved','production_started','correction_needed','ready'] | None = None,
order: int | None = Query(default=None, ge=1), offset: int = Query(default=0, ge=0),
limit: int = Query(default=20, ge=1, le=100), user=Depends(operator)):
"""The integration send log, newest first, filtered, by page with a total."""
clauses, params = [], []
if provider:
clauses.append('provider=%s'); params.append(provider)
if status == 'delivered':
clauses.append('delivered_at IS NOT NULL')
elif status == 'queued':
clauses.append('delivered_at IS NULL AND last_error IS NULL')
elif status == 'failing':
clauses.append('delivered_at IS NULL AND last_error IS NOT NULL')
if event:
clauses.append("payload->>'event'=%s"); params.append(event)
if order:
clauses.append("payload->>'number'=%s"); params.append(str(order))
where = ('WHERE ' + ' AND '.join(clauses)) if clauses else ''
with db.connect() as c:
rows = c.execute(f'SELECT * FROM dtf_local.outbox {where} ORDER BY id DESC LIMIT %s OFFSET %s',
(*params, limit, offset)).fetchall()
total = c.execute(f'SELECT count(*) AS n FROM dtf_local.outbox {where}', params).fetchone()['n']
return {'events': rows, 'total': total}
@router.get('/api/operator/quotes')
def quote_page(kind: Literal['pending','approved'], before_created_at: datetime | None = None,
before_id: UUID | None = None, offset: int = Query(default=0, ge=0),
limit: int = Query(default=50, ge=1, le=100), user=Depends(operator)):
"""Unpaid quotes, newest first: by cursor, or by page (offset) with a total."""
if (before_created_at is None) != (before_id is None):
raise HTTPException(422, 'Both quote cursor fields are required')
approved_filter = 'q.approved IS NULL' if kind == 'pending' else 'q.approved IS NOT NULL'
cursor = 'AND (q.created_at,q.id)<(%s,%s)' if before_created_at else ''
skip = 0 if before_created_at else offset
params = ((before_created_at,before_id) if before_created_at else ()) + (limit+1, skip)
with db.connect() as c:
rows = c.execute(f'''SELECT q.* FROM dtf_local.quotes q
LEFT JOIN dtf_local.orders o ON o.quote_id=q.id
WHERE o.id IS NULL AND {approved_filter} {cursor}
ORDER BY q.created_at DESC,q.id DESC LIMIT %s OFFSET %s''', params).fetchall()
total = c.execute(f'''SELECT count(*) AS n FROM dtf_local.quotes q
LEFT JOIN dtf_local.orders o ON o.quote_id=q.id
WHERE o.id IS NULL AND {approved_filter}''').fetchone()['n']
return {'quotes':[quote_view(c,row) for row in rows[:limit]],
'has_more':len(rows)>limit, 'total': total}
@router.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')
return quote_review.approve(c, row, body.items, user)
@router.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, 'O pedido mudou. Clique em Atualizar.')
if body.state == row['state']:
return row
back = BACK.get(row['state']) == body.state
if body.state not in TRANSITIONS[row['state']] and not back:
raise HTTPException(409, f"Não dá para ir de {STATES[row['state']]} para {STATES[body.state]}.")
if body.state == 'cor' and not body.reason.strip():
raise HTTPException(422, 'Informe o motivo da correção.')
if back and not body.reason.strip():
raise HTTPException(422, 'Informe por que o pedido está voltando de etapa.')
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, 'Aprove os arquivos finais de todos os itens antes de colocar na fila.')
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,back) VALUES(%s,%s,%s,%s,%s,%s)',
(uid,row['state'],body.state,user,body.reason,back))
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 and not back:
# "Production started" and "ready" reach the customer once per order,
# even if a mistaken move is undone and made again. Each correction
# is a new request, so it keeps one message per movement.
once = body.state in ('imp','fin')
event = {'order_id':str(uid), 'number':row['number'], 'event':events[body.state], 'reason':body.reason,
'customer_path': f'/portal.html?order={uid}'}
# Tiny needs the order (customer, pickup or delivery) and, once the
# sale reached it, its id, to set the situação without a search.
sale = c.execute("SELECT receipt FROM dtf_local.outbox WHERE event_key=%s AND delivered_at IS NOT NULL",
(f'{uid}:paid:tiny',)).fetchone()
tiny = {**event, 'order': row['snapshot'], 'tiny_id': ((sale or {}).get('receipt') or {}).get('tiny_id')}
for provider, payload in (('tiny', tiny), ('whatsapp', event)):
key = f'{uid}:{events[body.state]}:{provider}' if once else f'{uid}:{changed["version"]}:{provider}'
enqueue(c, key, provider, payload)
return changed
@router.post('/api/operator/orders/{uid}/print-files')
def regenerate(uid: UUID, user=Depends(operator)):
"""Queue generation again for items that failed or went to manual preparation,
and for orders paid before the generator existed."""
with db.connect() as c:
row = c.execute('SELECT id,snapshot,state FROM dtf_local.orders WHERE id=%s FOR UPDATE', (uid,)).fetchone()
if not row:
raise HTTPException(404, 'Order not found')
if row['state'] not in ('rec','tra'):
raise HTTPException(409, 'Print files are generated only before the order is queued')
if c.execute("SELECT 1 FROM dtf_local.order_files WHERE order_id=%s AND kind='correction' LIMIT 1", (uid,)).fetchone():
raise HTTPException(409, 'A customer correction replaced the original artwork; prepare the final file by hand')
queue_print_files(c, uid, len(row['snapshot']['items']), only_missing=True)
audit('print_file_requeued', order=str(uid), operator=user)
return c.execute('SELECT * FROM dtf_local.print_files WHERE order_id=%s ORDER BY item_index', (uid,)).fetchall()
@router.post('/api/operator/payment-events/{uid}/resolve')
def resolve_payment(uid: UUID, body: Resolution, user=Depends(operator)):
with db.connect() as c:
row = c.execute('''UPDATE dtf_local.payment_events SET resolved_at=now(), resolved_by=%s, resolution=%s
WHERE id=%s AND (outcome LIKE 'refused%%' OR outcome LIKE 'attention%%') AND resolved_at IS NULL RETURNING id''',
(user, body.note, uid)).fetchone()
if not row:
raise HTTPException(404, 'No open payment issue with this id')
audit('payment_issue_resolved', payment_event=str(uid), operator=user)
return {'ok': True}
def tiny_status():
if not tiny.configured():
return {'configured': False}
return {'configured': True, 'orders_enabled': os.environ.get('TINY_ADAPTER') == 'tiny',
**tiny.TinyAuth().status()}
@router.post('/api/operator/tiny/connect')
def tiny_connect(user=Depends(operator)):
"""Start the one-time authorisation of this system in the client's Tiny."""
if not tiny.configured():
raise HTTPException(503, 'Tiny application is not configured')
audit('tiny_connect_started', operator=user)
return {'url': tiny.TinyAuth().authorize_url(user)}
@router.post('/api/operator/tiny/test')
def tiny_test(user=Depends(operator)):
"""Read one order and one contact from Tiny. Creates nothing."""
if not tiny.configured():
raise HTTPException(503, 'Tiny application is not configured')
try:
results = tiny.check()
except tiny.TinyNotConnected as exc:
raise HTTPException(409, str(exc))
audit('tiny_tested', operator=user, ok=all(v == 'ok' for v in results.values()))
return {'ok': all(v == 'ok' for v in results.values()), 'results': results}
@router.post('/api/operator/mercadopago/check')
def mercadopago_check(user=Depends(operator)):
"""The Mercado Pago account behind the configured token; reads only."""
if payment.name != 'mercadopago':
raise HTTPException(409, 'Mercado Pago não está configurado')
from ..mercadopago_probe import check
audit('mercadopago_check', operator=user)
result = check(payment.access_token)
# Whether Mercado Pago's notifications reach this server and pass the
# signature: every accepted one is recorded, a refused one is audited.
with db.connect() as c:
received = c.execute('''SELECT count(*) AS n, max(received_at) AS last FROM dtf_local.payment_events
WHERE provider='mercadopago' AND received_at > now()-interval '24 hours' ''').fetchone()
last = c.execute('''SELECT received_at,status,outcome FROM dtf_local.payment_events
WHERE provider='mercadopago' ORDER BY received_at DESC LIMIT 1''').fetchone()
refused = c.execute('''SELECT count(*) AS n FROM dtf_local.security_events
WHERE event='payment_webhook_rejected' AND created_at > now()-interval '24 hours' ''').fetchone()
result['webhooks'] = {'accepted_24h': received['n'], 'refused_24h': refused['n'], 'last': last}
return result
@router.get('/api/operator/tiny/callback')
def tiny_callback(code: str = Query('', max_length=4096), state: str = Query('', max_length=128),
error: str = Query('', max_length=128)):
"""Tiny's redirect back. Cross-site, so the operator cookie is absent: the
single-use state an operator created is what authorises it. Tiny reports a
refusal (the operator declined, or offline access is not allowed for this
application) with `error` instead of a code."""
if not tiny.configured():
raise HTTPException(503, 'Tiny application is not configured')
try:
if error == 'invalid_scope':
retry = tiny.TinyAuth().without_offline(state)
audit('tiny_offline_refused')
return RedirectResponse(retry, status_code=303)
if error or not code:
raise tiny.TinyError(f'Tiny returned {error or "no code"}')
who = tiny.TinyAuth().complete(code, state)
except tiny.TinyError:
audit('tiny_connect_failed')
return RedirectResponse('/?tiny=failed', status_code=303)
audit('tiny_connected', operator=who)
return RedirectResponse('/?tiny=connected', status_code=303)
@router.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()
@router.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}