diff --git a/.env.example b/.env.example index a063da56..8d16580d 100644 --- a/.env.example +++ b/.env.example @@ -51,6 +51,11 @@ ROUTSTR_SECRET_KEY= # MINT_OPERATION_TIMEOUT_SECONDS=30 # MINT_MAX_CONCURRENCY=4 # MINT_RETRY_MAX_ATTEMPTS=3 +# Foreign token policy: reject, or swap into PRIMARY_MINT_URL over Lightning. +# FOREIGN_MINT_POLICY=reject +# FOREIGN_MINT_OPERATION_TIMEOUT_SECONDS=5 +# FOREIGN_MINT_MAX_CONCURRENCY=4 +# SWAP_RECONCILE_INTERVAL_SECONDS=60 # RECEIVE_LN_ADDRESS= # REFUND_SWEEP_CLAIM_TIMEOUT_SECONDS=900 diff --git a/docs/api/errors.md b/docs/api/errors.md index 5f151bb6..1764befd 100644 --- a/docs/api/errors.md +++ b/docs/api/errors.md @@ -193,7 +193,9 @@ granularity) on any of them. | `token_already_spent` | 400 | `cashu_token_already_spent` | No | The token was already redeemed. | | `invalid_token` | 400 | `invalid_cashu_token` | No | The token is malformed or cannot be decoded. | | `mint_error` | 422 | `cashu_token_swap_fees_exceed_amount` | No | Token value is too small to cover the mint's NUT-02 input fees. | -| `untrusted_mint` | 400 | `cashu_untrusted_source_mint` | No | The token was issued by a mint this node does not accept. Only the node's configured mints (`PRIMARY_MINT_URL` / `CASHU_MINTS`) are redeemable. | +| `untrusted_mint` | 400 | `cashu_untrusted_source_mint` | No | The token was issued by a mint this node does not accept. With the default `FOREIGN_MINT_POLICY=reject` only the node's configured mints (`PRIMARY_MINT_URL` / `CASHU_MINTS`) are redeemable. Bearer and X-Cashu payments always answer this for a foreign mint; only `/v1/wallet/topup` swaps when the operator enabled it. | +| `mint_error` | 422 | `cashu_foreign_mint_swap_failed` | No | Top-up only, `FOREIGN_MINT_POLICY=swap`: the token could not be swapped into the node's mint (fees exceed its value, unsupported unit, non-HTTPS mint URL, or the issuing mint refused the payment). Nothing was spent; the token is still yours. | +| `swap_pending` | 409 | `cashu_swap_pending` | No | Top-up only: the swap's Lightning payment was dispatched but the issuing mint has not confirmed it. Do **not** resend the token (its proofs may be spent). The balance is credited automatically once the payment is confirmed; poll `/v1/wallet/info`. | | `mint_unreachable` | 503 | `cashu_source_mint_unreachable` | **Yes** | The mint that issued the token could not be reached; it cannot be redeemed at another mint. | | `mint_rate_limited` | 503 | `cashu_mint_rate_limited` | **Yes** | The mint rate-limited the request; retry after the cooldown. | | `mint_timeout` | 503 | `cashu_mint_timeout` | **Yes** | The mint did not respond in time; retry later. | @@ -208,7 +210,8 @@ granularity) on any of them. retryable — the same token may work again later. Everything else is a permanent property of the token and must not be blindly retried. `untrusted_mint` is permanent: the node will never accept that mint until - an operator adds it to `CASHU_MINTS`. Use exponential backoff for the + an operator adds it to `CASHU_MINTS` or enables `FOREIGN_MINT_POLICY=swap` + for top-ups. Use exponential backoff for the 503 responses, and honor the mint's cooldown for `mint_rate_limited`. In particular, a `token_consumed` 500 means the mint already spent the token, so a retry would fail as `token_already_spent`. diff --git a/docs/provider/configuration.md b/docs/provider/configuration.md index 20668e60..e7c3b0a9 100644 --- a/docs/provider/configuration.md +++ b/docs/provider/configuration.md @@ -171,6 +171,30 @@ Setting `CASHU_MINTS` (env) or editing the list in the dashboard replaces this default entirely. An explicitly empty value leaves only the primary mint trusted. +#### Tokens from other mints + +By default a token issued by a mint outside this list is refused offline with +`cashu_untrusted_source_mint`; the node never contacts a mint it does not trust. + +`FOREIGN_MINT_POLICY=swap` lets `/v1/wallet/topup` accept such tokens by +melting them over Lightning into the primary mint. Bearer and X-Cashu payments +still refuse foreign mints (those paths run on every request and must not wait +on a third-party mint). Refunds of a key funded this way are swapped back to the +user's own mint, net of fees. Safeguards when enabled: + +- The token's mint URL must be HTTPS to a public address. +- Calls to the foreign mint get one attempt with a short deadline and share a + process-wide concurrency cap, so a dead or hostile mint can only stall its own + swap. They never run while the wallet lock is held. +- Fees are quoted before anything is spent; a token that cannot cover them is + refused with `cashu_foreign_mint_swap_failed` and stays spendable. +- Every swap is journaled in `cashu_swaps` before the Lightning leg. A timeout + answers `cashu_swap_pending`; a background reconciler credits or fails the row + once the mint confirms the outcome. + +Lightning routing fees and the mint's input fees are deducted from the amount +credited (and from the refund). Leftover fee reserve stays on the foreign mint. + ### Lightning Withdrawals Automatic profit withdrawal: @@ -234,6 +258,10 @@ Use environment variables for: | `MINT_OPERATION_TIMEOUT_SECONDS` | Per-attempt timeout for mint network calls | `30` | | `MINT_MAX_CONCURRENCY` | Concurrent operations allowed per mint (`0` disables the limit) | `4` | | `MINT_RETRY_MAX_ATTEMPTS` | Retries after a timeout or HTTP 429 (`0` disables retries) | `3` | +| `FOREIGN_MINT_POLICY` | `reject` refuses top-up tokens from unconfigured mints; `swap` melts them into the primary mint (see above) | `reject` | +| `FOREIGN_MINT_OPERATION_TIMEOUT_SECONDS` | Single-attempt deadline for calls to an unconfigured mint | `5` | +| `FOREIGN_MINT_MAX_CONCURRENCY` | Process-wide cap on in-flight calls to unconfigured mints | `4` | +| `SWAP_RECONCILE_INTERVAL_SECONDS` | How often unfinished swaps are re-checked against their mints | `60` | | `RECEIVE_LN_ADDRESS` | Lightning address for withdrawals | — | | `MIN_PAYOUT_SAT` | Min payout balance in sats (applies to all mints) | `210` | | `MAX_PAYOUT_SAT` | Maximum gross budget per periodic payout in sats, including fees (all mints) | `250000` | diff --git a/migrations/versions/b7e2c4d9f1a3_add_cashu_swaps_table.py b/migrations/versions/b7e2c4d9f1a3_add_cashu_swaps_table.py new file mode 100644 index 00000000..a2d38785 --- /dev/null +++ b/migrations/versions/b7e2c4d9f1a3_add_cashu_swaps_table.py @@ -0,0 +1,67 @@ +"""add cashu_swaps table + +Revision ID: b7e2c4d9f1a3 +Revises: a73d19b6c204 +Create Date: 2026-10-04 + +""" + +import sqlalchemy as sa +import sqlmodel +from alembic import op + +revision = "b7e2c4d9f1a3" +down_revision = "a73d19b6c204" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "cashu_swaps", + sa.Column("id", sqlmodel.sql.sqltypes.AutoString(), nullable=False), + sa.Column("direction", sqlmodel.sql.sqltypes.AutoString(), nullable=False), + sa.Column("status", sqlmodel.sql.sqltypes.AutoString(), nullable=False), + sa.Column( + "api_key_hashed_key", sqlmodel.sql.sqltypes.AutoString(), nullable=True + ), + sa.Column("refund_id", sqlmodel.sql.sqltypes.AutoString(), nullable=True), + sa.Column("token_hash", sqlmodel.sql.sqltypes.AutoString(), nullable=True), + sa.Column("source_mint", sqlmodel.sql.sqltypes.AutoString(), nullable=False), + sa.Column("source_unit", sqlmodel.sql.sqltypes.AutoString(), nullable=False), + sa.Column("source_amount", sa.Integer(), nullable=False), + sa.Column( + "destination_mint", sqlmodel.sql.sqltypes.AutoString(), nullable=False + ), + sa.Column( + "destination_unit", sqlmodel.sql.sqltypes.AutoString(), nullable=False + ), + sa.Column("destination_amount", sa.Integer(), nullable=False), + sa.Column("fee_reserve", sa.Integer(), nullable=False), + sa.Column("input_fees", sa.Integer(), nullable=False), + sa.Column("mint_quote_id", sqlmodel.sql.sqltypes.AutoString(), nullable=True), + sa.Column("melt_quote_id", sqlmodel.sql.sqltypes.AutoString(), nullable=True), + sa.Column("token", sqlmodel.sql.sqltypes.AutoString(), nullable=True), + sa.Column("error", sqlmodel.sql.sqltypes.AutoString(), nullable=True), + sa.Column("claimed_at", sa.Integer(), nullable=True), + sa.Column("created_at", sa.Integer(), nullable=False), + sa.Column("updated_at", sa.Integer(), nullable=False), + sa.ForeignKeyConstraint(["api_key_hashed_key"], ["api_keys.hashed_key"]), + sa.PrimaryKeyConstraint("id"), + ) + op.create_index("ix_cashu_swaps_status", "cashu_swaps", ["status"]) + op.create_index( + "ix_cashu_swaps_api_key_hashed_key", "cashu_swaps", ["api_key_hashed_key"] + ) + op.create_index("ix_cashu_swaps_refund_id", "cashu_swaps", ["refund_id"]) + op.create_index( + "ix_cashu_swaps_token_hash", "cashu_swaps", ["token_hash"], unique=True + ) + + +def downgrade() -> None: + op.drop_index("ix_cashu_swaps_token_hash", table_name="cashu_swaps") + op.drop_index("ix_cashu_swaps_refund_id", table_name="cashu_swaps") + op.drop_index("ix_cashu_swaps_api_key_hashed_key", table_name="cashu_swaps") + op.drop_index("ix_cashu_swaps_status", table_name="cashu_swaps") + op.drop_table("cashu_swaps") diff --git a/routstr/__init__.py b/routstr/__init__.py index 1bdd854c..9e6cf08f 100644 --- a/routstr/__init__.py +++ b/routstr/__init__.py @@ -1,3 +1,17 @@ -from .core.main import app as fastapi_app # noqa +"""Routstr application package.""" + +import os + +# cashu's settings loader reads the nearest .env with ``override=True`` during +# import. A library must not replace the node's already-configured process +# environment (notably CASHU_MINTS and networking settings), so contain that +# side effect while importing the application graph. +_environment_before_import = dict(os.environ) +try: + from .core.main import app as fastapi_app +finally: + os.environ.clear() + os.environ.update(_environment_before_import) + del _environment_before_import __all__ = ["fastapi_app"] diff --git a/routstr/balance.py b/routstr/balance.py index f887f5ae..cd67078f 100644 --- a/routstr/balance.py +++ b/routstr/balance.py @@ -20,10 +20,13 @@ from .core.db import ( ) from .core.logging import get_logger from .core.settings import settings +from .foreign_mint_swap import swap_enabled, swap_in_and_credit from .lightning import lightning_router from .wallet import ( + UntrustedSourceMintError, classify_redemption_error, credit_balance, + is_trusted_source_mint, recieve_token, token_mint_url, ) @@ -183,7 +186,14 @@ async def topup_wallet_endpoint( }, ) try: - amount_msats = await credit_balance(cashu_token, billing_key, session) + if source_mint != "unknown" and not is_trusted_source_mint(source_mint): + if not swap_enabled(): + raise UntrustedSourceMintError(f"Untrusted source mint: {source_mint}") + # Top-up is the only entry point that swaps: the caller is already + # waiting on a long operation here, unlike bearer auth or X-Cashu. + amount_msats = await swap_in_and_credit(cashu_token, billing_key, session) + else: + amount_msats = await credit_balance(cashu_token, billing_key, session) except Exception as e: # Shared taxonomy so top-up matches the bearer/X-Cashu paths (503 for an # unreachable mint, 422 for fee/swap failures, 400 for token faults). diff --git a/routstr/core/db.py b/routstr/core/db.py index c17a1093..f8c92f5c 100644 --- a/routstr/core/db.py +++ b/routstr/core/db.py @@ -628,6 +628,59 @@ class Refund(SQLModel, table=True): # type: ignore updated_at: int = Field(default_factory=lambda: int(time.time())) +# Swap rows the reconciler still owns: the melt was dispatched and its outcome +# or follow-up (mint, credit, token issue) is not final. +SWAP_OPEN_STATUSES = ("melting", "ambiguous", "melted", "minted", "issued") + + +class CashuSwap(SQLModel, table=True): # type: ignore + """Journal of one cross-mint swap, written before any Lightning payment. + + ``in`` swaps melt a token from a mint the operator does not trust into the + primary mint and credit an API key. ``out`` swaps melt owner proofs on the + primary mint to issue a refund token on the user's own mint. Every money + movement is recorded here first so a crash or timeout leaves a row the + reconciler can finish or fail, never an unknown balance. + """ + + __tablename__ = "cashu_swaps" + + id: str = Field(primary_key=True, default_factory=lambda: uuid.uuid4().hex) + direction: str = Field(description="in (token -> primary) or out (refund)") + status: str = Field( + default="melting", + index=True, + description=( + "melting, ambiguous, melted, minted, credited, issued, settled, failed" + ), + ) + api_key_hashed_key: str | None = Field( + default=None, foreign_key="api_keys.hashed_key", index=True + ) + refund_id: str | None = Field(default=None, index=True) + token_hash: str | None = Field( + default=None, + index=True, + unique=True, + description="sha256 of the incoming token", + ) + source_mint: str = Field() + source_unit: str = Field() + source_amount: int = Field(description="Gross amount leaving the source mint") + destination_mint: str = Field() + destination_unit: str = Field() + destination_amount: int = Field(description="Net amount minted at the destination") + fee_reserve: int = Field(default=0) + input_fees: int = Field(default=0) + mint_quote_id: str | None = Field(default=None) + melt_quote_id: str | None = Field(default=None) + token: str | None = Field(default=None, description="Issued token (out swaps)") + error: str | None = Field(default=None) + claimed_at: int | None = Field(default=None, description="Reconciler lease") + created_at: int = Field(default_factory=lambda: int(time.time())) + updated_at: int = Field(default_factory=lambda: int(time.time())) + + async def store_cashu_transaction( token: str, amount: int, diff --git a/routstr/core/main.py b/routstr/core/main.py index cd4cab90..e2ed8c7d 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -18,6 +18,7 @@ from ..auth import ( ) from ..balance import balance_router, deprecated_wallet_router from ..cashu_compat import install_cashu_httpx_shim +from ..foreign_mint_swap import periodic_swap_reconcile from ..lightning import ( lightning_router, periodic_invoice_watcher, @@ -75,6 +76,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: auto_topup_task = None refund_sweep_task = None refund_reconcile_task = None + swap_reconcile_task = None routstr_fee_task = None invoice_watcher_task = None @@ -163,6 +165,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: auto_topup_task = asyncio.create_task(periodic_auto_topup()) refund_sweep_task = asyncio.create_task(periodic_refund_sweep()) refund_reconcile_task = asyncio.create_task(periodic_refund_reconcile()) + swap_reconcile_task = asyncio.create_task(periodic_swap_reconcile()) routstr_fee_task = asyncio.create_task(periodic_routstr_fee_payout()) invoice_watcher_task = asyncio.create_task(periodic_invoice_watcher()) @@ -208,6 +211,8 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: refund_sweep_task.cancel() if refund_reconcile_task is not None: refund_reconcile_task.cancel() + if swap_reconcile_task is not None: + swap_reconcile_task.cancel() if routstr_fee_task is not None: routstr_fee_task.cancel() if invoice_watcher_task is not None: @@ -243,6 +248,8 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: tasks_to_wait.append(refund_sweep_task) if refund_reconcile_task is not None: tasks_to_wait.append(refund_reconcile_task) + if swap_reconcile_task is not None: + tasks_to_wait.append(swap_reconcile_task) if routstr_fee_task is not None: tasks_to_wait.append(routstr_fee_task) if invoice_watcher_task is not None: diff --git a/routstr/core/settings.py b/routstr/core/settings.py index d77d6fc1..0e71ae1c 100644 --- a/routstr/core/settings.py +++ b/routstr/core/settings.py @@ -6,7 +6,7 @@ import os import secrets import time from datetime import datetime, timezone -from typing import Any +from typing import Any, Literal from pydantic.v1 import BaseModel, BaseSettings, Field from sqlmodel.ext.asyncio.session import AsyncSession @@ -100,6 +100,26 @@ class Settings(BaseSettings): mint_max_concurrency: int = Field(default=4, ge=0, env="MINT_MAX_CONCURRENCY") # Max retries when a mint returns 429 or times out (exponential backoff). mint_retry_max_attempts: int = Field(default=3, ge=0, env="MINT_RETRY_MAX_ATTEMPTS") + # What to do with a top-up token issued by a mint outside primary_mint / + # cashu_mints. "reject" refuses it offline. "swap" melts it over Lightning + # into the primary mint (and refunds back the same way) under a separate, + # short budget so a hostile or dead mint can never hold wallet state. + foreign_mint_policy: Literal["reject", "swap"] = Field( + default="reject", env="FOREIGN_MINT_POLICY" + ) + # Single-attempt deadline for any call to a mint the operator did not + # configure. No retries: the sender chose that mint, not the operator. + foreign_mint_operation_timeout_seconds: float = Field( + default=5.0, gt=0, env="FOREIGN_MINT_OPERATION_TIMEOUT_SECONDS" + ) + # Process-wide cap on in-flight foreign-mint calls across all such mints, so + # rotating hostnames cannot multiply the per-mint budget. + foreign_mint_max_concurrency: int = Field( + default=4, ge=1, env="FOREIGN_MINT_MAX_CONCURRENCY" + ) + swap_reconcile_interval_seconds: int = Field( + default=60, gt=0, env="SWAP_RECONCILE_INTERVAL_SECONDS" + ) # Pricing # Default behavior: derive pricing from MODELS diff --git a/routstr/foreign_mint_swap.py b/routstr/foreign_mint_swap.py new file mode 100644 index 00000000..c986ebc9 --- /dev/null +++ b/routstr/foreign_mint_swap.py @@ -0,0 +1,947 @@ +"""Cross-mint swaps for tokens from mints the operator did not configure. + +The original swap path was removed in cc55868e because it contacted the +sender's mint while holding the process-wide wallet lock. This version keeps +three things apart that were mixed before: + +* **Trust and destination checks** run offline first. The mint URL inside the + token is unauthenticated input and gets the same SSRF treatment as any other + client-supplied URL. +* **Foreign-mint I/O** runs outside ``wallet_operation_guard`` under its own + budget: one attempt, a short deadline, a process-wide concurrency cap and a + per-mint file lock. A dead or hostile mint can stall its own swap, nothing + else. +* **Wallet mutation** (minting on a trusted mint, crediting a key) runs under + the guard as before, but only after the Lightning leg has settled. + +Every money movement is journaled in ``cashu_swaps`` before it is dispatched, +so a crash or timeout leaves a row the reconciler can finish or fail. +""" + +import asyncio +import fcntl +import hashlib +import os +import re +import time +from contextlib import asynccontextmanager +from typing import Any, AsyncGenerator, Awaitable, Callable + +from cashu.core.base import MeltQuote, MintQuote, Proof, Token +from cashu.wallet.helpers import deserialize_token_from_string +from sqlalchemy.exc import IntegrityError +from sqlmodel import col, select, update + +from . import wallet as _wallet_module +from .core import db, get_logger +from .core.db import SWAP_OPEN_STATUSES, ApiKey, AsyncSession, CashuSwap, Refund +from .core.settings import settings +from .mint import ( + MINT_TRANSPORT_COOLDOWN_SECONDS, + MINT_TRANSPORT_EXCEPTIONS, + MintRateGuard, + mint_cooldown_remaining, + run_mint_operation, +) +from .net_guard import BlockedDestinationError, assert_public_https_origin +from .wallet import ( + Bolt11PaymentAmbiguous, + Bolt11PaymentNotAttempted, + Bolt11PaymentPlan, + ForeignMintSwapError, + ForeignMintUnavailableError, + SwapPendingError, + TokenConsumedError, + Wallet, + _apply_credit_locked, + _check_bolt11_payment_status_locked, + _execute_bolt11_payment, + _wallet_operation_depth, + get_proofs_per_mint_and_unit, + get_wallet, + resolve_trusted_source_mint, + wallet_operation_guard, +) + +logger = get_logger(__name__) + +RECONCILE_BATCH_LIMIT = 100 +_UNITS = ("sat", "msat") + +_MINT_ERROR_CODE_RE = re.compile(r"\(Code: (\d+)\)") +_FOREIGN_FAILURE_EXCEPTIONS: tuple[type[BaseException], ...] = ( + asyncio.TimeoutError, + *MINT_TRANSPORT_EXCEPTIONS, +) + + +def swap_enabled() -> bool: + return settings.foreign_mint_policy.strip().lower() == "swap" + + +def refund_destination_mint(key: ApiKey) -> str | None: + """The user's own mint to refund to, when it is foreign and swaps are on.""" + mint = key.refund_mint_url + if not mint or not swap_enabled() or resolve_trusted_source_mint(mint): + return None + return mint + + +# --- foreign-mint budget --------------------------------------------------- + + +_foreign_slots: asyncio.Semaphore | None = None +_foreign_slots_capacity = 0 + + +def _slots() -> asyncio.Semaphore: + global _foreign_slots, _foreign_slots_capacity + capacity = settings.foreign_mint_max_concurrency + if _foreign_slots is None or _foreign_slots_capacity != capacity: + _foreign_slots = asyncio.Semaphore(capacity) + _foreign_slots_capacity = capacity + return _foreign_slots + + +def _assert_outside_wallet_guard(what: str) -> None: + # A programming error, not a runtime condition: this is the exact shape of + # the DoS that got the feature removed. + if _wallet_operation_depth.get(): + raise RuntimeError(f"{what} must not run under wallet_operation_guard") + + +@asynccontextmanager +async def foreign_mint_lock(mint_url: str) -> AsyncGenerator[None, None]: + """Serialize all work against one foreign mint across worker processes. + + The foreign wallet's secret-derivation counter lives in the shared wallet + database, so two processes minting change on the same foreign mint would + collide. Bounded wait: a slot that does not free up within the foreign + budget is treated like an unreachable mint, not queued behind. + """ + _assert_outside_wallet_guard("foreign_mint_lock") + digest = hashlib.sha256(mint_url.strip().lower().encode()).hexdigest()[:16] + # Resolved at call time so tests that relocate the wallet lock move this too. + path = ( + _wallet_module._WALLET_OPERATION_LOCK.parent / f".routstr-foreign-{digest}.lock" + ) + path.parent.mkdir(parents=True, exist_ok=True) + fd = os.open(path, os.O_CREAT | os.O_RDWR, 0o600) + deadline = time.monotonic() + settings.foreign_mint_operation_timeout_seconds + acquired = False + try: + while True: + try: + fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + acquired = True + break + except BlockingIOError: + if time.monotonic() >= deadline: + raise ForeignMintUnavailableError( + "Another swap against this mint is still in progress" + ) from None + await asyncio.sleep(0.05) + yield + finally: + if acquired: + fcntl.flock(fd, fcntl.LOCK_UN) + os.close(fd) + + +async def run_foreign_mint_operation( + factory: Callable[[], Awaitable[Any]], *, mint_url: str, op_name: str +) -> Any: + """One bounded attempt against a mint the operator did not configure. + + No retries and no queueing: if every slot is busy or the mint is cooling + down, fail now. A transport failure or timeout puts the mint on the normal + transport cooldown so a flood naming the same dead mint is refused offline. + """ + _assert_outside_wallet_guard(op_name) + if mint_cooldown_remaining(mint_url) > 0: + raise ForeignMintUnavailableError("Issuing mint is cooling down") + slots = _slots() + if slots.locked(): + raise ForeignMintUnavailableError("Foreign-mint budget exhausted; retry later") + await slots.acquire() + try: + return await asyncio.wait_for( + factory(), timeout=settings.foreign_mint_operation_timeout_seconds + ) + except _FOREIGN_FAILURE_EXCEPTIONS as error: + MintRateGuard.get(mint_url).apply_cooldown( + MINT_TRANSPORT_COOLDOWN_SECONDS, reason="transport" + ) + logger.warning( + "Foreign mint operation failed", + extra={ + "event": "foreign_mint_operation_failed", + "op_name": op_name, + "mint_url": mint_url, + "error_type": type(error).__name__, + }, + ) + raise ForeignMintUnavailableError( + f"Issuing mint did not answer {op_name} in time" + ) from error + finally: + slots.release() + + +# --- amounts ---------------------------------------------------------------- + + +def _convert(amount: int, from_unit: str, to_unit: str) -> int: + msats = amount * 1000 if from_unit == "sat" else amount + return msats // 1000 if to_unit == "sat" else msats + + +def _melt_definitively_failed(error: BaseException) -> bool: + """The mint authoritatively rejected the Lightning payment; proofs are unspent.""" + message = str(error).strip() + return message.lower() == "could not pay invoice." or "(Code: 20004)" in message + + +def _melt_rejected_inputs(error: BaseException) -> bool: + """The mint refused the melt before paying because the inputs fell short. + + 11005 is the registered "Transaction is not balanced" code (cdk). 11000 is + nutshell's generic TransactionError and only counts alongside the + "not enough inputs" detail text. + """ + message = str(error) + match = _MINT_ERROR_CODE_RE.search(message) + code = match.group(1) if match else None + shortfall_text = "not enough inputs" in message.lower() + return code == "11005" or (code in (None, "11000") and shortfall_text) + + +def _state_name(response: object) -> str: + raw = getattr(response, "state", None) + if raw is None: + return "paid" if getattr(response, "paid", None) is True else "" + return str(raw).lower().rsplit(".", 1)[-1] + + +# --- journal ---------------------------------------------------------------- + + +async def _save(swap: CashuSwap) -> None: + async with db.create_session() as session: + session.add(swap) + await session.commit() + + +async def _update(swap: CashuSwap, **values: Any) -> None: + values.setdefault("updated_at", int(time.time())) + async with db.create_session() as session: + await session.exec( # type: ignore[call-overload] + update(CashuSwap).where(col(CashuSwap.id) == swap.id).values(**values) + ) + await session.commit() + for name, value in values.items(): + setattr(swap, name, value) + + +async def _prior_swap_for_token(token_hash: str) -> CashuSwap | None: + async with db.create_session() as session: + result = await session.exec( + select(CashuSwap) + .where(CashuSwap.token_hash == token_hash) + .order_by(col(CashuSwap.created_at).desc()) + ) + return result.first() + + +def _raise_for_prior_swap(prior: CashuSwap) -> None: + if prior.status in SWAP_OPEN_STATUSES: + raise SwapPendingError("A swap for this token is already in progress") + if prior.status == "failed": + raise ForeignMintSwapError( + "A prior swap for this token failed; the token was not spent" + ) + raise ValueError("Cashu token already spent") + + +# --- inbound: foreign token -> primary mint -> API key credit ------------- + + +async def _load_foreign_proofs(wallet: Wallet, token_obj: Token) -> list[Proof]: + await run_foreign_mint_operation( + wallet.load_mint_keysets, mint_url=token_obj.mint, op_name="swap_load_keysets" + ) + try: + await wallet.activate_keyset() + except Exception as error: + raise ForeignMintSwapError("Issuing mint has no active keyset") from error + proofs = token_obj.proofs + try: + await wallet._expand_short_keyset_ids(proofs) + except (KeyError, ValueError) as error: + raise ForeignMintSwapError( + "Cashu token references an unknown or ambiguous keyset" + ) from error + try: + wallet.verify_proofs_dleq(proofs) + except Exception as error: + raise ValueError("Invalid Cashu token: DLEQ proof failed") from error + return proofs + + +async def _quote_pair( + dest_wallet: Wallet, + dest_mint: str, + source_wallet: Wallet, + source_mint: str, + amount: int, +) -> tuple[MintQuote, MeltQuote]: + mint_quote = await run_mint_operation( + lambda: dest_wallet.request_mint(amount), + op_name="swap_request_mint", + mint_url=dest_mint, + retry_timeouts=False, + ) + melt_quote = await run_foreign_mint_operation( + lambda: source_wallet.melt_quote(mint_quote.request), + mint_url=source_mint, + op_name="swap_melt_quote", + ) + return mint_quote, melt_quote + + +async def swap_in_and_credit( + cashu_token: str, key: ApiKey, session: AsyncSession +) -> int: + """Melt a foreign-mint token into the primary mint and credit ``key``. + + Returns the credited msats. Raises before anything is spent for every + refusal (``ForeignMintSwapError``, ``ForeignMintUnavailableError``, + ``ValueError``) and ``SwapPendingError`` once the melt was dispatched but + not confirmed. + """ + if not swap_enabled(): + raise ForeignMintSwapError("Foreign-mint swaps are disabled on this node") + token_obj = deserialize_token_from_string(cashu_token) + source_mint = str(token_obj.mint) + if resolve_trusted_source_mint(source_mint) is not None: + raise ValueError("Token is from a trusted mint; redeem it directly") + source_unit = str(token_obj.unit) + dest_unit = settings.primary_mint_unit + dest_mint = settings.primary_mint + if source_unit not in _UNITS or dest_unit not in _UNITS or not dest_mint: + raise ForeignMintSwapError("Unsupported token unit for swap") + if key.refund_currency is not None and key.refund_currency != dest_unit: + raise ValueError( + "Cashu token unit does not match the API key liability unit: " + f"expected {key.refund_currency}, got {dest_unit}" + ) + try: + await assert_public_https_origin(source_mint) + except BlockedDestinationError as error: + raise ForeignMintSwapError(str(error)) from error + + token_hash = hashlib.sha256(cashu_token.encode()).hexdigest() + prior = await _prior_swap_for_token(token_hash) + if prior is not None: + _raise_for_prior_swap(prior) + + source_amount = int(token_obj.amount) + swap: CashuSwap + async with foreign_mint_lock(source_mint): + # The first check fails fast. This second check closes the in-process + # race for requests that were already waiting on the per-mint lock; + # the unique token hash below closes it across worker processes. + prior = await _prior_swap_for_token(token_hash) + if prior is not None: + _raise_for_prior_swap(prior) + + source_wallet = await get_wallet(source_mint, source_unit, load=False) + proofs = await _load_foreign_proofs(source_wallet, token_obj) + input_fees = source_wallet.get_fees_for_proofs(proofs) + dest_wallet = await get_wallet(dest_mint, dest_unit, load_proofs=False) + + # Round one quotes the whole token minus input fees; the melt quote + # then tells us the real Lightning fee reserve. Round two, if needed, + # re-quotes for what is left. No open-ended retry loop: a mint whose + # second quote still does not fit is refused with nothing spent. + gross = _convert(source_amount - input_fees, source_unit, dest_unit) + if gross <= 0: + raise ForeignMintSwapError( + "Token value does not cover the mint's input fees" + ) + mint_quote, melt_quote = await _quote_pair( + dest_wallet, dest_mint, source_wallet, source_mint, gross + ) + needed = melt_quote.amount + melt_quote.fee_reserve + input_fees + if needed > source_amount: + net = _convert( + source_amount - input_fees - melt_quote.fee_reserve, + source_unit, + dest_unit, + ) + if net <= 0: + raise ForeignMintSwapError("Token value does not cover swap fees") + mint_quote, melt_quote = await _quote_pair( + dest_wallet, dest_mint, source_wallet, source_mint, net + ) + needed = melt_quote.amount + melt_quote.fee_reserve + input_fees + if needed > source_amount: + raise ForeignMintSwapError("Token value does not cover swap fees") + minted_amount = net + else: + minted_amount = gross + + swap = CashuSwap( + direction="in", + status="melting", + api_key_hashed_key=key.hashed_key, + token_hash=token_hash, + token=cashu_token, + source_mint=source_mint, + source_unit=source_unit, + source_amount=source_amount, + destination_mint=dest_mint, + destination_unit=dest_unit, + destination_amount=minted_amount, + fee_reserve=int(melt_quote.fee_reserve), + input_fees=input_fees, + mint_quote_id=mint_quote.quote, + melt_quote_id=melt_quote.quote, + ) + try: + await _save(swap) + except IntegrityError: + prior = await _prior_swap_for_token(token_hash) + if prior is not None: + _raise_for_prior_swap(prior) + raise + logger.info( + "Cross-mint swap dispatching melt", + extra={ + "event": "cashu_swap_melting", + "swap_id": swap.id, + "source_mint": source_mint, + "destination_mint": dest_mint, + "source_amount": source_amount, + "minted_amount": minted_amount, + "fee_reserve": melt_quote.fee_reserve, + "input_fees": input_fees, + }, + ) + + try: + response = await run_foreign_mint_operation( + lambda: source_wallet.melt( + proofs=proofs, + invoice=mint_quote.request, + fee_reserve_sat=melt_quote.fee_reserve, + quote_id=melt_quote.quote, + ), + mint_url=source_mint, + op_name="swap_melt", + ) + except ForeignMintUnavailableError as error: + # Dispatched, outcome unknown: the Lightning payment may still land. + await _update(swap, status="ambiguous", error=str(error)) + raise SwapPendingError("Source melt outcome unknown") from error + except Exception as error: + if _melt_definitively_failed(error) or _melt_rejected_inputs(error): + await _update(swap, status="failed", error=str(error)) + raise ForeignMintSwapError( + "Issuing mint refused the Lightning payment" + ) from error + if "already spent" in str(error).lower(): + await _update(swap, status="failed", error=str(error)) + raise ValueError("Cashu token already spent") from error + await _update(swap, status="ambiguous", error=str(error)) + raise SwapPendingError("Source melt outcome unknown") from error + + state = _state_name(response) + if state == "unpaid": + await _update(swap, status="failed", error="melt reported unpaid") + raise ForeignMintSwapError("Issuing mint did not pay the swap invoice") + if state != "paid": + await _update(swap, status="ambiguous", error=f"melt state {state!r}") + raise SwapPendingError("Source melt is still pending") + await _update(swap, status="melted", error=None) + + return await _finish_swap_in(swap, key=key, session=session) + + +async def _mint_with_recovery( + wallet: Wallet, + amount: int, + quote_id: str, + *, + mint_url: str, + foreign: bool, +) -> list[Proof]: + """Mint for a paid quote; recover proofs the mint already signed once. + + An earlier attempt may have signed outputs at the mint but died before the + local derivation counter advanced. Restoring the keyset recovers those + proofs instead of crediting money the wallet does not hold. + """ + await wallet.load_proofs(reload=True) + + def proofs_for_quote() -> list[Proof]: + proofs = [ + proof + for proof in wallet.proofs + if getattr(proof, "mint_id", None) == quote_id + ] + return proofs if sum(proof.amount for proof in proofs) == amount else [] + + # A prior attempt may have completed remotely and in the wallet DB before + # the swap journal commit. Cashu records the quote id on every minted proof, + # which is the idempotency key we need to resume without minting or crediting + # a different set of proofs. + existing = proofs_for_quote() + if existing: + return existing + before = wallet.available_balance.amount + + async def do_mint() -> list[Proof]: + return await wallet.mint(amount, quote_id=quote_id) + + try: + if foreign: + return await run_foreign_mint_operation( + do_mint, mint_url=mint_url, op_name="swap_mint_foreign" + ) + return await run_mint_operation( + do_mint, + op_name="swap_mint_on_destination", + mint_url=mint_url, + retry_timeouts=False, + ) + except Exception as error: + text = str(error).lower() + if "11003" not in text and "outputs already signed" not in text: + raise + logger.warning( + "Swap mint outputs already signed; recovering orphaned proofs", + extra={"mint_url": mint_url, "quote_id": quote_id, "amount": amount}, + ) + for keyset_id in list(wallet.keysets): + await wallet.restore_tokens_for_keyset(keyset_id, to=1, batch=25) + await wallet.load_proofs(reload=True) + recovered_for_quote = proofs_for_quote() + if recovered_for_quote: + return recovered_for_quote + gained = wallet.available_balance.amount - before + if gained < amount: + raise TokenConsumedError( + f"Swap recovery restored {gained} of {amount}; manual reconciliation required" + ) from error + recovered = [p for p in wallet.proofs if not p.reserved] + try: + # offline: never ask the mint to split here, the budget is spent. + picked, _ = await wallet.select_to_send( + recovered, amount, set_reserved=False, offline=True + ) + except Exception as selection_error: + raise TokenConsumedError( + "Swap recovery restored proofs but none match the swapped amount; " + "manual reconciliation required" + ) from selection_error + return picked + + +async def _finish_swap_in( + swap: CashuSwap, + *, + key: ApiKey | None = None, + session: AsyncSession | None = None, +) -> int: + """Mint on the trusted destination and credit the key, under the guard.""" + async with wallet_operation_guard(): + if swap.status == "melted": + dest_wallet = await get_wallet(swap.destination_mint, swap.destination_unit) + try: + await _mint_with_recovery( + dest_wallet, + swap.destination_amount, + str(swap.mint_quote_id), + mint_url=swap.destination_mint, + foreign=False, + ) + except TokenConsumedError as error: + await _update(swap, error=str(error)) + raise + except Exception as error: + # Invoice is paid; the quote stays mintable. Leave the row for + # the reconciler rather than losing track of settled money. + await _update(swap, error=str(error)) + logger.error( + "Swap mint on destination failed after a paid melt", + extra={"swap_id": swap.id, "error": str(error)}, + ) + raise SwapPendingError( + "Destination mint failed; retrying later" + ) from error + await _update(swap, status="minted", error=None) + + if swap.status != "minted": + raise SwapPendingError(f"Swap is {swap.status}") + + if session is None or key is None: + async with db.create_session() as own_session: + own_key = await own_session.get(ApiKey, swap.api_key_hashed_key) + if own_key is None: + await _update(swap, status="failed", error="api key missing") + logger.critical( + "Swapped funds have no API key to credit", + extra={"swap_id": swap.id, "amount": swap.destination_amount}, + ) + raise TokenConsumedError("API key vanished before swap credit") + credited = await _apply_credit_locked( + own_key, + own_session, + amount=swap.destination_amount, + unit=swap.destination_unit, + mint_url=swap.destination_mint, + token=str(swap.token), + refund_mint_url=swap.source_mint, + swap_id=swap.id, + ) + else: + credited = await _apply_credit_locked( + key, + session, + amount=swap.destination_amount, + unit=swap.destination_unit, + mint_url=swap.destination_mint, + token=str(swap.token), + refund_mint_url=swap.source_mint, + swap_id=swap.id, + ) + swap.status = "credited" + swap.error = None + logger.info( + "Cross-mint swap credited", + extra={ + "event": "cashu_swap_completed", + "swap_id": swap.id, + "source_mint": swap.source_mint, + "destination_mint": swap.destination_mint, + "credited_msats": credited, + }, + ) + return credited + + +# --- outbound: primary mint -> user's mint (refund) ------------------------- + + +async def swap_out_for_refund( + session: AsyncSession, refund: Refund, destination_mint: str +) -> bool: + """Pay a refund as a token on the user's own (foreign) mint. + + Owner proofs on the primary mint pay a mint quote on the user's mint; the + user receives the net amount after the Lightning fee reserve and input + fees. The ``Refund`` claim carries the melt quote so the existing refund + reconciler can hold or release the balance; the swap row carries the rest. + """ + from . import refund as refund_module + from .payment.lnurl import MeltOutcomeAmbiguousError, MeltUnpaidError + + unit = refund.unit + amount = refund.amount_msats // 1000 if unit == "sat" else refund.amount_msats + primary = settings.primary_mint + try: + await assert_public_https_origin(destination_mint) + except BlockedDestinationError as error: + raise ForeignMintSwapError(str(error)) from error + + async with foreign_mint_lock(destination_mint): + dest_wallet = await get_wallet(destination_mint, unit, load=False) + await run_foreign_mint_operation( + dest_wallet.load_mint_keysets, + mint_url=destination_mint, + op_name="refund_swap_load_keysets", + ) + try: + await dest_wallet.activate_keyset() + except Exception as error: + raise ForeignMintSwapError("Refund mint has no active keyset") from error + + source_wallet = await get_wallet(primary, unit) + async with wallet_operation_guard(): + await source_wallet.load_proofs(reload=True) + proofs = get_proofs_per_mint_and_unit( + source_wallet, primary, unit, not_reserved=True + ) + if sum(p.amount for p in proofs) < amount: + raise ValueError("Primary mint balance cannot cover this refund") + selection, _ = await source_wallet.select_to_send( + proofs, amount, set_reserved=False, include_fees=True + ) + input_fees = source_wallet.get_fees_for_proofs(selection) + + async def quotes(mint_amount: int) -> tuple[MintQuote, MeltQuote]: + mint_quote = await run_foreign_mint_operation( + lambda: dest_wallet.request_mint(mint_amount), + mint_url=destination_mint, + op_name="refund_swap_request_mint", + ) + melt_quote = await run_mint_operation( + lambda: source_wallet.melt_quote(mint_quote.request), + op_name="refund_swap_melt_quote", + mint_url=primary, + retry_timeouts=False, + ) + return mint_quote, melt_quote + + mint_quote, melt_quote = await quotes(amount) + net = amount - melt_quote.fee_reserve - input_fees + if net <= 0: + raise ForeignMintSwapError("Refund amount does not cover swap fees") + if net < amount: + mint_quote, melt_quote = await quotes(net) + if melt_quote.amount + melt_quote.fee_reserve + input_fees > amount: + raise ForeignMintSwapError("Refund amount does not cover swap fees") + + swap = CashuSwap( + direction="out", + status="melting", + api_key_hashed_key=refund.api_key_hashed_key, + refund_id=refund.id, + source_mint=primary, + source_unit=unit, + source_amount=amount, + destination_mint=destination_mint, + destination_unit=unit, + destination_amount=net, + fee_reserve=int(melt_quote.fee_reserve), + input_fees=input_fees, + mint_quote_id=mint_quote.quote, + melt_quote_id=melt_quote.quote, + ) + await _save(swap) + await refund_module.record_quote(refund, melt_quote.quote, primary) + + async with wallet_operation_guard(): + await source_wallet.load_proofs(reload=True) + proofs = get_proofs_per_mint_and_unit( + source_wallet, primary, unit, not_reserved=True + ) + plan = Bolt11PaymentPlan( + mint_quote.request, source_wallet, proofs, melt_quote, primary, unit + ) + try: + await _execute_bolt11_payment(plan) + except Bolt11PaymentNotAttempted as error: + await _update(swap, status="failed", error=str(error)) + raise MeltUnpaidError(str(error)) from error + except Bolt11PaymentAmbiguous as error: + await _update(swap, status="ambiguous", error=str(error)) + await refund_module.hold(session, refund, melt_quote.quote) + raise MeltOutcomeAmbiguousError(str(error)) from error + await _update(swap, status="melted", error=None) + + try: + token = await _issue_refund_token(swap, dest_wallet) + except Exception as error: + await _update(swap, error=str(error)) + await refund_module.hold(session, refund, melt_quote.quote) + logger.error( + "Refund swap paid but minting on the user's mint failed; held", + extra={"swap_id": swap.id, "refund_id": refund.id, "error": str(error)}, + ) + raise MeltOutcomeAmbiguousError( + "Refund token could not be minted yet" + ) from error + + refund.token = token + refund.mint_url = destination_mint + settled = await refund_module.settle( + session, refund, token=token, mint_url=destination_mint + ) + await _update(swap, status="settled", error=None) + return settled + + +async def _issue_refund_token(swap: CashuSwap, dest_wallet: Wallet) -> str: + """Mint the paid quote on the user's mint and hand the proofs over as a token.""" + new_proofs = await _mint_with_recovery( + dest_wallet, + swap.destination_amount, + str(swap.mint_quote_id), + mint_url=swap.destination_mint, + foreign=True, + ) + token = await dest_wallet.serialize_proofs( + new_proofs, include_dleq=False, legacy=False, memo=None + ) + await dest_wallet.set_reserved_for_send(new_proofs, reserved=True) + await _update(swap, status="issued", token=token, error=None) + return token + + +# --- reconciler --------------------------------------------------------------- + + +async def _lease(swap_id: str, now: int, cutoff: int) -> bool: + async with db.create_session() as session: + result = await session.exec( # type: ignore[call-overload] + update(CashuSwap) + .where(col(CashuSwap.id) == swap_id) + .where(col(CashuSwap.status).in_(SWAP_OPEN_STATUSES)) + .where( + col(CashuSwap.claimed_at).is_(None) + | (col(CashuSwap.claimed_at) < cutoff) + ) + .values(claimed_at=now) + ) + await session.commit() + return bool(result.rowcount) + + +async def _source_melt_state(swap: CashuSwap) -> str: + """Ask the foreign mint what became of an inbound swap's melt.""" + async with foreign_mint_lock(swap.source_mint): + wallet = await get_wallet(swap.source_mint, swap.source_unit, load=False) + quote = await run_foreign_mint_operation( + lambda: wallet.get_melt_quote(str(swap.melt_quote_id)), + mint_url=swap.source_mint, + op_name="swap_reconcile_melt_quote", + ) + return _state_name(quote) if quote is not None else "unknown" + + +async def _reconcile_in(swap: CashuSwap, now: int) -> None: + if swap.status in ("melting", "ambiguous"): + if swap.updated_at > now - settings.refund_claim_timeout_seconds: + return + state = await _source_melt_state(swap) + if state == "paid": + await _update(swap, status="melted", error=None) + elif state == "unpaid": + # The mint says it never paid, so the sender still holds the proofs. + await _update(swap, status="failed", error="melt unpaid at the mint") + return + else: + logger.warning( + "Inbound swap melt still unresolved", + extra={"swap_id": swap.id, "melt_state": state}, + ) + return + if swap.status in ("melted", "minted"): + try: + await _finish_swap_in(swap) + except SwapPendingError as error: + logger.warning( + "Inbound swap not finished yet", + extra={"swap_id": swap.id, "error": str(error)}, + ) + + +async def _reconcile_out(swap: CashuSwap, now: int) -> None: + from . import refund as refund_module + + async with db.create_session() as session: + refund = await session.get(Refund, swap.refund_id) + if refund is None: + await _update(swap, status="failed", error="refund claim missing") + return + + if swap.status in ("melting", "ambiguous"): + if swap.updated_at > now - settings.refund_claim_timeout_seconds: + return + async with wallet_operation_guard(): + state = await _check_bolt11_payment_status_locked( + swap.source_mint, swap.source_unit, str(swap.melt_quote_id) + ) + if state == "paid": + await _update(swap, status="melted", error=None) + elif state == "unpaid": + await _update(swap, status="failed", error="melt unpaid at the mint") + async with db.create_session() as session: + await refund_module.release(session, refund) + return + else: + logger.warning( + "Refund swap melt still unresolved", + extra={"swap_id": swap.id, "melt_state": state}, + ) + return + + if swap.status == "melted": + async with foreign_mint_lock(swap.destination_mint): + dest_wallet = await get_wallet( + swap.destination_mint, swap.destination_unit, load=False + ) + await run_foreign_mint_operation( + dest_wallet.load_mint_keysets, + mint_url=swap.destination_mint, + op_name="refund_swap_reconcile_keysets", + ) + await dest_wallet.activate_keyset() + await _issue_refund_token(swap, dest_wallet) + + if swap.status == "issued": + token = str(swap.token) + async with db.create_session() as session: + await refund_module.settle( + session, refund, token=token, mint_url=swap.destination_mint + ) + refund.token = token + refund.mint_url = swap.destination_mint + await refund_module._record_cashu_payout(refund) + await _update(swap, status="settled", error=None) + + +async def reconcile_swaps_once() -> None: + """Finish or fail swap rows whose request died before a final state.""" + now = int(time.time()) + cutoff = now - settings.refund_claim_timeout_seconds + async with db.create_session() as session: + result = await session.exec( + select(CashuSwap) + .where(col(CashuSwap.status).in_(SWAP_OPEN_STATUSES)) + .where( + col(CashuSwap.claimed_at).is_(None) + | (col(CashuSwap.claimed_at) < cutoff) + ) + .order_by(col(CashuSwap.created_at)) + .limit(RECONCILE_BATCH_LIMIT) + ) + stale = list(result.all()) + + for swap in stale: + if not await _lease(swap.id, now, cutoff): + continue + try: + if swap.direction == "in": + await _reconcile_in(swap, now) + else: + await _reconcile_out(swap, now) + except Exception as error: + logger.error( + "Swap reconciliation failed", + extra={ + "swap_id": swap.id, + "direction": swap.direction, + "error": str(error), + "error_type": type(error).__name__, + }, + exc_info=True, + ) + finally: + await _update(swap, claimed_at=None) + + +async def periodic_swap_reconcile() -> None: + while True: + await asyncio.sleep(settings.swap_reconcile_interval_seconds) + try: + await reconcile_swaps_once() + except asyncio.CancelledError: + raise + except Exception as error: + logger.error( + "Swap reconcile loop error", + extra={"error": str(error), "error_type": type(error).__name__}, + ) diff --git a/routstr/net_guard.py b/routstr/net_guard.py new file mode 100644 index 00000000..6a155497 --- /dev/null +++ b/routstr/net_guard.py @@ -0,0 +1,61 @@ +"""Destination checks for URLs a client chose and the node will connect to.""" + +import asyncio +import ipaddress +import socket +from urllib.parse import urlsplit + + +class BlockedDestinationError(ValueError): + """The URL points at something the node must not connect to.""" + + +def is_blocked_address(address: str) -> bool: + """Allow only globally reachable addresses (RFC 6890).""" + try: + ip = ipaddress.ip_address(address) + except ValueError: + return True + if isinstance(ip, ipaddress.IPv6Address): + # An embedded v4 address would otherwise smuggle a rejected target past + # the v6 checks. + for embedded in (ip.ipv4_mapped, ip.sixtofour): + if embedded is not None: + return is_blocked_address(str(embedded)) + return not ip.is_global or ip.is_multicast + + +async def assert_public_https_origin(url: str) -> None: + """Reject a client-supplied origin unless it is HTTPS to a public host. + + Used for mints named inside incoming Cashu tokens: the token is + unauthenticated input, so without this the node would open connections to + whatever address the sender wrote into it. HTTPS is required because the + host is only verified by name here; certificate validation binds the later + connection to the same name. + """ + parts = urlsplit(url.strip()) + if parts.scheme != "https": + raise BlockedDestinationError("mint URL must use https") + if parts.username is not None or parts.password is not None: + raise BlockedDestinationError("mint URL must not carry credentials") + if parts.query or parts.fragment: + raise BlockedDestinationError("mint URL must not carry a query or fragment") + host = parts.hostname + if not host: + raise BlockedDestinationError("mint URL has no host") + try: + port = parts.port or 443 + except ValueError as error: + raise BlockedDestinationError("mint URL port is invalid") from error + try: + infos = await asyncio.get_running_loop().getaddrinfo( + host, port, proto=socket.IPPROTO_TCP + ) + except socket.gaierror as error: + raise BlockedDestinationError("mint host did not resolve") from error + if not infos: + raise BlockedDestinationError("mint host did not resolve") + for info in infos: + if is_blocked_address(str(info[4][0])): + raise BlockedDestinationError("mint host resolves to a blocked address") diff --git a/routstr/payment/helpers.py b/routstr/payment/helpers.py index 4582031b..3e21e130 100644 --- a/routstr/payment/helpers.py +++ b/routstr/payment/helpers.py @@ -1,6 +1,5 @@ import asyncio import base64 -import ipaddress import json import math import socket @@ -26,6 +25,7 @@ from ..core.error_scope import ( from ..core.exceptions import UpstreamError from ..core.redaction import redact_org_ids from ..core.settings import settings +from ..net_guard import is_blocked_address as _is_blocked_address from ..wallet import ( UntrustedSourceMintError, classify_redemption_error, @@ -396,21 +396,6 @@ def _get_image_dimensions(image_data: bytes) -> tuple[int, int]: return (512, 512) -def _is_blocked_address(address: str) -> bool: - """Allow only globally reachable addresses (RFC 6890).""" - try: - ip = ipaddress.ip_address(address) - except ValueError: - return True - if isinstance(ip, ipaddress.IPv6Address): - # An embedded v4 address would otherwise smuggle a rejected target past - # the v6 checks. - for embedded in (ip.ipv4_mapped, ip.sixtofour): - if embedded is not None: - return _is_blocked_address(str(embedded)) - return not ip.is_global or ip.is_multicast - - async def _validated_fetch_target(url: str) -> tuple[str, str]: """Return the URL to request and its ``Host`` header. diff --git a/routstr/refund.py b/routstr/refund.py index 2aee20e8..eb1a2d6f 100644 --- a/routstr/refund.py +++ b/routstr/refund.py @@ -5,6 +5,7 @@ import time from typing import Any import httpx +from cashu.wallet.helpers import deserialize_token_from_string from fastapi import HTTPException from sqlalchemy.exc import IntegrityError from sqlmodel import col, func, select, update @@ -50,11 +51,25 @@ def amount_in_unit(amount_msats: int, unit: str) -> int: def refund_mint(key: ApiKey) -> str: + """Trusted mint the payout is drawn from. + + A foreign refund mint (a key funded by a swapped-in token) is paid from the + primary mint and swapped back; see :func:`swap_destination`. + """ if key.refund_mint_url and key.refund_mint_url in settings.cashu_mints: return key.refund_mint_url return settings.primary_mint +def swap_destination(key: ApiKey, method: str) -> str | None: + """The user's own mint to swap a Cashu refund to, if that applies.""" + from .foreign_mint_swap import refund_destination_mint + + if method != "cashu": + return None + return refund_destination_mint(key) + + async def validate_lightning_destination(destination: str) -> None: try: await get_lnurl_data(destination) @@ -85,6 +100,10 @@ async def open_claim( ) ) created_at = max(int(time.time()), (latest.one() or 0) + 1) + if destination is None: + # A cashu refund to the user's own foreign mint records that mint as + # the destination; the payout is still drawn from a trusted mint. + destination = swap_destination(key, method) refund = Refund( api_key_hashed_key=key.hashed_key, method=method, @@ -306,16 +325,32 @@ async def latest_terminal(session: AsyncSession, key: ApiKey) -> Refund | None: return await _latest_with_status(session, key, ("paid",)) +def _delivered_amount_msats(refund: Refund) -> int: + """Actual bearer-token value, which can be net of cross-mint fees.""" + if not refund.token: + return refund.amount_msats + try: + token = deserialize_token_from_string(refund.token) + except Exception: + logger.error( + "paid refund token could not be decoded for amount reporting", + extra={"refund_id": refund.id}, + ) + return refund.amount_msats + return int(token.amount) * 1000 if str(token.unit) == "sat" else int(token.amount) + + def describe(refund: Refund) -> dict[str, str]: body: dict[str, str] = {"refund_id": refund.id, "status": refund.status} if refund.token: body["token"] = refund.token if refund.destination: body["recipient"] = refund.destination + delivered_msats = _delivered_amount_msats(refund) if refund.unit == "sat": - body["sats"] = str(refund.amount_msats // 1000) + body["sats"] = str(delivered_msats // 1000) else: - body["msats"] = str(refund.amount_msats) + body["msats"] = str(delivered_msats) return body @@ -349,6 +384,10 @@ async def _pay_lightning(session: AsyncSession, refund: Refund) -> bool: async def _pay_cashu(session: AsyncSession, refund: Refund) -> bool: amount = amount_in_unit(refund.amount_msats, refund.unit) await renew_lease(session, refund) + if refund.destination: + from .foreign_mint_swap import swap_out_for_refund + + return await swap_out_for_refund(session, refund, refund.destination) token = await send_token(amount, refund.unit, refund.mint_url) # From here the token is bearer money: keep it on the claim so a failed # settle withholds the balance instead of restoring it. @@ -363,7 +402,7 @@ async def _record_cashu_payout(refund: Refund) -> None: try: await store_cashu_transaction( token=str(refund.token), - amount=amount_in_unit(refund.amount_msats, refund.unit), + amount=amount_in_unit(_delivered_amount_msats(refund), refund.unit), unit=refund.unit, mint_url=refund.mint_url, typ="out", @@ -514,6 +553,10 @@ async def _reconcile(refund: Refund, now: int) -> None: async with create_session() as session: await settle(session, refund) return + if refund.destination and refund.quote_id: + # A swap to the user's mint: its journal row owns the outcome and + # the swap reconciler settles or releases this claim from there. + return # No quote to query for cashu; withhold the balance and alert once. async with create_session() as session: if await _close(session, refund, require_no_token=True, status="stuck"): diff --git a/routstr/wallet.py b/routstr/wallet.py index 1f2a9e0e..0b280067 100644 --- a/routstr/wallet.py +++ b/routstr/wallet.py @@ -24,7 +24,6 @@ from sqlmodel import col, select, update from .cashu_compat import install_cashu_httpx_shim from .checkstate import filter_unspent_proofs from .core import db, get_logger -from .core.db import store_cashu_transaction_with_retry as store_cashu_transaction from .core.settings import settings from .mint import ( MINT_TRANSPORT_EXCEPTIONS, @@ -208,6 +207,30 @@ class UntrustedSourceMintError(ValueError): """The token names a mint outside primary_mint/cashu_mints.""" +class ForeignMintSwapError(ValueError): + """A cross-mint swap was refused before any proof was spent. + + The token is still fully usable by its holder: fees exceeded its value, the + unit is unsupported, or the issuing mint rejected the quote. + """ + + +class ForeignMintUnavailableError(MintConnectionError): + """The issuing mint did not answer within the foreign-mint budget. + + Nothing was spent. Unlike trusted mints there is no retry: the sender, not + the operator, picked this mint. + """ + + +class SwapPendingError(Exception): + """The swap's Lightning leg was dispatched but its outcome is not yet known. + + The token must not be retried: its proofs may already be spent. The journal + row keeps the quote ids and the reconciler credits or fails it later. + """ + + class TokenConsumedError(Exception): """A failure that happened AFTER the token's proofs were spent (melt succeeded, or redemption already returned) — e.g. minting on the primary @@ -314,6 +337,28 @@ def classify_redemption_error( "Cashu token was issued by a mint this node does not accept", "cashu_untrusted_source_mint", ) + if isinstance(error, SwapPendingError): + return ( + "swap_pending", + 409, + "Cross-mint swap was dispatched and is awaiting confirmation; do not " + "resend this token, the balance is credited once the mint confirms", + "cashu_swap_pending", + ) + if isinstance(error, ForeignMintSwapError): + return ( + "mint_error", + 422, + "Cashu token cannot be swapped into this node's mint; nothing was spent", + "cashu_foreign_mint_swap_failed", + ) + if isinstance(error, ForeignMintUnavailableError): + return ( + "mint_unreachable", + 503, + "The mint that issued this Cashu token did not answer in time; retry later", + "cashu_source_mint_unreachable", + ) if is_mint_rate_limited(error): return ( "mint_rate_limited", @@ -1090,97 +1135,18 @@ async def _credit_balance_locked( if isinstance(key.refund_currency, str) else None, ) - original_amount = amount - original_unit = unit logger.info( "credit_balance: Token redeemed successfully", extra={"amount": amount, "unit": unit, "mint_url": mint_url}, ) - - if unit == "sat": - amount = _sats_to_msats(amount) - logger.info( - "credit_balance: Converted to msat", extra={"amount_msat": amount} - ) - - # Guard against zero/negative redemptions (empty or dust tokens, or - # swap-to-primary-mint amounts that net to <= 0 after fees). Raising here - # — before the UPDATE/commit below — leaves any freshly-created, still - # uncommitted ApiKey row to be rolled back when the request session - # closes, instead of persisting an orphan key with balance 0. - if amount <= 0: - logger.error( - "credit_balance: Redeemed amount is zero or negative; refusing to credit", - extra={"amount": amount, "unit": unit, "mint_url": mint_url}, - ) - raise ValueError( - f"Redeemed token amount must be positive, got {amount} msats" - ) - - logger.info( - "credit_balance: Updating balance", - extra={"old_balance": key.balance, "credit_amount": amount}, - ) - - # The token is already redeemed (spent) here, so any crediting failure - # is post-redemption and non-retryable — surface it as TokenConsumedError - # (a key that vanished mid-flight, or an unexpected DB fault), never a - # retryable/token-error taxonomy. - try: - # Atomic UPDATE to prevent race conditions during concurrent topups. - updates: dict[str, object] = { - "balance": db.ApiKey.balance + amount, - } - # Legacy keys may predate refund provenance. Pin them to the - # destination used for this credit before exposing the balance. - if key.refund_mint_url is None: - updates["refund_mint_url"] = mint_url - if key.refund_currency is None: - updates["refund_currency"] = unit - stmt = ( - update(db.ApiKey) - .where(col(db.ApiKey.hashed_key) == key.hashed_key) - .values(**updates) - ) - result = await session.exec(stmt) # type: ignore[call-overload] - # If pruning removed this key after redemption, do not commit a no-op - # balance update and pretend the top-up succeeded. - if (getattr(result, "rowcount", 0) or 0) == 0: - raise TokenConsumedError( - "Token redeemed but the API key disappeared before the " - "credit could be recorded" - ) - await session.commit() - await session.refresh(key) - # refresh() starts a read transaction; release it before the - # transaction-history write opens its own session below. - await session.commit() - except TokenConsumedError: - raise - except Exception as db_error: - raise TokenConsumedError( - "Token redeemed but crediting the balance failed" - ) from db_error - - logger.info( - "credit_balance: Balance updated successfully", - extra={"new_balance": key.balance}, - ) - - await store_cashu_transaction( - token=cashu_token, - amount=original_amount, - unit=original_unit, + return await _apply_credit_locked( + key, + session, + amount=amount, + unit=unit, mint_url=mint_url, - typ="in", - source="apikey", - api_key_hashed_key=key.hashed_key, + token=cashu_token, ) - logger.debug( - "Cashu token successfully redeemed and stored", - extra={"amount": amount, "unit": unit, "mint_url": mint_url}, - ) - return amount except Exception as e: classification = classify_redemption_error(e) expected_codes = { @@ -1206,6 +1172,104 @@ async def _credit_balance_locked( raise +async def _apply_credit_locked( + key: db.ApiKey, + session: db.AsyncSession, + *, + amount: int, + unit: str, + mint_url: str, + token: str, + refund_mint_url: str | None = None, + swap_id: str | None = None, +) -> int: + """Atomically credit a redeemed amount and record its ledger row. + + ``amount`` is in ``unit``. ``refund_mint_url`` overrides the mint pinned as + the key's refund destination when the key has none yet. When ``swap_id`` is + present, the same transaction also claims the swap's ``minted`` state, so + reconciliation can never apply one minted quote twice. + """ + original_amount = amount + if unit == "sat": + amount = _sats_to_msats(amount) + logger.info("credit_balance: Converted to msat", extra={"amount_msat": amount}) + + if amount <= 0: + logger.error( + "credit_balance: Redeemed amount is zero or negative; refusing to credit", + extra={"amount": amount, "unit": unit, "mint_url": mint_url}, + ) + raise ValueError(f"Redeemed token amount must be positive, got {amount} msats") + + logger.info( + "credit_balance: Updating balance", + extra={"old_balance": key.balance, "credit_amount": amount}, + ) + + try: + updates: dict[str, object] = {"balance": db.ApiKey.balance + amount} + if key.refund_mint_url is None: + updates["refund_mint_url"] = refund_mint_url or mint_url + if key.refund_currency is None: + updates["refund_currency"] = unit + result = await session.exec( # type: ignore[call-overload] + update(db.ApiKey) + .where(col(db.ApiKey.hashed_key) == key.hashed_key) + .values(**updates) + ) + if (getattr(result, "rowcount", 0) or 0) == 0: + raise TokenConsumedError( + "Token redeemed but the API key disappeared before the " + "credit could be recorded" + ) + + if swap_id is not None: + claimed = await session.exec( # type: ignore[call-overload] + update(db.CashuSwap) + .where(col(db.CashuSwap.id) == swap_id) + .where(col(db.CashuSwap.status) == "minted") + .values(status="credited", error=None, updated_at=int(time.time())) + ) + if (getattr(claimed, "rowcount", 0) or 0) != 1: + raise TokenConsumedError( + "Swapped funds were already credited or their journal vanished" + ) + + session.add( + db.CashuTransaction( + token=token, + amount=original_amount, + unit=unit, + mint_url=mint_url, + type="in", + source="apikey", + api_key_hashed_key=key.hashed_key, + ) + ) + await session.flush() + await session.refresh(key) + await session.commit() + except TokenConsumedError: + await session.rollback() + raise + except Exception as db_error: + await session.rollback() + raise TokenConsumedError( + "Token redeemed but crediting the balance failed" + ) from db_error + + logger.info( + "credit_balance: Balance updated successfully", + extra={"new_balance": key.balance}, + ) + logger.debug( + "Cashu token successfully redeemed and stored", + extra={"amount": amount, "unit": unit, "mint_url": mint_url}, + ) + return amount + + _wallets: dict[str, Wallet] = {} # Proofs require a shorter refresh interval than remote mint metadata. _wallet_last_load: dict[str, float] = {} diff --git a/tests/unit/test_balance.py b/tests/unit/test_balance.py index cb262d70..8ffdec9f 100644 --- a/tests/unit/test_balance.py +++ b/tests/unit/test_balance.py @@ -481,27 +481,25 @@ async def test_credit_balance_stores_apikey_transaction_history() -> None: session.exec = AsyncMock(return_value=_update_result(1)) session.commit = AsyncMock() session.rollback = AsyncMock() + session.flush = AsyncMock() session.refresh = AsyncMock() - with ( - patch( - "routstr.wallet.recieve_token", - AsyncMock(return_value=(100, "sat", "https://mint.example")), - ), - patch("routstr.wallet.store_cashu_transaction", AsyncMock()) as mock_store, + with patch( + "routstr.wallet.recieve_token", + AsyncMock(return_value=(100, "sat", "https://mint.example")), ): amount = await credit_balance("cashuAtopup_token", key, session) assert amount == 100_000 - mock_store.assert_awaited_once() - call_kwargs = mock_store.call_args.kwargs - assert call_kwargs["typ"] == "in" - assert call_kwargs["source"] == "apikey" - assert call_kwargs["api_key_hashed_key"] == key.hashed_key - assert call_kwargs["amount"] == 100 - assert call_kwargs["unit"] == "sat" - assert call_kwargs["token"] == "cashuAtopup_token" - assert call_kwargs["mint_url"] == "https://mint.example" + stored = session.add.call_args.args[0] + assert isinstance(stored, CashuTransaction) + assert stored.type == "in" + assert stored.source == "apikey" + assert stored.api_key_hashed_key == key.hashed_key + assert stored.amount == 100 + assert stored.unit == "sat" + assert stored.token == "cashuAtopup_token" + assert stored.mint_url == "https://mint.example" @pytest.mark.asyncio diff --git a/tests/unit/test_foreign_mint_swap.py b/tests/unit/test_foreign_mint_swap.py new file mode 100644 index 00000000..4f59ceaf --- /dev/null +++ b/tests/unit/test_foreign_mint_swap.py @@ -0,0 +1,864 @@ +"""Cross-mint swaps run outside the wallet lock, under a bounded budget, and +leave a journal row for every Lightning leg they dispatch.""" + +import asyncio +import time +from contextlib import asynccontextmanager +from pathlib import Path +from types import SimpleNamespace +from typing import Any, AsyncGenerator +from unittest.mock import AsyncMock, Mock, patch + +import pytest +from cashu.core.base import MeltQuoteState +from sqlalchemy.exc import IntegrityError +from sqlalchemy.ext.asyncio import AsyncEngine, create_async_engine +from sqlalchemy.pool import StaticPool +from sqlmodel import SQLModel, select +from sqlmodel.ext.asyncio.session import AsyncSession + +from routstr import foreign_mint_swap as fms +from routstr import refund as refund_module +from routstr import wallet +from routstr.core import db +from routstr.core.db import ApiKey, CashuSwap, CashuTransaction, Refund +from routstr.core.settings import settings +from routstr.mint import MintRateGuard, mint_cooldown_remaining +from routstr.payment.lnurl import MeltOutcomeAmbiguousError +from routstr.wallet import ( + Bolt11PaymentAmbiguous, + ForeignMintSwapError, + ForeignMintUnavailableError, + SwapPendingError, + TokenConsumedError, + wallet_operation_guard, +) + +PRIMARY = "https://primary.example" +FOREIGN = "https://foreign.example" +KEY_HASH = "a" * 64 + + +@pytest.fixture +async def engine( + monkeypatch: pytest.MonkeyPatch, tmp_path: Path +) -> AsyncGenerator[AsyncEngine, None]: + engine = create_async_engine( + "sqlite+aiosqlite://", + poolclass=StaticPool, + connect_args={"check_same_thread": False}, + ) + async with engine.begin() as connection: + await connection.run_sync(SQLModel.metadata.create_all) + + @asynccontextmanager + async def create_session() -> AsyncGenerator[AsyncSession, None]: + async with AsyncSession(engine, expire_on_commit=False) as session: + yield session + + monkeypatch.setattr(db, "create_session", create_session) + monkeypatch.setattr(refund_module, "create_session", create_session) + monkeypatch.setattr(wallet, "_WALLET_OPERATION_LOCK", tmp_path / "op.lock") + monkeypatch.setattr(settings, "primary_mint", PRIMARY) + monkeypatch.setattr(settings, "primary_mint_unit", "sat") + monkeypatch.setattr(settings, "cashu_mints", [PRIMARY]) + monkeypatch.setattr(settings, "foreign_mint_policy", "swap") + monkeypatch.setattr(settings, "foreign_mint_operation_timeout_seconds", 0.2) + monkeypatch.setattr(settings, "foreign_mint_max_concurrency", 4) + monkeypatch.setattr(fms, "_foreign_slots", None) + MintRateGuard._guards.clear() + try: + yield engine + finally: + await engine.dispose() + + +@pytest.fixture +async def session(engine: AsyncEngine) -> AsyncGenerator[AsyncSession, None]: + async with AsyncSession(engine, expire_on_commit=False) as session: + yield session + + +async def _make_key(session: AsyncSession, **fields: Any) -> ApiKey: + fields.setdefault("balance", 0) + key = ApiKey(hashed_key=KEY_HASH, **fields) + session.add(key) + await session.commit() + await session.refresh(key) + return key + + +def _proof(amount: int) -> SimpleNamespace: + return SimpleNamespace(amount=amount, reserved=False, secret=f"s{amount}", id="k") + + +def _token(amount: int = 1000, mint: str = FOREIGN) -> SimpleNamespace: + return SimpleNamespace( + mint=mint, unit="sat", amount=amount, keysets=["k"], proofs=[_proof(amount)] + ) + + +class _ForeignWallet: + """Fake wallet for the sender's mint; records that the guard was not held.""" + + def __init__(self, fee_reserve: int = 5, melt_state: Any = MeltQuoteState.paid): + self.fee_reserve = fee_reserve + self.melt_state = melt_state + self.guard_depth_seen: list[int] = [] + self.load_mint_keysets = AsyncMock(side_effect=self._observe) + self.activate_keyset = AsyncMock() + self._expand_short_keyset_ids = AsyncMock() + self.verify_proofs_dleq = Mock() + self.get_fees_for_proofs = Mock(return_value=1) + self.melt_quote = AsyncMock(side_effect=self._melt_quote) + self.melt = AsyncMock(side_effect=self._melt) + self.get_melt_quote = AsyncMock( + return_value=SimpleNamespace(state=MeltQuoteState.paid) + ) + self.request_mint = AsyncMock(side_effect=self._request_mint) + self.mint = AsyncMock(side_effect=self._mint) + self.serialize_proofs = AsyncMock(return_value="cashuBrefund") + self.set_reserved_for_send = AsyncMock() + self.load_proofs = AsyncMock() + self.available_balance = SimpleNamespace(amount=0) + self.keysets: dict[str, Any] = {} + self.proofs: list[Any] = [] + + async def _observe(self, *args: Any, **kwargs: Any) -> None: + self.guard_depth_seen.append(wallet._wallet_operation_depth.get()) + + async def _melt_quote(self, invoice: str, amount_msat: int | None = None) -> Any: + await self._observe() + amount = int(invoice.rsplit(":", 1)[1]) + return SimpleNamespace( + quote=f"melt-{amount}", amount=amount, fee_reserve=self.fee_reserve + ) + + async def _melt(self, **kwargs: Any) -> Any: + await self._observe() + if isinstance(self.melt_state, BaseException): + raise self.melt_state + if self.melt_state == "hang": + await asyncio.sleep(5) + return SimpleNamespace(state=self.melt_state) + + async def _request_mint(self, amount: int, memo: str | None = None) -> Any: + await self._observe() + return SimpleNamespace(quote=f"mint-{amount}", request=f"lnbc:{amount}") + + async def _mint(self, amount: int, quote_id: str, split: Any = None) -> list[Any]: + await self._observe() + return [_proof(amount)] + + +class _PrimaryWallet: + def __init__(self, proofs: list[Any] | None = None, fee_reserve: int = 2): + self.proofs = proofs or [_proof(500), _proof(500)] + self.fee_reserve = fee_reserve + self.request_mint = AsyncMock(side_effect=self._request_mint) + self.mint = AsyncMock(return_value=[_proof(1)]) + self.load_proofs = AsyncMock() + self.available_balance = SimpleNamespace(amount=0) + self.keysets: dict[str, Any] = {} + self.select_to_send = AsyncMock(side_effect=self._select) + self.get_fees_for_proofs = Mock(return_value=0) + self.melt_quote = AsyncMock(side_effect=self._melt_quote) + + async def _request_mint(self, amount: int, memo: str | None = None) -> Any: + return SimpleNamespace(quote=f"mint-{amount}", request=f"lnbc:{amount}") + + async def _select(self, proofs: Any, amount: int, **kwargs: Any) -> Any: + return proofs, 0 + + async def _melt_quote(self, invoice: str, amount_msat: int | None = None) -> Any: + amount = int(invoice.rsplit(":", 1)[1]) + return SimpleNamespace( + quote=f"melt-{amount}", amount=amount, fee_reserve=self.fee_reserve + ) + + +@asynccontextmanager +async def _swap_env( + foreign: _ForeignWallet, + primary: _PrimaryWallet, + token: SimpleNamespace, +) -> AsyncGenerator[None, None]: + wallets = {FOREIGN: foreign, PRIMARY: primary} + + async def get_wallet(mint_url: str, unit: str = "sat", **kwargs: Any) -> Any: + return wallets[mint_url] + + async def run_mint_operation(factory: Any, **kwargs: Any) -> Any: + return await factory() + + with ( + patch.object(fms, "deserialize_token_from_string", return_value=token), + patch.object(fms, "assert_public_https_origin", AsyncMock()), + patch.object(fms, "get_wallet", get_wallet), + patch.object(fms, "run_mint_operation", run_mint_operation), + patch.object( + fms, + "get_proofs_per_mint_and_unit", + lambda w, m, u, not_reserved=False: list(w.proofs), + ), + ): + yield + + +async def _swap_rows(session: AsyncSession) -> list[CashuSwap]: + # Rows are written through other sessions; read with a fresh one so the + # test session's identity map cannot serve stale copies. + async with AsyncSession(session.bind, expire_on_commit=False) as fresh: + return list((await fresh.exec(select(CashuSwap))).all()) + + +async def _refund_row(session: AsyncSession, refund_id: str) -> Refund: + async with AsyncSession(session.bind, expire_on_commit=False) as fresh: + row = await fresh.get(Refund, refund_id) + assert row is not None + return row + + +def _mint_recovered() -> None: + """The melt timeout put the mint on cooldown; the reconciler runs later.""" + MintRateGuard._guards.clear() + + +# --- policy and budget ------------------------------------------------------ + + +def test_swap_is_off_by_default() -> None: + from routstr.core.settings import Settings + + assert Settings.__fields__["foreign_mint_policy"].default == "reject" + + +@pytest.mark.asyncio +async def test_swap_in_refused_when_policy_is_reject(engine: AsyncEngine) -> None: + settings.foreign_mint_policy = "reject" + key = ApiKey(hashed_key=KEY_HASH) + with patch.object(fms, "deserialize_token_from_string") as parse: + with pytest.raises(ForeignMintSwapError): + await fms.swap_in_and_credit("cashuA", key, Mock()) + parse.assert_not_called() + + +@pytest.mark.asyncio +async def test_foreign_operation_is_single_attempt_with_cooldown( + engine: AsyncEngine, +) -> None: + calls = 0 + + async def slow() -> None: + nonlocal calls + calls += 1 + await asyncio.sleep(5) + + started = time.monotonic() + with pytest.raises(ForeignMintUnavailableError): + await fms.run_foreign_mint_operation(slow, mint_url=FOREIGN, op_name="t") + assert time.monotonic() - started < 2 + assert calls == 1 + assert mint_cooldown_remaining(FOREIGN) > 0 + # Cooling down: refused offline without touching the mint again. + with pytest.raises(ForeignMintUnavailableError): + await fms.run_foreign_mint_operation(slow, mint_url=FOREIGN, op_name="t") + assert calls == 1 + + +@pytest.mark.asyncio +async def test_foreign_operation_refuses_to_run_under_wallet_guard( + engine: AsyncEngine, +) -> None: + async with wallet_operation_guard(): + with pytest.raises(RuntimeError): + await fms.run_foreign_mint_operation( + AsyncMock(), mint_url=FOREIGN, op_name="t" + ) + with pytest.raises(RuntimeError): + async with fms.foreign_mint_lock(FOREIGN): + pass + + +@pytest.mark.asyncio +async def test_foreign_budget_is_global_and_fails_fast(engine: AsyncEngine) -> None: + settings.foreign_mint_max_concurrency = 1 + release = asyncio.Event() + + async def hold() -> None: + await release.wait() + + holder = asyncio.create_task( + fms.run_foreign_mint_operation(hold, mint_url=FOREIGN, op_name="hold") + ) + await asyncio.sleep(0.01) + other_mint_called = False + + async def other() -> None: + nonlocal other_mint_called + other_mint_called = True + + # A different hostname does not get its own budget. + with pytest.raises(ForeignMintUnavailableError): + await fms.run_foreign_mint_operation( + other, mint_url="https://other.example", op_name="o" + ) + assert not other_mint_called + release.set() + await holder + + +@pytest.mark.asyncio +async def test_foreign_mint_lock_waits_bounded(engine: AsyncEngine) -> None: + entered = asyncio.Event() + release = asyncio.Event() + + async def hold() -> None: + async with fms.foreign_mint_lock(FOREIGN): + entered.set() + await release.wait() + + holder = asyncio.create_task(hold()) + await entered.wait() + with pytest.raises(ForeignMintUnavailableError): + async with fms.foreign_mint_lock(FOREIGN): + pass + release.set() + await holder + + +# --- inbound swap ----------------------------------------------------------- + + +@pytest.mark.asyncio +async def test_swap_in_rejects_non_https_mint_before_any_contact( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + get_wallet = AsyncMock() + with ( + patch.object( + fms, + "deserialize_token_from_string", + return_value=_token(mint="http://foreign.example"), + ), + patch.object(fms, "get_wallet", get_wallet), + ): + with pytest.raises(ForeignMintSwapError): + await fms.swap_in_and_credit("cashuAhttp", key, session) + get_wallet.assert_not_awaited() + assert await _swap_rows(session) == [] + + +@pytest.mark.asyncio +async def test_swap_in_happy_path_credits_net_and_pins_refund_mint( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet(fee_reserve=5) + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + credited = await fms.swap_in_and_credit("cashuAswap", key, session) + + # 1000 sat token, 1 sat input fee, 5 sat fee reserve -> 994 sat minted. + assert credited == 994_000 + assert [c.args[0] for c in primary.request_mint.await_args_list] == [999, 994] + foreign.melt.assert_awaited_once() + assert foreign.melt.await_args is not None + assert foreign.melt.await_args.kwargs["quote_id"] == "melt-994" + primary.mint.assert_awaited_once_with(994, quote_id="mint-994") + # Every call to the sender's mint ran with the wallet guard released. + assert foreign.guard_depth_seen and set(foreign.guard_depth_seen) == {0} + + await session.refresh(key) + assert key.balance == 994_000 + assert key.refund_mint_url == FOREIGN + (row,) = await _swap_rows(session) + assert (row.status, row.direction, row.destination_amount) == ( + "credited", + "in", + 994, + ) + assert row.fee_reserve == 5 and row.input_fees == 1 + ledger = (await session.exec(select(CashuTransaction))).all() + assert [(t.type, t.amount, t.mint_url) for t in ledger] == [("in", 994, PRIMARY)] + + +@pytest.mark.asyncio +async def test_swap_in_accepts_proofs_from_rotated_keysets( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet(fee_reserve=1) + primary = _PrimaryWallet() + token = _token() + token.keysets = ["old", "new"] + token.proofs = [_proof(400), _proof(600)] + + async with _swap_env(foreign, primary, token): + credited = await fms.swap_in_and_credit("cashuArotated", key, session) + + assert credited == 998_000 + foreign.get_fees_for_proofs.assert_called_once_with(token.proofs) + assert foreign.melt.await_count == 1 + + +@pytest.mark.asyncio +async def test_token_hash_is_unique_across_swap_journals(engine: AsyncEngine) -> None: + first = CashuSwap( + direction="in", + status="failed", + token_hash="same-token", + source_mint=FOREIGN, + source_unit="sat", + source_amount=100, + destination_mint=PRIMARY, + destination_unit="sat", + destination_amount=0, + ) + second = CashuSwap( + direction="in", + status="melting", + token_hash="same-token", + source_mint=FOREIGN, + source_unit="sat", + source_amount=100, + destination_mint=PRIMARY, + destination_unit="sat", + destination_amount=0, + ) + await fms._save(first) + with pytest.raises(IntegrityError): + await fms._save(second) + + +@pytest.mark.asyncio +async def test_swap_in_fee_shortfall_spends_nothing( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet(fee_reserve=2000) + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + with pytest.raises(ForeignMintSwapError): + await fms.swap_in_and_credit("cashuAsmall", key, session) + foreign.melt.assert_not_awaited() + assert await _swap_rows(session) == [] + await session.refresh(key) + assert key.balance == 0 + + +@pytest.mark.asyncio +async def test_swap_in_melt_timeout_is_journaled_as_ambiguous( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet(melt_state="hang") + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + with pytest.raises(SwapPendingError): + await fms.swap_in_and_credit("cashuAhang", key, session) + primary.mint.assert_not_awaited() + (row,) = await _swap_rows(session) + assert row.status == "ambiguous" + assert row.melt_quote_id == "melt-994" + await session.refresh(key) + assert key.balance == 0 + + +@pytest.mark.asyncio +async def test_swap_in_mint_refusal_fails_without_consuming_token( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet( + melt_state=Exception( + "Mint Error: not enough inputs provided for melt. Provided: 999, needed: 1004 (Code: 11000)" + ) + ) + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + with pytest.raises(ForeignMintSwapError): + await fms.swap_in_and_credit("cashuArefused", key, session) + (row,) = await _swap_rows(session) + assert row.status == "failed" + + +@pytest.mark.asyncio +async def test_swap_in_rejects_replayed_token( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet() + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + await fms.swap_in_and_credit("cashuAonce", key, session) + with pytest.raises(ValueError, match="already spent"): + await fms.swap_in_and_credit("cashuAonce", key, session) + foreign.melt.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_reconciler_credits_ambiguous_swap_once_mint_confirms_paid( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet(melt_state="hang") + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + with pytest.raises(SwapPendingError): + await fms.swap_in_and_credit("cashuAlate", key, session) + (row,) = await _swap_rows(session) + await fms._update(row, updated_at=int(time.time()) - 10_000) + _mint_recovered() + await fms.reconcile_swaps_once() + + foreign.get_melt_quote.assert_awaited_once_with("melt-994") + primary.mint.assert_awaited_once_with(994, quote_id="mint-994") + await session.refresh(key) + assert key.balance == 994_000 + assert key.refund_mint_url == FOREIGN + (row,) = await _swap_rows(session) + assert row.status == "credited" and row.claimed_at is None + + +@pytest.mark.asyncio +async def test_minted_swap_credit_is_atomic_and_cannot_repeat( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + swap = CashuSwap( + direction="in", + status="minted", + api_key_hashed_key=key.hashed_key, + token="cashuAatomic", + token_hash="atomic", + source_mint=FOREIGN, + source_unit="sat", + source_amount=1000, + destination_mint=PRIMARY, + destination_unit="sat", + destination_amount=998, + ) + await fms._save(swap) + + credited = await fms._finish_swap_in(swap, key=key, session=session) + assert credited == 998_000 + swap.status = "minted" # stale worker copy after the first transaction + with pytest.raises(TokenConsumedError, match="already credited"): + await fms._finish_swap_in(swap, key=key, session=session) + + await session.refresh(key) + assert key.balance == 998_000 + rows = list((await session.exec(select(CashuTransaction))).all()) + assert len(rows) == 1 + (stored,) = await _swap_rows(session) + assert stored.status == "credited" + + +@pytest.mark.asyncio +async def test_mint_recovery_reuses_proofs_tagged_with_quote() -> None: + proof = SimpleNamespace(amount=998, reserved=False, mint_id="mint-998") + mint = AsyncMock(side_effect=AssertionError("must not mint the quote twice")) + fake_wallet: Any = SimpleNamespace( + proofs=[proof], + load_proofs=AsyncMock(), + available_balance=SimpleNamespace(amount=998), + mint=mint, + ) + + recovered = await fms._mint_with_recovery( + fake_wallet, + 998, + "mint-998", + mint_url=PRIMARY, + foreign=False, + ) + + assert recovered == [proof] + mint.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_reconciler_fails_swap_the_mint_reports_unpaid( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet(melt_state="hang") + foreign.get_melt_quote.return_value = SimpleNamespace(state=MeltQuoteState.unpaid) + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + with pytest.raises(SwapPendingError): + await fms.swap_in_and_credit("cashuAunpaid", key, session) + (row,) = await _swap_rows(session) + await fms._update(row, updated_at=int(time.time()) - 10_000) + _mint_recovered() + await fms.reconcile_swaps_once() + + primary.mint.assert_not_awaited() + (row,) = await _swap_rows(session) + assert row.status == "failed" + await session.refresh(key) + assert key.balance == 0 + + +@pytest.mark.asyncio +async def test_reconciler_leaves_fresh_rows_alone( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session) + foreign = _ForeignWallet(melt_state="hang") + primary = _PrimaryWallet() + async with _swap_env(foreign, primary, _token(1000)): + with pytest.raises(SwapPendingError): + await fms.swap_in_and_credit("cashuAfresh", key, session) + await fms.reconcile_swaps_once() + foreign.get_melt_quote.assert_not_awaited() + + +# --- refund back to the user's mint ------------------------------------------- + + +def test_refund_destination_requires_foreign_mint_and_swap_policy( + engine: AsyncEngine, +) -> None: + foreign_key = ApiKey(hashed_key=KEY_HASH, refund_mint_url=FOREIGN) + assert fms.refund_destination_mint(foreign_key) == FOREIGN + assert fms.refund_destination_mint(ApiKey(hashed_key=KEY_HASH)) is None + assert ( + fms.refund_destination_mint( + ApiKey(hashed_key=KEY_HASH, refund_mint_url=PRIMARY) + ) + is None + ) + settings.foreign_mint_policy = "reject" + assert fms.refund_destination_mint(foreign_key) is None + + +async def _open_cashu_refund(session: AsyncSession, key: ApiKey) -> Refund: + return await refund_module.open_claim( + session, key, method="cashu", destination=None + ) + + +@pytest.mark.asyncio +async def test_open_claim_records_users_mint_as_destination( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=FOREIGN) + claim = await _open_cashu_refund(session, key) + assert claim.destination == FOREIGN + assert claim.mint_url == PRIMARY + + +@pytest.mark.asyncio +async def test_refund_swaps_back_to_users_mint_net_of_fees( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=FOREIGN) + claim = await _open_cashu_refund(session, key) + foreign = _ForeignWallet() + primary = _PrimaryWallet(fee_reserve=2) + execute = AsyncMock(return_value=(998, PRIMARY, "sat")) + async with _swap_env(foreign, primary, _token()): + with ( + patch.object(fms, "_execute_bolt11_payment", execute), + patch.object( + refund_module, + "deserialize_token_from_string", + return_value=SimpleNamespace(amount=998, unit="sat"), + ), + ): + result = await refund_module.execute(session, claim) + + assert result["status"] == "paid" + assert result["token"] == "cashuBrefund" + assert result["recipient"] == FOREIGN + assert result["sats"] == "998" + # 1000 sat refund, 2 sat fee reserve, 0 input fees -> 998 sat on the user's mint. + assert [c.args[0] for c in foreign.request_mint.await_args_list] == [1000, 998] + foreign.mint.assert_awaited_once_with(998, quote_id="mint-998") + assert execute.await_args is not None + plan = execute.await_args.args[0] + assert (plan.mint_url, plan.quote.quote) == (PRIMARY, "melt-998") + assert foreign.guard_depth_seen and set(foreign.guard_depth_seen) == {0} + + refreshed = await _refund_row(session, claim.id) + assert (refreshed.status, refreshed.token, refreshed.mint_url) == ( + "paid", + "cashuBrefund", + FOREIGN, + ) + (row,) = await _swap_rows(session) + assert (row.direction, row.status, row.destination_amount) == ( + "out", + "settled", + 998, + ) + assert row.refund_id == claim.id + + +@pytest.mark.asyncio +async def test_reconciler_settles_a_persisted_issued_refund( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=FOREIGN) + claim = await _open_cashu_refund(session, key) + swap = CashuSwap( + direction="out", + status="issued", + api_key_hashed_key=key.hashed_key, + refund_id=claim.id, + source_mint=PRIMARY, + source_unit="sat", + source_amount=1000, + destination_mint=FOREIGN, + destination_unit="sat", + destination_amount=998, + token="cashuBpersisted", + ) + await fms._save(swap) + + with patch.object(refund_module, "_record_cashu_payout", AsyncMock()) as record: + await fms.reconcile_swaps_once() + + refreshed = await _refund_row(session, claim.id) + assert (refreshed.status, refreshed.token) == ("paid", "cashuBpersisted") + (stored,) = await _swap_rows(session) + assert stored.status == "settled" and stored.claimed_at is None + record.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_refund_swap_ambiguous_melt_withholds_balance( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=FOREIGN) + claim = await _open_cashu_refund(session, key) + foreign = _ForeignWallet() + primary = _PrimaryWallet(fee_reserve=2) + execute = AsyncMock(side_effect=Bolt11PaymentAmbiguous("melt did not return")) + async with _swap_env(foreign, primary, _token()): + with patch.object(fms, "_execute_bolt11_payment", execute): + with pytest.raises(Exception) as exc_info: + await refund_module.execute(session, claim) + assert getattr(exc_info.value, "status_code", None) == 502 + foreign.mint.assert_not_awaited() + refreshed = await _refund_row(session, claim.id) + assert (refreshed.status, refreshed.quote_id) == ("ambiguous", "melt-998") + (row,) = await _swap_rows(session) + assert row.status == "ambiguous" + await session.refresh(key) + assert key.balance == 0 + + +@pytest.mark.asyncio +async def test_refund_swap_reconciler_finishes_after_melt_paid( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=FOREIGN) + claim = await _open_cashu_refund(session, key) + foreign = _ForeignWallet() + primary = _PrimaryWallet(fee_reserve=2) + execute = AsyncMock(side_effect=Bolt11PaymentAmbiguous("melt did not return")) + async with _swap_env(foreign, primary, _token()): + with patch.object(fms, "_execute_bolt11_payment", execute): + with pytest.raises(Exception): + await refund_module.execute(session, claim) + (row,) = await _swap_rows(session) + await fms._update(row, updated_at=int(time.time()) - 10_000) + with patch.object( + fms, "_check_bolt11_payment_status_locked", AsyncMock(return_value="paid") + ): + await fms.reconcile_swaps_once() + + foreign.mint.assert_awaited_once_with(998, quote_id="mint-998") + refreshed = await _refund_row(session, claim.id) + assert (refreshed.status, refreshed.token) == ("paid", "cashuBrefund") + (row,) = await _swap_rows(session) + assert row.status == "settled" + + +@pytest.mark.asyncio +async def test_refund_swap_reconciler_releases_balance_when_unpaid( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=FOREIGN) + claim = await _open_cashu_refund(session, key) + foreign = _ForeignWallet() + primary = _PrimaryWallet(fee_reserve=2) + execute = AsyncMock(side_effect=Bolt11PaymentAmbiguous("melt did not return")) + async with _swap_env(foreign, primary, _token()): + with patch.object(fms, "_execute_bolt11_payment", execute): + with pytest.raises(Exception): + await refund_module.execute(session, claim) + (row,) = await _swap_rows(session) + await fms._update(row, updated_at=int(time.time()) - 10_000) + with patch.object( + fms, "_check_bolt11_payment_status_locked", AsyncMock(return_value="unpaid") + ): + await fms.reconcile_swaps_once() + + foreign.mint.assert_not_awaited() + refreshed = await _refund_row(session, claim.id) + assert refreshed.status == "failed" + await session.refresh(key) + assert key.balance == 1_000_000 + (row,) = await _swap_rows(session) + assert row.status == "failed" + + +@pytest.mark.asyncio +async def test_refund_reconciler_does_not_mark_swap_claims_stuck( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=FOREIGN) + claim = await _open_cashu_refund(session, key) + await refund_module.record_quote(claim, "melt-998", PRIMARY) + await refund_module._reconcile(claim, int(time.time()) + 10_000) + refreshed = await _refund_row(session, claim.id) + assert refreshed.status == "pending" + + +@pytest.mark.asyncio +async def test_trusted_refund_mint_still_pays_directly( + engine: AsyncEngine, session: AsyncSession +) -> None: + key = await _make_key(session, balance=1_000_000, refund_mint_url=PRIMARY) + claim = await _open_cashu_refund(session, key) + assert claim.destination is None + send_token = AsyncMock(return_value="cashuBdirect") + swap = AsyncMock() + with ( + patch.object(refund_module, "send_token", send_token), + patch.object(fms, "swap_out_for_refund", swap), + patch.object(refund_module, "token_mint_url", lambda t, f: PRIMARY), + ): + result = await refund_module.execute(session, claim) + assert result["token"] == "cashuBdirect" + swap.assert_not_awaited() + + +# --- classification ----------------------------------------------------------- + + +def _status_and_code(error: Exception) -> tuple[int, str]: + classified = wallet.classify_redemption_error(error) + assert classified is not None + return classified[1], classified[3] + + +def test_swap_errors_have_dedicated_codes() -> None: + assert _status_and_code(ForeignMintSwapError("x")) == ( + 422, + "cashu_foreign_mint_swap_failed", + ) + assert _status_and_code(SwapPendingError("x")) == (409, "cashu_swap_pending") + assert _status_and_code(ForeignMintUnavailableError("x")) == ( + 503, + "cashu_source_mint_unreachable", + ) + + +def test_ambiguous_melt_error_type_is_lnurl_ambiguous() -> None: + assert issubclass(MeltOutcomeAmbiguousError, Exception) diff --git a/tests/unit/test_provider_slugs.py b/tests/unit/test_provider_slugs.py index 0d446fc1..32d4957e 100644 --- a/tests/unit/test_provider_slugs.py +++ b/tests/unit/test_provider_slugs.py @@ -14,6 +14,24 @@ from routstr.core.provider_slugs import ( from routstr.upstream.helpers import _seed_providers_from_settings +def _isolate_provider_environment(monkeypatch: pytest.MonkeyPatch) -> None: + for env_key in ( + "OPENAI_API_KEY", + "ANTHROPIC_API_KEY", + "OPENROUTER_API_KEY", + "GROQ_API_KEY", + "PERPLEXITY_API_KEY", + "FIREWORKS_API_KEY", + "XAI_API_KEY", + "DEEPSEEK_API_KEY", + "TINFOIL_API_KEY", + "TYPESAFE_API_KEY", + "OLLAMA_BASE_URL", + "OLLAMA_API_KEY", + ): + monkeypatch.delenv(env_key, raising=False) + + @pytest.mark.asyncio async def test_allocate_unique_provider_slug_is_deterministic_with_suffixes() -> None: engine = create_async_engine("sqlite+aiosqlite:///:memory:") @@ -81,6 +99,7 @@ async def test_seed_providers_from_settings_sets_deterministic_slug( async with engine.begin() as conn: await conn.run_sync(SQLModel.metadata.create_all) + _isolate_provider_environment(monkeypatch) monkeypatch.setenv("OPENAI_API_KEY", "seeded-openai-key") class SettingsStub: @@ -108,6 +127,7 @@ async def test_seed_providers_from_settings_keeps_slug_stable_on_reseed( async with engine.begin() as conn: await conn.run_sync(SQLModel.metadata.create_all) + _isolate_provider_environment(monkeypatch) monkeypatch.setenv("OPENAI_API_KEY", "seeded-openai-key") class SettingsStub: diff --git a/tests/unit/test_settings.py b/tests/unit/test_settings.py index 61f657ac..35011cf0 100644 --- a/tests/unit/test_settings.py +++ b/tests/unit/test_settings.py @@ -1,6 +1,9 @@ import json import os +import subprocess +import sys from pathlib import Path +from typing import Any import pytest from pydantic.v1 import ValidationError @@ -63,6 +66,29 @@ def test_payout_settings_have_sensible_defaults() -> None: assert s.payout_interval_seconds == 900 +def test_foreign_mint_policy_rejects_typos() -> None: + bad_policy: Any = "swpa" + with pytest.raises(ValidationError): + Settings(foreign_mint_policy=bad_policy) + + +def test_cashu_import_cannot_override_operator_environment() -> None: + env = dict(os.environ) + env["CASHU_MINTS"] = "https://mint.operator.example" + result = subprocess.run( + [ + sys.executable, + "-c", + "import os, routstr; print(os.environ['CASHU_MINTS'])", + ], + check=True, + capture_output=True, + text=True, + env=env, + ) + assert result.stdout.splitlines()[-1] == "https://mint.operator.example" + + def test_database_pool_defaults_provide_concurrency_headroom() -> None: s = Settings() assert s.database_pool_size == 10 diff --git a/tests/unit/test_wallet.py b/tests/unit/test_wallet.py index 30b71f09..be50987c 100644 --- a/tests/unit/test_wallet.py +++ b/tests/unit/test_wallet.py @@ -762,8 +762,8 @@ async def test_credit_balance() -> None: "routstr.wallet.recieve_token", return_value=(1000, "sat", "http://mint:3338"), ): - with patch("routstr.wallet.store_cashu_transaction", AsyncMock()): - amount = await credit_balance(token_str, mock_key, mock_session) + mock_session.add = Mock() + amount = await credit_balance(token_str, mock_key, mock_session) assert amount == 1000000 # converted to msat assert mock_key.balance == 6000000 # Should be updated after refresh # Verify atomic operations were used @@ -784,12 +784,9 @@ async def test_concurrent_duplicate_token_credits_exactly_once() -> None: ValueError("Mint Error: proofs already spent (Code: 11001)"), ] ) - store = AsyncMock() + session.add = Mock() - with ( - patch("routstr.wallet.recieve_token", receive), - patch("routstr.wallet.store_cashu_transaction", store), - ): + with patch("routstr.wallet.recieve_token", receive): results = await asyncio.gather( credit_balance("cashuAduplicate", key, session), credit_balance("cashuAduplicate", key, session), @@ -801,7 +798,7 @@ async def test_concurrent_duplicate_token_credits_exactly_once() -> None: classified = classify_redemption_error(failure) assert classified is not None and classified[3] == "cashu_token_already_spent" assert session.exec.await_count == 1 - store.assert_awaited_once() + assert session.add.call_count == 1 @pytest.mark.asyncio @@ -819,9 +816,9 @@ async def test_credit_balance_redeems_on_token_mint_not_key_mint() -> None: mock_session.exec.return_value.rowcount = 1 receive = AsyncMock(return_value=(1000, "sat", key_mint)) + mock_session.add = Mock() with patch("routstr.wallet.recieve_token", receive): - with patch("routstr.wallet.store_cashu_transaction", AsyncMock()): - await credit_balance("cashuAtoken", mock_key, mock_session) + await credit_balance("cashuAtoken", mock_key, mock_session) receive.assert_awaited_once_with("cashuAtoken", destination_unit="sat") @@ -969,36 +966,41 @@ async def test_credit_balance_msat_unit_not_converted() -> None: "routstr.wallet.recieve_token", return_value=(1_000_000, "msat", "http://mint:3338"), ): - with patch("routstr.wallet.store_cashu_transaction", AsyncMock()): - amount = await credit_balance("cashuAtest", mock_key, mock_session) + mock_session.add = Mock() + amount = await credit_balance("cashuAtest", mock_key, mock_session) assert amount == 1_000_000 assert mock_session.commit.called @pytest.mark.asyncio -async def test_credit_balance_propagates_audit_store_failure_after_credit() -> None: - """A final transaction-history failure propagates after committing credit.""" - mock_key = Mock() - mock_key.balance = 0 - mock_key.hashed_key = "test_hash" +async def test_credit_balance_rolls_back_when_audit_row_cannot_flush() -> None: + """Balance and history stay atomic when the ledger write fails.""" + mock_key = Mock( + balance=0, + hashed_key="test_hash", + refund_mint_url=None, + refund_currency=None, + ) mock_session = AsyncMock() + mock_session.exec.return_value.rowcount = 1 + mock_session.add = Mock() + mock_session.flush.side_effect = Exception("history table locked") from routstr.core.settings import settings - with patch.object(settings, "cashu_mints", ["http://mint:3338"]): - with patch( + with ( + patch.object(settings, "cashu_mints", ["http://mint:3338"]), + patch( "routstr.wallet.recieve_token", return_value=(1000, "sat", "http://mint:3338"), - ): - with patch( - "routstr.wallet.store_cashu_transaction", - side_effect=Exception("history table locked"), - ): - with pytest.raises(Exception, match="history table locked"): - await credit_balance("cashuAtest", mock_key, mock_session) + ), + pytest.raises(Exception, match="crediting the balance failed"), + ): + await credit_balance("cashuAtest", mock_key, mock_session) - assert mock_session.commit.called + mock_session.commit.assert_not_awaited() + mock_session.rollback.assert_awaited_once() # --- Mint-unreachable classification (is_mint_connection_error) ---------------