← Files paytechARCHIVED FILE
skills/psp-payments/references/code-examples-python.md
35.1 KB · Oct 3, 2026 · 06:35 UTC
# Python / FastAPI — PSP integration code
Rules and the failures they prevent: `references/integration-patterns.md`. When the project outgrows
the baseline (several instances, real concurrency on one order, crash-during-POST, sweep jobs, plus
the session-per-request and no-transaction-over-I/O traps): `references/hardening-concurrency.md`.
The code below is the **baseline** level.
**Adapt, don't transplant:** reuse the project's HTTP client, ORM, logger, config and test
framework. Field names, endpoints and states are fixed by the API; everything else is yours.
## 1. PSP client
```python
# psp/client.py
import os
from decimal import Decimal
from typing import Any
import httpx
from pydantic import BaseModel, ConfigDict, ValidationError
# Documented state enum: CHECKOUT | PENDING | AWAITING_APPROVAL | AUTHORIZED | COMPLETED
# | DECLINED | CANCELLED. Deliberately NOT a Literal on the field below — see _one().
class PaymentResult(BaseModel):
model_config = ConfigDict(extra="ignore") # the API sends many more fields, all optional
id: str
state: str # str, never Literal[...]: see _one()
referenceId: str | None = None
paymentType: str | None = None # DEPOSIT | WITHDRAWAL | REFUND
amount: Decimal | None = None # decimal MAJOR units: 10.01 == 10.01 GBP
currency: str | None = None
redirectUrl: str | None = None
parentPaymentId: str | None = None
errorCode: str | None = None
errorMessage: str | None = None
class PspTimeout(RuntimeError):
"""Call did not complete: the outcome is UNKNOWN — never treat it as a failure."""
class PspApiError(RuntimeError):
def __init__(self, status: int, body: Any) -> None:
super().__init__(f"PSP HTTP {status}")
self.status, self.body = status, body
def env(name: str) -> str:
value = os.environ.get(name)
if not value:
raise RuntimeError(f"{name} is not set") # fail fast at startup
return value
def _one(payload: Any) -> PaymentResult:
"""Validate INSIDE the fail-safe boundary. `PaymentResult.model_validate(...)` at the call site
runs *after* the request was sent, and a ValidationError is neither PspTimeout nor PspApiError:
it escapes the checkout's `except` clause and leaves the attempt claimed with nothing to resolve
it. Same reason `state` is a plain `str` — a state the API adds later must reach the transition
whitelist (which ignores what it does not know), and in the webhook handler a ValidationError is
a 500, so the PSP redelivers that event forever instead of it being acked and ignored."""
try:
return PaymentResult.model_validate(payload)
except ValidationError as exc:
raise PspTimeout(f"undecodable payment body: {exc.error_count()} error(s)") from exc
class PspClient:
def __init__(self, client: httpx.AsyncClient | None = None) -> None:
# An injected client must carry base_url, the auth header AND a timeout: with `timeout=None`
# a POST can hang forever, so no PspTimeout is ever raised and nothing reconciles.
self._http = client or httpx.AsyncClient(
base_url=env("PSP_API_URL").rstrip("/"),
timeout=httpx.Timeout(connect=5.0, read=30.0, write=10.0, pool=5.0),
headers={"Authorization": f"Bearer {env('PSP_API_KEY')}", # never log client.headers
"Content-Type": "application/json",
"User-Agent": "psp-integration/1.0"}) # some WL hosts (WAF) 403 the default client UA
async def _call(self, method: str, url: str, json: dict | None = None) -> Any:
# Classify FAIL-SAFE: only a real HTTP status tells you what the PSP did.
# Everything else must become PspTimeout so it reaches the reconcile path —
# letting a raw exception escape leaves the attempt claimed with nothing to
# resolve it, which wedges the order behind a permanent 409.
# Catch httpx.HTTPError, NOT TransportError: TimeoutException is one subclass of
# TransportError, and TransportError itself misses two endings that happen AFTER the
# request went out — DecodingError (a proxy answering with a broken Content-Encoding)
# and TooManyRedirects. Both are RequestError. HTTPStatusError cannot occur here
# because raise_for_status() is never called, so no real status is ever mislabelled.
try:
r = await self._http.request(method, url, json=json)
except httpx.HTTPError as exc:
raise PspTimeout(f"{method} {url}: {type(exc).__name__}") from exc
if r.is_error: # status is KNOWN: 4xx/5xx, never retry a 4xx
try:
body: Any = r.json() # structured error body...
except ValueError:
body = r.text[:500] # ...unless a proxy answered HTML
raise PspApiError(r.status_code, body)
try:
return r.json()["result"] # responses are {timestamp, status, result}
except (ValueError, KeyError, TypeError) as exc: # non-JSON, or JSON that is not the
raise PspTimeout( # envelope (a bare list/scalar subscripted
f"{method} {url}: undecodable 2xx body") from exc # by "result" -> TypeError)
async def create_deposit(self, *, amount: Decimal, currency: str, reference_id: str,
return_url: str, webhook_url: str, customer: dict | None = None,
billing_address: dict | None = None) -> PaymentResult:
body = {"paymentType": "DEPOSIT", "amount": float(amount), "currency": currency,
"referenceId": reference_id, "returnUrl": return_url, "webhookUrl": webhook_url,
"customer": customer, "billingAddress": billing_address}
return _one(await self._call(
"POST", "/api/v1/payments", {k: v for k, v in body.items() if v is not None}))
async def create_refund(self, *, parent_payment_id: str, amount: Decimal, currency: str,
reference_id: str) -> PaymentResult:
# No refund endpoint: a REFUND is a new payment linked via parentPaymentId.
return _one(await self._call("POST", "/api/v1/payments", {
"paymentType": "REFUND", "parentPaymentId": parent_payment_id,
"amount": float(amount), "currency": currency, "referenceId": reference_id}))
async def get_payment(self, payment_id: str) -> PaymentResult:
return _one(await self._call("GET", f"/api/v1/payments/{payment_id}"))
async def find_by_reference_id(self, reference_id: str) -> list[PaymentResult]: # reconciliation
result = await self._call("GET", f"/api/v1/payments?referenceId.eq={reference_id}")
# `result or []`: an empty match can arrive as `null`, and iterating None on the ONE path
# that must never fail unclassified — reconciliation — escapes as a raw TypeError.
return [_one(p) for p in (result or [])]
```
## 2. Webhook endpoint — the raw-body recipe
```python
# psp/signature.py
import base64, hashlib, hmac
from psp.client import env
_SIGNING_KEY = env("PSP_SIGNING_KEY").encode()
def signature_valid(raw_body: bytes, header: str | None) -> bool:
if not header:
return False
mac = hmac.new(_SIGNING_KEY, raw_body, hashlib.sha256).digest()
# BYTES, not str: compare_digest("…", "é") raises TypeError ("comparing strings with non-ASCII
# characters is not supported"), so one crafted header turns an unauthenticated 401 into a 500.
presented = header.strip().encode()
# Encoding (hex vs base64) is NOT documented — accept both, then pin the one your sandbox
# sends and delete the other branch (authentication.md §3). compare_digest = constant time.
return (hmac.compare_digest(mac.hex().encode(), presented.lower())
or hmac.compare_digest(base64.b64encode(mac), presented))
```
```python
# api/webhooks.py
import logging
from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Request, Response
from db import get_session # YOUR session-per-request dependency: one AsyncSession
from orders.transition import apply_webhook # per request, never one shared across requests
from psp.client import PaymentResult
from psp.signature import signature_valid
router, log = APIRouter(), logging.getLogger(__name__)
@router.post("/webhooks/psp")
async def psp_webhook(request: Request, background: BackgroundTasks,
session=Depends(get_session)) -> Response:
# RAW bytes. Do NOT declare a Pydantic body parameter (e.g. `event: PaymentResult`) on this
# endpoint: FastAPI would consume and re-parse the stream, and re-serialised JSON never
# reproduces the byte sequence the HMAC was computed over.
raw: bytes = await request.body()
if not signature_valid(raw, request.headers.get("Signature")):
log.warning("psp webhook signature mismatch (%d bytes)", len(raw)) # never log the key
raise HTTPException(status_code=401, detail="invalid signature")
event = PaymentResult.model_validate_json(raw) # parse only after verifying
outcome = await apply_webhook(session, event, background)
log.info("psp webhook id=%s state=%s outcome=%s", event.id, event.state, outcome)
return Response(status_code=200) # always 2xx once verified
```
## 3. Order state transition (idempotent, DB-guarded)
```python
# orders/transition.py (SQLAlchemy 2.x async, PostgreSQL)
from sqlalchemy import func, or_, select, update
from sqlalchemy.dialects.postgresql import insert
from db import orders, webhook_event # YOUR tables/mapped classes — integration-patterns.md schema
from orders.fulfilment import fulfil_order # YOUR heavy post-payment work: e-mails, ledger, …
_STATE_MAP = {"COMPLETED": "PAID", "AUTHORIZED": "AUTHORIZED", # whitelist: integration-patterns.md
"DECLINED": "PAYMENT_FAILED", "CANCELLED": "PAYMENT_FAILED"}
_ALLOWED_FROM = {"PAID": ("AWAITING_PAYMENT", "PROCESSING", "AUTHORIZED"),
"AUTHORIZED": ("AWAITING_PAYMENT", "PROCESSING"),
"PAYMENT_FAILED": ("AWAITING_PAYMENT", "PROCESSING", "AUTHORIZED")}
async def apply_webhook(session, event, background) -> str:
next_status = _STATE_MAP.get(event.state)
if next_status is None:
return "ignored" # non-final / unknown state: change nothing
async with session.begin():
# INBOX claim, not a tombstone: `do update ... where processed_at is null` takes over (and
# row-locks) a receipt that was recorded but never applied, so concurrent redeliveries
# serialise here; no row back means it was already applied = a real duplicate.
claim = (await session.execute(insert(webhook_event)
.values(payment_id=event.id, state=event.state)
.on_conflict_do_update(index_elements=["payment_id", "state"],
set_={"received_at": func.now()},
where=webhook_event.c.processed_at.is_(None))
.returning(webhook_event.c.id))).first()
if claim is None:
return "duplicate"
match = [orders.c.psp_payment_id == event.id]
if event.referenceId:
match.append(orders.c.order_ref == event.referenceId)
# Conditional UPDATE: re-application and any downgrade of a final status match 0 rows. Book
# amount/currency FROM THE PAYLOAD (the final amount may differ from the requested one).
values = {"status": next_status, "psp_payment_id": event.id, "updated_at": func.now(),
"error_code": event.errorCode, "error_message": event.errorMessage}
if event.amount is not None:
values |= {"paid_amount": event.amount, "paid_currency": event.currency}
res = await session.execute(update(orders)
.where(or_(*match), orders.c.status.in_(_ALLOWED_FROM[next_status]))
.values(**values).returning(orders.c.id))
row = res.first()
if row is None:
known = (await session.execute(select(orders.c.id).where(or_(*match)))).first()
# The webhook can beat the create-payment response. Commit the RECEIPT but leave
# processed_at NULL: marking it processed here loses the event forever, because the
# redelivery would be dismissed as a duplicate and the order would never transition.
if known is None:
return "unknown" # ack with 200 so the PSP stops redelivering
# 0 rows with a known order = already past this transition: the receipt is settled either way.
await session.execute(update(webhook_event)
.where(webhook_event.c.id == claim.id)
.values(processed_at=func.now(),
processed_reason="applied" if row is not None else "duplicate"))
if row is None:
return "duplicate"
if next_status == "PAID":
background.add_task(fulfil_order, row.id) # heavy work, after the 200
return "applied"
```
## 4. Creation and refund idempotency
```python
# orders/checkout.py
import uuid
from decimal import Decimal
from sqlalchemy import select, text, update
from sqlalchemy.dialects.postgresql import insert
from db import orders, psp_attempt, psp_refund_attempt # YOUR tables/mapped classes
from psp.client import PspApiError, PspClient, PspTimeout
class CheckoutInProgress(RuntimeError):
"""The attempt state is churning -> HTTP 409 + Retry-After."""
class PaymentOutcomeUnknown(RuntimeError):
"""Reconciliation ran and is STILL inconclusive -> HTTP 409 + Retry-After. The attempt stays
claimed, so the next call reconciles again: never a dead end."""
class CheckoutFailed(RuntimeError):
"""The attempt is FAILED and the order is free again: the client may start a NEW checkout."""
class RefundOutcomeUnknown(RuntimeError):
"""A committed refund attempt is unresolved -> HTTP 409; reconcile it, never start a new one."""
_ACTIVE = ("IN_FLIGHT", "READY") # FAILED excluded: it must never block a new attempt
_ATTEMPT = (psp_attempt.c.id, psp_attempt.c.order_id, psp_attempt.c.reference_id,
psp_attempt.c.state, psp_attempt.c.redirect_url)
async def _claim_attempt(session, order_id: int):
"""Atomic get-or-create, committed BEFORE the PSP call. Returns (attempt, owner): only the caller
that INSERTED the row may POST."""
for _ in range(2): # 2nd pass: the active attempt turned FAILED between the two statements
async with session.begin():
fresh = (await session.execute(insert(psp_attempt)
.values(order_id=order_id, reference_id=f"order-{order_id}-{uuid.uuid4()}",
state="IN_FLIGHT")
.on_conflict_do_nothing(
index_elements=["order_id"],
index_where=text("state in ('IN_FLIGHT','READY')"))
.returning(*_ATTEMPT))).first()
if fresh is not None:
return fresh, True
active = (await session.execute(select(*_ATTEMPT).where(
psp_attempt.c.order_id == order_id,
psp_attempt.c.state.in_(_ACTIVE)))).first()
if active is not None:
return active, False
raise CheckoutInProgress(order_id)
async def start_checkout(session, psp: PspClient, order_id: int,
amount: Decimal, currency: str) -> str | None:
attempt, owner = await _claim_attempt(session, order_id)
if not owner:
return await _join_attempt(session, psp, attempt) # a non-owner never POSTs, never mints
try:
# Outside session.begin(): the claim is committed, and no transaction is held over the call.
payment = await psp.create_deposit(
amount=amount, currency=currency, reference_id=attempt.reference_id,
return_url="https://shop.example/return/{id}/{referenceId}/{state}/{type}",
webhook_url="https://shop.example/webhooks/psp",
customer={"referenceId": f"customer_{order_id}"})
except (PspTimeout, PspApiError) as exc:
if isinstance(exc, PspApiError) and exc.status < 500:
await _set_attempt_state(session, attempt.id, "FAILED") # confirmed: nothing created
raise # ...so the order is free
# Timeout or ambiguous 5xx: the payment may exist. Reconcile, never a second POST.
return await resolve_attempt(session, psp, attempt)
return await _promote_ready(session, attempt, payment)
async def _join_attempt(session, psp: PspClient, attempt):
"""READY -> the stored URL. Still IN_FLIGHT -> reconcile, so a 409 always follows real progress.
(One extra GET per concurrent click; the cheaper UNKNOWN split: hardening-concurrency.md §1.)"""
if attempt.state == "READY":
return attempt.redirect_url
return await resolve_attempt(session, psp, attempt)
async def resolve_attempt(session, psp: PspClient, attempt):
"""The only way out of an unresolved attempt: GET by the PERSISTED referenceId."""
found = await psp.find_by_reference_id(attempt.reference_id)
if not found: # stays claimed: retried by the next call
raise PaymentOutcomeUnknown(attempt.reference_id)
if found[0].state in ("DECLINED", "CANCELLED"):
await _set_attempt_state(session, attempt.id, "FAILED") # frees the order for a NEW attempt
raise CheckoutFailed(found[0].errorCode or found[0].state)
return await _promote_ready(session, attempt, found[0])
async def _promote_ready(session, attempt, payment) -> str | None:
async with session.begin():
await session.execute(update(psp_attempt).where(psp_attempt.c.id == attempt.id)
.values(state="READY", psp_payment_id=payment.id, redirect_url=payment.redirectUrl))
await session.execute(update(orders) # AWAITING_PAYMENT only: a webhook may have won
.where(orders.c.id == attempt.order_id, orders.c.status == "AWAITING_PAYMENT")
.values(status="PROCESSING", psp_payment_id=payment.id))
return payment.redirectUrl # CHECKOUT, not paid; None once the payment moved past checkout
async def _set_attempt_state(session, attempt_id: int, state: str) -> None:
async with session.begin():
await session.execute(update(psp_attempt)
.where(psp_attempt.c.id == attempt_id).values(state=state))
_REFUND = (psp_refund_attempt.c.id, psp_refund_attempt.c.order_id,
psp_refund_attempt.c.reference_id, psp_refund_attempt.c.amount,
psp_refund_attempt.c.state, psp_refund_attempt.c.psp_payment_id)
async def refund_order(session, psp: PspClient, order_id: int, refund_key: str,
amount: Decimal, currency: str):
"""Idempotent per (order_id, refund_key) — refund_key identifies ONE logical refund. ONLY the
caller that INSERTED the attempt may POST: referenceId is NOT an idempotency key at the PSP, so
a second POST for the same refund is a second payout."""
attempt, owner, parent_id = await reserve_refund(session, order_id, refund_key, amount, currency)
if attempt.psp_payment_id:
return await psp.get_payment(attempt.psp_payment_id) # DONE/FAILED: pure replay
if not owner:
# Someone else's (or an earlier crashed) attempt: reconcile or 409, never POST.
return await _reconcile_refund(session, psp, attempt)
try:
result = await psp.create_refund(parent_payment_id=parent_id, amount=amount,
currency=currency, reference_id=attempt.reference_id)
except PspTimeout:
# Outcome UNKNOWN. The row is already committed: reconcile by ITS referenceId.
return await _reconcile_refund(session, psp, attempt)
await _settle_refund(session, attempt, result)
return result
async def reserve_refund(session, order_id, refund_key, amount, currency): # public: tests use it
"""Returns (attempt, owner, parent_payment_id). ONE commit for the attempt AND the amount
reservation, before the PSP call. A row that already existed belongs to another call, so its
amount must NOT be reserved a second time and its caller must NOT POST."""
async with session.begin():
order = (await session.execute(select(orders.c.psp_payment_id)
.where(orders.c.id == order_id))).first()
if order is None:
raise ValueError("unknown order")
row = (await session.execute(insert(psp_refund_attempt)
.values(order_id=order_id, refund_key=refund_key, amount=amount, currency=currency,
reference_id=f"refund-{order_id}-{uuid.uuid4()}", state="IN_FLIGHT")
.on_conflict_do_nothing(index_elements=["order_id", "refund_key"])
.returning(*_REFUND))).first()
if row is None: # the row exists: reuse ITS referenceId, reserve nothing
row = (await session.execute(select(*_REFUND).where(
psp_refund_attempt.c.order_id == order_id,
psp_refund_attempt.c.refund_key == refund_key))).first()
return row, False, order.psp_payment_id
reserved = await session.execute(update(orders) # refund only the remainder
.where(orders.c.id == order_id, orders.c.status == "PAID",
orders.c.refunded_amount + amount <= orders.c.paid_amount)
.values(refunded_amount=orders.c.refunded_amount + amount))
if reserved.rowcount == 0:
raise ValueError("refund exceeds remaining refundable amount")
return row, True, order.psp_payment_id
async def _reconcile_refund(session, psp: PspClient, attempt):
# GET by the PERSISTED referenceId. A fresh referenceId here is a SECOND payout.
found = await psp.find_by_reference_id(attempt.reference_id)
if not found:
raise RefundOutcomeUnknown(f"refund {attempt.reference_id} unresolved; retry reconciliation")
await _settle_refund(session, attempt, found[0])
return found[0]
async def _settle_refund(session, attempt, result):
failed = result.state in ("DECLINED", "CANCELLED")
async with session.begin():
# State-conditional: the owner and a reconciler can settle the same attempt, and releasing
# the reservation twice would inflate the refundable amount.
done = await session.execute(update(psp_refund_attempt)
.where(psp_refund_attempt.c.id == attempt.id,
psp_refund_attempt.c.state == "IN_FLIGHT")
.values(state="FAILED" if failed else "DONE", psp_payment_id=result.id))
if failed and done.rowcount: # only a CONFIRMED failure gives the amount back, not a timeout
await session.execute(update(orders).where(orders.c.id == attempt.order_id)
.values(refunded_amount=orders.c.refunded_amount - attempt.amount))
```
## 5. Tests (pytest + respx + httpx.ASGITransport)
```python
import asyncio, base64, hashlib, hmac, json
from decimal import Decimal
import httpx, pytest, respx
from orders.checkout import (PaymentOutcomeUnknown, RefundOutcomeUnknown,
refund_order, reserve_refund, start_checkout)
from psp.client import PspApiError, PspClient, PspTimeout
BASE = "https://sandbox.psp.invalid" # sandbox stub; never production credentials
PAY = f"{BASE}/api/v1/payments"
SIGNING_KEY = b"test-signing-key"
BODY = json.dumps({"id": "pay1", "referenceId": "order-1-a", "state": "COMPLETED",
"amount": 10.01, "currency": "GBP"}).encode()
CHECKOUT_1 = {"id": "pay1", "state": "CHECKOUT", "redirectUrl": "https://checkout.example/pay1"}
REFUND_5 = {"state": "COMPLETED", "paymentType": "REFUND", "amount": 5.00, "currency": "GBP"}
def sign(body: bytes, encoding: str = "hex") -> str:
mac = hmac.new(SIGNING_KEY, body, hashlib.sha256).digest()
return mac.hex() if encoding == "hex" else base64.b64encode(mac).decode()
def wrapped(result) -> httpx.Response: # every response is {timestamp, status, result}
return httpx.Response(200, json={"status": 200, "result": result})
@respx.mock
async def test_successful_deposit(psp: PspClient):
route = respx.post(PAY).mock(return_value=wrapped(CHECKOUT_1))
p = await psp.create_deposit(amount=Decimal("10.01"), currency="GBP", reference_id="order-1-a",
return_url="https://shop.example/r",
webhook_url="https://shop.example/webhooks/psp")
assert p.state == "CHECKOUT" and p.redirectUrl == CHECKOUT_1["redirectUrl"] # created, not paid
sent = json.loads(route.calls.last.request.content)
assert sent["paymentType"] == "DEPOSIT" and sent["amount"] == 10.01
assert route.calls.last.request.headers["authorization"].startswith("Bearer ")
@respx.mock
async def test_decline_is_http_200_with_declined_state(psp: PspClient):
respx.get(f"{PAY}/pay2").mock(return_value=wrapped(
{"id": "pay2", "state": "DECLINED", "errorCode": "4.01"}))
p = await psp.get_payment("pay2")
assert (p.state, p.errorCode) == ("DECLINED", "4.01")
@respx.mock
async def test_timeout_means_unknown_outcome(psp: PspClient):
respx.post(PAY).mock(side_effect=httpx.ReadTimeout("boom"))
with pytest.raises(PspTimeout):
await psp.create_deposit(amount=Decimal("1"), currency="GBP", reference_id="order-9-a",
return_url="x", webhook_url="y")
@respx.mock
async def test_every_ending_without_a_status_is_unknown_never_a_raw_exception(psp: PspClient):
# None of these is a TransportError and none carries a usable status, yet all of them happen
# AFTER the request went out. Each must be PspTimeout; a raw exception here would escape
# start_checkout's except clause and leave the attempt claimed with nothing to resolve it.
for ending in ({"side_effect": httpx.DecodingError("broken Content-Encoding")}, # a proxy
{"return_value": httpx.Response(200, text="<html>proxy</html>")}, # undecodable
{"return_value": httpx.Response(200, json=[])}, # 2xx, but not the envelope
{"return_value": wrapped({"state": "CHECKOUT"})}, # no id: the model rejects it
{"return_value": wrapped(None)}): # result: null
respx.post(PAY).mock(**ending)
with pytest.raises(PspTimeout):
await psp.create_deposit(amount=Decimal("1"), currency="GBP", reference_id="order-9-a",
return_url="x", webhook_url="y")
respx.get(f"{PAY}/pay8").mock(return_value=wrapped({"id": "pay8", "state": "ADDED_IN_v2"}))
assert (await psp.get_payment("pay8")).state == "ADDED_IN_v2" # a state the API added later is
# accepted and left to the transition whitelist, NOT a validation error
@pytest.mark.parametrize("encoding", ["hex", "base64"])
async def test_valid_webhook_marks_order_paid(client: httpx.AsyncClient, db, encoding):
r = await client.post("/webhooks/psp", content=BODY, headers={
"Signature": sign(BODY, encoding), "Content-Type": "application/json"})
assert r.status_code == 200
order = await fetch_order(db, "order-1-a")
assert order.status == "PAID"
assert order.paid_amount == Decimal("10.01") # from the payload, not from the request
async def test_invalid_signature_rejected_and_order_untouched(client, db):
r = await client.post("/webhooks/psp", content=BODY,
headers={"Signature": "deadbeef", "Content-Type": "application/json"})
assert r.status_code == 401
assert (await fetch_order(db, "order-1-a")).status == "AWAITING_PAYMENT"
async def test_duplicate_webhook_is_a_no_op(client, db):
headers = {"Signature": sign(BODY), "Content-Type": "application/json"}
assert (await client.post("/webhooks/psp", content=BODY, headers=headers)).status_code == 200
before = await fetch_order(db, "order-1-a")
assert (await client.post("/webhooks/psp", content=BODY, headers=headers)).status_code == 200
after = await fetch_order(db, "order-1-a")
assert (after.status, after.updated_at) == (before.status, before.updated_at)
assert await count_events(db, "pay1", "COMPLETED") == 1
async def test_webhook_before_the_order_link_is_not_swallowed(client, db):
# referenceId is the ATTEMPT's reference (not order_ref) and psp_payment_id is not stored yet:
# the real create-payment/webhook race. The receipt must stay unprocessed.
early = json.dumps({"id": "pay7", "referenceId": "order-1-9f2c", "state": "COMPLETED",
"amount": 10.01, "currency": "GBP"}).encode()
headers = {"Signature": sign(early), "Content-Type": "application/json"}
assert (await client.post("/webhooks/psp", content=early, headers=headers)).status_code == 200
assert await fetch_processed_at(db, "pay7", "COMPLETED") is None
await link_payment(db, "order-1-a", "pay7") # the create response lands late
assert (await client.post("/webhooks/psp", content=early, headers=headers)).status_code == 200
assert (await fetch_order(db, "order-1-a")).status == "PAID" # redelivery still applies it
# The tests below need a REAL PostgreSQL (Testcontainers) and independent sessions: the guarantees
# rest on ON CONFLICT, partial indexes and committed transactions, which no in-memory/mocked DB
# reproduces. Sharing one AsyncSession across asyncio.gather serialises the race away.
@respx.mock
async def test_concurrent_start_checkout_creates_exactly_one_payment(session_factory, psp, db):
route = respx.post(PAY).mock(return_value=wrapped(CHECKOUT_1))
respx.get(url__startswith=f"{PAY}?referenceId.eq=").mock( # the loser reconciles...
return_value=wrapped([])) # ...and finds nothing yet
async def attempt():
async with session_factory() as s:
return await start_checkout(s, psp, 1, Decimal("10.01"), "GBP")
outcomes = await asyncio.gather(attempt(), attempt(), return_exceptions=True)
deferred = [o for o in outcomes if isinstance(o, PaymentOutcomeUnknown)] # loser -> 409, no POST
assert len(deferred) + len([o for o in outcomes if isinstance(o, str)]) == 2
assert route.call_count == 1 # exactly ONE payment created
assert await count_open_attempts(db, 1) == 1 # ONE attempt, ONE referenceId
# Baseline recovery path: no UNKNOWN state and no sweep job — the NEXT call reconciles.
@respx.mock
async def test_checkout_timeout_then_empty_reconciliation_recovers_later(session, psp, db):
respx.post(PAY).mock(side_effect=httpx.ReadTimeout("boom"))
respx.get(url__startswith=f"{PAY}?referenceId.eq=").mock(return_value=wrapped([])) # not yet
with pytest.raises(PaymentOutcomeUnknown):
await start_checkout(session, psp, 2, Decimal("10.01"), "GBP")
stuck = await fetch_active_attempt(db, 2)
assert stuck.state == "IN_FLIGHT" # still claimed, not a permanent 409
respx.reset()
posts = respx.post(PAY)
respx.get(f"{PAY}?referenceId.eq={stuck.reference_id}").mock(return_value=wrapped(
[{"id": "pay2", "state": "CHECKOUT", "redirectUrl": "https://checkout.example/pay2"}]))
assert (await start_checkout(session, psp, 2, Decimal("10.01"), "GBP")
== "https://checkout.example/pay2") # the later call reconciles and serves it
assert (await fetch_active_attempt(db, 2)).state == "READY"
assert posts.call_count == 0 # ONE referenceId, no second POST
@respx.mock
async def test_a_failed_attempt_lets_a_new_checkout_start(session, psp, db):
respx.post(PAY).mock(return_value=httpx.Response( # confirmed refusal, nothing created
400, json={"status": 400, "errorCode": "2.01"}))
with pytest.raises(PspApiError):
await start_checkout(session, psp, 3, Decimal("1.00"), "GBP")
assert await fetch_active_attempt(db, 3) is None # FAILED sits outside the partial index
respx.reset()
respx.post(PAY).mock(return_value=wrapped(
{"id": "pay3", "state": "CHECKOUT", "redirectUrl": "https://checkout.example/pay3"}))
assert (await start_checkout(session, psp, 3, Decimal("1.00"), "GBP")
== "https://checkout.example/pay3") # a NEW attempt, a NEW referenceId
assert await count_attempts(db, 3) == 2
@respx.mock
async def test_refund_timeout_then_retry_does_not_refund_twice(session, psp, db):
respx.post(PAY).mock(side_effect=httpx.ReadTimeout("boom"))
respx.get(url__startswith=f"{PAY}?referenceId.eq=").mock(return_value=wrapped([])) # not yet
with pytest.raises(RefundOutcomeUnknown):
await refund_order(session, psp, 1, "rk-1", Decimal("5.00"), "GBP")
stuck = await fetch_refund_attempt(db, 1, "rk-1")
assert stuck.state == "IN_FLIGHT" # the attempt row SURVIVED
assert (await fetch_order(db, "order-1-a")).refunded_amount == Decimal("5.00") # commit survived
# The PSP had processed it; only the client timed out. The retry must reconcile the SAME
# referenceId, never POST a second REFUND.
respx.reset()
posted = respx.post(PAY)
respx.get(f"{PAY}?referenceId.eq={stuck.reference_id}").mock(
return_value=wrapped([{"id": "rf1", **REFUND_5}]))
result = await refund_order(session, psp, 1, "rk-1", Decimal("5.00"), "GBP")
assert result.id == "rf1" and posted.call_count == 0 # no second payout
assert (await fetch_order(db, "order-1-a")).refunded_amount == Decimal("5.00") # reserved once
@respx.mock
async def test_concurrent_refunds_with_the_same_key_post_once(session_factory, psp, db):
posts = respx.post(PAY).mock(return_value=wrapped({"id": "rf9", **REFUND_5}))
respx.get(url__startswith=f"{PAY}?referenceId.eq=").mock( # the non-owner reconciles...
return_value=wrapped([])) # ...and finds nothing yet
respx.get(f"{PAY}/rf9").mock(return_value=wrapped({"id": "rf9", **REFUND_5}))
async def attempt():
async with session_factory() as s:
return await refund_order(s, psp, 1, "rk-9", Decimal("5.00"), "GBP")
outcomes = await asyncio.gather(attempt(), attempt(), return_exceptions=True)
deferred = [o for o in outcomes if isinstance(o, RefundOutcomeUnknown)] # non-owner -> 409
assert len(deferred) + len([o for o in outcomes if not isinstance(o, Exception)]) == 2
assert posts.call_count == 1 # ONE payout, never two
assert (await fetch_order(db, "order-1-a")).refunded_amount == Decimal("5.00") # reserved once
@respx.mock
async def test_crash_after_an_accepted_refund_post_does_not_post_again(session, psp, db):
# The crash state: attempt committed IN_FLIGHT, amount reserved, the PSP already holds the payment.
attempt, _, _ = await reserve_refund(session, 1, "rk-2", Decimal("5.00"), "GBP")
posts = respx.post(PAY)
respx.get(f"{PAY}?referenceId.eq={attempt.reference_id}").mock(
return_value=wrapped([{"id": "rf2", **REFUND_5}]))
result = await refund_order(session, psp, 1, "rk-2", Decimal("5.00"), "GBP")
assert result.id == "rf2" and posts.call_count == 0 # reconciled, never re-POSTed
assert (await fetch_refund_attempt(db, 1, "rk-2")).state == "DONE"
```
The replay-job test (a stale `AUTHORIZED` receipt closed as `superseded`) exercises the hardened
variant: `references/hardening-concurrency.md` §2.
SHA-256: 02a1c64cee63bf1968a320fe07c148ad50ff3af192039524e80fbdd828fd8599