"""Private durable recovery for x402 2.24.0. Dependencies: rfc8785, cryptography. Use one client and one new file in a user-owned 0700 directory per purchase. Never feed a retained credential through an automatic payment wrapper. Two verified endings exist. An answer: interpret label and coverage before acting. A terminal no-charge non-answer (a signed receipt whose state is in TERMINAL_NO_CHARGE_STATES with billing.charged "no", on HTTP 200 or 503): nothing was charged, never retry the credential, archive the record; a new purchase later is allowed only when hints.new_quote_allowed is true. Plain-text ingress 503s and explicitly transient service errors are retryable. """ from __future__ import annotations import argparse import base64 import hashlib import json import math from datetime import timezone from email.utils import parsedate_to_datetime import os import re from pathlib import Path import stat import tempfile import time from urllib.parse import urlsplit import rfc8785 from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey # Receipt states that end a purchase with nothing charged. With a VERIFIED # receipt and billing.charged "no" they are final: polling cannot change them. TERMINAL_NO_CHARGE_STATES = ('service_failed', 'refused', 'closed_no_charge', 'frozen_unsettled') PENDING_MESSAGE = ('Still pending. Your payment record is saved. Run recover again later ' 'with the same record. Do not pay again.') class PendingRecovery(TimeoutError): """The polling budget ended; the retained credential must be reused.""" def canonical_request(body, accepted): q = accepted['extra']['seconded'] return rfc8785.dumps({'v': 1, 'input': body['input'], 'product': body['product'], 'product_schema_version': q['product_schema_version'], 'predicate_version': q['predicate_version'], 'tier': q['tier'], 'price_atomic': accepted['amount'], 'asset': accepted['network'] + '/erc20:' + accepted['asset'], 'network': accepted['network'], 'scheme': accepted['scheme'], 'pay_to': accepted['payTo'], 'billing_mode': 'paid'}).decode() def normalized_authorization(authorization): if set(authorization) != {'from', 'to', 'value', 'validAfter', 'validBefore', 'nonce'}: raise ValueError('authorization_fields') for key, size in (('from', 40), ('to', 40), ('nonce', 64)): if not isinstance(authorization[key], str) or not re.fullmatch('0x[0-9a-fA-F]{' + str(size) + '}', authorization[key]): raise ValueError('authorization_hex') auth = {k: authorization[k].lower() for k in ('from', 'to', 'nonce')} for key in ('value', 'validAfter', 'validBefore'): raw = authorization[key] if not ((type(raw) is int and 0 <= raw <= 2**53-1) or (isinstance(raw, str) and re.fullmatch('0|[1-9][0-9]*', raw))): raise ValueError('authorization_integer') if not 0 <= int(raw) < 2**256: raise ValueError('authorization_uint256') auth[key] = str(int(raw)) if int(auth['validAfter']) >= int(auth['validBefore']): raise ValueError('authorization_window') return auth def valid_association(value): return (isinstance(value, dict) and set(value) == {'v', 'sha256'} and type(value['v']) is int and value['v'] == 1 and isinstance(value['sha256'], str) and re.fullmatch('[0-9a-f]{64}', value['sha256']) is not None) def purchase_association(body, accepted, authorization): auth = normalized_authorization(authorization) terms = json.loads(json.dumps(accepted)) terms['extra'].pop('ticket', None) value = {'authorization': auth, 'accepted': terms, 'request': json.loads(canonical_request(body, accepted))} return {'v': 1, 'sha256': hashlib.sha256(b'seconded-purchase-association/v1\0' + rfc8785.dumps(value)).hexdigest()} def guard_new_payment(path): path = Path(path) if os.path.lexists(path) or os.path.lexists(str(path) + '.lock'): raise FileExistsError('purchase_record_exists_recover_same_credential') def private_parent(path): parent = Path(path).parent info = parent.lstat() if not stat.S_ISDIR(info.st_mode) or info.st_uid != os.getuid() or info.st_mode & 0o077: raise PermissionError('recovery_directory_must_be_private') return parent def durable_write(path, record): """Serialize writers; publish only after file fsync, then fsync the directory. A crash may leave a lock or a record: both refuse a second payment. Never auto-remove an orphan lock; inspect the purchase before choosing a new file. """ path = Path(path) parent = private_parent(path) guard_new_payment(path) lock = str(path) + '.lock' lockfd = os.open(lock, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) temporary = None try: if os.path.lexists(path): raise FileExistsError('purchase_record_exists') fd, temporary = tempfile.mkstemp(prefix='.recovery-', dir=parent) with os.fdopen(fd, 'w', encoding='utf-8') as stream: os.fchmod(stream.fileno(), 0o600) json.dump(record, stream, ensure_ascii=False, allow_nan=False) stream.flush() os.fsync(stream.fileno()) os.rename(temporary, path) temporary = None directory = os.open(parent, os.O_RDONLY | os.O_DIRECTORY) try: os.fsync(directory) finally: os.close(directory) finally: os.close(lockfd) if temporary is not None: os.unlink(temporary) os.unlink(lock) def capture(path, url, body, accepted, payment_signature): payload = json.loads(base64.b64decode(payment_signature, validate=True)) if payload['x402Version'] != 2 or payload['accepted'] != accepted: raise ValueError('payment_record_mismatch') normalized_authorization(payload['payload']['authorization']) canonical = canonical_request(body, accepted) quote = accepted['extra']['seconded'] legacy = 'commitment_salt' in quote if legacy: salt = bytes.fromhex(quote['commitment_salt']) if len(salt) != 32 or hashlib.sha256(b'seconded-request/v1' + salt + canonical.encode()).hexdigest() != quote['request_commitment']: raise ValueError('request_commitment_mismatch') record = {'url': url, 'body': body, 'accepted': accepted, 'quote': quote, 'payment_signature': payment_signature, 'canonical_request': canonical, 'payment_signature_sha256': hashlib.sha256(payment_signature.encode()).hexdigest(), 'state': 'unresolved'} if not legacy: record.update(format='seconded-door-recovery/v2', purchase_association=purchase_association(body, accepted, payload['payload']['authorization'])) durable_write(path, record) return record def install_hooks(client, path, url, body): """Install on a dedicated x402 async client before its first request.""" from x402.http import encode_payment_signature_header # Snapshot caller-owned input before asynchronous hooks can run. body = json.loads(json.dumps(body)) async def before(context): guard_new_payment(path) async def after(context): if context.payment_required.resource.url != url: raise ValueError('resource_url_mismatch') accepted = context.selected_requirements.model_dump(by_alias=True, exclude_none=True) capture(path, url, body, accepted, encode_payment_signature_header(context.payment_payload)) client.on_before_payment_creation(before) client.on_after_payment_creation(after) return client def load_record(path): private_parent(path) fd = os.open(path, os.O_RDONLY | os.O_NOFOLLOW) with os.fdopen(fd) as stream: info = os.fstat(stream.fileno()) if not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() or info.st_mode & 0o077: raise PermissionError('recovery_record_must_be_private') return json.load(stream) def verify(record, receipt, keys, chain_check=None): """Verify association from the record alone; optional chain_check proves tx logs. 'verified' means receipt authenticity and association, not a passing verdict. The caller must interpret answer, coverage and state separately. """ try: required_record = {'url', 'body', 'payment_signature', 'accepted', 'quote', 'canonical_request', 'state', 'payment_signature_sha256'} if not required_record <= record.keys() or record['state'] not in ('unresolved', 'resolved'): return 'unverified' if not isinstance(record['url'], str) or not urlsplit(record['url']).netloc: return 'unverified' quoted = record['quote'] new = record.get('format') == 'seconded-door-recovery/v2' if record.get('format') not in (None, 'seconded-door-recovery/v2'): return 'unverified' for field in (() if new else ('offer_id', 'quote_id', 'check_id')): if not re.fullmatch('[0-9a-f]{32}', quoted[field]): return 'unverified' for field in (() if new else ('commitment_salt', 'request_commitment')): if not re.fullmatch('[0-9a-f]{64}', quoted[field]): return 'unverified' for field in ('expires_at', 'max_valid_before', 'min_remaining_s'): if type(quoted[field]) is not int or quoted[field] <= 0: return 'unverified' a = record['accepted'] if (a['scheme'] != 'exact' or not re.fullmatch('eip155:[0-9]+', a['network']) or not re.fullmatch('[0-9]+', a['amount']) or int(a['amount']) <= 0 or type(a['maxTimeoutSeconds']) is not int or a['maxTimeoutSeconds'] <= 0 or not a['extra']['name'] or not a['extra']['version'] or not re.fullmatch('[0-9a-f]{64}', a['extra']['ticket'])): return 'unverified' envelope = receipt['envelope'] kid = receipt['key_id'] if kid != envelope['key_id'] or type(envelope['v']) is not int or envelope['v'] not in (1, 2, 3): return 'unverified' key, = [k for k in keys['keys'] if k['key_id'] == kid and k['algorithm'] == 'Ed25519'] signature = base64.urlsafe_b64decode(receipt['sig'] + '==') Ed25519PublicKey.from_public_bytes(bytes.fromhex(key['public_key_hex'])).verify( signature, f"SECONDED-RECEIPT/v{envelope['v']}\0".encode() + rfc8785.dumps(envelope)) accepted, q = record['accepted'], record['quote'] header = record['payment_signature'] if hashlib.sha256(header.encode()).hexdigest() != record['payment_signature_sha256']: return 'unverified' payload = json.loads(base64.b64decode(header, validate=True)) normalized_authorization(payload['payload']['authorization']) if payload['x402Version'] != 2 or payload['accepted'] != accepted or accepted['extra']['seconded'] != q: return 'unverified' canonical = canonical_request(record['body'], accepted) if canonical != record['canonical_request']: return 'unverified' if new: association = purchase_association(record['body'], accepted, payload['payload']['authorization']) if not (valid_association(record['purchase_association']) and valid_association(envelope.get('purchase_association')) and association == record['purchase_association'] == envelope.get('purchase_association') and all(envelope[k] == q[k] for k in ('tier', 'product_schema_version', 'predicate_version')) and envelope['product'] == record['body']['product']): return 'unverified' else: salt = bytes.fromhex(q['commitment_salt']) commitment = hashlib.sha256(b'seconded-request/v1' + salt + canonical.encode()).hexdigest() if len(salt) != 32 or commitment != q['request_commitment']: return 'unverified' if envelope.get('purchase_association') is not None: if (not valid_association(envelope['purchase_association']) or envelope['purchase_association'] != purchase_association(record['body'], accepted, payload['payload']['authorization'])): return 'unverified' elif not (commitment == envelope['request_commitment'] and envelope['check_id'] == q['check_id']): return 'unverified' auth, billing = payload['payload']['authorization'], envelope['billing'] if (not re.fullmatch('0x[0-9a-fA-F]{64}', auth['nonce']) or not re.fullmatch('0x[0-9a-fA-F]{130}', payload['payload']['signature']) or not 0 <= int(auth['validAfter']) < int(auth['validBefore'])): return 'unverified' if not (billing['network'] == accepted['network'] and billing['asset'] == accepted['network'] + '/erc20:' + accepted['asset'] and billing['amount_atomic'] == str(auth['value']) == accepted['amount'] and billing['payer'].lower() == auth['from'].lower() and billing['pay_to'].lower() == auth['to'].lower() == accepted['payTo'].lower()): return 'unverified' if chain_check is not None and billing['tx'] is not None and chain_check(billing, auth) is not True: return 'unverified' return 'verified' except Exception: return 'unverified' def _terminal_envelope(envelope): return ('answer' not in envelope and envelope.get('state') in TERMINAL_NO_CHARGE_STATES and isinstance(envelope.get('billing'), dict) and envelope['billing'].get('charged') == 'no') def _terminal_reply(record, reply, keys): receipt = reply.get('receipt') return (isinstance(receipt, dict) and isinstance(receipt.get('envelope'), dict) and verify(record, receipt, keys) == 'verified' and _terminal_envelope(receipt['envelope'])) _TRANSIENT_ERRORS = {'store_unavailable', 'internal_error', 'verification_unavailable'} _PENDING_STATES = {'running', 'settling', 'delayed', 'unavailable'} def _poll_delay(response, reply, wall_time): hints = reply.get('hints', {}) if not isinstance(hints, dict): raise ValueError('invalid_retry_delay_retain_record') raw = response.headers.get('Retry-After') delay = 0.0 if raw not in (None, ''): try: delay = float(raw) except (TypeError, ValueError): try: date = parsedate_to_datetime(raw) if date.tzinfo is None: date = date.replace(tzinfo=timezone.utc) delay = max(0.0, date.timestamp() - wall_time()) except Exception: raise ValueError('invalid_retry_delay_retain_record') from None hint = hints.get('poll_after_s') if hint is None: hint = 0 if isinstance(hints.get('poll_after_s'), bool): raise ValueError('invalid_retry_delay_retain_record') try: hint = float(hint) except (TypeError, ValueError, OverflowError): raise ValueError('invalid_retry_delay_retain_record') from None if not all(math.isfinite(value) and value >= 0 for value in (delay, hint)): raise ValueError('invalid_retry_delay_retain_record') return max(delay, hint, 1.0) def recover(record, keys, *, attempts=1800, budget_seconds=1800, send=None, sleep=time.sleep, clock=time.monotonic, wall_time=time.time, timeout=None, now=None): """Bounded same-credential polling; no signer, redirect or payment wrapper. Site defaults retain the 30-minute budget and PendingRecovery message. timeout/now are service-compatible aliases for budget_seconds/clock. Verified terminal no-charge replies still support the buyer's archive flow. """ timeout = budget_seconds if timeout is None else timeout now = clock if now is None else now if (type(attempts) is not int or attempts < 0 or type(timeout) not in (int, float) or not math.isfinite(timeout) or timeout <= 0): raise ValueError('invalid_recovery_budget_retain_record') parsed = urlsplit(record['url']) if parsed.scheme != 'https' and not (parsed.scheme == 'http' and parsed.hostname in ('127.0.0.1', '::1')): raise ValueError('https_or_loopback_required') deadline = now() + timeout url, header = record['url'], record['payment_signature'] request_body = json.dumps(record['body'], ensure_ascii=False, allow_nan=False) if send is None: import httpx send = lambda url, **kw: httpx.post(url, follow_redirects=False, timeout=min(35, max(0.001, deadline - now())), **kw) for attempt in range(attempts): if now() >= deadline: raise PendingRecovery(PENDING_MESSAGE) try: response = send(url, json=json.loads(request_body), headers={ 'PAYMENT-SIGNATURE': header, 'Prefer': 'wait=25'}) except Exception: raise ValueError('transport_error_retain_record') from None if now() >= deadline: raise PendingRecovery(PENDING_MESSAGE) status = response.status_code if status not in (200, 202, 503): raise ValueError('unverified_or_non_answer_reply_retain_record') try: reply = response.json() except Exception: # Uvicorn overload happens before the app and returns plain text. if status != 503: raise ValueError('invalid_reply_shape_retain_record') from None reply = {} if now() >= deadline: raise PendingRecovery(PENDING_MESSAGE) if not isinstance(reply, dict): raise ValueError('invalid_reply_shape_retain_record') if status == 503: hints, error = reply.get('hints'), reply.get('error') ingress = 'error' not in reply and 'receipt' not in reply transient = (isinstance(error, str) and error in _TRANSIENT_ERRORS and isinstance(hints, dict) and hints.get('do_not_resign') is True and hints.get('new_quote_allowed') is False) if not (ingress or transient): if _terminal_reply(record, reply, keys): if now() >= deadline: raise PendingRecovery(PENDING_MESSAGE) return reply raise ValueError('unverified_or_non_answer_reply_retain_record') elif status == 200 or 'receipt' in reply: receipt = reply.get('receipt') if not (isinstance(receipt, dict) and isinstance(receipt.get('envelope'), dict) and verify(record, receipt, keys) == 'verified'): raise ValueError('unverified_or_non_answer_reply_retain_record') envelope = receipt['envelope'] if status == 200: if _terminal_envelope(envelope): if now() >= deadline: raise PendingRecovery(PENDING_MESSAGE) return reply if (envelope.get('state') in ('included', 'final', 'released') and isinstance(envelope.get('answer'), dict) and envelope['answer']): if now() >= deadline: raise PendingRecovery(PENDING_MESSAGE) return reply raise ValueError('unverified_or_non_answer_reply_retain_record') if (not isinstance(envelope.get('state'), str) or envelope['state'] not in _PENDING_STATES or 'answer' in envelope or 'payment-response' in {key.lower() for key in response.headers}): raise ValueError('unverified_or_non_answer_reply_retain_record') delay = _poll_delay(response, reply, wall_time) if attempt + 1 >= attempts or now() + delay >= deadline: raise PendingRecovery(PENDING_MESSAGE) sleep(delay) raise PendingRecovery(PENDING_MESSAGE) def archive_terminal(path, reply): """After a VERIFIED terminal no-charge receipt only: move the record aside. Frees the path for a deliberate NEW purchase (fresh 402, new record file) and keeps the record plus the terminal reply as evidence. Never deletes. """ path = Path(path) private_parent(path) archived = Path(str(path) + '.terminal-' + time.strftime('%Y%m%dT%H%M%SZ', time.gmtime())) fd = os.open(str(archived) + '.reply', os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600) with os.fdopen(fd, 'w', encoding='utf-8') as stream: json.dump(reply, stream, ensure_ascii=False, allow_nan=False) stream.flush() os.fsync(stream.fileno()) os.rename(path, archived) return archived def main(): p = argparse.ArgumentParser(description=__doc__) p.add_argument('record') p.add_argument('keys', help='public /v1/keys JSON obtained from the trusted service origin') p.add_argument('--receipt', help='verify a saved signed receipt without making requests') args = p.parse_args() record = load_record(args.record) keys = json.loads(Path(args.keys).read_text()) if args.receipt: result = verify(record, json.loads(Path(args.receipt).read_text()), keys) print(result) return 0 if result == 'verified' else 1 try: body = recover(record, keys) except PendingRecovery as error: print(error) return 1 envelope = body['receipt']['envelope'] if 'answer' in envelope: print('verified answer recovered; interpret label and coverage before acting') return 0 archived = archive_terminal(args.record, body) guidance = ('a deliberate new purchase later, from a fresh 402 and a new record file, is allowed' if body.get('hints', {}).get('new_quote_allowed') else 'do not start a new purchase for this input yet') print('terminal no-charge result: state=%s reason=%s. Nothing was charged. ' 'Never retry this credential; %s. Record archived at %s' % (envelope.get('state'), envelope.get('reason'), guidance, archived)) return 0 if __name__ == '__main__': try: raise SystemExit(main()) except Exception: raise SystemExit('recovery_failed_retain_private_record') from None