From a2e2a5c662427362069aef00f6b29436c7c9025c Mon Sep 17 00:00:00 2001 From: 9qeklajc Date: Tue, 4 Aug 2026 01:13:20 +0200 Subject: [PATCH] clean up --- routstr/upstream/auto_topup.py | 188 ++++++++++++------ routstr/wallet.py | 85 ++++---- tests/unit/test_auto_topup.py | 37 ++++ ui/components/provider-card.tsx | 35 +++- .../provider-form-dialog-content.tsx | 8 +- .../providers/PPQAutoTopupSettings.tsx | 69 +++++-- 6 files changed, 296 insertions(+), 126 deletions(-) diff --git a/routstr/upstream/auto_topup.py b/routstr/upstream/auto_topup.py index f6e46d59..fb1cc914 100644 --- a/routstr/upstream/auto_topup.py +++ b/routstr/upstream/auto_topup.py @@ -29,6 +29,7 @@ from ..wallet import ( release_token_reservation, send_token, token_mint_url, + wallet_operation_guard, ) from .ppqai import PPQAIUpstreamProvider from .routstr import RoutstrUpstreamProvider @@ -51,6 +52,12 @@ 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 = 1000 async def periodic_auto_topup() -> None: @@ -822,6 +829,28 @@ async def _mark_ppq_reconcile( ) +async def _ppq_spent_last_24h_sats() -> int: + """Total sats committed to PPQ top-ups in the last 24 hours. + + Counts every payment audit row, including in-flight and ambiguous ones: + for spend-cap purposes an unresolved payment must be assumed spent. + """ + cutoff = int(time.time()) - 24 * 60 * 60 + async with create_session() as session: + rows = ( + await session.exec( + select(CashuTransaction.amount, CashuTransaction.unit).where( + col(CashuTransaction.source) == "ppq_auto_topup", + col(CashuTransaction.type) == "out", + col(CashuTransaction.created_at) >= cutoff, + ) + ) + ).all() + return sum( + amount if unit == "sat" else math.ceil(amount / 1000) for amount, unit 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"]) @@ -855,6 +884,19 @@ async def _check_and_topup_ppq(row: UpstreamProviderRow, settings: dict) -> None ) return + spent_24h_usd = await _ppq_spent_last_24h_sats() * 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 @@ -882,21 +924,6 @@ async def _check_and_topup_ppq(row: UpstreamProviderRow, settings: dict) -> None invoice_expires_at //= 1000 if invoice_expires_at <= now: raise ValueError("PPQ returned an expired Lightning invoice") - - 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), - 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. @@ -910,61 +937,96 @@ async def _check_and_topup_ppq(row: UpstreamProviderRow, settings: dict) -> None ) 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}, + # 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: + 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), + unit=plan.unit, + mint_url=plan.mint_url, ) - 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( + 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( - "Unexpected failure while paying PPQ invoice; payment requires reconciliation", - extra={"provider_id": row.id, "invoice_id": topup.invoice_id}, - exc_info=True, - ) - raise + 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 diff --git a/routstr/wallet.py b/routstr/wallet.py index 064080cd..cbd2174c 100644 --- a/routstr/wallet.py +++ b/routstr/wallet.py @@ -19,19 +19,14 @@ from pydantic_core import PydanticUndefined from sqlmodel import col, select, update from .core import db, get_logger -from .core.db import store_cashu_transaction_with_retry as store_cashu_transaction +from .core.db import \ + store_cashu_transaction_with_retry as store_cashu_transaction from .core.settings import settings -from .mint import ( - MINT_TRANSPORT_COOLDOWN_SECONDS, - MINT_TRANSPORT_EXCEPTIONS, - MintRateGuard, - MintRateLimitedError, - fail_fast_mint_operations, - is_mint_rate_limited, - mint_cooldown_reason, - mint_cooldown_remaining, - run_mint_operation, -) +from .mint import (MINT_TRANSPORT_COOLDOWN_SECONDS, MINT_TRANSPORT_EXCEPTIONS, + MintRateGuard, MintRateLimitedError, + fail_fast_mint_operations, is_mint_rate_limited, + mint_cooldown_reason, mint_cooldown_remaining, + run_mint_operation) from .payment.lnurl import raw_send_to_lnurl # Backwards-compatible aliases for callers/tests that imported the former @@ -396,9 +391,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( @@ -560,11 +553,13 @@ 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) - ) + 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"])) + ( + 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") @@ -577,7 +572,17 @@ async def prepare_bolt11_payment(invoice: str) -> Bolt11PaymentPlan: 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]] = [] @@ -631,7 +636,7 @@ async def prepare_bolt11_payment(invoice: str) -> Bolt11PaymentPlan: 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( @@ -648,7 +653,15 @@ async def execute_bolt11_payment(plan: Bolt11PaymentPlan) -> tuple[int, str, str 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: @@ -684,8 +697,8 @@ async def execute_bolt11_payment(plan: Bolt11PaymentPlan) -> tuple[int, str, str try: await asyncio.shield( plan.wallet.set_reserved_for_melt( - selected, reserved=True, quote_id=plan.quote.quote - ) + selected, reserved=True, quote_id=plan.quote.quote + ) ) except BaseException: logger.critical( @@ -1972,9 +1985,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 ) @@ -2088,9 +2099,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: @@ -2329,11 +2338,8 @@ async def periodic_refund_sweep() -> None: async def periodic_routstr_fee_payout() -> None: - from .auth import ( - ROUTSTR_FEE_DEFAULT_PAYOUT, - ROUTSTR_FEE_PAYOUT_INTERVAL_SECONDS, - ROUTSTR_LN_ADDRESS, - ) + from .auth import (ROUTSTR_FEE_DEFAULT_PAYOUT, + ROUTSTR_FEE_PAYOUT_INTERVAL_SECONDS, ROUTSTR_LN_ADDRESS) if not ROUTSTR_LN_ADDRESS: logger.info("ROUTSTR_LN_ADDRESS not set, skipping fee payout") @@ -2434,16 +2440,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) @@ -2469,3 +2469,4 @@ async def send_to_lnurl(amount: int, unit: str, mint: str, address: str) -> int: # def refund_partial(self, amount: int) -> None: # raise NotImplementedError + diff --git a/tests/unit/test_auto_topup.py b/tests/unit/test_auto_topup.py index f74af477..c5018991 100644 --- a/tests/unit/test_auto_topup.py +++ b/tests/unit/test_auto_topup.py @@ -530,6 +530,43 @@ async def test_ppq_auto_topup_skips_when_balance_meets_threshold() -> None: 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), + ), + # 1_000_000 sats * 0.001 USD/sat = 1000 USD, the daily cap: the next + # 10 USD top-up must be refused. + patch( + "routstr.upstream.auto_topup._ppq_spent_last_24h_sats", + AsyncMock(return_value=1_000_000), + ), + 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() diff --git a/ui/components/provider-card.tsx b/ui/components/provider-card.tsx index 3cdce61d..943fac88 100644 --- a/ui/components/provider-card.tsx +++ b/ui/components/provider-card.tsx @@ -30,6 +30,7 @@ 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 { ApiError } from '@/lib/api/client'; import { useMutation, useQuery, useQueryClient } from '@tanstack/react-query'; import { useState } from 'react'; import { toast } from 'sonner'; @@ -93,9 +94,11 @@ export function ProviderCard({ const queryClient = useQueryClient(); const [isKeyModalOpen, setIsKeyModalOpen] = useState(false); const [isReleaseDialogOpen, setIsReleaseDialogOpen] = useState(false); - // Snapshot of the claim as it looked when the admin opened the dialog. - // The mutation sends this token, never the live query data: a background - // refetch must not swap in a state the admin never reviewed. + // 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 ); @@ -103,7 +106,7 @@ export function ProviderCard({ const isRoutstr = provider.provider_type === 'routstr'; const isPPQ = provider.provider_type === 'ppqai'; - const { data: ppqAutoTopupState } = useQuery({ + const { data: ppqAutoTopupState, isError: ppqStateFetchFailed } = useQuery({ queryKey: ['ppq-auto-topup-state', provider.id], queryFn: () => AdminService.getPPQAutoTopupState(provider.id), enabled: isPPQ, @@ -137,12 +140,21 @@ export function ProviderCard({ toast.success('PPQ auto top-up claim released'); }, onError: (error: Error) => { - // A 409 usually means the claim changed since it was reviewed. queryClient.invalidateQueries({ queryKey: ['ppq-auto-topup-state', provider.id], }); - setIsReleaseDialogOpen(false); - setReviewedState(null); + if (error instanceof ApiError && error.status === 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}`); }, }); @@ -200,6 +212,15 @@ export function ProviderCard({ : 'Auto top-up needs review'} )} + {isPPQ && ppqStateFetchFailed && ( + + + Top-up status unavailable + + )} {provider.base_url} diff --git a/ui/components/provider-form-dialog-content.tsx b/ui/components/provider-form-dialog-content.tsx index 5ef761a4..b3872b63 100644 --- a/ui/components/provider-form-dialog-content.tsx +++ b/ui/components/provider-form-dialog-content.tsx @@ -12,6 +12,7 @@ import { DialogTitle, } from '@/components/ui/dialog'; import { ProviderFormFields } from '@/components/provider-form-fields'; +import { ppqAutoTopupSettingsInvalid } from '@/components/providers/PPQAutoTopupSettings'; interface ProviderFormDialogContentProps { mode: 'create' | 'edit'; @@ -52,6 +53,11 @@ export function ProviderFormDialogContent({ isSubmitting, availableMints, }: ProviderFormDialogContentProps) { + // The server re-validates these bounds; this only stops submitting a form + // whose inline errors are already visible. + const hasInvalidSettings = + formData.provider_type === 'ppqai' && + ppqAutoTopupSettingsInvalid(formData.provider_settings || {}); return ( @@ -80,7 +86,7 @@ export function ProviderFormDialogContent({