diff --git a/docs/provider/configuration.md b/docs/provider/configuration.md index 11cc2934..68433cc1 100644 --- a/docs/provider/configuration.md +++ b/docs/provider/configuration.md @@ -48,6 +48,68 @@ Connect to your AI provider(s): | **Upstream URL** | API endpoint (e.g., `https://api.openai.com/v1`) | | **API Key** | Your provider's API key | +### PPQ Auto Top-up + +PPQ providers can automatically purchase more credits when their USD balance +falls below a configured threshold. Configure this per provider in the Admin +Dashboard by editing a **PPQ.AI** provider and opening **PPQ Auto Top-up**. +There are no environment variables for this feature. + +#### Requirements + +Before enabling auto top-up, make sure that: + +- the PPQ provider has a valid API key; +- at least one trusted Cashu mint is configured; +- the node wallet has enough **node-owned** funds at one mint to pay the + Lightning invoice; client balances are never used; and +- the node has a current BTC/USD price for validating the invoice amount. + +| Setting | Description | +| ------- | ----------- | +| **Enable Auto Top-up** | Enables automatic PPQ credit purchases for this provider. | +| **When credits are below (USD)** | Starts a top-up when the reported PPQ balance is below this positive USD value. | +| **Purchase this amount (USD)** | Amount of PPQ credit to buy per top-up. Must be a whole number from **1 to 500 USD**. | + +For example, a threshold of `5` and purchase amount of `20` buys 20 USD of +credit when the PPQ balance drops below 5 USD. + +#### How it works + +The worker checks eligible providers approximately once per minute. When the +balance is below the threshold, it: + +1. verifies the node has enough owner funds before creating an invoice; +2. requests a USD-denominated Lightning top-up invoice from PPQ; +3. rejects expired, mismatched, or unexpectedly expensive invoices (more than + 10% above the local BTC/USD estimate); +4. pays from the configured Cashu mint with sufficient owner funds; and +5. waits for PPQ to confirm that the credit settled. + +Only one attempt can be active for a provider. An attempt that was active at +the start of a cycle suppresses another top-up for that entire cycle, even if +PPQ reports it settled immediately. This prevents a temporarily stale PPQ +balance from causing a duplicate purchase. + +Completed PPQ payments appear in the dashboard transaction history with source +`ppq_auto_topup`. The payment record is separate from the internal claim used +to prevent concurrent attempts. + +#### Payment recovery + +If the Cashu mint paid the invoice but PPQ settlement cannot be confirmed, the +provider card shows **Auto top-up needs review**. A payment still owned by a +running worker is shown as **Paying invoice** and cannot be released. + +Before choosing **Release top-up**, manually verify both PPQ and the Cashu mint. +Release the claim only when the previous Lightning payment is definitively +unable to settle. Releasing an ambiguous payment allows the next cycle to try +again and can therefore cause a duplicate top-up. + +Disabling auto top-up prevents new purchases, but the node continues to +reconcile an already active payment until it reaches a safe terminal state or +requires operator review. + ### Node Identity How your node appears to clients: diff --git a/routstr/core/admin.py b/routstr/core/admin.py index 0b03cdc3..604ed948 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -862,6 +862,33 @@ class UpstreamProviderUpdateBySlug(BaseModel): provider_settings: dict | None = None +async def _active_ppq_claim_in_session(session: AsyncSession, provider_id: int) -> bool: + """Check for an active claim inside the caller's transaction. + + Must share the transaction of whatever destructive write it is guarding — + a check in its own session leaves a window for a worker to create the + claim between the check and the commit. + """ + from ..upstream.auto_topup import _ppq_state_id_for_provider + + claim = await session.get(CashuTransaction, _ppq_state_id_for_provider(provider_id)) + return claim is not None and not claim.collected and not claim.swept + + +def _require_valid_ppq_auto_topup( + provider_type: str, settings: dict | None +) -> None: + """Reject PPQ auto top-up settings the worker would later refuse.""" + if provider_type != "ppqai": + return + + from ..upstream.auto_topup import validate_ppq_auto_topup_settings + + problem = validate_ppq_auto_topup_settings(settings) + if problem is not None: + raise HTTPException(status_code=400, detail=problem) + + async def _apply_provider_update( session: AsyncSession, provider: UpstreamProviderRow, @@ -873,6 +900,29 @@ async def _apply_provider_update( await _ensure_unique_slug(session, validated, exclude_id=provider.id) provider.slug = validated + provider_type_changed = ( + payload.provider_type is not None + and payload.provider_type != provider.provider_type + ) + ppq_type_changed = provider_type_changed and ( + provider.provider_type == "ppqai" or payload.provider_type == "ppqai" + ) + if ( + provider_type_changed + and provider.provider_type == "ppqai" + and provider.id is not None + and await _active_ppq_claim_in_session(session, provider.id) + ): + # Changing the type would orphan the claim: the PPQ endpoints refuse + # non-ppqai providers, so nobody could ever inspect or release it. + raise HTTPException( + status_code=409, + detail=( + "This provider has an active PPQ auto top-up claim. Release " + "it before changing the provider type" + ), + ) + if payload.provider_type is not None: provider.provider_type = payload.provider_type if payload.base_url is not None: @@ -885,6 +935,41 @@ async def _apply_provider_update( provider.enabled = payload.enabled if payload.provider_fee is not None: provider.provider_fee = payload.provider_fee + + # Auto-top-up fields have provider-specific units and meaning. Reusing + # enabled Routstr settings for PPQ (or vice versa) can silently reinterpret + # sats as USD, so a type change must provide settings for the new type. + if ( + ppq_type_changed + and payload.provider_settings is None + and provider.provider_settings + ): + try: + stored_settings = json.loads(provider.provider_settings) + except (json.JSONDecodeError, TypeError): + stored_settings = None + if isinstance(stored_settings, dict) and stored_settings.get("auto_topup"): + raise HTTPException( + status_code=400, + detail=( + "Changing provider type requires explicit auto-top-up " + "settings because the units are provider-specific" + ), + ) + + # Validate against the effective type and effective settings. + effective_settings = payload.provider_settings + if effective_settings is None and payload.provider_type is not None: + try: + effective_settings = ( + json.loads(provider.provider_settings) + if provider.provider_settings + else None + ) + except (json.JSONDecodeError, TypeError): + effective_settings = None + if effective_settings is not None: + _require_valid_ppq_auto_topup(provider.provider_type, effective_settings) if payload.provider_settings is not None: provider.provider_settings = json.dumps(payload.provider_settings) @@ -924,6 +1009,10 @@ async def create_upstream_provider( else: slug = await allocate_unique_provider_slug(session, payload.provider_type) + _require_valid_ppq_auto_topup( + payload.provider_type, payload.provider_settings + ) + provider = UpstreamProviderRow( slug=slug, provider_type=payload.provider_type, @@ -1015,6 +1104,25 @@ async def delete_upstream_provider(provider_id: str) -> dict[str, object]: async with create_session() as session: provider = await _get_upstream_provider_by_ref(session, provider_id) deleted_id = _provider_pk(provider) + + # Checked inside the delete transaction: the worker's claim creation + # re-reads the provider inside its own transaction, so these two + # writes serialise — either the claim lands first and this 409s, or + # the delete lands first and the worker refuses to claim. + if provider.provider_type == "ppqai" and await _active_ppq_claim_in_session( + session, deleted_id + ): + # Deleting now would orphan the claim and any funds it tracks: + # the PPQ endpoints 404 without the provider row, so the claim + # could never again be inspected or released. + raise HTTPException( + status_code=409, + detail=( + "This provider has an active PPQ auto top-up claim. " + "Resolve and release it before deleting the provider" + ), + ) + await session.delete(provider) await session.commit() await reinitialize_upstreams() @@ -1623,6 +1731,78 @@ async def get_log_dates_api(request: Request) -> dict[str, object]: return {"dates": dates} +_PPQ_RELEASE_ERRORS = { + "no_active_claim": "No active PPQ claim to release", + "stale_state": ("The claim changed since it was reviewed; reload and check again"), + "payment_in_flight": ( + "A Lightning payment is still in flight for this claim. Wait for it to " + "finish or expire before releasing" + ), + "claim_changed": ( + "The claim changed while the release was being applied; reload and check again" + ), +} + + +class ReleasePPQAutoTopupRequest(BaseModel): + confirmed_safe_to_retry: bool + # Echoes the state_token the admin reviewed — the claim's full versioned + # state, not just its operation id. Any change since the review (a new + # attempt, a phase change, a renewed lease) fails the match, so the + # release cannot land on a state the admin never saw. + state_token: str | None = None + + +async def _require_ppq_provider(provider_id: int) -> UpstreamProviderRow: + async with create_session() as session: + provider = await session.get(UpstreamProviderRow, provider_id) + if provider is None: + raise HTTPException(status_code=404, detail="Provider not found") + if provider.provider_type != "ppqai": + raise HTTPException(status_code=400, detail="Provider is not PPQ") + return provider + + +@admin_router.get( + "/api/upstream-providers/{provider_id}/ppq-auto-topup", + dependencies=[Depends(require_admin_api)], +) +async def get_ppq_auto_topup_api(provider_id: int) -> dict[str, object]: + await _require_ppq_provider(provider_id) + from ..upstream.auto_topup import get_ppq_auto_topup_state + + return {"ok": True, **await get_ppq_auto_topup_state(provider_id)} + + +@admin_router.post( + "/api/upstream-providers/{provider_id}/ppq-auto-topup/release", + dependencies=[Depends(require_admin_api)], +) +async def release_ppq_auto_topup_api( + provider_id: int, payload: ReleasePPQAutoTopupRequest +) -> dict[str, object]: + await _require_ppq_provider(provider_id) + if not payload.confirmed_safe_to_retry: + raise HTTPException( + status_code=400, + detail="Confirm the Lightning payment outcome is safe before releasing", + ) + + from ..upstream.auto_topup import release_ppq_auto_topup_state + + outcome = await release_ppq_auto_topup_state( + provider_id, state_token=payload.state_token + ) + if not outcome.released: + raise HTTPException(status_code=409, detail=_PPQ_RELEASE_ERRORS[outcome.reason]) + + logger.warning( + "Admin released PPQ auto top-up claim after manual reconciliation", + extra={"provider_id": provider_id, "state_token": payload.state_token}, + ) + return {"ok": True, "released": True} + + @admin_router.get("/api/transactions", dependencies=[Depends(require_admin_api)]) async def get_transactions_api( type: str | None = None, @@ -1635,7 +1815,11 @@ async def get_transactions_api( async with create_session() as session: from sqlmodel import col, func - base = select(CashuTransaction) + # Hide only the deterministic PPQ claim-lock rows. Append-only PPQ + # payment rows remain visible as the audit trail for irreversible melts. + base = select(CashuTransaction).where( + ~col(CashuTransaction.id).like("ppq-auto-topup-%") + ) if type: base = base.where(CashuTransaction.type == type) if source: diff --git a/routstr/upstream/auto_topup.py b/routstr/upstream/auto_topup.py index c824d614..c8eadb6e 100644 --- a/routstr/upstream/auto_topup.py +++ b/routstr/upstream/auto_topup.py @@ -1,7 +1,13 @@ import asyncio import json +import math +import time +import typing +import uuid -from sqlmodel import select +from sqlalchemy.exc import IntegrityError +from sqlmodel import col, or_, select, update +from sqlmodel.ext.asyncio.session import AsyncSession from ..core import get_logger from ..core.db import ( @@ -12,23 +18,50 @@ from ..core.db import ( from ..core.db import ( store_cashu_transaction_with_retry as store_cashu_transaction, ) -from ..wallet import release_token_reservation, send_token, token_mint_url +from ..payment.price import sats_usd_price +from ..wallet import ( + Bolt11PaymentAmbiguous, + Bolt11PaymentNotAttempted, + check_bolt11_payment_status, + execute_bolt11_payment, + maximum_owner_cashu_balance_sats, + prepare_bolt11_payment, + release_token_reservation, + send_token, + token_mint_url, + wallet_operation_guard, +) +from .ppqai import PPQAIUpstreamProvider from .routstr import RoutstrUpstreamProvider logger = get_logger(__name__) # Check every 60 seconds AUTO_TOPUP_INTERVAL_SECONDS = 60 +# Claim lifecycle. "claimed" holds the slot while the invoice is being created +# and priced; nothing has been spent yet, so it is always safe to release. +# "in_flight" means proofs are committed to a mint and the outcome is unknown. +# "reconcile" means the worker gave up and an admin must decide. +PPQ_PHASE_CLAIMED = "claimed" +PPQ_PHASE_IN_FLIGHT = "in_flight" +PPQ_PHASE_RECONCILE = "reconcile" +PPQ_PHASES = frozenset({PPQ_PHASE_CLAIMED, PPQ_PHASE_IN_FLIGHT, PPQ_PHASE_RECONCILE}) +PPQ_SETTLEMENT_ATTEMPTS = 5 +PPQ_SETTLEMENT_POLL_SECONDS = 2 +PPQ_PENDING_TTL_SECONDS = 15 * 60 +PPQ_MAX_INVOICE_PREMIUM = 1.10 +PPQ_MIN_TOPUP_USD = 1 +PPQ_MAX_TOPUP_USD = 500 +# Rolling 24h ceiling on total PPQ auto top-up spend across all providers. +# Independent of PPQ's own balance endpoint: if that endpoint is buggy or +# compromised and keeps reporting a below-threshold balance, this cap bounds +# the damage instead of letting the worker drain the owner's mint funds one +# per-transaction-capped payment at a time. +PPQ_MAX_DAILY_TOPUP_USD = 300 async def periodic_auto_topup() -> None: - """Background task that monitors Routstr provider balances and auto-tops up when below threshold. - - For each Routstr provider with auto_topup enabled in provider_settings: - 1. Checks the upstream balance via get_balance() - 2. If balance < topup_threshold, creates a cashu token from the configured mint - 3. Sends the token to the upstream provider via topup() - """ + """Background task that monitors Routstr and PPQ provider balances.""" # Wait for initial startup to complete await asyncio.sleep(30) logger.info("Auto top-up worker started") @@ -47,15 +80,21 @@ async def periodic_auto_topup() -> None: async def _run_auto_topup_cycle() -> None: """Single cycle: check all eligible providers and top up if needed.""" + reconciled_ppq_provider_ids = await _reconcile_all_ppq_claims() + async with create_session() as session: query = select(UpstreamProviderRow).where( - UpstreamProviderRow.provider_type == "routstr", + col(UpstreamProviderRow.provider_type).in_(["routstr", "ppqai"]), UpstreamProviderRow.enabled == True, # noqa: E712 ) result = await session.exec(query) providers = result.all() for row in providers: + # Do not immediately retry a PPQ claim reconciled this cycle: PPQ's + # balance endpoint may lag its invoice status and cause a duplicate. + if row.id is not None and row.id in reconciled_ppq_provider_ids: + continue try: await _check_and_topup(row) except Exception as e: @@ -70,8 +109,98 @@ async def _run_auto_topup_cycle() -> None: ) +async def _reconcile_all_ppq_claims() -> set[int]: + """Reconcile active PPQ claims and return providers suppressed this cycle.""" + async with create_session() as session: + rows = ( + await session.exec( + select(UpstreamProviderRow).where( + col(UpstreamProviderRow.provider_type) == "ppqai" + ) + ) + ).all() + + active_provider_ids: set[int] = set() + for row in rows: + try: + if row.id is None: + continue + state = await get_ppq_auto_topup_state(row.id) + if not state.get("active"): + continue + active_provider_ids.add(row.id) + provider = PPQAIUpstreamProvider.from_db_row(row) if row.api_key else None + await _reconcile_ppq_state(row, provider) + except Exception as e: + logger.error( + "PPQ claim reconciliation failed", + extra={"provider_id": row.id, "error": str(e)}, + ) + return active_provider_ids + + +def _invalid_ppq_number(value: object, *, integer: bool = False) -> bool: + if isinstance(value, bool) or not isinstance(value, (int, float)): + return True + try: + number = float(value) + except OverflowError: + return True + return ( + not math.isfinite(number) + or number <= 0 + or (integer and not number.is_integer()) + ) + + +def validate_ppq_auto_topup_settings(settings: dict | None) -> str | None: + """Return why enabled PPQ auto top-up settings are invalid, if anything.""" + if not settings or not settings.get("auto_topup"): + return None + + threshold = settings.get("topup_threshold") + amount = settings.get("topup_amount_limit") + if _invalid_ppq_number(threshold): + return "PPQ auto top-up threshold must be a positive number" + if _invalid_ppq_number(amount, integer=True): + return "PPQ auto top-up amount must be a positive whole number" + amount_usd = int(typing.cast(int | float, amount)) + if not PPQ_MIN_TOPUP_USD <= amount_usd <= PPQ_MAX_TOPUP_USD: + return ( + f"PPQ auto top-up amount must be between {PPQ_MIN_TOPUP_USD} " + f"and {PPQ_MAX_TOPUP_USD} USD" + ) + return None + + +async def _check_and_topup_ppq_from_row(row: UpstreamProviderRow) -> None: + settings: dict = {} + if row.provider_settings: + try: + settings = json.loads(row.provider_settings) + except (json.JSONDecodeError, TypeError): + return + if not isinstance(settings, dict) or not settings.get("auto_topup"): + return + + problem = validate_ppq_auto_topup_settings(settings) + if problem is not None: + logger.warning( + "PPQ auto top-up enabled but its configuration is invalid", + extra={"provider_id": row.id, "problem": problem}, + ) + return + if not row.api_key: + return + await _check_and_topup_ppq(row, settings) + + async def _check_and_topup(row: UpstreamProviderRow) -> None: """Check a single provider's balance and top up if below threshold.""" + if row.provider_type == "ppqai": + await _check_and_topup_ppq_from_row(row) + return + # Parse provider settings settings: dict = {} if row.provider_settings: @@ -217,3 +346,811 @@ async def _check_and_topup(row: UpstreamProviderRow) -> None: "new_balance_approx": balance + amount, }, ) + + +def _ppq_state_id(row: UpstreamProviderRow) -> str: + if row.id is None: + raise ValueError("PPQ auto top-up requires a persisted provider row") + return _ppq_state_id_for_provider(row.id) + + +def _ppq_state_id_for_provider(provider_id: int | str) -> str: + return f"ppq-auto-topup-{provider_id}" + + +def _ppq_payment_id(operation_id: str) -> str: + return f"ppq-payment-{operation_id}" + + +class PPQClaim(typing.NamedTuple): + operation_id: str + # Worker lease, not the BOLT11 invoice expiry: the invoice can expire + # while melt() is still running, and only the lease says whether the + # owning worker can still be alive. + lease_expires_at: int + phase: str + invoice_id: str + # Cashu melt quote id, "none" until a payment plan exists. This is what + # lets an ambiguous payment be reconciled against the mint later. + quote_id: str + + +def _ppq_request_id( + operation_id: str, + lease_expires_at: int, + phase: str, + invoice_id: str, + quote_id: str = "none", +) -> str: + return f"ppq:{operation_id}:{lease_expires_at}:{phase}:{invoice_id}:{quote_id}" + + +def _parse_ppq_request_id(request_id: str | None) -> PPQClaim | None: + parts = (request_id or "").split(":", 5) + if len(parts) != 6 or parts[0] != "ppq" or parts[3] not in PPQ_PHASES: + return None + try: + lease_expires_at = int(parts[2]) + except (TypeError, ValueError): + return None + return PPQClaim(parts[1], lease_expires_at, parts[3], parts[4], parts[5]) + + +def _ppq_claim_is_releasable(claim: PPQClaim | None, now: float) -> bool: + """Whether an admin may sweep this claim without risking a double payment. + + A claim in ``in_flight`` is owned by a worker that is somewhere between + reserving proofs and hearing back from the mint. Releasing it there frees + the next cycle to pay a second invoice, which is exactly what the claim + exists to prevent — so it is releasable only once its lease has passed, + which means the owning worker died rather than that it is still working. + """ + if claim is None: + # Corrupt state cannot be reasoned about and has no operation id to + # fence on. Leaving it unreleasable would strand the provider forever. + return True + if claim.phase == PPQ_PHASE_IN_FLIGHT: + return now >= claim.lease_expires_at + return True + + +async def get_ppq_auto_topup_state(provider_id: int) -> dict[str, object]: + """Return admin-safe state for a provider's durable PPQ claim.""" + async with create_session() as session: + transaction = await session.get( + CashuTransaction, _ppq_state_id_for_provider(provider_id) + ) + if transaction is None or transaction.collected or transaction.swept: + return {"active": False} + + claim = _parse_ppq_request_id(transaction.request_id) + return { + "active": True, + # The exact version of the claim the admin is looking at. A release + # must echo it back verbatim: the operation id alone is stable across + # phase changes, so it cannot distinguish "the state I reviewed" from + # "the same attempt after its payment turned ambiguous". + "state_token": transaction.request_id, + "operation_id": claim.operation_id if claim else None, + "phase": claim.phase if claim else None, + "releasable": _ppq_claim_is_releasable(claim, time.time()), + "expires_at": claim.lease_expires_at if claim else None, + "invoice_id": ( + claim.invoice_id if claim and claim.invoice_id != "pending" else None + ), + "created_at": transaction.created_at, + "amount": transaction.amount, + "unit": transaction.unit, + "mint_url": transaction.mint_url, + "malformed": claim is None, + } + + +class PPQReleaseOutcome(typing.NamedTuple): + released: bool + reason: str + + +async def release_ppq_auto_topup_state( + provider_id: int, *, state_token: str | None +) -> PPQReleaseOutcome: + """Force-release an active claim after an admin reconciles payment status. + + ``state_token`` is the ``state_token`` the caller read from + :func:`get_ppq_auto_topup_state` — the claim's full ``request_id``. Two + things are enforced with it. The claim must not be inside an unexpired + ``in_flight`` lease, because a worker between reserving proofs and hearing + from the mint still owns the outcome. And the row must be byte-identical + to the one the caller reviewed: the update fences on the whole token, so + any phase change, lease renewal, or new attempt since the review fails the + write instead of sweeping a state the admin never saw. + """ + state_id = _ppq_state_id_for_provider(provider_id) + async with create_session() as session: + transaction = await session.get(CashuTransaction, state_id) + + if transaction is None or transaction.collected or transaction.swept: + return PPQReleaseOutcome(False, "no_active_claim") + + if transaction.request_id != state_token: + return PPQReleaseOutcome(False, "stale_state") + + claim = _parse_ppq_request_id(transaction.request_id) + if not _ppq_claim_is_releasable(claim, time.time()): + return PPQReleaseOutcome(False, "payment_in_flight") + + async with create_session() as session: + result = await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where( + col(CashuTransaction.id) == state_id, + col(CashuTransaction.source).in_( + ["ppq_auto_topup", "ppq_auto_topup_claim"] + ), + # Fence on the token itself, not the row read above: a change + # between the read and this write must lose the race. + col(CashuTransaction.request_id) == state_token, + col(CashuTransaction.collected) == False, # noqa: E712 + col(CashuTransaction.swept) == False, # noqa: E712 + ) + .values(swept=True) + ) + if (getattr(result, "rowcount", 0) or 0) == 1: + claim = _parse_ppq_request_id(state_token) + if claim is not None: + await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where( + col(CashuTransaction.id) == _ppq_payment_id(claim.operation_id) + ) + .values(swept=True) + ) + await session.commit() + return PPQReleaseOutcome(True, "released") + await session.rollback() + return PPQReleaseOutcome(False, "claim_changed") + + +async def _set_ppq_state_terminal( + row: UpstreamProviderRow, + operation_id: str, + *, + collected: bool, + swept: bool, +) -> bool: + """Finish a PPQ attempt only if this worker still owns the claim.""" + async with create_session() as session: + result = await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where( + col(CashuTransaction.id) == _ppq_state_id(row), + col(CashuTransaction.request_id).like(f"ppq:{operation_id}:%"), + col(CashuTransaction.collected) == False, # noqa: E712 + col(CashuTransaction.swept) == False, # noqa: E712 + ) + .values(collected=collected, swept=swept) + ) + updated = (getattr(result, "rowcount", 0) or 0) == 1 + if updated: + await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where(col(CashuTransaction.id) == _ppq_payment_id(operation_id)) + .values(collected=collected, swept=swept) + ) + await session.commit() + else: + await session.rollback() + return updated + + +async def _reconcile_ppq_state( + row: UpstreamProviderRow, provider: PPQAIUpstreamProvider | None +) -> bool: + """Return True while a prior PPQ attempt must suppress a new payment. + + ``provider`` may be ``None`` when the API key is gone: PPQ settlement + cannot be polled then, but mint-side reconciliation still runs. + """ + async with create_session() as session: + transaction = await session.get(CashuTransaction, _ppq_state_id(row)) + if transaction is None or transaction.collected or transaction.swept: + return False + + claim = _parse_ppq_request_id(transaction.request_id) + if claim is None: + logger.critical( + "Malformed PPQ auto top-up state; suppressing duplicate payment", + extra={"provider_id": row.id}, + ) + return True + + if claim.invoice_id != "pending": + if provider is not None and await provider.check_topup_status(claim.invoice_id): + if not await _set_ppq_state_terminal( + row, claim.operation_id, collected=True, swept=False + ): + # An admin release won the race against a settlement that + # turned out to have succeeded. The next cycle may pay again; + # the balance check is the only remaining guard, so shout. + logger.critical( + "PPQ invoice settled but its claim was already released; " + "a duplicate top-up is possible on the next cycle", + extra={ + "provider_id": row.id, + "invoice_id": claim.invoice_id, + }, + ) + return True + # PPQ has not credited the invoice. Ask the mint what became of the + # melt — the durable reconciliation path for a payment whose worker + # died or whose melt call never returned. cashu settles the wallet + # database as a side effect: "unpaid" releases the reserved proofs. + if ( + claim.quote_id != "none" + and transaction.mint_url + and time.time() >= claim.lease_expires_at + ): + status = await check_bolt11_payment_status( + transaction.mint_url, transaction.unit, claim.quote_id + ) + if status == "unpaid": + # Provably never paid, funds recovered — safe to retry. + released = await _set_ppq_state_terminal( + row, claim.operation_id, collected=False, swept=True + ) + if released: + logger.warning( + "PPQ auto top-up melt was never paid; claim released", + extra={ + "provider_id": row.id, + "invoice_id": claim.invoice_id, + }, + ) + return not released + # "paid" means the mint paid but PPQ has not credited yet: keep + # waiting on PPQ. "pending"/"unknown" stay locked for the admin. + return True + if time.time() < claim.lease_expires_at: + return True + + released = await _set_ppq_state_terminal( + row, claim.operation_id, collected=False, swept=True + ) + return not released + + +async def _ppq_provider_is_claimable( + session: AsyncSession, provider_id: int | None +) -> bool: + """Re-read the provider inside the claim transaction. + + SQLite serialises write transactions, so checking here — rather than + trusting the row the cycle loaded earlier — means a concurrent provider + deletion or type change either commits before us (we see it and refuse) + or after us (its own claim check sees our claim and refuses). Without + this the worker could create a claim for a provider that no longer + exists, orphaning it forever. + """ + if provider_id is None: + return False + current = await session.get(UpstreamProviderRow, provider_id) + return current is not None and current.provider_type == "ppqai" + + +async def _claim_ppq_topup(row: UpstreamProviderRow) -> str | None: + """Acquire a durable, ownership-fenced per-provider claim.""" + state_id = _ppq_state_id(row) + operation_id = uuid.uuid4().hex + expires_at = int(time.time()) + PPQ_PENDING_TTL_SECONDS + request_id = _ppq_request_id(operation_id, expires_at, PPQ_PHASE_CLAIMED, "pending") + + async with create_session() as session: + if not await _ppq_provider_is_claimable(session, row.id): + return None + existing = await session.get(CashuTransaction, state_id) + if existing is not None: + result = await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where( + col(CashuTransaction.id) == state_id, + or_( + CashuTransaction.collected == True, # noqa: E712 + CashuTransaction.swept == True, # noqa: E712 + ), + ) + .values( + token="pending", + amount=0, + unit="sat", + mint_url=None, + request_id=request_id, + collected=False, + swept=False, + created_at=int(time.time()), + source="ppq_auto_topup_claim", + ) + ) + await session.commit() + if (getattr(result, "rowcount", 0) or 0) != 1: + return None + return operation_id + + try: + async with create_session() as session: + # Same fencing as the update path: the provider must still exist + # inside the transaction that creates the claim. + if not await _ppq_provider_is_claimable(session, row.id): + return None + session.add( + CashuTransaction( + id=state_id, + token="pending", + amount=0, + unit="sat", + type="out", + request_id=request_id, + collected=False, + source="ppq_auto_topup_claim", + ) + ) + await session.commit() + except IntegrityError: + return None + return operation_id + + +async def _record_ppq_invoice( + row: UpstreamProviderRow, + operation_id: str, + *, + invoice: str, + invoice_id: str, + quote_id: str, + amount: int, + amount_usd: int, + unit: str, + mint_url: str, +) -> int: + """Move the claim to in_flight and return its fresh worker lease. + + The lease is minted here, at the start of the payment, and deliberately + not derived from the BOLT11 invoice's expiry: the invoice can expire + while melt() is still running, and the lease answers a different question + — can the worker that owns this claim still be alive? + """ + lease_expires_at = int(time.time()) + PPQ_PENDING_TTL_SECONDS + async with create_session() as session: + result = await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where( + col(CashuTransaction.id) == _ppq_state_id(row), + col(CashuTransaction.request_id).like( + f"ppq:{operation_id}:%:{PPQ_PHASE_CLAIMED}:pending:none" + ), + col(CashuTransaction.collected) == False, # noqa: E712 + col(CashuTransaction.swept) == False, # noqa: E712 + ) + .values( + token=invoice, + # Moves the claim to in_flight: from here the proofs are + # committed and an admin may not release it until the lease + # runs out. + request_id=_ppq_request_id( + operation_id, + lease_expires_at, + PPQ_PHASE_IN_FLIGHT, + invoice_id, + quote_id, + ), + amount=amount, + unit=unit, + mint_url=mint_url, + ) + ) + if (getattr(result, "rowcount", 0) or 0) != 1: + await session.rollback() + raise RuntimeError("PPQ auto top-up claim ownership was lost") + session.add( + CashuTransaction( + id=_ppq_payment_id(operation_id), + # Do not expose the raw BOLT11 through the transaction API. + # The USD amount is stamped here so the daily spend cap can + # aggregate what each payment was worth when it was made, + # independent of later BTC price moves. + token=f"ppq-invoice:{invoice_id}:usd:{amount_usd}", + amount=amount, + unit=unit, + type="out", + request_id=invoice_id, + mint_url=mint_url, + collected=False, + swept=False, + source="ppq_auto_topup", + ) + ) + await session.commit() + return lease_expires_at + + +async def _record_ppq_payment_spent(operation_id: str, paid_amount: int) -> None: + """Persist the irreversible wallet spend before polling PPQ settlement.""" + async with create_session() as session: + result = await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where(col(CashuTransaction.id) == _ppq_payment_id(operation_id)) + .values(amount=paid_amount) + ) + await session.commit() + if (getattr(result, "rowcount", 0) or 0) != 1: + raise RuntimeError("PPQ payment audit row is missing") + + +async def _mark_ppq_reconcile( + row: UpstreamProviderRow, + operation_id: str, + lease_expires_at: int, + invoice_id: str, + quote_id: str, +) -> None: + """Move an in_flight claim to reconcile so an admin may release it. + + Without this an ambiguous payment stays in_flight, and in_flight is only + releasable once its lease expires — the admin would have to wait out the + lease before they could act on an alert that already fired. + """ + async with create_session() as session: + result = await session.exec( # type: ignore[call-overload] + update(CashuTransaction) + .where( + col(CashuTransaction.id) == _ppq_state_id(row), + col(CashuTransaction.request_id) + == _ppq_request_id( + operation_id, + lease_expires_at, + PPQ_PHASE_IN_FLIGHT, + invoice_id, + quote_id, + ), + col(CashuTransaction.collected) == False, # noqa: E712 + col(CashuTransaction.swept) == False, # noqa: E712 + ) + .values( + request_id=_ppq_request_id( + operation_id, + lease_expires_at, + PPQ_PHASE_RECONCILE, + invoice_id, + quote_id, + ) + ) + ) + await session.commit() + if (getattr(result, "rowcount", 0) or 0) != 1: + logger.warning( + "Could not flag the PPQ claim for reconciliation; " + "it is no longer owned by this attempt", + extra={"provider_id": row.id, "invoice_id": invoice_id}, + ) + + +def _ppq_payment_usd(amount: int, unit: str, token: str, price: float) -> float: + """USD value of one PPQ payment audit row. + + Prefers the USD amount stamped into the token when the payment was + recorded: converting stored sats at today's price would undercount past + spend whenever the BTC price has fallen since. Falls back to a current + price conversion for rows recorded before the stamp existed. + """ + marker = ":usd:" + if marker in token: + try: + return float(token.rsplit(marker, 1)[1]) + except ValueError: + pass + sats = amount if unit == "sat" else math.ceil(amount / 1000) + return sats * price + + +async def _ppq_spent_last_24h_usd(price: float) -> float: + """Total USD committed to PPQ top-ups in the last 24 hours. + + Counts in-flight and ambiguous payments — for spend-cap purposes an + unresolved payment must be assumed spent — but not rows marked + ``collected=False, swept=True``, which record payments the mint provably + never attempted; those must not starve future top-ups for a day. + """ + cutoff = int(time.time()) - 24 * 60 * 60 + async with create_session() as session: + rows = ( + await session.exec( + select( + CashuTransaction.amount, + CashuTransaction.unit, + CashuTransaction.token, + ).where( + col(CashuTransaction.source) == "ppq_auto_topup", + col(CashuTransaction.type) == "out", + col(CashuTransaction.created_at) >= cutoff, + or_( + col(CashuTransaction.collected) == True, # noqa: E712 + col(CashuTransaction.swept) == False, # noqa: E712 + ), + ) + ) + ).all() + return sum( + _ppq_payment_usd(amount, unit, token, price) for amount, unit, token in rows + ) + + +async def _check_and_topup_ppq(row: UpstreamProviderRow, settings: dict) -> None: + threshold_usd = float(settings["topup_threshold"]) + amount_usd = int(settings["topup_amount_limit"]) + provider = PPQAIUpstreamProvider.from_db_row(row) + if provider is None or await _reconcile_ppq_state(row, provider): + return + + balance = await provider.get_balance() + if balance is None or not math.isfinite(balance) or balance < 0: + logger.warning( + "Could not fetch a valid PPQ balance for auto top-up", + extra={"provider_id": row.id}, + ) + return + if balance >= threshold_usd: + return + + # Perform local pricing and owner-funds checks before asking PPQ to create + # an invoice. The exact mint quote still has to be checked afterward, but + # predictable local failures should not leave abandoned PPQ invoices. + price = sats_usd_price() + minimum_invoice_sats = math.ceil(amount_usd / price) + max_invoice_sats = math.ceil(minimum_invoice_sats * PPQ_MAX_INVOICE_PREMIUM) + if await maximum_owner_cashu_balance_sats() < minimum_invoice_sats: + logger.warning( + "PPQ auto top-up skipped because no mint has enough owner funds", + extra={ + "provider_id": row.id, + "minimum_invoice_sats": minimum_invoice_sats, + }, + ) + return + + # Cheap early check to avoid claim and invoice churn; the authoritative + # re-check happens under the wallet guard just before payment, where no + # concurrent worker can move the total. + spent_24h_usd = await _ppq_spent_last_24h_usd(price) + if spent_24h_usd + amount_usd > PPQ_MAX_DAILY_TOPUP_USD: + logger.critical( + "PPQ auto top-up skipped: rolling 24h spend cap reached", + extra={ + "provider_id": row.id, + "spent_24h_usd": round(spent_24h_usd, 2), + "topup_usd": amount_usd, + "daily_cap_usd": PPQ_MAX_DAILY_TOPUP_USD, + }, + ) + return + + operation_id = await _claim_ppq_topup(row) + if operation_id is None: + return + + logger.info( + "PPQ auto top-up triggered", + extra={ + "provider_id": row.id, + "balance_usd": balance, + "threshold_usd": threshold_usd, + "topup_usd": amount_usd, + }, + ) + + try: + topup = await provider.initiate_topup(amount_usd) + if topup.currency.upper() != "USD" or topup.amount != amount_usd: + raise ValueError("PPQ top-up response amount or currency does not match") + + now = int(time.time()) + # The invoice's own expiry is a pre-payment sanity check only; the + # claim's lease is minted separately in _record_ppq_invoice. + invoice_expires_at = topup.expires_at or now + PPQ_PENDING_TTL_SECONDS + if invoice_expires_at > 10**12: + invoice_expires_at //= 1000 + if invoice_expires_at <= now: + raise ValueError("PPQ returned an expired Lightning invoice") + except Exception: + # Nothing has been paid yet, so the claim can be handed back. If the + # release does not land, the claim is no longer ours to reason about. + if not await _set_ppq_state_terminal( + row, operation_id, collected=False, swept=True + ): + logger.warning( + "Could not release the PPQ auto top-up claim after a " + "pre-payment failure; it is owned by another attempt", + extra={"provider_id": row.id}, + ) + raise + + # Hold the wallet guard across planning and execution: the plan snapshots + # live proof state, and another worker process reserving or spending those + # proofs between the two calls would invalidate the snapshot. + async with wallet_operation_guard(): + try: + # Authoritative daily-cap check: the early check above is raceable + # across worker processes, but here the guard serializes every + # payment, so the total cannot move between this read and the + # melt. + spent_24h_usd = await _ppq_spent_last_24h_usd(price) + if spent_24h_usd + amount_usd > PPQ_MAX_DAILY_TOPUP_USD: + logger.critical( + "PPQ auto top-up aborted: rolling 24h spend cap reached", + extra={ + "provider_id": row.id, + "spent_24h_usd": round(spent_24h_usd, 2), + "topup_usd": amount_usd, + "daily_cap_usd": PPQ_MAX_DAILY_TOPUP_USD, + }, + ) + raise ValueError("PPQ auto top-up daily spend cap reached") + + plan = await prepare_bolt11_payment(topup.payment_request) + if plan.maximum_spend_sats > max_invoice_sats: + raise ValueError("PPQ Lightning invoice exceeds the USD spending cap") + + lease_expires_at = await _record_ppq_invoice( + row, + operation_id, + invoice=topup.payment_request, + invoice_id=topup.invoice_id, + quote_id=str(plan.quote.quote), + amount=int(plan.quote.amount + plan.quote.fee_reserve), + amount_usd=amount_usd, + unit=plan.unit, + mint_url=plan.mint_url, + ) + except Exception: + # Nothing has been paid yet, so the claim can be handed back. If + # the release does not land, the claim is no longer ours to reason + # about. + if not await _set_ppq_state_terminal( + row, operation_id, collected=False, swept=True + ): + logger.warning( + "Could not release the PPQ auto top-up claim after a " + "pre-payment failure; it is owned by another attempt", + extra={"provider_id": row.id}, + ) + raise + + try: + paid_amount, mint_url, unit = await execute_bolt11_payment(plan) + except Bolt11PaymentNotAttempted: + # The mint's own answer rules out a settlement and any reserved + # proofs were handed back, so this claim is safe to retry next + # cycle. + if not await _set_ppq_state_terminal( + row, operation_id, collected=False, swept=True + ): + logger.warning( + "Could not release the PPQ auto top-up claim after a " + "payment that was never attempted; it is owned by another " + "attempt", + extra={"provider_id": row.id}, + ) + logger.warning( + "PPQ Lightning payment was not attempted; claim released for retry", + extra={"provider_id": row.id, "invoice_id": topup.invoice_id}, + exc_info=True, + ) + raise + except Bolt11PaymentAmbiguous: + await _mark_ppq_reconcile( + row, + operation_id, + lease_expires_at, + topup.invoice_id, + str(plan.quote.quote), + ) + logger.critical( + "PPQ auto top-up payment outcome is ambiguous; claim remains locked until admin reconciliation", + extra={ + "provider_id": row.id, + "invoice_id": topup.invoice_id, + "admin_action": f"POST /admin/api/upstream-providers/{row.id}/ppq-auto-topup/release", + }, + exc_info=True, + ) + raise + except BaseException: + # Cancellation or an unexpected error after execution began is + # also ambiguous. Preserve the claim before propagating it. + await asyncio.shield( + _mark_ppq_reconcile( + row, + operation_id, + lease_expires_at, + topup.invoice_id, + str(plan.quote.quote), + ) + ) + logger.critical( + "Unexpected failure while paying PPQ invoice; payment requires reconciliation", + extra={"provider_id": row.id, "invoice_id": topup.invoice_id}, + exc_info=True, + ) + raise + + try: + # The melt is irreversible now. Persist the actual spend before any + # fallible PPQ status call so operators retain an audit trail. + await _record_ppq_payment_spent(operation_id, paid_amount) + settled = False + for attempt in range(PPQ_SETTLEMENT_ATTEMPTS): + if await provider.check_topup_status(topup.invoice_id): + settled = True + break + if attempt + 1 < PPQ_SETTLEMENT_ATTEMPTS: + await asyncio.sleep(PPQ_SETTLEMENT_POLL_SECONDS) + except BaseException as error: + await asyncio.shield( + _mark_ppq_reconcile( + row, + operation_id, + lease_expires_at, + topup.invoice_id, + str(plan.quote.quote), + ) + ) + logger.critical( + "PPQ Lightning payment completed but settlement polling failed; reconciliation required", + extra={ + "provider_id": row.id, + "invoice_id": topup.invoice_id, + "error": str(error), + "error_type": type(error).__name__, + }, + exc_info=True, + ) + if isinstance(error, asyncio.CancelledError): + raise + return + + if not settled: + # The payment left the wallet, so the claim must stay. Flag it for + # reconciliation: _reconcile_ppq_state keeps polling PPQ, and an admin + # can step in without waiting out the lease. + await _mark_ppq_reconcile( + row, + operation_id, + lease_expires_at, + topup.invoice_id, + str(plan.quote.quote), + ) + logger.critical( + "PPQ Lightning payment completed but credit settlement is unconfirmed", + extra={"provider_id": row.id, "invoice_id": topup.invoice_id}, + ) + return + + if not await _set_ppq_state_terminal( + row, operation_id, collected=True, swept=False + ): + # Something else finished this claim while the payment was in flight — + # an admin release, most likely. The next cycle is now free to pay + # again, so surface it rather than completing quietly. + logger.critical( + "PPQ auto top-up settled but its claim was already released; " + "a duplicate top-up is possible on the next cycle", + extra={"provider_id": row.id, "invoice_id": topup.invoice_id}, + ) + logger.info( + "PPQ auto top-up completed", + extra={ + "provider_id": row.id, + "topup_usd": amount_usd, + "cashu_paid": paid_amount, + "cashu_unit": unit, + "mint_url": mint_url, + }, + ) diff --git a/routstr/upstream/ppqai.py b/routstr/upstream/ppqai.py index 6cc26cb2..50ab4532 100644 --- a/routstr/upstream/ppqai.py +++ b/routstr/upstream/ppqai.py @@ -441,7 +441,7 @@ class PPQAIUpstreamProvider(BaseUpstreamProvider): """ data = await self.check_balance() balance = data.get("balance") - if isinstance(balance, (int, float)): + if isinstance(balance, (int, float)) and not isinstance(balance, bool): return float(balance) return None diff --git a/routstr/wallet.py b/routstr/wallet.py index 486ed8d1..24b146ab 100644 --- a/routstr/wallet.py +++ b/routstr/wallet.py @@ -6,11 +6,12 @@ import time import typing from contextlib import asynccontextmanager from contextvars import ContextVar +from dataclasses import dataclass from pathlib import Path from typing import AsyncGenerator, TypedDict import httpx -from cashu.core.base import MeltQuoteState, MintQuote, Proof, Token +from cashu.core.base import MeltQuote, MeltQuoteState, MintQuote, Proof, Token from cashu.core.mint_info import MintInfo as _CashuMintInfo from cashu.wallet.helpers import deserialize_token_from_string from cashu.wallet.wallet import Wallet as _CashuWallet @@ -105,6 +106,11 @@ def _msats_to_sats(amount: int) -> int: return amount // 1000 +def _msats_to_sats_ceil(amount: int) -> int: + """Round liabilities up so fractional sats are never treated as owner funds.""" + return (amount + 999) // 1000 + + def _mints_to_inspect() -> list[str]: """Return configured mints plus the primary mint, without duplicates.""" mint_urls = list(settings.cashu_mints) @@ -390,9 +396,7 @@ async def _recieve_token_locked( else list(dict.fromkeys([settings.primary_mint, *settings.cashu_mints])) ) output_unit = ( - token_obj.unit - if token_obj.mint in destinations - else settings.primary_mint_unit + token_obj.unit if token_obj.mint in destinations else settings.primary_mint_unit ) if destination_unit is not None and output_unit != destination_unit: raise ValueError( @@ -493,6 +497,275 @@ async def send_token(amount: int, unit: str, mint_url: str | None = None) -> str return token +class Bolt11PaymentNotAttempted(Exception): + """The invoice was definitively not paid, so the attempt can be retried. + + Raised only where the mint's own answer rules out a settlement: coin + selection never reached ``melt``, or ``melt`` returned an explicit unpaid + state. Any proofs reserved along the way are released before this is + raised. + """ + + +class Bolt11PaymentAmbiguous(Exception): + """The payment may or may not have settled, so it must not be retried. + + Raised when ``melt`` errored, timed out, or came back pending. The selected + proofs stay reserved: the mint may still complete the payment with them, + and spending them elsewhere would be a double spend. + """ + + +@dataclass +class Bolt11PaymentPlan: + invoice: str + wallet: Wallet + proofs: list[Proof] + quote: MeltQuote + mint_url: str + unit: str + + @property + def invoice_amount_sats(self) -> int: + amount = int(self.quote.amount) + return amount if self.unit == "sat" else (amount + 999) // 1000 + + @property + def maximum_spend_sats(self) -> int: + maximum = ( + int(self.quote.amount) + + int(self.quote.fee_reserve) + + int(self.wallet.get_fees_for_proofs(self.proofs)) + ) + return maximum if self.unit == "sat" else (maximum + 999) // 1000 + + +async def _owner_balance_for_mint_and_unit( + mint_url: str, unit: str, proofs_balance: int +) -> int: + """Return spendable node-owned funds without crossing user liabilities.""" + async with db.create_session() as session: + # Refund mint is a preference, not funding provenance. Mirror payout's + # conservative rule and protect the full liability at every mint. + user_liability = await db.total_user_liability(session) + # API-key balances are stored in msats. Cashu ``sat`` proofs are not. + if unit == "sat": + user_liability = _msats_to_sats_ceil(user_liability) + return max(0, proofs_balance - user_liability) + + +async def maximum_owner_cashu_balance_sats() -> int: + """Return the largest conservatively owner-funded mint/unit balance.""" + details, _, _, _ = await fetch_all_balances() + async with db.create_session() as session: + liability_sats = _msats_to_sats_ceil(await db.total_user_liability(session)) + balances = [ + ( + detail["wallet_balance"] + if detail["unit"] == "sat" + else _msats_to_sats(detail["wallet_balance"]) + ) + - liability_sats + for detail in details + if not detail.get("error") + ] + return max([0, *balances]) + + +async def prepare_bolt11_payment(invoice: str) -> Bolt11PaymentPlan: + """Choose the sufficiently funded configured mint with most owner funds. + + Candidate discovery reads balances, user liabilities, and melt quotes. Coin + selection, which may split proofs, is deferred until the winner is known. + + Runs under ``wallet_operation_guard``: the plan snapshots live proof state, + which another worker process could otherwise mutate mid-read. Callers that + go on to execute the plan should hold the guard across both calls so the + snapshot stays valid. + """ + async with wallet_operation_guard(): + return await _prepare_bolt11_payment(invoice) + + +async def _prepare_bolt11_payment(invoice: str) -> Bolt11PaymentPlan: + mint_urls = list(dict.fromkeys([*settings.cashu_mints, settings.primary_mint])) + candidates: list[tuple[int, Wallet, list[Proof], MeltQuote, str, str]] = [] + failures: list[dict[str, str]] = [] + evaluated = 0 + + for mint_url in mint_urls: + if not mint_url: + continue + for unit in ("sat", "msat"): + try: + # force_reload: the guard's flock only serializes access — a + # cached wallet can still hold proof state from before another + # process's reservation landed on disk. + wallet = await get_wallet(mint_url, unit, force_reload=True) + proofs = get_proofs_per_mint_and_unit( + wallet, mint_url, unit, not_reserved=True + ) + proofs = await slow_filter_spend_proofs(proofs, wallet) + proofs_balance = sum(proof.amount for proof in proofs) + if proofs_balance <= 0: + evaluated += 1 + continue + + quote = await wallet.melt_quote(invoice=invoice) + evaluated += 1 + # select_to_send runs with include_fees=True, so the input fee + # has to be part of sufficiency too. Without it a mint passes + # this filter and then fails coin selection. + required = ( + quote.amount + + quote.fee_reserve + + wallet.get_fees_for_proofs(proofs) + ) + owner_balance = await _owner_balance_for_mint_and_unit( + mint_url, unit, proofs_balance + ) + if owner_balance < required: + continue + owner_balance_msats = ( + owner_balance * 1000 if unit == "sat" else owner_balance + ) + candidates.append( + (owner_balance_msats, wallet, proofs, quote, mint_url, unit) + ) + except Exception as e: + failures.append({"mint_url": mint_url, "unit": unit, "error": str(e)}) + logger.debug( + "Cashu mint cannot fund BOLT11 invoice", + extra={"mint_url": mint_url, "unit": unit, "error": str(e)}, + ) + + if not candidates: + if failures: + logger.warning( + "No Cashu mint could fund the BOLT11 invoice", + extra={"evaluated": evaluated, "failures": failures}, + ) + if evaluated == 0 and failures: + raise RuntimeError("Every configured Cashu mint refused the payment") + raise ValueError( + "No configured Cashu mint has enough balance after user liabilities to pay invoice" + ) + + _, wallet, proofs, quote, mint_url, unit = max(candidates, key=lambda item: item[0]) + return Bolt11PaymentPlan(invoice, wallet, proofs, quote, mint_url, unit) + + +async def execute_bolt11_payment(plan: Bolt11PaymentPlan) -> tuple[int, str, str]: + """Execute a prepared payment, separating retryable from ambiguous failure. + + Raises ``Bolt11PaymentNotAttempted`` when the invoice provably did not + settle, and ``Bolt11PaymentAmbiguous`` when the outcome is unknown. Callers + may safely retry the first and must never retry the second. + + Runs under ``wallet_operation_guard``: coin selection and reservation must + not race another worker process spending the same proofs. + """ + async with wallet_operation_guard(): + return await _execute_bolt11_payment(plan) + + +async def _execute_bolt11_payment(plan: Bolt11PaymentPlan) -> tuple[int, str, str]: + # Select unreserved, mirroring send_token: a selection failure must not + # strand proofs that were never handed to the mint. + try: + selected, _ = await plan.wallet.select_to_send( + plan.proofs, + plan.quote.amount + plan.quote.fee_reserve, + set_reserved=False, + include_fees=True, + ) + except Exception as e: + raise Bolt11PaymentNotAttempted(f"Coin selection failed: {e}") from e + + await plan.wallet.set_reserved_for_send(selected, reserved=True) + + try: + result = await asyncio.wait_for( + plan.wallet.melt( + proofs=selected, + invoice=plan.invoice, + fee_reserve_sat=plan.quote.fee_reserve, + quote_id=plan.quote.quote, + ), + timeout=60, + ) + except BaseException as e: + # The mint may still be settling with these proofs, so they must stay + # reserved — but cashu's melt() un-reserves them itself on a mint + # transport error, the exact ambiguous case. Re-reserve with the melt + # quote id, not as a send: get_melt_quote() finds the proofs to settle + # by melt_id, so a send-style reservation would strand them — paid + # proofs never invalidated, unpaid ones never released. BaseException + # includes task cancellation after the melt was submitted. + try: + await asyncio.shield( + plan.wallet.set_reserved_for_melt( + selected, reserved=True, quote_id=plan.quote.quote + ) + ) + except BaseException: + logger.critical( + "Could not re-reserve proofs after an ambiguous melt", + extra={"mint_url": plan.mint_url, "quote_id": plan.quote.quote}, + ) + if isinstance(e, asyncio.CancelledError): + raise + raise Bolt11PaymentAmbiguous(f"Cashu melt did not return: {e}") from e + + raw_state = getattr(result, "state", None) + state = str(raw_state).lower().rsplit(".", 1)[-1] if raw_state is not None else "" + if state == "paid" or getattr(result, "paid", None) is True: + change = getattr(result, "change", None) or [] + paid = sum(proof.amount for proof in selected) - sum( + int(item.amount) for item in change + ) + return paid, plan.mint_url, plan.unit + + if state == "unpaid": + # The mint is telling us it did not pay, so the proofs are ours again. + await plan.wallet.set_reserved_for_send(selected, reserved=False) + raise Bolt11PaymentNotAttempted("Cashu mint reported the melt as unpaid") + + raise Bolt11PaymentAmbiguous( + f"Cashu melt did not reach a final state: {state or 'unknown'}" + ) + + +async def check_bolt11_payment_status(mint_url: str, unit: str, quote_id: str) -> str: + """Ask the mint what became of an earlier melt attempt. + + Returns ``"paid"``, ``"unpaid"``, ``"pending"``, or ``"unknown"``. This is + the durable reconciliation path for an ambiguous payment: cashu's + ``get_melt_quote`` also settles the wallet database — invalidating the + proofs on ``paid`` and releasing their reservation on ``unpaid`` — so a + caller that sees ``"unpaid"`` may safely retry with the same funds. + + Runs under ``wallet_operation_guard`` because of that side effect: it + mutates proof state and must not race other processes' wallet operations. + """ + try: + async with wallet_operation_guard(): + wallet = await get_wallet(mint_url, unit, force_reload=True) + quote = await wallet.get_melt_quote(quote_id) + except Exception as e: + logger.warning( + "Could not query the mint for a melt quote's status", + extra={"mint_url": mint_url, "quote_id": quote_id, "error": str(e)}, + ) + return "unknown" + if quote is None: + return "unknown" + state = str(getattr(quote, "state", "")).lower().rsplit(".", 1)[-1] + if state in ("paid", "unpaid", "pending"): + return state + return "unknown" + + async def release_token_reservation(token: str) -> None: """Release a token that was created locally but never handed off.""" async with wallet_operation_guard(): @@ -1724,9 +1997,7 @@ async def fetch_all_balances( try: async with mint_check_limit: - wallet = await get_wallet( - mint_url, unit, retry_on_rate_limit=False - ) + wallet = await get_wallet(mint_url, unit, retry_on_rate_limit=False) proofs = get_proofs_per_mint_and_unit( wallet, mint_url, unit, not_reserved=True ) @@ -1840,9 +2111,7 @@ async def _payout_mint_and_unit(mint_url: str, unit: str) -> None: # snapshot up to 30s stale from another process's reservation, so the # cross-process lock is only safe with a fresh reload. wallet = await get_wallet(mint_url, unit, force_reload=True) - proofs = get_proofs_per_mint_and_unit( - wallet, mint_url, unit, not_reserved=True - ) + proofs = get_proofs_per_mint_and_unit(wallet, mint_url, unit, not_reserved=True) proofs = await slow_filter_spend_proofs(proofs, wallet) await asyncio.sleep(5) except Exception as e: @@ -1948,6 +2217,12 @@ async def _refund_sweep_once(cutoff: int) -> None: db.CashuTransaction.type == "out", db.CashuTransaction.collected == False, # noqa: E712 db.CashuTransaction.swept == False, # noqa: E712 + # PPQ rows describe a Lightning spend or claim lock, not a + # refundable Cashu token. Preserve legacy rows without a source. + col(db.CashuTransaction.source).is_(None) + | col(db.CashuTransaction.source).notin_( + ["ppq_auto_topup", "ppq_auto_topup_claim"] + ), db.CashuTransaction.created_at < cutoff, claim_available, ) @@ -2180,16 +2455,10 @@ async def periodic_routstr_fee_payout() -> None: async def send_to_lnurl(amount: int, unit: str, mint: str, address: str) -> int: async with wallet_operation_guard(): - mint = await find_trusted_mint_with_funds( - amount, unit, mint, force_reload=True - ) + mint = await find_trusted_mint_with_funds(amount, unit, mint, force_reload=True) wallet = await get_wallet(mint, unit) - available = get_proofs_per_mint_and_unit( - wallet, mint, unit, not_reserved=True - ) - proofs, _ = await wallet.select_to_send( - available, amount, set_reserved=True - ) + available = get_proofs_per_mint_and_unit(wallet, mint, unit, not_reserved=True) + proofs, _ = await wallet.select_to_send(available, amount, set_reserved=True) return await raw_send_to_lnurl(wallet, proofs, address, unit) diff --git a/tests/integration/test_ppq_auto_topup_claim.py b/tests/integration/test_ppq_auto_topup_claim.py new file mode 100644 index 00000000..433389cd --- /dev/null +++ b/tests/integration/test_ppq_auto_topup_claim.py @@ -0,0 +1,602 @@ +"""Real-database tests for the PPQ auto top-up claim lifecycle. + +These exercise the claim against actual SQL rather than mocked sessions, +because the guarantees under test are all about what the database will and +will not let two concurrent writers do. +""" + +import time +from typing import Any +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest +from sqlmodel import select + +from routstr.core.db import CashuTransaction, create_session +from routstr.upstream.auto_topup import ( + PPQ_PHASE_CLAIMED, + PPQ_PHASE_IN_FLIGHT, + PPQ_PHASE_RECONCILE, + _claim_ppq_topup, + _ppq_payment_id, + _ppq_payment_usd, + _ppq_request_id, + _ppq_spent_last_24h_usd, + _ppq_state_id_for_provider, + _record_ppq_invoice, + _set_ppq_state_terminal, + get_ppq_auto_topup_state, + release_ppq_auto_topup_state, +) + +pytestmark = pytest.mark.asyncio + + +def _row(provider_id: int = 1) -> MagicMock: + row = MagicMock() + row.id = provider_id + return row + + +async def _seed_provider(provider_id: int = 1, slug: str = "ppq") -> None: + """Claim creation is fenced on the provider row existing; seed it.""" + from routstr.core.db import UpstreamProviderRow + + async with create_session() as session: + session.add( + UpstreamProviderRow( + id=provider_id, + slug=slug, + provider_type="ppqai", + base_url="https://api.ppq.ai", + api_key="secret", + enabled=True, + ) + ) + await session.commit() + + +async def _state_row(provider_id: int = 1) -> CashuTransaction | None: + async with create_session() as session: + return await session.get( + CashuTransaction, _ppq_state_id_for_provider(provider_id) + ) + + +async def _seed_claim( + provider_id: int, + phase: str, + invoice_id: str, + lease_expires_at: int, + quote_id: str = "quote-1", +) -> str: + """Seed a claim row and return its state token (the full request_id).""" + token = _ppq_request_id( + "operation-1", lease_expires_at, phase, invoice_id, quote_id + ) + async with create_session() as session: + session.add( + CashuTransaction( + id=_ppq_state_id_for_provider(provider_id), + token="lnbc-invoice", + amount=102, + unit="sat", + type="out", + request_id=token, + mint_url="https://mint.test", + collected=False, + source="ppq_auto_topup", + ) + ) + await session.commit() + return token + + +async def test_second_claim_is_refused_while_the_first_is_active( + patched_db_engine: Any, +) -> None: + await _seed_provider() + assert await _claim_ppq_topup(_row()) is not None + # The whole point of the claim: a concurrent cycle must not get one. + assert await _claim_ppq_topup(_row()) is None + + async with create_session() as session: + rows = (await session.exec(select(CashuTransaction))).all() + assert len(rows) == 1 + + +async def test_claim_is_reusable_once_the_previous_attempt_finished( + patched_db_engine: Any, +) -> None: + await _seed_provider() + first = await _claim_ppq_topup(_row()) + assert first is not None + assert await _set_ppq_state_terminal(_row(), first, collected=True, swept=False) + + second = await _claim_ppq_topup(_row()) + assert second is not None and second != first + + +async def test_recording_the_invoice_moves_the_claim_in_flight( + patched_db_engine: Any, +) -> None: + await _seed_provider() + operation_id = await _claim_ppq_topup(_row()) + assert operation_id is not None + + state = await get_ppq_auto_topup_state(1) + assert state["phase"] == PPQ_PHASE_CLAIMED + assert state["releasable"] is True + assert state["invoice_id"] is None + + lease = await _record_ppq_invoice( + _row(), + operation_id, + invoice="lnbc-invoice", + invoice_id="invoice-1", + quote_id="quote-1", + amount=102, + amount_usd=10, + unit="sat", + mint_url="https://mint.test", + ) + assert lease > int(time.time()) + + state = await get_ppq_auto_topup_state(1) + assert state["phase"] == PPQ_PHASE_IN_FLIGHT + assert state["invoice_id"] == "invoice-1" + # A payment is committed to a mint, so an admin must not sweep it. + assert state["releasable"] is False + # The raw BOLT11 invoice must never reach the admin API. + assert "token" not in state + + +async def test_release_refuses_an_in_flight_claim(patched_db_engine: Any) -> None: + token = await _seed_claim( + 1, PPQ_PHASE_IN_FLIGHT, "invoice-1", int(time.time()) + 900 + ) + + outcome = await release_ppq_auto_topup_state(1, state_token=token) + + assert outcome.released is False + assert outcome.reason == "payment_in_flight" + row = await _state_row() + assert row is not None and row.swept is False + + +async def test_release_refuses_a_stale_state_token(patched_db_engine: Any) -> None: + await _seed_claim(1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) + 900) + + outcome = await release_ppq_auto_topup_state(1, state_token="ppq:stale:token") + + assert outcome.released is False + assert outcome.reason == "stale_state" + row = await _state_row() + assert row is not None and row.swept is False + + +async def test_release_accepts_a_reconcile_claim(patched_db_engine: Any) -> None: + token = await _seed_claim( + 1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) + 900 + ) + + outcome = await release_ppq_auto_topup_state(1, state_token=token) + + assert outcome.released is True + row = await _state_row() + assert row is not None and row.swept is True + + +async def test_expired_in_flight_claim_becomes_releasable( + patched_db_engine: Any, +) -> None: + # A worker that died mid-payment must not lock the provider forever. + token = await _seed_claim(1, PPQ_PHASE_IN_FLIGHT, "invoice-1", int(time.time()) - 1) + + assert (await get_ppq_auto_topup_state(1))["releasable"] is True + outcome = await release_ppq_auto_topup_state(1, state_token=token) + assert outcome.released is True + + +async def test_release_reports_no_active_claim_once_swept( + patched_db_engine: Any, +) -> None: + token = await _seed_claim( + 1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) + 900 + ) + assert (await release_ppq_auto_topup_state(1, state_token=token)).released + + outcome = await release_ppq_auto_topup_state(1, state_token=token) + assert outcome.released is False + assert outcome.reason == "no_active_claim" + + +async def test_terminal_write_fails_after_the_claim_was_released( + patched_db_engine: Any, +) -> None: + """The symptom an admin release leaves behind for the owning worker.""" + token = await _seed_claim( + 1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) + 900 + ) + assert (await release_ppq_auto_topup_state(1, state_token=token)).released + + assert ( + await _set_ppq_state_terminal( + _row(), "operation-1", collected=True, swept=False + ) + is False + ) + + +async def test_ppq_claim_rows_are_excluded_from_the_admin_transaction_list( + patched_db_engine: Any, +) -> None: + from routstr.core.admin import get_transactions_api + + await _seed_provider() + await _claim_ppq_topup(_row()) + async with create_session() as session: + session.add( + CashuTransaction( + id="real-transaction", + token="cashuAreal", + amount=50, + unit="sat", + type="out", + source="x-cashu", + ) + ) + await session.commit() + + result = await get_transactions_api() + + ids = {t["id"] for t in result["transactions"]} # type: ignore[index,union-attr] + assert "real-transaction" in ids + assert _ppq_state_id_for_provider(1) not in ids + + +async def test_ppq_payment_audit_row_is_visible_and_survives_next_claim( + patched_db_engine: Any, +) -> None: + from routstr.core.admin import get_transactions_api + + await _seed_provider() + operation_id = await _claim_ppq_topup(_row()) + assert operation_id is not None + await _record_ppq_invoice( + _row(), + operation_id, + invoice="lnbc-secret-invoice", + invoice_id="invoice-1", + quote_id="quote-1", + amount=102, + amount_usd=10, + unit="sat", + mint_url="https://mint.test", + ) + assert await _set_ppq_state_terminal( + _row(), operation_id, collected=True, swept=False + ) + + result = await get_transactions_api(source="ppq_auto_topup") + transactions = result["transactions"] + assert len(transactions) == 1 + audit = transactions[0] + assert audit["id"] == _ppq_payment_id(operation_id) + assert audit["token"] == "ppq-invoice:invoice-1:usd:10" + assert audit["collected"] is True + assert "lnbc-secret-invoice" not in audit["token"] + + # Reusing the deterministic claim lock must not overwrite history. + assert await _claim_ppq_topup(_row()) is not None + async with create_session() as session: + assert await session.get(CashuTransaction, audit["id"]) is not None + + +async def test_reconcile_settles_a_recorded_invoice(patched_db_engine: Any) -> None: + from routstr.upstream.auto_topup import _reconcile_ppq_state + + await _seed_claim(1, PPQ_PHASE_IN_FLIGHT, "invoice-1", int(time.time()) + 900) + provider = MagicMock() + provider.check_topup_status = AsyncMock(return_value=True) + + # Still suppresses this cycle, but the claim is now finished. + assert await _reconcile_ppq_state(_row(), provider) is True + + row = await _state_row() + assert row is not None and row.collected is True + + +async def test_stale_token_from_before_a_phase_change_cannot_release( + patched_db_engine: Any, +) -> None: + """The blocker scenario: admin reviews `claimed`, payment turns ambiguous. + + The operation id is identical in both states, so an id-based fence would + let the stale confirmation land. The full state token must not. + """ + await _seed_provider() + operation_id = await _claim_ppq_topup(_row()) + assert operation_id is not None + reviewed = await get_ppq_auto_topup_state(1) + assert reviewed["phase"] == PPQ_PHASE_CLAIMED + + # Worker records the invoice: same operation, new phase, proofs committed. + await _record_ppq_invoice( + _row(), + operation_id, + invoice="lnbc-invoice", + invoice_id="invoice-1", + quote_id="quote-1", + amount=102, + amount_usd=10, + unit="sat", + mint_url="https://mint.test", + ) + + outcome = await release_ppq_auto_topup_state( + 1, state_token=str(reviewed["state_token"]) + ) + assert outcome.released is False + assert outcome.reason == "stale_state" + row = await _state_row() + assert row is not None and row.swept is False + + +async def test_concurrent_claims_only_one_wins(patched_db_engine: Any) -> None: + import asyncio + + await _seed_provider() + + results = await asyncio.gather( + *(_claim_ppq_topup(_row()) for _ in range(5)), return_exceptions=True + ) + winners = [r for r in results if isinstance(r, str)] + assert len(winners) == 1 + + async with create_session() as session: + rows = (await session.exec(select(CashuTransaction))).all() + assert len(rows) == 1 + + +async def test_reconcile_releases_claim_when_mint_reports_unpaid( + patched_db_engine: Any, +) -> None: + from routstr.upstream.auto_topup import _reconcile_ppq_state + + # Lease expired, PPQ never credited: only the mint's own "unpaid" answer + # may hand the claim back. + await _seed_claim(1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) - 1) + provider = MagicMock() + provider.check_topup_status = AsyncMock(return_value=False) + + with patch( + "routstr.upstream.auto_topup.check_bolt11_payment_status", + AsyncMock(return_value="unpaid"), + ) as status: + suppressed = await _reconcile_ppq_state(_row(), provider) + + status.assert_awaited_once_with("https://mint.test", "sat", "quote-1") + assert suppressed is False + row = await _state_row() + assert row is not None and row.swept is True + + +async def test_reconcile_keeps_claim_when_mint_answer_is_not_final( + patched_db_engine: Any, +) -> None: + from routstr.upstream.auto_topup import _reconcile_ppq_state + + await _seed_claim(1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) - 1) + provider = MagicMock() + provider.check_topup_status = AsyncMock(return_value=False) + + for answer in ("paid", "pending", "unknown"): + with patch( + "routstr.upstream.auto_topup.check_bolt11_payment_status", + AsyncMock(return_value=answer), + ): + assert await _reconcile_ppq_state(_row(), provider) is True + row = await _state_row() + assert row is not None and row.swept is False, answer + + +async def test_release_endpoint_maps_refusals_to_409(patched_db_engine: Any) -> None: + from fastapi import HTTPException + + from routstr.core.admin import ( + ReleasePPQAutoTopupRequest, + release_ppq_auto_topup_api, + ) + + provider_row = MagicMock() + provider_row.provider_type = "ppqai" + + token = await _seed_claim( + 1, PPQ_PHASE_IN_FLIGHT, "invoice-1", int(time.time()) + 900 + ) + + with patch( + "routstr.core.admin._require_ppq_provider", + AsyncMock(return_value=provider_row), + ): + with pytest.raises(HTTPException) as excinfo: + await release_ppq_auto_topup_api( + 1, + ReleasePPQAutoTopupRequest( + confirmed_safe_to_retry=True, state_token=token + ), + ) + assert excinfo.value.status_code == 409 + assert "in flight" in excinfo.value.detail + + with pytest.raises(HTTPException) as excinfo: + await release_ppq_auto_topup_api( + 1, + ReleasePPQAutoTopupRequest( + confirmed_safe_to_retry=True, state_token="ppq:wrong" + ), + ) + assert excinfo.value.status_code == 409 + assert "changed since" in excinfo.value.detail + + +async def test_provider_delete_is_blocked_by_an_active_claim( + patched_db_engine: Any, +) -> None: + from fastapi import HTTPException + + from routstr.core.admin import delete_upstream_provider + from routstr.core.db import UpstreamProviderRow + + async with create_session() as session: + session.add( + UpstreamProviderRow( + id=1, + slug="ppq", + provider_type="ppqai", + base_url="https://api.ppq.ai", + api_key="secret", + enabled=True, + ) + ) + await session.commit() + await _seed_claim(1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) + 900) + + with pytest.raises(HTTPException) as excinfo: + await delete_upstream_provider("1") + assert excinfo.value.status_code == 409 + + # Provider must still exist. + async with create_session() as session: + assert await session.get(UpstreamProviderRow, 1) is not None + + +async def test_claim_is_refused_when_the_provider_row_is_gone( + patched_db_engine: Any, +) -> None: + """The worker's half of the delete race: no provider row, no claim.""" + assert await _claim_ppq_topup(_row()) is None + + async with create_session() as session: + rows = (await session.exec(select(CashuTransaction))).all() + assert rows == [] + + +async def test_claim_is_refused_after_a_provider_type_change( + patched_db_engine: Any, +) -> None: + from routstr.core.db import UpstreamProviderRow + + await _seed_provider() + async with create_session() as session: + provider = await session.get(UpstreamProviderRow, 1) + assert provider is not None + provider.provider_type = "openai" + session.add(provider) + await session.commit() + + assert await _claim_ppq_topup(_row()) is None + + +async def test_disabled_provider_with_claim_still_reconciles( + patched_db_engine: Any, +) -> None: + """A claim tracks committed money; eligibility must not stop reconciling.""" + from routstr.core.db import UpstreamProviderRow + from routstr.upstream.auto_topup import _reconcile_all_ppq_claims + + await _seed_provider() + async with create_session() as session: + provider = await session.get(UpstreamProviderRow, 1) + assert provider is not None + provider.enabled = False + session.add(provider) + await session.commit() + await _seed_claim(1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) + 900) + + ppq = MagicMock() + ppq.check_topup_status = AsyncMock(return_value=True) + with patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=ppq, + ): + await _reconcile_all_ppq_claims() + + row = await _state_row() + assert row is not None and row.collected is True + + +async def test_claim_without_api_key_still_reconciles_via_the_mint( + patched_db_engine: Any, +) -> None: + from routstr.core.db import UpstreamProviderRow + from routstr.upstream.auto_topup import _reconcile_all_ppq_claims + + await _seed_provider() + async with create_session() as session: + provider = await session.get(UpstreamProviderRow, 1) + assert provider is not None + provider.api_key = "" + session.add(provider) + await session.commit() + # Lease expired, so the mint may be consulted. + await _seed_claim(1, PPQ_PHASE_RECONCILE, "invoice-1", int(time.time()) - 1) + + with patch( + "routstr.upstream.auto_topup.check_bolt11_payment_status", + AsyncMock(return_value="unpaid"), + ) as status: + await _reconcile_all_ppq_claims() + + # No API key: PPQ was never polled, but the mint was, and its definitive + # "unpaid" released the claim. + status.assert_awaited_once() + row = await _state_row() + assert row is not None and row.swept is True + + +def test_ppq_payment_usd_prefers_stamped_amount() -> None: + # Stamped rows must not move with the BTC price. + assert _ppq_payment_usd(102, "sat", "ppq-invoice:a:usd:10", 0.5) == 10.0 + + +def test_ppq_payment_usd_falls_back_to_current_price() -> None: + # Rows recorded before the stamp existed convert sats at today's price. + assert _ppq_payment_usd(2000, "sat", "ppq-invoice:legacy", 0.001) == 2.0 + assert _ppq_payment_usd(2_000_000, "msat", "ppq-invoice:legacy", 0.001) == 2.0 + + +def test_ppq_payment_usd_survives_malformed_stamp() -> None: + assert _ppq_payment_usd(3000, "sat", "ppq-invoice:x:usd:oops", 0.001) == 3.0 + + +async def test_daily_spend_ignores_provably_unattempted_payments( + patched_db_engine: Any, +) -> None: + def _payment( + id_: str, token: str, collected: bool, swept: bool + ) -> CashuTransaction: + return CashuTransaction( + id=id_, + token=token, + amount=1, + unit="sat", + type="out", + source="ppq_auto_topup", + collected=collected, + swept=swept, + ) + + async with create_session() as session: + # Settled, in-flight, and provably-unattempted payments plus a + # pre-stamp row: only the unattempted one must be excluded. + session.add(_payment("pay-usd-1", "ppq-invoice:a:usd:100", True, False)) + session.add(_payment("pay-usd-2", "ppq-invoice:b:usd:50", False, False)) + session.add(_payment("pay-usd-3", "ppq-invoice:c:usd:25", False, True)) + legacy = _payment("pay-usd-4", "ppq-invoice:legacy", True, False) + legacy.amount = 2000 + session.add(legacy) + await session.commit() + + assert await _ppq_spent_last_24h_usd(0.001) == 152.0 diff --git a/tests/integration/test_wallet_melt_restart.py b/tests/integration/test_wallet_melt_restart.py new file mode 100644 index 00000000..5bb607fa --- /dev/null +++ b/tests/integration/test_wallet_melt_restart.py @@ -0,0 +1,115 @@ +"""Restart reconciliation for ambiguous melts, against a real cashu wallet DB. + +The ambiguous-melt path in ``execute_bolt11_payment`` re-reserves proofs with +``set_reserved_for_melt(..., quote_id=...)`` after cashu's ``melt()`` clears +both the reservation and the ``melt_id`` on a transport error. These tests +prove, on cashu's actual sqlite store rather than mocks, that the recovery +survives a process restart: a fresh wallet instance on the same database can +still find the proofs by ``melt_id`` — the lookup ``get_melt_quote()`` uses to +invalidate them on "paid" or release them on "unpaid". +""" + +from pathlib import Path + +import pytest +from cashu.core.base import Proof +from cashu.wallet import crud +from cashu.wallet.wallet import Wallet + +pytestmark = pytest.mark.asyncio + +QUOTE_ID = "quote-restart-1" + + +def _proof(secret: str, amount: int = 64) -> Proof: + return Proof( + id="009a1f293253e41e", + amount=amount, + secret=secret, + C="02bc9097997d81afb2cc7346b5e4345a9346bd2a506eb7958598a72f0cf85163ea", + ) + + +async def _wallet(db_dir: Path) -> Wallet: + # with_db builds the instance and runs migrations locally; nothing here + # talks to a mint. + return await Wallet.with_db("https://mint.test", str(db_dir)) + + +async def _seed_ambiguous_melt(wallet: Wallet) -> list[Proof]: + """Reproduce the exact sequence of an ambiguous melt failure. + + 1. Proofs exist and are selected for a melt. + 2. cashu's melt() reserves them with the quote id, then hits a transport + error and rolls that back — reservation gone, melt_id gone. + 3. Our recovery in execute_bolt11_payment re-reserves with the quote id. + """ + proofs = [_proof("secret-a"), _proof("secret-b", amount=32)] + for proof in proofs: + await crud.store_proof(proof, db=wallet.db) + + await wallet.set_reserved_for_melt(proofs, reserved=True, quote_id=QUOTE_ID) + # cashu's `except` block in melt(): + await wallet.set_reserved_for_melt(proofs, reserved=False, quote_id=None) + # our recovery: + await wallet.set_reserved_for_melt(proofs, reserved=True, quote_id=QUOTE_ID) + return proofs + + +async def test_melt_recovery_is_findable_by_quote_after_restart( + tmp_path: Path, +) -> None: + wallet = await _wallet(tmp_path) + await _seed_ambiguous_melt(wallet) + + # "Restart": a brand-new wallet on the same database file, as after a + # process crash between the melt and any reconciliation. + restarted = await _wallet(tmp_path) + found = await crud.get_proofs(db=restarted.db, melt_id=QUOTE_ID) + + # This is get_melt_quote()'s own lookup. If it comes back empty, a "paid" + # answer can never invalidate these proofs and an "unpaid" answer can + # never release them — the strand the send-style re-reserve caused. + assert sorted(p.secret for p in found) == ["secret-a", "secret-b"] + assert all(p.reserved for p in found) + assert all(p.melt_id == QUOTE_ID for p in found) + + +async def test_send_style_reservation_would_not_be_reconcilable( + tmp_path: Path, +) -> None: + """The defect the fix removed, demonstrated on the real store.""" + wallet = await _wallet(tmp_path) + proofs = [_proof("secret-send")] + for proof in proofs: + await crud.store_proof(proof, db=wallet.db) + + await wallet.set_reserved_for_melt(proofs, reserved=True, quote_id=QUOTE_ID) + await wallet.set_reserved_for_melt(proofs, reserved=False, quote_id=None) + # The old recovery: reserve as a send, no quote association. + await wallet.set_reserved_for_send(proofs, reserved=True) + + restarted = await _wallet(tmp_path) + found = await crud.get_proofs(db=restarted.db, melt_id=QUOTE_ID) + assert found == [] # reconciliation would never see these proofs + + +async def test_unpaid_reconciliation_releases_recovered_proofs_after_restart( + tmp_path: Path, +) -> None: + """The full recovery arc: crash, restart, mint says unpaid, funds usable.""" + wallet = await _wallet(tmp_path) + await _seed_ambiguous_melt(wallet) + + restarted = await _wallet(tmp_path) + found = await crud.get_proofs(db=restarted.db, melt_id=QUOTE_ID) + assert len(found) == 2 + + # What get_melt_quote() does on an "unpaid" answer. + await restarted.set_reserved_for_melt(found, reserved=False, quote_id=None) + + released = await crud.get_proofs(db=restarted.db, melt_id=QUOTE_ID) + assert released == [] + all_proofs = await crud.get_proofs(db=restarted.db) + assert len(all_proofs) == 2 + assert all(not p.reserved for p in all_proofs) # spendable again diff --git a/tests/unit/test_auto_topup.py b/tests/unit/test_auto_topup.py index 8aff3550..0db7408f 100644 --- a/tests/unit/test_auto_topup.py +++ b/tests/unit/test_auto_topup.py @@ -4,7 +4,29 @@ from unittest.mock import AsyncMock, MagicMock, patch import pytest from routstr.core.db import CashuTransaction -from routstr.upstream.auto_topup import _check_and_topup +from routstr.upstream.auto_topup import ( + _check_and_topup, + _parse_ppq_request_id, + _run_auto_topup_cycle, + validate_ppq_auto_topup_settings, +) +from routstr.upstream.ppqai import PPQAIUpstreamProvider +from routstr.wallet import Bolt11PaymentAmbiguous, Bolt11PaymentNotAttempted + + +def test_ppq_claim_parser_rejects_invalid_expiry() -> None: + assert ( + _parse_ppq_request_id("ppq:operation:not-a-timestamp:claimed:invoice:none") + is None + ) + + +@pytest.mark.asyncio +async def test_ppq_balance_rejects_boolean_api_value() -> None: + provider = PPQAIUpstreamProvider("secret") + provider.check_balance = AsyncMock(return_value={"balance": False}) # type: ignore[method-assign] + + assert await provider.get_balance() is None def _row() -> MagicMock: @@ -12,6 +34,7 @@ def _row() -> MagicMock: row.id = "provider-1" row.base_url = "https://provider.test" row.api_key = "secret" + row.provider_type = "routstr" row.provider_settings = json.dumps( { "auto_topup": True, @@ -151,3 +174,563 @@ async def test_auto_topup_does_not_send_untracked_token() -> None: reclaim.assert_awaited_once_with("cashu-token") provider.topup.assert_not_awaited() + + +def _ppq_row() -> MagicMock: + row = MagicMock() + row.id = "ppq-provider-1" + row.base_url = "https://api.ppq.ai" + row.api_key = "secret" + row.provider_type = "ppqai" + row.provider_settings = json.dumps( + { + "auto_topup": True, + "topup_threshold": 5.0, + "topup_amount_limit": 10, + } + ) + return row + + +@pytest.mark.asyncio +async def test_ppq_auto_topup_pays_invoice_and_confirms_settlement() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=2.5) + provider.initiate_topup = AsyncMock( + return_value=MagicMock( + invoice_id="invoice-1", + payment_request="lnbc-invoice", + amount=10, + currency="USD", + expires_at=None, + ) + ) + provider.check_topup_status = AsyncMock(return_value=True) + plan = MagicMock() + plan.invoice_amount_sats = 100 + plan.maximum_spend_sats = 102 + plan.quote.amount = 100 + plan.quote.fee_reserve = 2 + plan.mint_url = "https://mint-rich.test" + plan.unit = "sat" + row = _ppq_row() + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch( + "routstr.upstream.auto_topup._claim_ppq_topup", + AsyncMock(return_value="operation-1"), + ), + patch( + "routstr.upstream.auto_topup.maximum_owner_cashu_balance_sats", + AsyncMock(return_value=10_000), + ), + patch( + "routstr.upstream.auto_topup._ppq_spent_last_24h_usd", + AsyncMock(return_value=0.0), + ), + patch( + "routstr.upstream.auto_topup.prepare_bolt11_payment", + AsyncMock(return_value=plan), + ) as prepare, + patch( + "routstr.upstream.auto_topup.execute_bolt11_payment", + AsyncMock(return_value=(101, "https://mint-rich.test", "sat")), + ) as execute, + patch("routstr.upstream.auto_topup._record_ppq_invoice", AsyncMock()) as record, + patch( + "routstr.upstream.auto_topup._record_ppq_payment_spent", AsyncMock() + ) as record_spent, + patch( + "routstr.upstream.auto_topup._set_ppq_state_terminal", AsyncMock() + ) as terminal, + patch("routstr.upstream.auto_topup.sats_usd_price", return_value=0.001), + ): + await _check_and_topup(row) + + provider.initiate_topup.assert_awaited_once_with(10) + prepare.assert_awaited_once_with("lnbc-invoice") + execute.assert_awaited_once_with(plan) + record.assert_awaited_once() + record_spent.assert_awaited_once_with("operation-1", 101) + provider.check_topup_status.assert_awaited_once_with("invoice-1") + terminal.assert_awaited_once_with(row, "operation-1", collected=True, swept=False) + + +@pytest.mark.asyncio +async def test_ppq_ambiguous_melt_keeps_claim_and_emits_critical_alert() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=2.5) + provider.initiate_topup = AsyncMock( + return_value=MagicMock( + invoice_id="invoice-1", + payment_request="lnbc-invoice", + amount=10, + currency="USD", + expires_at=None, + ) + ) + plan = MagicMock(maximum_spend_sats=102, mint_url="https://mint.test", unit="sat") + plan.quote.amount = 100 + plan.quote.fee_reserve = 2 + row = _ppq_row() + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch( + "routstr.upstream.auto_topup._claim_ppq_topup", + AsyncMock(return_value="operation-1"), + ), + patch( + "routstr.upstream.auto_topup.maximum_owner_cashu_balance_sats", + AsyncMock(return_value=10_000), + ), + patch( + "routstr.upstream.auto_topup._ppq_spent_last_24h_usd", + AsyncMock(return_value=0.0), + ), + patch( + "routstr.upstream.auto_topup.prepare_bolt11_payment", + AsyncMock(return_value=plan), + ), + patch( + "routstr.upstream.auto_topup.execute_bolt11_payment", + AsyncMock(side_effect=Bolt11PaymentAmbiguous("ambiguous melt")), + ), + patch( + "routstr.upstream.auto_topup._record_ppq_invoice", + AsyncMock(return_value=2_000_000_000), + ), + patch( + "routstr.upstream.auto_topup._mark_ppq_reconcile", AsyncMock() + ) as reconcile_mark, + patch( + "routstr.upstream.auto_topup._set_ppq_state_terminal", AsyncMock() + ) as terminal, + patch("routstr.upstream.auto_topup.sats_usd_price", return_value=0.001), + patch("routstr.upstream.auto_topup.logger.critical") as critical, + ): + with pytest.raises(Bolt11PaymentAmbiguous, match="ambiguous melt"): + await _check_and_topup(row) + + # The claim is never released — it moves to reconcile for the admin. + terminal.assert_not_awaited() + reconcile_mark.assert_awaited_once() + critical.assert_called_once() + assert "admin reconciliation" in critical.call_args.args[0] + + +@pytest.mark.asyncio +async def test_ppq_payment_not_attempted_releases_claim_for_retry() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=2.5) + provider.initiate_topup = AsyncMock( + return_value=MagicMock( + invoice_id="invoice-1", + payment_request="lnbc-invoice", + amount=10, + currency="USD", + expires_at=None, + ) + ) + plan = MagicMock(maximum_spend_sats=102, mint_url="https://mint.test", unit="sat") + plan.quote.amount = 100 + plan.quote.fee_reserve = 2 + plan.quote.quote = "quote-1" + terminal = AsyncMock(return_value=True) + row = _ppq_row() + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch( + "routstr.upstream.auto_topup.maximum_owner_cashu_balance_sats", + AsyncMock(return_value=10_000), + ), + patch( + "routstr.upstream.auto_topup._ppq_spent_last_24h_usd", + AsyncMock(return_value=0.0), + ), + patch( + "routstr.upstream.auto_topup._claim_ppq_topup", + AsyncMock(return_value="operation-1"), + ), + patch( + "routstr.upstream.auto_topup.prepare_bolt11_payment", + AsyncMock(return_value=plan), + ), + patch( + "routstr.upstream.auto_topup._record_ppq_invoice", + AsyncMock(return_value=2_000_000_000), + ), + patch( + "routstr.upstream.auto_topup.execute_bolt11_payment", + AsyncMock(side_effect=Bolt11PaymentNotAttempted("unpaid")), + ), + patch("routstr.upstream.auto_topup._set_ppq_state_terminal", terminal), + patch( + "routstr.upstream.auto_topup._mark_ppq_reconcile", AsyncMock() + ) as reconcile, + patch("routstr.upstream.auto_topup.sats_usd_price", return_value=0.001), + pytest.raises(Bolt11PaymentNotAttempted, match="unpaid"), + ): + await _check_and_topup(row) + + terminal.assert_awaited_once_with(row, "operation-1", collected=False, swept=True) + reconcile.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_ppq_status_error_after_payment_marks_reconcile_and_alerts() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=2.5) + provider.initiate_topup = AsyncMock( + return_value=MagicMock( + invoice_id="invoice-1", + payment_request="lnbc-invoice", + amount=10, + currency="USD", + expires_at=None, + ) + ) + provider.check_topup_status = AsyncMock(side_effect=RuntimeError("PPQ 502")) + plan = MagicMock(maximum_spend_sats=102, mint_url="https://mint.test", unit="sat") + plan.quote.amount = 100 + plan.quote.fee_reserve = 2 + plan.quote.quote = "quote-1" + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch( + "routstr.upstream.auto_topup.maximum_owner_cashu_balance_sats", + AsyncMock(return_value=10_000), + ), + patch( + "routstr.upstream.auto_topup._ppq_spent_last_24h_usd", + AsyncMock(return_value=0.0), + ), + patch( + "routstr.upstream.auto_topup._claim_ppq_topup", + AsyncMock(return_value="operation-1"), + ), + patch( + "routstr.upstream.auto_topup.prepare_bolt11_payment", + AsyncMock(return_value=plan), + ), + patch( + "routstr.upstream.auto_topup._record_ppq_invoice", + AsyncMock(return_value=2_000_000_000), + ), + patch( + "routstr.upstream.auto_topup.execute_bolt11_payment", + AsyncMock(return_value=(101, "https://mint.test", "sat")), + ), + patch( + "routstr.upstream.auto_topup._record_ppq_payment_spent", AsyncMock() + ) as spent, + patch( + "routstr.upstream.auto_topup._mark_ppq_reconcile", AsyncMock() + ) as reconcile, + patch( + "routstr.upstream.auto_topup._set_ppq_state_terminal", AsyncMock() + ) as terminal, + patch("routstr.upstream.auto_topup.sats_usd_price", return_value=0.001), + patch("routstr.upstream.auto_topup.logger.critical") as critical, + ): + await _check_and_topup(_ppq_row()) + + spent.assert_awaited_once_with("operation-1", 101) + reconcile.assert_awaited_once() + terminal.assert_not_awaited() + assert "settlement polling failed" in critical.call_args.args[0] + + +@pytest.mark.asyncio +async def test_ppq_preflight_funding_check_happens_before_invoice_creation() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=2.5) + provider.initiate_topup = AsyncMock() + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch( + "routstr.upstream.auto_topup.maximum_owner_cashu_balance_sats", + AsyncMock(return_value=1), + ), + patch("routstr.upstream.auto_topup._claim_ppq_topup", AsyncMock()) as claim, + patch("routstr.upstream.auto_topup.sats_usd_price", return_value=0.001), + ): + await _check_and_topup(_ppq_row()) + + provider.initiate_topup.assert_not_awaited() + claim.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_active_claim_at_cycle_start_suppresses_topup_for_whole_cycle() -> None: + row = _ppq_row() + row.id = 1 + session = AsyncMock() + result = MagicMock() + result.all.return_value = [row] + session.exec.return_value = result + context = MagicMock() + context.__aenter__ = AsyncMock(return_value=session) + context.__aexit__ = AsyncMock(return_value=None) + + with ( + patch( + "routstr.upstream.auto_topup._reconcile_all_ppq_claims", + AsyncMock(return_value={1}), + ), + patch("routstr.upstream.auto_topup.create_session", return_value=context), + patch("routstr.upstream.auto_topup._check_and_topup", AsyncMock()) as check, + ): + await _run_auto_topup_cycle() + + check.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_ppq_auto_topup_skips_when_balance_meets_threshold() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=5.0) + provider.initiate_topup = AsyncMock() + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + ): + await _check_and_topup(_ppq_row()) + + provider.initiate_topup.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_ppq_auto_topup_skips_when_daily_spend_cap_reached() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=2.5) + provider.initiate_topup = AsyncMock() + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch( + "routstr.upstream.auto_topup.maximum_owner_cashu_balance_sats", + AsyncMock(return_value=10_000_000), + ), + # 1000 USD already spent, exactly the daily cap: the next 10 USD + # top-up must be refused. + patch( + "routstr.upstream.auto_topup._ppq_spent_last_24h_usd", + AsyncMock(return_value=1000.0), + ), + patch( + "routstr.upstream.auto_topup._claim_ppq_topup", + AsyncMock(), + ) as claim, + patch("routstr.upstream.auto_topup.sats_usd_price", return_value=0.001), + ): + await _check_and_topup(_ppq_row()) + + claim.assert_not_awaited() + provider.initiate_topup.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_ppq_pending_attempt_suppresses_duplicate_topup() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock() + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=True), + ), + ): + await _check_and_topup(_ppq_row()) + + provider.get_balance.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_ppq_auto_topup_rejects_non_finite_balance() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=float("nan")) + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch("routstr.upstream.auto_topup._claim_ppq_topup", AsyncMock()) as claim, + ): + await _check_and_topup(_ppq_row()) + + claim.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_settled_topup_alerts_when_its_claim_was_already_released() -> None: + provider = MagicMock() + provider.get_balance = AsyncMock(return_value=2.5) + provider.initiate_topup = AsyncMock( + return_value=MagicMock( + invoice_id="invoice-1", + payment_request="lnbc-invoice", + amount=10, + currency="USD", + expires_at=None, + ) + ) + provider.check_topup_status = AsyncMock(return_value=True) + plan = MagicMock() + plan.maximum_spend_sats = 102 + plan.quote.amount = 100 + plan.quote.fee_reserve = 2 + plan.mint_url = "https://mint-rich.test" + plan.unit = "sat" + + with ( + patch( + "routstr.upstream.auto_topup.PPQAIUpstreamProvider.from_db_row", + return_value=provider, + ), + patch( + "routstr.upstream.auto_topup._reconcile_ppq_state", + AsyncMock(return_value=False), + ), + patch( + "routstr.upstream.auto_topup._claim_ppq_topup", + AsyncMock(return_value="operation-1"), + ), + patch( + "routstr.upstream.auto_topup.prepare_bolt11_payment", + AsyncMock(return_value=plan), + ), + patch( + "routstr.upstream.auto_topup.maximum_owner_cashu_balance_sats", + AsyncMock(return_value=10_000), + ), + patch( + "routstr.upstream.auto_topup._ppq_spent_last_24h_usd", + AsyncMock(return_value=0.0), + ), + patch( + "routstr.upstream.auto_topup.execute_bolt11_payment", + AsyncMock(return_value=(101, "https://mint-rich.test", "sat")), + ), + patch("routstr.upstream.auto_topup._record_ppq_invoice", AsyncMock()), + patch("routstr.upstream.auto_topup._record_ppq_payment_spent", AsyncMock()), + patch( + "routstr.upstream.auto_topup._set_ppq_state_terminal", + AsyncMock(return_value=False), + ), + patch("routstr.upstream.auto_topup.sats_usd_price", return_value=0.001), + patch("routstr.upstream.auto_topup.logger") as log, + ): + await _check_and_topup(_ppq_row()) + + assert any( + "claim was already released" in call.args[0] + for call in log.critical.call_args_list + ) + + +@pytest.mark.parametrize( + ("settings", "expected"), + [ + ({"auto_topup": False, "topup_threshold": -1}, None), + ( + {"auto_topup": True, "topup_threshold": 5, "topup_amount_limit": 10}, + None, + ), + ( + {"auto_topup": True, "topup_threshold": None, "topup_amount_limit": 10}, + "threshold", + ), + ( + {"auto_topup": True, "topup_threshold": 5, "topup_amount_limit": 0.5}, + "whole number", + ), + ( + {"auto_topup": True, "topup_threshold": 5, "topup_amount_limit": 5000}, + "between", + ), + ( + {"auto_topup": True, "topup_threshold": True, "topup_amount_limit": 10}, + "threshold", + ), + ], +) +def test_ppq_auto_topup_settings_validation( + settings: dict, expected: str | None +) -> None: + problem = validate_ppq_auto_topup_settings(settings) + if expected is None: + assert problem is None + else: + assert problem is not None and expected in problem + + +def test_ppq_auto_topup_settings_validation_survives_huge_json_integers() -> None: + # json.loads happily produces integers past float range; float() raises + # OverflowError there instead of returning inf. + problem = validate_ppq_auto_topup_settings( + {"auto_topup": True, "topup_threshold": 10**400, "topup_amount_limit": 10} + ) + assert problem is not None and "threshold" in problem diff --git a/tests/unit/test_wallet.py b/tests/unit/test_wallet.py index 800e4fb9..54e8d4dd 100644 --- a/tests/unit/test_wallet.py +++ b/tests/unit/test_wallet.py @@ -4,7 +4,7 @@ import json import socket from collections.abc import AsyncIterator, Generator from contextlib import asynccontextmanager -from unittest.mock import AsyncMock, Mock, patch +from unittest.mock import AsyncMock, MagicMock, Mock, patch import httpx import pytest @@ -12,13 +12,17 @@ from cashu.core.base import MeltQuoteState from routstr.core.db import ApiKey from routstr.wallet import ( + Bolt11PaymentAmbiguous, + Bolt11PaymentNotAttempted, MintConnectionError, TokenConsumedError, _is_mint_rate_limited, classify_redemption_error, credit_balance, + execute_bolt11_payment, get_balance, is_mint_connection_error, + prepare_bolt11_payment, recieve_token, send, send_token, @@ -239,9 +243,7 @@ async def test_recieve_token_uses_only_requested_destination_mint() -> None: ) assert result == (99, "sat", destination) - swap.assert_awaited_once_with( - token, source_wallet, destination_mints=[destination] - ) + swap.assert_awaited_once_with(token, source_wallet, destination_mints=[destination]) @pytest.mark.asyncio @@ -492,9 +494,7 @@ async def test_send_refreshes_reservations_inside_wallet_guard() -> None: ): assert await send(1000, "sat", mint) == (1000, "token") - wallet.set_reserved_for_send.assert_awaited_once_with( - [proof], reserved=True - ) + wallet.set_reserved_for_send.assert_awaited_once_with([proof], reserved=True) @pytest.mark.asyncio @@ -888,9 +888,7 @@ def _make_swap_mocks( quote=f"melt_quote_{invoice}", amount=invoice, fee_reserve=_next_fee() ) ) - mock_token_wallet.melt = AsyncMock( - return_value=Mock(state=MeltQuoteState.paid) - ) + mock_token_wallet.melt = AsyncMock(return_value=Mock(state=MeltQuoteState.paid)) return mock_token, mock_token_wallet, mock_primary_wallet @@ -1804,6 +1802,220 @@ async def test_swap_melt_transport_error_is_never_reported_reusable() -> None: assert mock_token_wallet.melt.call_count == 1 +@pytest.mark.asyncio +async def test_execute_bolt11_payment_rejects_unpaid_melt_state() -> None: + plan = MagicMock() + plan.proofs = [MagicMock(amount=110)] + plan.quote.amount = 100 + plan.quote.fee_reserve = 10 + plan.quote.quote = "quote-1" + plan.invoice = "lnbc-invoice" + plan.wallet.select_to_send = AsyncMock(return_value=(plan.proofs, 0)) + plan.wallet.set_reserved_for_send = AsyncMock() + plan.wallet.melt = AsyncMock(return_value=MagicMock(state="UNPAID", change=[])) + + with pytest.raises(Bolt11PaymentNotAttempted): + await execute_bolt11_payment(plan) + + # An explicit unpaid answer means the proofs are ours again. + plan.wallet.set_reserved_for_send.assert_awaited_with(plan.proofs, reserved=False) + + +@pytest.mark.asyncio +async def test_execute_bolt11_payment_accepts_legacy_paid_response() -> None: + plan = MagicMock() + plan.proofs = [MagicMock(amount=110)] + plan.quote.amount = 100 + plan.quote.fee_reserve = 10 + plan.quote.quote = "quote-1" + plan.invoice = "lnbc-invoice" + plan.mint_url = "https://mint.test" + plan.unit = "sat" + plan.wallet.select_to_send = AsyncMock(return_value=(plan.proofs, 0)) + plan.wallet.set_reserved_for_send = AsyncMock() + plan.wallet.melt = AsyncMock( + return_value=MagicMock(state=None, paid=True, change=[]) + ) + + assert await execute_bolt11_payment(plan) == ( + 110, + "https://mint.test", + "sat", + ) + + +@pytest.mark.asyncio +async def test_execute_bolt11_payment_keeps_proofs_reserved_when_melt_errors() -> None: + plan = MagicMock() + plan.proofs = [MagicMock(amount=110)] + plan.quote.amount = 100 + plan.quote.fee_reserve = 10 + plan.quote.quote = "quote-1" + plan.invoice = "lnbc-invoice" + plan.wallet.select_to_send = AsyncMock(return_value=(plan.proofs, 0)) + plan.wallet.set_reserved_for_send = AsyncMock() + plan.wallet.set_reserved_for_melt = AsyncMock() + plan.wallet.melt = AsyncMock(side_effect=TimeoutError("no answer")) + + with pytest.raises(Bolt11PaymentAmbiguous): + await execute_bolt11_payment(plan) + + # The mint may still settle with these proofs. cashu's own melt() + # un-reserves them on a mint transport error, so the ambiguous path must + # re-reserve — and it must do so with the melt quote id, because + # get_melt_quote() finds the proofs to settle by melt_id. + plan.wallet.set_reserved_for_melt.assert_awaited_once_with( + plan.proofs, reserved=True, quote_id="quote-1" + ) + + +@pytest.mark.asyncio +async def test_execute_bolt11_payment_does_not_reserve_when_selection_fails() -> None: + plan = MagicMock() + plan.proofs = [MagicMock(amount=110)] + plan.quote.amount = 100 + plan.quote.fee_reserve = 10 + plan.wallet.select_to_send = AsyncMock(side_effect=ValueError("insufficient")) + plan.wallet.set_reserved_for_send = AsyncMock() + plan.wallet.melt = AsyncMock() + + with pytest.raises(Bolt11PaymentNotAttempted): + await execute_bolt11_payment(plan) + + plan.wallet.set_reserved_for_send.assert_not_awaited() + plan.wallet.melt.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_prepare_bolt11_payment_counts_input_fees_in_sufficiency() -> None: + from routstr.core.settings import settings + + wallet = MagicMock() + wallet.proofs = [MagicMock(amount=105)] + wallet.melt_quote = AsyncMock( + return_value=MagicMock(amount=100, fee_reserve=2, quote="quote-1") + ) + # Balance covers amount + fee_reserve (102) but not the 5 sat input fee. + wallet.get_fees_for_proofs = Mock(return_value=5) + + async def get_wallet(mint_url: str, unit: str = "sat", **_: object) -> MagicMock: + if unit == "msat": + raise ValueError("unit unsupported") + return wallet + + with ( + patch.object(settings, "cashu_mints", ["https://only.test"]), + patch.object(settings, "primary_mint", "https://only.test"), + patch("routstr.wallet.get_wallet", side_effect=get_wallet), + patch( + "routstr.wallet.get_proofs_per_mint_and_unit", + side_effect=lambda wallet, *args, **kwargs: wallet.proofs, + ), + patch( + "routstr.wallet.slow_filter_spend_proofs", + side_effect=lambda proofs, wallet: proofs, + ), + pytest.raises(ValueError, match="enough balance"), + ): + await prepare_bolt11_payment("lnbc-invoice") + + +@pytest.mark.asyncio +async def test_prepare_bolt11_payment_does_not_spend_user_liabilities() -> None: + from routstr.core.settings import settings + + wallet = MagicMock() + wallet.proofs = [MagicMock(amount=500)] + wallet.melt_quote = AsyncMock( + return_value=MagicMock(amount=100, fee_reserve=2, quote="quote-1") + ) + wallet.get_fees_for_proofs = Mock(return_value=0) + + async def get_wallet(mint_url: str, unit: str = "sat", **_: object) -> MagicMock: + if unit == "msat": + raise ValueError("unit unsupported") + return wallet + + with ( + patch.object(settings, "cashu_mints", ["https://only.test"]), + patch.object(settings, "primary_mint", "https://only.test"), + patch("routstr.wallet.get_wallet", side_effect=get_wallet), + patch( + "routstr.wallet.get_proofs_per_mint_and_unit", + side_effect=lambda wallet, *args, **kwargs: wallet.proofs, + ), + patch( + "routstr.wallet.slow_filter_spend_proofs", + side_effect=lambda proofs, wallet: proofs, + ), + patch( + "routstr.wallet._owner_balance_for_mint_and_unit", + AsyncMock(return_value=90), + ), + pytest.raises(ValueError, match="user liabilities"), + ): + await prepare_bolt11_payment("lnbc-invoice") + + +@pytest.mark.asyncio +async def test_prepare_bolt11_payment_rounds_user_liability_up_to_whole_sats() -> None: + from routstr.core.settings import settings + + wallet = MagicMock() + wallet.proofs = [MagicMock(amount=100)] + wallet.melt_quote = AsyncMock( + return_value=MagicMock(amount=1, fee_reserve=0, quote="quote-1") + ) + wallet.get_fees_for_proofs = Mock(return_value=0) + + async def get_wallet(mint_url: str, unit: str = "sat", **_: object) -> MagicMock: + if unit == "msat": + raise ValueError("unit unsupported") + return wallet + + with ( + patch.object(settings, "cashu_mints", ["https://only.test"]), + patch.object(settings, "primary_mint", "https://only.test"), + patch("routstr.wallet.get_wallet", side_effect=get_wallet), + patch( + "routstr.wallet.get_proofs_per_mint_and_unit", + side_effect=lambda wallet, *args, **kwargs: wallet.proofs, + ), + patch( + "routstr.wallet.slow_filter_spend_proofs", + side_effect=lambda proofs, wallet: proofs, + ), + patch( + "routstr.wallet.db.total_user_liability", + AsyncMock(return_value=99_999), + ), + pytest.raises(ValueError, match="user liabilities"), + ): + await prepare_bolt11_payment("lnbc-invoice") + + +@pytest.mark.asyncio +async def test_execute_bolt11_payment_rereserves_when_cancelled() -> None: + plan = MagicMock() + plan.proofs = [MagicMock(amount=110)] + plan.quote.amount = 100 + plan.quote.fee_reserve = 10 + plan.quote.quote = "quote-1" + plan.invoice = "lnbc-invoice" + plan.mint_url = "https://mint.test" + plan.wallet.select_to_send = AsyncMock(return_value=(plan.proofs, 0)) + plan.wallet.set_reserved_for_send = AsyncMock() + plan.wallet.set_reserved_for_melt = AsyncMock() + plan.wallet.melt = AsyncMock(side_effect=asyncio.CancelledError()) + + with pytest.raises(asyncio.CancelledError): + await execute_bolt11_payment(plan) + + plan.wallet.set_reserved_for_melt.assert_awaited_once_with( + plan.proofs, reserved=True, quote_id="quote-1" + ) + + # --------------------------------------------------------------------------- # Per-mint adaptive guard + _mint_operation factory/retry # --------------------------------------------------------------------------- @@ -2014,9 +2226,7 @@ async def test_default_timeout_allows_retry_after_rate_limit_cooldown() -> None: response = httpx.Response(429, request=request) operation = AsyncMock( side_effect=[ - httpx.HTTPStatusError( - "rate limited", request=request, response=response - ), + httpx.HTTPStatusError("rate limited", request=request, response=response), "ok", ] ) diff --git a/ui/components/provider-card.tsx b/ui/components/provider-card.tsx index 851dda85..44f05183 100644 --- a/ui/components/provider-card.tsx +++ b/ui/components/provider-card.tsx @@ -1,3 +1,5 @@ +import { AdminService } from '@/lib/api/services/admin'; +import type { PPQAutoTopupState } from '@/lib/api/services/admin'; import type { AdminModel, ProviderModels, @@ -20,12 +22,16 @@ import { Trash2, Key, RotateCcw, + AlertTriangle, + Unlock, + Loader2, } from 'lucide-react'; import { ProviderBalance } from '@/components/provider-balance'; import { ProviderModelsPanel } from '@/components/provider-models-panel'; import { RoutstrCreateKeySection } from '@/components/providers/RoutstrCreateKeySection'; import { RoutstrProviderService } from '@/lib/api/services/routstr-provider'; -import { useMutation, useQueryClient } from '@tanstack/react-query'; +import { getErrorStatus } from '@/lib/api/client'; +import { useMutation, useQuery, useQueryClient } from '@tanstack/react-query'; import { useState } from 'react'; import { toast } from 'sonner'; import { cn } from '@/lib/utils'; @@ -36,6 +42,16 @@ import { DialogHeader, DialogTitle, } from '@/components/ui/dialog'; +import { + AlertDialog, + AlertDialogAction, + AlertDialogCancel, + AlertDialogContent, + AlertDialogDescription, + AlertDialogFooter, + AlertDialogHeader, + AlertDialogTitle, +} from '@/components/ui/alert-dialog'; interface ProviderCardProps { provider: UpstreamProvider; @@ -77,8 +93,71 @@ export function ProviderCard({ }: ProviderCardProps) { const queryClient = useQueryClient(); const [isKeyModalOpen, setIsKeyModalOpen] = useState(false); + const [isReleaseDialogOpen, setIsReleaseDialogOpen] = useState(false); + // The claim as the query cache held it when the admin opened the dialog. + // The mutation sends this token rather than re-reading the query at submit + // time: a background refetch after the dialog opened must not swap in a + // state the admin never saw. The server rejects a stale token with a 409, + // which is the authoritative guard. + const [reviewedState, setReviewedState] = useState( + null + ); const hasDetails = Boolean(provider.api_version) || isExpanded; const isRoutstr = provider.provider_type === 'routstr'; + const isPPQ = provider.provider_type === 'ppqai'; + + const { data: ppqAutoTopupState, isError: ppqStateFetchFailed } = useQuery({ + queryKey: ['ppq-auto-topup-state', provider.id], + queryFn: () => AdminService.getPPQAutoTopupState(provider.id), + enabled: isPPQ, + refetchInterval: 30000, + }); + + // A claim the server will not let us release: a worker is between reserving + // proofs and hearing back from the mint, and sweeping it would let the next + // cycle pay a second invoice. + const isPPQPaymentInFlight = + Boolean(ppqAutoTopupState?.active) && + ppqAutoTopupState?.releasable === false; + + const openReleaseDialog = () => { + setReviewedState(ppqAutoTopupState ?? null); + setIsReleaseDialogOpen(true); + }; + + const releasePPQMutation = useMutation({ + mutationFn: () => + AdminService.releasePPQAutoTopup( + provider.id, + reviewedState?.state_token ?? null + ), + onSuccess: () => { + queryClient.invalidateQueries({ + queryKey: ['ppq-auto-topup-state', provider.id], + }); + setIsReleaseDialogOpen(false); + setReviewedState(null); + toast.success('PPQ auto top-up claim released'); + }, + onError: (error: Error) => { + queryClient.invalidateQueries({ + queryKey: ['ppq-auto-topup-state', provider.id], + }); + if (getErrorStatus(error) === 409) { + // The claim changed since it was reviewed; the stale snapshot is + // useless, so force a fresh review. + setIsReleaseDialogOpen(false); + setReviewedState(null); + toast.error( + 'PPQ claim changed since it was reviewed; reopen to see the new state' + ); + return; + } + // Transient failure: keep the dialog and the reviewed snapshot so the + // admin can retry without re-navigating. + toast.error(`Failed to release PPQ claim: ${error.message}`); + }, + }); const refundMutation = useMutation({ mutationFn: () => RoutstrProviderService.refundBalance(provider.id), @@ -113,6 +192,35 @@ export function ProviderCard({ > {provider.enabled ? 'Enabled' : 'Disabled'} + {ppqAutoTopupState?.active && ( + + {isPPQPaymentInFlight ? ( + + ) : ( + + )} + {isPPQPaymentInFlight + ? 'Paying invoice' + : 'Auto top-up needs review'} + + )} + {isPPQ && ppqStateFetchFailed && ( + + + Top-up status unavailable + + )} {provider.base_url} @@ -153,6 +261,19 @@ export function ProviderCard({ )} + {isPPQ && ppqAutoTopupState?.active && !isPPQPaymentInFlight && ( + + )} + {isRoutstr && provider.api_key && (