feat: accept payment notifications, once, from a verified sender
There was no inbound payment path at all: a button called a fake synchronously and wrote an order. A real provider does the opposite — it charges, then tells us, repeatedly, out of order, and sometimes long afterwards. POST /api/payments/webhook verifies the signature before the body is parsed, so an unsigned or tampered delivery is refused and recorded without touching an order. Verified deliveries are stored under the provider's own event id with a unique constraint, and applied inside the same transaction that marks them processed: a repeat is a no-op, a crash is retried rather than half-applied. An approval whose amount disagrees with the reviewed quote does not become an order. Underpayment would ship artwork nobody paid for, and overpayment means something a person should look at. Order creation moved to app/payments.py so the webhook and the local development checkout share one implementation and cannot drift. That also closes 3.5: the charge happens inside the transaction that persists the order, rather than before it. The adapter contract is create/verify/parse. FakePayment implements it with a real HMAC scheme so the whole path is exercised now, by tests/payment_test.py: unsigned, tampered, underpaid, duplicate, re-sent, unknown reference, and non-approved statuses. Connecting Mercado Pago is one adapter; no service code changes. PAYMENT_WEBHOOK_SECRET is optional in production on purpose. Required would break the next Portainer render, and a guessable default would be worse than either: with no secret configured the adapter verifies nothing and therefore accepts nothing, which is the right state until a provider is connected. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -1,6 +1,9 @@
|
||||
"""Local-only composition root. No production provider implementations/imports."""
|
||||
import hashlib
|
||||
import hmac
|
||||
import json
|
||||
import os
|
||||
from typing import Protocol
|
||||
from typing import Mapping, NamedTuple, Protocol
|
||||
from urllib.parse import urlparse
|
||||
import boto3
|
||||
from botocore.config import Config
|
||||
@@ -41,10 +44,81 @@ def require_runtime():
|
||||
# Compatibility alias for local-only callers outside the active runtime.
|
||||
require_local = require_runtime
|
||||
|
||||
class PaymentEvent(NamedTuple):
|
||||
"""One provider notification, normalised.
|
||||
|
||||
`event_id` identifies the delivery and makes it idempotent. `reference` is
|
||||
our quote id, echoed back by the provider. `amount_cents` is what the
|
||||
provider says was actually paid, which the service compares against the
|
||||
approved total before it will create an order.
|
||||
"""
|
||||
event_id: str
|
||||
reference: str
|
||||
status: str # 'approved' | 'rejected' | 'pending' | 'refunded'
|
||||
amount_cents: int | None
|
||||
raw: dict
|
||||
|
||||
|
||||
class PaymentAdapter(Protocol):
|
||||
def pay(self, quote_id: str, total_cents: int) -> dict: ...
|
||||
def create(self, quote_id: str, total_cents: int, customer: dict) -> dict:
|
||||
"""Start a payment. Must be idempotent on quote_id: a retry after a
|
||||
timeout has to return the existing payment, never charge twice."""
|
||||
|
||||
def verify(self, headers: Mapping[str, str], body: bytes) -> bool:
|
||||
"""Whether this delivery genuinely came from the provider."""
|
||||
|
||||
def parse(self, body: bytes) -> PaymentEvent | None:
|
||||
"""Normalise a verified delivery, or None if it is not about a payment."""
|
||||
|
||||
|
||||
class FakePayment:
|
||||
"""Local stand-in with a real signature scheme, so the webhook path is
|
||||
exercised end to end rather than waiting for a provider account.
|
||||
|
||||
Signs the body with HMAC-SHA256 under PAYMENT_WEBHOOK_SECRET. A real adapter
|
||||
replaces verify() and parse() with the provider's own scheme; nothing else in
|
||||
the service changes.
|
||||
"""
|
||||
|
||||
header = 'x-payment-signature'
|
||||
|
||||
def _secret(self) -> bytes | None:
|
||||
secret = os.environ.get('PAYMENT_WEBHOOK_SECRET', '')
|
||||
return secret.encode() if secret else None
|
||||
|
||||
def create(self, quote_id: str, total_cents: int, customer: dict) -> dict:
|
||||
return {'provider': 'fake', 'id': f'local-{quote_id}',
|
||||
'status': 'pending', 'total_cents': total_cents}
|
||||
|
||||
def sign(self, body: bytes) -> str:
|
||||
secret = self._secret()
|
||||
if secret is None:
|
||||
raise RuntimeError('PAYMENT_WEBHOOK_SECRET is not configured')
|
||||
return hmac.new(secret, body, hashlib.sha256).hexdigest()
|
||||
|
||||
def verify(self, headers, body: bytes) -> bool:
|
||||
# No configured secret means nothing can be verified, so nothing is
|
||||
# accepted. A guessable default would let anyone forge an approval and
|
||||
# create an order that was never paid for.
|
||||
if self._secret() is None:
|
||||
return False
|
||||
supplied = headers.get(self.header) or headers.get(self.header.title()) or ''
|
||||
return hmac.compare_digest(supplied, self.sign(body))
|
||||
|
||||
def parse(self, body: bytes):
|
||||
try:
|
||||
data = json.loads(body)
|
||||
except ValueError:
|
||||
return None
|
||||
if not isinstance(data, dict) or 'event_id' not in data:
|
||||
return None
|
||||
return PaymentEvent(event_id=str(data['event_id']),
|
||||
reference=str(data.get('reference', '')),
|
||||
status=str(data.get('status', 'pending')),
|
||||
amount_cents=data.get('amount_cents'),
|
||||
raw=data)
|
||||
|
||||
# The local development checkout still needs a direct "it is paid" path.
|
||||
def pay(self, quote_id: str, total_cents: int) -> dict:
|
||||
return {'provider': 'fake', 'id': f'local-{quote_id}',
|
||||
'status': 'paid', 'total_cents': total_cents}
|
||||
|
||||
@@ -1,41 +1,38 @@
|
||||
"""Paid orders. Local development payment only; no provider is wired yet."""
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from uuid import UUID, uuid4
|
||||
"""Local development checkout.
|
||||
|
||||
The real path is the provider webhook. This exists so the local stack can reach
|
||||
a paid order without a provider account, and it goes through the same service so
|
||||
the two cannot drift apart.
|
||||
"""
|
||||
from uuid import UUID
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException
|
||||
from psycopg.types.json import Jsonb
|
||||
|
||||
from .. import payments
|
||||
from ..core import db
|
||||
from ..core.auth import owner
|
||||
from ..core.models import Pay
|
||||
from ..runtime import ENVIRONMENT, enqueue, payment, upload_row
|
||||
from ..scanning import require_clean
|
||||
from ..runtime import ENVIRONMENT, payment
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
||||
@router.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
|
||||
try:
|
||||
quote = payments.approved_quote(c, body.quote_id, session_id)
|
||||
except payments.PaymentRefused as refusal:
|
||||
# The quote may already be paid; that is not a refusal.
|
||||
existing = c.execute('SELECT * FROM dtf_local.orders WHERE quote_id=%s',
|
||||
(body.quote_id,)).fetchone()
|
||||
if existing:
|
||||
return existing
|
||||
raise HTTPException(404 if 'not found' in str(refusal) else 409, str(refusal))
|
||||
# The charge happens inside the transaction that persists the order, so a
|
||||
# failure to record it cannot leave a customer charged without an order.
|
||||
receipt = payment.pay(str(body.quote_id), quote['approved']['total_cents'])
|
||||
order, _ = payments.create_order(c, quote, receipt)
|
||||
return order
|
||||
|
||||
53
app/api/payments.py
Normal file
53
app/api/payments.py
Normal file
@@ -0,0 +1,53 @@
|
||||
"""The provider's callback.
|
||||
|
||||
Unauthenticated by necessity — a payment provider has no session — so the
|
||||
signature is the only thing standing between this endpoint and an attacker
|
||||
creating orders. It is verified before the body is parsed, let alone acted on,
|
||||
and an unverified delivery is recorded and refused rather than retried.
|
||||
"""
|
||||
from fastapi import APIRouter, HTTPException, Request
|
||||
|
||||
from .. import payments
|
||||
from ..core import db
|
||||
from ..core.auth import audit, client_ip, rate_limit
|
||||
from ..runtime import payment
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
# Generous: a provider legitimately retries, and a signature check is cheap.
|
||||
# This exists so an unsigned flood cannot keep the database busy.
|
||||
WEBHOOK_LIMIT = 600
|
||||
|
||||
|
||||
@router.post('/api/payments/webhook')
|
||||
async def webhook(request: Request):
|
||||
rate_limit('payment-webhook', client_ip(request), WEBHOOK_LIMIT, 900)
|
||||
body = await request.body()
|
||||
|
||||
if not payment.verify(request.headers, body):
|
||||
audit('payment_webhook_rejected', ip=client_ip(request), reason='signature')
|
||||
raise HTTPException(403, 'Invalid signature')
|
||||
|
||||
event = payment.parse(body)
|
||||
if event is None:
|
||||
# Verified, so genuinely from the provider, but not about a payment.
|
||||
# Acknowledge it: refusing would make the provider retry for ever.
|
||||
return {'status': 'ignored'}
|
||||
|
||||
with db.connect() as c:
|
||||
stored = payments.record(c, event_provider(), event)
|
||||
if stored is None:
|
||||
# Already delivered. Acknowledge without acting again.
|
||||
return {'status': 'duplicate'}
|
||||
outcome = payments.apply(c, event)
|
||||
c.execute('UPDATE dtf_local.payment_events SET processed_at=now(), outcome=%s WHERE id=%s',
|
||||
(outcome, stored['id']))
|
||||
|
||||
# audit()'s own first parameter is named `event`, so the id goes under another key.
|
||||
audit('payment_webhook_applied', payment_event=event.event_id,
|
||||
status=event.status, outcome=outcome)
|
||||
return {'status': 'applied', 'outcome': outcome}
|
||||
|
||||
|
||||
def event_provider():
|
||||
return getattr(payment, 'name', payment.__class__.__name__.replace('Payment', '').lower() or 'fake')
|
||||
@@ -13,7 +13,7 @@ from starlette.middleware.trustedhost import TrustedHostMiddleware
|
||||
from .core import db
|
||||
from .core.auth import audit, client_ip
|
||||
from .runtime import ALLOWED_HOSTS, ALLOWED_ORIGINS, storage
|
||||
from .api import artwork, customer, health, operator, orders, quotes, uploads
|
||||
from .api import artwork, customer, health, operator, orders, payments, quotes, uploads
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
@@ -44,5 +44,5 @@ async def safe_headers(request, call_next):
|
||||
|
||||
|
||||
# Order is not significant: no two routers declare the same path.
|
||||
for module in (health, uploads, quotes, orders, operator, customer, artwork):
|
||||
for module in (health, uploads, quotes, orders, payments, operator, customer, artwork):
|
||||
app.include_router(module.router)
|
||||
|
||||
103
app/payments.py
Normal file
103
app/payments.py
Normal file
@@ -0,0 +1,103 @@
|
||||
"""Turning a payment into an order, once.
|
||||
|
||||
A provider may deliver the same notification several times, out of order, or
|
||||
long after the fact. None of that may produce a second order, a second charge,
|
||||
or a second WhatsApp message. Every delivery is recorded under the provider's
|
||||
own event id and applied inside one transaction, so a duplicate is a no-op and a
|
||||
crash mid-way is retried rather than half-applied.
|
||||
|
||||
Order creation lives here rather than in a route because two paths reach it: the
|
||||
webhook, and the local development checkout. They must agree.
|
||||
"""
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from uuid import UUID, uuid4
|
||||
|
||||
from psycopg.types.json import Jsonb
|
||||
|
||||
from .core.auth import audit
|
||||
from .runtime import enqueue, upload_row
|
||||
from .scanning import require_clean
|
||||
|
||||
QUOTE_VALID_HOURS = 24
|
||||
|
||||
|
||||
class PaymentRefused(Exception):
|
||||
"""The payment cannot become an order, with a reason worth recording."""
|
||||
|
||||
|
||||
def approved_quote(c, quote_id, owner=None):
|
||||
"""The reviewed quote behind a payment, or a refusal explaining why not."""
|
||||
sql = 'SELECT * FROM dtf_local.quotes WHERE id=%s' + (' AND owner=%s' if owner else '')
|
||||
row = c.execute(sql + ' FOR UPDATE', (quote_id, owner) if owner else (quote_id,)).fetchone()
|
||||
if not row:
|
||||
raise PaymentRefused('quote not found')
|
||||
if not row['approved']:
|
||||
raise PaymentRefused('quote was never reviewed')
|
||||
if row['approved_at'] < datetime.now(timezone.utc) - timedelta(hours=QUOTE_VALID_HOURS):
|
||||
raise PaymentRefused('quote expired before payment')
|
||||
return row
|
||||
|
||||
|
||||
def create_order(c, quote, payment):
|
||||
"""Create the order for a reviewed quote, or return the one already there.
|
||||
|
||||
Returns (order, created). The caller decides what to do about a duplicate;
|
||||
the important part is that asking twice cannot produce two orders, because
|
||||
orders.quote_id is unique and this runs inside the caller's transaction.
|
||||
"""
|
||||
existing = c.execute('SELECT * FROM dtf_local.orders WHERE quote_id=%s', (quote['id'],)).fetchone()
|
||||
if existing:
|
||||
return existing, False
|
||||
|
||||
approved = quote['approved']
|
||||
for item in approved['items']:
|
||||
for upload_id in item['uploads']:
|
||||
require_clean(upload_row(c, UUID(upload_id), quote['owner']))
|
||||
|
||||
order = c.execute(
|
||||
'INSERT INTO dtf_local.orders(id,quote_id,owner,snapshot,payment) VALUES(%s,%s,%s,%s,%s) RETURNING *',
|
||||
(uuid4(), quote['id'], quote['owner'], Jsonb(approved), Jsonb(payment))).fetchone()
|
||||
for provider in ('tiny', 'whatsapp'):
|
||||
enqueue(c, f"{order['id']}:paid:{provider}", provider,
|
||||
{'order_id': str(order['id']), 'number': order['number'],
|
||||
'event': 'payment_approved', 'order': approved})
|
||||
return order, True
|
||||
|
||||
|
||||
def record(c, provider, event):
|
||||
"""Store a delivery. Returns None if this exact event was already seen."""
|
||||
inserted = c.execute(
|
||||
'''INSERT INTO dtf_local.payment_events(id,provider,event_id,reference,status,amount_cents,payload)
|
||||
VALUES(%s,%s,%s,%s,%s,%s,%s) ON CONFLICT(provider,event_id) DO NOTHING RETURNING *''',
|
||||
(uuid4(), provider, event.event_id, event.reference, event.status,
|
||||
event.amount_cents, Jsonb(event.raw))).fetchone()
|
||||
return inserted
|
||||
|
||||
|
||||
def apply(c, event):
|
||||
"""Act on a payment notification. Returns the outcome recorded against it."""
|
||||
if event.status != 'approved':
|
||||
return f'ignored: {event.status}'
|
||||
|
||||
try:
|
||||
quote_id = UUID(event.reference)
|
||||
except (ValueError, AttributeError):
|
||||
return 'refused: reference is not a quote id'
|
||||
|
||||
try:
|
||||
quote = approved_quote(c, quote_id)
|
||||
except PaymentRefused as refusal:
|
||||
return f'refused: {refusal}'
|
||||
|
||||
# The provider is the authority on what was paid, and the reviewed quote is
|
||||
# the authority on what was owed. If they disagree, no order is created:
|
||||
# underpayment would ship artwork that was not paid for, and overpayment
|
||||
# means something is wrong that a person should look at.
|
||||
expected = quote['approved']['total_cents']
|
||||
if event.amount_cents is not None and event.amount_cents != expected:
|
||||
audit('payment_amount_mismatch', quote=str(quote_id),
|
||||
expected_cents=expected, paid_cents=event.amount_cents)
|
||||
return f'refused: paid {event.amount_cents} but quote total is {expected}'
|
||||
|
||||
order, created = create_order(c, quote, {'provider': 'webhook', **event.raw})
|
||||
return f"order {order['number']}" + ('' if created else ' (already existed)')
|
||||
@@ -64,6 +64,13 @@ CREATE TABLE IF NOT EXISTS dtf_local.operators (
|
||||
password_hash text NOT NULL, active boolean NOT NULL DEFAULT true,
|
||||
created_at timestamptz NOT NULL DEFAULT now(), last_login_at timestamptz
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS dtf_local.payment_events (
|
||||
id uuid PRIMARY KEY, provider text NOT NULL, event_id text NOT NULL,
|
||||
reference text, status text NOT NULL, amount_cents bigint,
|
||||
payload jsonb NOT NULL, received_at timestamptz NOT NULL DEFAULT now(),
|
||||
processed_at timestamptz, outcome text,
|
||||
UNIQUE(provider, event_id)
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS dtf_local.security_events (
|
||||
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
|
||||
event text NOT NULL, details jsonb NOT NULL, created_at timestamptz NOT NULL DEFAULT now()
|
||||
@@ -118,3 +125,9 @@ CREATE INDEX IF NOT EXISTS operator_sessions_username ON dtf_local.operator_sess
|
||||
|
||||
-- security_status reads recent events; the worker prunes old ones by age.
|
||||
CREATE INDEX IF NOT EXISTS security_events_created ON dtf_local.security_events(created_at);
|
||||
|
||||
-- The webhook looks an event up by provider and id on every delivery, and the
|
||||
-- unique constraint already indexes that pair. Only the unprocessed sweep needs
|
||||
-- its own index, and it stays the size of the backlog.
|
||||
CREATE INDEX IF NOT EXISTS payment_events_unprocessed ON dtf_local.payment_events(received_at)
|
||||
WHERE processed_at IS NULL;
|
||||
|
||||
Reference in New Issue
Block a user