diff --git a/migrations/versions/e4c7a1b9d520_add_direction_to_lightning_invoices.py b/migrations/versions/e4c7a1b9d520_add_direction_to_lightning_invoices.py new file mode 100644 index 00000000..a0a0964d --- /dev/null +++ b/migrations/versions/e4c7a1b9d520_add_direction_to_lightning_invoices.py @@ -0,0 +1,25 @@ +"""Add direction to lightning_invoices + +Revision ID: e4c7a1b9d520 +Revises: 3a0fbd387f10 +Create Date: 2026-09-20 00:00:00.000000 +""" + +import sqlalchemy as sa +from alembic import op + +revision = "e4c7a1b9d520" +down_revision = "3a0fbd387f10" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.add_column( + "lightning_invoices", + sa.Column("direction", sa.String(), nullable=False, server_default="in"), + ) + + +def downgrade() -> None: + op.drop_column("lightning_invoices", "direction") diff --git a/routstr/core/admin.py b/routstr/core/admin.py index cc140f6c..365e2f68 100644 --- a/routstr/core/admin.py +++ b/routstr/core/admin.py @@ -2051,6 +2051,7 @@ async def get_transactions_api( async def get_lightning_invoices_api( status: str | None = None, purpose: str | None = None, + direction: str | None = None, search: str | None = None, limit: int = 50, offset: int = 0, @@ -2063,6 +2064,8 @@ async def get_lightning_invoices_api( base = base.where(LightningInvoice.status == status) if purpose: base = base.where(LightningInvoice.purpose == purpose) + if direction: + base = base.where(LightningInvoice.direction == direction) if search: pattern = f"%{search}%" base = base.where( diff --git a/routstr/core/db.py b/routstr/core/db.py index eb4ba2fe..78019c32 100644 --- a/routstr/core/db.py +++ b/routstr/core/db.py @@ -509,14 +509,17 @@ class LightningInvoice(SQLModel, table=True): # type: ignore status: str = Field( default="pending", description=( - "pending, settlement_pending, paid, expired, cancelled, " + "pending, settlement_pending, paid, failed, expired, cancelled, " "reconciliation_required" ), ) api_key_hash: str | None = Field( default=None, description="Associated API key hash for topup operations" ) - purpose: str = Field(description="create or topup") + direction: str = Field( + default="in", description="in for incoming invoices, out for payouts" + ) + purpose: str = Field(description="create, topup or payout") mint_url: str | None = Field( default=None, description="Mint URL where the quote was created (fallback tracking)", @@ -1005,6 +1008,80 @@ async def complete_routstr_fee_payout( return result.rowcount == 1 +async def record_lightning_payout( + session: AsyncSession, + *, + quote_id: str, + bolt11: str, + amount_sats: int, + mint_url: str, + destination: str, +) -> None: + """Record a dispatched payout so it shows up in the Lightning history.""" + session.add( + LightningInvoice( + id=uuid.uuid4().hex, + bolt11=bolt11, + amount_sats=amount_sats, + description=f"Payout to {destination}", + payment_hash=quote_id, + status="pending", + direction="out", + purpose="payout", + mint_url=mint_url, + # Payouts settle or fail at the mint; they never expire on our side. + expires_at=int(time.time()), + ) + ) + await session.commit() + + +async def settle_lightning_payout( + session: AsyncSession, + quote_id: str, + *, + status: str, + amount_sats: int | None = None, +) -> None: + result = await session.exec( + select(LightningInvoice) + .where(col(LightningInvoice.payment_hash) == quote_id) + .where(col(LightningInvoice.direction) == "out") + ) + payout = result.first() + if payout is None: + logger.warning( + "No Lightning payout history row for quote", + extra={"quote_id": quote_id, "status": status}, + ) + return + payout.status = status + if status == "paid": + payout.paid_at = int(time.time()) + if amount_sats is not None: + payout.amount_sats = amount_sats + session.add(payout) + await session.commit() + + +UNSETTLED_PAYOUT_STATUSES = ("pending", "reconciliation_required") + + +async def list_unsettled_lightning_payouts( + session: AsyncSession, mint_url: str, *, created_before: int +) -> list[LightningInvoice]: + """Payout rows whose mint outcome was never written back to history.""" + result = await session.exec( + select(LightningInvoice) + .where(col(LightningInvoice.direction) == "out") + .where(col(LightningInvoice.mint_url) == mint_url) + .where(col(LightningInvoice.status).in_(UNSETTLED_PAYOUT_STATUSES)) + .where(col(LightningInvoice.created_at) < created_before) + .order_by(col(LightningInvoice.created_at)) + ) + return list(result.all()) + + async def total_user_liability(db_session: AsyncSession) -> int: """Return all outstanding user funds in millisatoshis. diff --git a/routstr/lightning.py b/routstr/lightning.py index b3ece68f..732f1847 100644 --- a/routstr/lightning.py +++ b/routstr/lightning.py @@ -443,7 +443,8 @@ async def get_invoice_status( structured_errors: bool = Depends(_uses_v2_errors), ) -> InvoiceStatusResponse: invoice = await session.get(LightningInvoice, invoice_id) - if not invoice: + # Payout rows (direction="out") are operator history, never user invoices. + if not invoice or invoice.direction != "in": raise _invoice_error( 404, "Invoice not found", @@ -486,7 +487,9 @@ async def recover_invoice( structured_errors: bool = Depends(_uses_v2_errors), ) -> InvoiceStatusResponse: result = await session.exec( - select(LightningInvoice).where(LightningInvoice.bolt11 == request.bolt11) + select(LightningInvoice) + .where(LightningInvoice.bolt11 == request.bolt11) + .where(col(LightningInvoice.direction) == "in") ) invoice = result.first() @@ -942,6 +945,7 @@ async def _expire_overdue_invoices(now: int) -> int: expired = await expiry_session.exec( # type: ignore[call-overload] update(LightningInvoice) .where( + col(LightningInvoice.direction) == "in", col(LightningInvoice.status) == "pending", col(LightningInvoice.expires_at) < now, ) @@ -959,13 +963,17 @@ async def _process_invoice_watch_batch(session: AsyncSession, prev_now: int) -> logger.info("Expired overdue invoices", extra={"invoice_count": swept}) settling = await session.exec( select(LightningInvoice) - .where(col(LightningInvoice.status) == "settlement_pending") + .where( + col(LightningInvoice.direction) == "in", + col(LightningInvoice.status) == "settlement_pending", + ) .order_by(col(LightningInvoice.created_at)) .limit(INVOICE_WATCH_BATCH_LIMIT // 2) ) unpaid = await session.exec( select(LightningInvoice) .where( + col(LightningInvoice.direction) == "in", col(LightningInvoice.status) == "pending", col(LightningInvoice.expires_at) >= now, ) @@ -975,6 +983,7 @@ async def _process_invoice_watch_batch(session: AsyncSession, prev_now: int) -> recoverable = await session.exec( select(LightningInvoice) .where( + col(LightningInvoice.direction) == "in", col(LightningInvoice.status) == "expired", col(LightningInvoice.expires_at) > now - INVOICE_EXPIRY_GRACE_SECONDS, ) diff --git a/routstr/payment/lnurl.py b/routstr/payment/lnurl.py index 5f8545b0..808de593 100644 --- a/routstr/payment/lnurl.py +++ b/routstr/payment/lnurl.py @@ -325,7 +325,7 @@ async def raw_send_to_lnurl( unit: str, amount: int | None = None, *, - on_melt_quote: Callable[[str], Awaitable[None]] | None = None, + on_melt_quote: Callable[[str, str], Awaitable[None]] | None = None, ) -> int: """Send funds to an LNURL address. @@ -413,7 +413,7 @@ async def raw_send_to_lnurl( raise LNURLError("Cashu melt fees exceed the requested gross amount") if on_melt_quote is not None: - await on_melt_quote(melt_quote_resp.quote) + await on_melt_quote(melt_quote_resp.quote, bolt11_invoice) assert selected_proofs is not None proofs = selected_proofs diff --git a/routstr/wallet.py b/routstr/wallet.py index 806484ef..2c592b3b 100644 --- a/routstr/wallet.py +++ b/routstr/wallet.py @@ -36,7 +36,7 @@ from .mint import ( mint_cooldown_remaining, run_mint_operation, ) -from .payment.lnurl import raw_send_to_lnurl +from .payment.lnurl import MeltUnpaidError, raw_send_to_lnurl # cashu 0.20.x passes the `proxies` kwarg httpx removed in 0.28; see the module # docstring. Installed at import so no mint call can run before the patch. @@ -1529,6 +1529,102 @@ async def fetch_all_balances( ) +PAYOUT_HISTORY_STALE_SECONDS = 600 + + +async def _record_payout_history( + *, + quote_id: str, + bolt11: str, + amount_sats: int, + mint_url: str, + destination: str, +) -> None: + """Best-effort history insert; a history failure must never block a payout.""" + try: + async with db.create_session() as session: + await db.record_lightning_payout( + session, + quote_id=quote_id, + bolt11=bolt11, + amount_sats=amount_sats, + mint_url=mint_url, + destination=destination, + ) + except Exception as e: + logger.error( + "Failed to record Lightning payout history", + extra={ + "error": str(e), + "error_type": type(e).__name__, + "quote_id": quote_id, + "mint_url": mint_url, + }, + ) + + +async def _reconcile_stale_payout_history(mint_url: str, unit: str) -> None: + """Resolve payout rows left pending by a crash or an ambiguous melt. + + Runs under ``wallet_operation_guard``. Only writes what the mint asserts + (paid/unpaid); quotes still pending or unreachable are left for later. + """ + try: + cutoff = int(time.time()) - PAYOUT_HISTORY_STALE_SECONDS + async with db.create_session() as session: + stale = await db.list_unsettled_lightning_payouts( + session, mint_url, created_before=cutoff + ) + for payout in stale: + quote_state = await _check_bolt11_payment_status_locked( + mint_url, unit, payout.payment_hash + ) + if quote_state == "paid": + await _settle_payout_history(payout.payment_hash, status="paid") + elif quote_state == "unpaid": + await _settle_payout_history(payout.payment_hash, status="failed") + else: + continue + logger.info( + "Reconciled stale Lightning payout history", + extra={ + "quote_id": payout.payment_hash, + "mint_url": mint_url, + "quote_state": quote_state, + }, + ) + except Exception as e: + logger.error( + "Failed to reconcile Lightning payout history", + extra={ + "error": str(e), + "error_type": type(e).__name__, + "mint_url": mint_url, + }, + ) + + +async def _settle_payout_history( + quote_id: str, *, status: str, amount_sats: int | None = None +) -> None: + """Best-effort history update after the external payment outcome is known.""" + try: + async with db.create_session() as session: + await db.settle_lightning_payout( + session, quote_id, status=status, amount_sats=amount_sats + ) + except Exception as e: + logger.error( + "Failed to update Lightning payout history", + extra={ + "error": str(e), + "error_type": type(e).__name__, + "quote_id": quote_id, + "status": status, + }, + ) + + async def _payout_mint_and_unit(mint_url: str, unit: str) -> None: """Send only conservatively proven owner funds for one wallet.""" try: @@ -1582,13 +1678,49 @@ async def _payout_mint_and_unit(mint_url: str, unit: str) -> None: ) if available_balance > min_amount: payout_amount = min(available_balance, max_amount) - amount_received = await raw_send_to_lnurl( - wallet, - proofs, - settings.receive_ln_address, - unit, - amount=payout_amount, - ) + payout_quote_id: str | None = None + + async def record_payout(quote_id: str, bolt11: str) -> None: + nonlocal payout_quote_id + payout_quote_id = quote_id + await _record_payout_history( + quote_id=quote_id, + bolt11=bolt11, + amount_sats=( + payout_amount + if unit == "sat" + else _msats_to_sats(payout_amount) + ), + mint_url=mint_url, + destination=settings.receive_ln_address, + ) + + try: + amount_received = await raw_send_to_lnurl( + wallet, + proofs, + settings.receive_ln_address, + unit, + amount=payout_amount, + on_melt_quote=record_payout, + ) + except Exception as e: + if payout_quote_id is not None: + await _settle_payout_history( + payout_quote_id, + status=( + "failed" + if isinstance(e, MeltUnpaidError) + else "reconciliation_required" + ), + ) + raise + if payout_quote_id is not None: + await _settle_payout_history( + payout_quote_id, + status="paid", + amount_sats=_msats_to_sats(amount_received), + ) logger.info( "Payout sent successfully", extra={ @@ -1635,6 +1767,7 @@ async def periodic_payout() -> None: # Proof mutation, liability observation, and sending are one # cross-process critical section. Credits take the same lock. async with wallet_operation_guard(): + await _reconcile_stale_payout_history(mint_url, unit) await _payout_mint_and_unit(mint_url, unit) except Exception as e: logger.error( @@ -1861,6 +1994,7 @@ async def periodic_routstr_fee_payout() -> None: payout_unit, ) if completed: + await _settle_payout_history(payout_quote_id, status="paid") logger.info( "Routstr fee payout reconciled as paid", extra={"payout_quote_id": payout_quote_id}, @@ -1875,6 +2009,9 @@ async def periodic_routstr_fee_payout() -> None: payout_unit, ) if restored: + await _settle_payout_history( + payout_quote_id, status="failed" + ) logger.warning( "Routstr fee payout reconciled as unpaid and restored for retry", extra={"payout_quote_id": payout_quote_id}, @@ -1905,7 +2042,7 @@ async def periodic_routstr_fee_payout() -> None: attempt_quote_id: str | None = None - async def checkpoint_quote(quote_id: str) -> None: + async def checkpoint_quote(quote_id: str, bolt11: str) -> None: nonlocal attempt_quote_id async with db.create_session() as session: checkpointed = await db.reset_routstr_fee( @@ -1918,6 +2055,13 @@ async def periodic_routstr_fee_payout() -> None: if not checkpointed: raise _RoutstrFeePayoutAlreadyClaimed attempt_quote_id = quote_id + await _record_payout_history( + quote_id=quote_id, + bolt11=bolt11, + amount_sats=accumulated_sats, + mint_url=settings.primary_mint, + destination=ROUTSTR_LN_ADDRESS, + ) try: amount_received = await raw_send_to_lnurl( @@ -1944,6 +2088,9 @@ async def periodic_routstr_fee_payout() -> None: extra={"payout_in_progress_msats": paid_msats}, exc_info=isinstance(e, Exception), ) + await _settle_payout_history( + attempt_quote_id, status="reconciliation_required" + ) if not isinstance(e, Exception): raise continue @@ -1964,6 +2111,9 @@ async def periodic_routstr_fee_payout() -> None: extra={"payout_in_progress_msats": paid_msats}, exc_info=isinstance(e, Exception), ) + await _settle_payout_history( + attempt_quote_id, status="reconciliation_required" + ) if not isinstance(e, Exception): raise continue @@ -1972,8 +2122,17 @@ async def periodic_routstr_fee_payout() -> None: "Routstr fee payout sent but checkpoint was not completed; awaiting quote reconciliation", extra={"payout_in_progress_msats": paid_msats}, ) + await _settle_payout_history( + attempt_quote_id, status="reconciliation_required" + ) continue + await _settle_payout_history( + attempt_quote_id, + status="paid", + amount_sats=_msats_to_sats(amount_received), + ) + logger.info( "Routstr fee payout sent", extra={ @@ -1990,8 +2149,8 @@ async def periodic_routstr_fee_payout() -> None: def _quote_callback( notify: Callable[[str, str], Awaitable[None]], mint: str -) -> Callable[[str], Awaitable[None]]: - async def callback(quote_id: str) -> None: +) -> Callable[[str, str], Awaitable[None]]: + async def callback(quote_id: str, _bolt11: str) -> None: await notify(quote_id, mint) return callback diff --git a/tests/integration/test_lightning_invoice_constraints.py b/tests/integration/test_lightning_invoice_constraints.py index 7370ebe4..f715e67d 100644 --- a/tests/integration/test_lightning_invoice_constraints.py +++ b/tests/integration/test_lightning_invoice_constraints.py @@ -17,8 +17,10 @@ import pytest from cashu.core.base import Proof from sqlalchemy import inspect from sqlalchemy.ext.asyncio import AsyncEngine +from sqlmodel import select from sqlmodel.ext.asyncio.session import AsyncSession +from routstr.core import db from routstr.core.db import ApiKey, LightningInvoice from routstr.lightning import _create_api_key_record @@ -81,6 +83,41 @@ async def test_invoice_persists_validity_date( assert stored.validity_date == expiry +@pytest.mark.asyncio +async def test_outgoing_payout_history_is_recorded_and_settled( + integration_session: AsyncSession, +) -> None: + await db.record_lightning_payout( + integration_session, + quote_id="payout-quote", + bolt11="lnbc1payout", + amount_sats=1_000, + mint_url="https://mint.test", + destination="owner@example.com", + ) + + result = await integration_session.exec( + select(LightningInvoice).where(LightningInvoice.payment_hash == "payout-quote") + ) + payout = result.one() + assert payout.direction == "out" + assert payout.purpose == "payout" + assert payout.status == "pending" + assert payout.mint_url == "https://mint.test" + + await db.settle_lightning_payout( + integration_session, + "payout-quote", + status="paid", + amount_sats=995, + ) + await integration_session.refresh(payout) + + assert payout.status == "paid" + assert payout.amount_sats == 995 + assert payout.paid_at is not None + + # --------------------------------------------------------------------------- # Propagation to ApiKey # --------------------------------------------------------------------------- diff --git a/tests/integration/test_lightning_settlement.py b/tests/integration/test_lightning_settlement.py index 28803315..1a9a28e4 100644 --- a/tests/integration/test_lightning_settlement.py +++ b/tests/integration/test_lightning_settlement.py @@ -5,6 +5,7 @@ from unittest.mock import AsyncMock, Mock, patch import pytest from cashu.core.base import Proof +from fastapi import HTTPException from sqlalchemy.ext.asyncio import AsyncEngine from sqlmodel import col, update from sqlmodel.ext.asyncio.session import AsyncSession @@ -13,12 +14,15 @@ from routstr.core.db import ApiKey, LightningInvoice from routstr.lightning import ( INVOICE_EXPIRY_GRACE_SECONDS, INVOICE_WATCH_BATCH_LIMIT, + InvoiceRecoverRequest, _expire_invoice_if_authoritatively_unpaid, _expire_overdue_invoices, _finalize_invoice_settlement, _InvoiceSettlement, _process_invoice_watch_batch, check_invoice_payment, + get_invoice_status, + recover_invoice, ) @@ -58,7 +62,8 @@ async def test_invoice_read_transaction_closes_before_external_mint_io( return wallet with patch( - "routstr.lightning.get_wallet", side_effect=get_wallet_without_open_db_transaction + "routstr.lightning.get_wallet", + side_effect=get_wallet_without_open_db_transaction, ): await check_invoice_payment(stored, integration_session) @@ -191,9 +196,7 @@ async def test_failed_final_commit_rolls_back_claim_and_credit_for_retry( assert unchanged.balance == 100_000 async with AsyncSession(integration_engine, expire_on_commit=False) as retry: - settled, _ = await _finalize_invoice_settlement( - snapshot, retry, 1_700_000_001 - ) + settled, _ = await _finalize_invoice_settlement(snapshot, retry, 1_700_000_001) assert settled async with AsyncSession(integration_engine, expire_on_commit=False) as verify: @@ -302,9 +305,7 @@ async def test_expiry_cas_cannot_overwrite_concurrent_paid_invoice( assert result.rowcount == 1 await paid.commit() - expired = await _expire_invoice_if_authoritatively_unpaid( - stale, caller, True - ) + expired = await _expire_invoice_if_authoritatively_unpaid(stale, caller, True) assert expired is False assert stale.status == "paid" @@ -381,8 +382,9 @@ async def test_sweep_expires_only_overdue_pending_invoices( overdue = _lightning_invoice(expires_at=now - 1) fresh = _lightning_invoice(expires_at=now + 3600) settling = _lightning_invoice(expires_at=now - 1, status="settlement_pending") + outgoing = _lightning_invoice(expires_at=now - 1, direction="out") async with AsyncSession(integration_engine, expire_on_commit=False) as seed: - seed.add_all([overdue, fresh, settling]) + seed.add_all([overdue, fresh, settling, outgoing]) await seed.commit() await _expire_overdue_invoices(now) @@ -392,6 +394,7 @@ async def test_sweep_expires_only_overdue_pending_invoices( (overdue, "expired"), (fresh, "pending"), (settling, "settlement_pending"), + (outgoing, "pending"), ): stored = await verify.get(LightningInvoice, invoice.id) assert stored is not None @@ -409,8 +412,11 @@ async def test_watch_batch_expires_overdue_invoices_and_keeps_settling_rows( settling = _lightning_invoice( expires_at=now - 86_400, created_at=now - 86_400, status="settlement_pending" ) + outgoing = _lightning_invoice( + expires_at=now + 3600, created_at=now, direction="out" + ) async with AsyncSession(integration_engine, expire_on_commit=False) as seed: - seed.add_all([overdue, fresh, settling]) + seed.add_all([overdue, fresh, settling, outgoing]) await seed.commit() polled: list[str] = [] @@ -425,6 +431,7 @@ async def test_watch_batch_expires_overdue_invoices_and_keeps_settling_rows( assert fresh.id in polled assert settling.id in polled + assert outgoing.id not in polled async with AsyncSession(integration_engine, expire_on_commit=False) as verify: stored = await verify.get(LightningInvoice, overdue.id) @@ -618,3 +625,34 @@ async def test_recovery_tail_cannot_starve_owed_or_live_invoices( assert len(polled) == INVOICE_WATCH_BATCH_LIMIT assert {inv.id for inv in settling} <= set(polled) assert {inv.id for inv in fresh} <= set(polled) + + +@pytest.mark.asyncio +async def test_public_invoice_endpoints_ignore_payout_rows( + integration_engine: AsyncEngine, + patched_db_engine: None, +) -> None: + """A payout's bolt11/id must not let /recover or /status touch the row.""" + payout = _lightning_invoice( + direction="out", + purpose="payout", + expires_at=int(time.time()) - 1, + ) + async with AsyncSession(integration_engine, expire_on_commit=False) as seed: + seed.add(payout) + await seed.commit() + + async with AsyncSession(integration_engine, expire_on_commit=False) as session: + with pytest.raises(HTTPException) as recover_error: + await recover_invoice( + InvoiceRecoverRequest(bolt11=payout.bolt11), session, False + ) + with pytest.raises(HTTPException) as status_error: + await get_invoice_status(payout.id, session, False) + assert recover_error.value.status_code == 404 + assert status_error.value.status_code == 404 + + async with AsyncSession(integration_engine, expire_on_commit=False) as verify: + stored = await verify.get(LightningInvoice, payout.id) + assert stored is not None + assert stored.status == "pending" diff --git a/tests/unit/test_fee_payout_crash_safety.py b/tests/unit/test_fee_payout_crash_safety.py index c6fc3073..b60f42a0 100644 --- a/tests/unit/test_fee_payout_crash_safety.py +++ b/tests/unit/test_fee_payout_crash_safety.py @@ -1,5 +1,5 @@ import asyncio -from collections.abc import AsyncGenerator +from collections.abc import AsyncGenerator, Generator from contextlib import asynccontextmanager from types import SimpleNamespace from unittest.mock import AsyncMock, Mock, patch @@ -29,6 +29,21 @@ def _session_context(session: Mock) -> _SessionContext: return _SessionContext(session) +@pytest.fixture(autouse=True) +def _mock_lightning_payout_history() -> Generator[ + tuple[AsyncMock, AsyncMock], None, None +]: + with ( + patch( + "routstr.wallet.db.record_lightning_payout", new_callable=AsyncMock + ) as record, + patch( + "routstr.wallet.db.settle_lightning_payout", new_callable=AsyncMock + ) as settle, + ): + yield record, settle + + @pytest.mark.asyncio async def test_fee_payout_checkpoint_is_atomic_and_durable() -> None: engine = create_async_engine("sqlite+aiosqlite://") @@ -178,7 +193,7 @@ async def test_fee_payout_prepares_wallet_then_checkpoints_before_sending() -> N async def send(*_args: object, **kwargs: object) -> int: checkpoint_quote = kwargs["on_melt_quote"] - await checkpoint_quote("quote-1") # type: ignore[operator] + await checkpoint_quote("quote-1", "lnbc1payout") # type: ignore[operator] events.append("send") return 5 @@ -256,7 +271,7 @@ async def test_fee_payout_lost_checkpoint_race_does_not_send() -> None: async def send(*_args: object, **kwargs: object) -> int: checkpoint_quote = kwargs["on_melt_quote"] - await checkpoint_quote("quote-1") # type: ignore[operator] + await checkpoint_quote("quote-1", "lnbc1payout") # type: ignore[operator] await dispatched() return 5 @@ -289,7 +304,9 @@ async def test_fee_payout_lost_checkpoint_race_does_not_send() -> None: @pytest.mark.asyncio -async def test_fee_payout_finalizes_a_paid_unresolved_quote_without_resending() -> None: +async def test_fee_payout_finalizes_a_paid_unresolved_quote_without_resending( + _mock_lightning_payout_history: tuple[AsyncMock, AsyncMock], +) -> None: session = Mock() fee = SimpleNamespace( accumulated_msats=10_000, @@ -331,10 +348,14 @@ async def test_fee_payout_finalizes_a_paid_unresolved_quote_without_resending() ) restore.assert_not_awaited() send.assert_not_awaited() + _, settle = _mock_lightning_payout_history + settle.assert_awaited_once_with(session, "quote-1", status="paid", amount_sats=None) @pytest.mark.asyncio -async def test_fee_payout_restores_only_an_unpaid_quote_and_retries() -> None: +async def test_fee_payout_restores_only_an_unpaid_quote_and_retries( + _mock_lightning_payout_history: tuple[AsyncMock, AsyncMock], +) -> None: session = Mock() unresolved_fee = SimpleNamespace( accumulated_msats=10_000, @@ -351,7 +372,9 @@ async def test_fee_payout_restores_only_an_unpaid_quote_and_retries() -> None: ) async def send(*_args: object, **kwargs: object) -> int: - await kwargs["on_melt_quote"]("quote-2") # type: ignore[index,operator] + await kwargs["on_melt_quote"]( # type: ignore[index,operator] + "quote-2", "lnbc1payout" + ) return 15 with ( @@ -398,6 +421,8 @@ async def test_fee_payout_restores_only_an_unpaid_quote_and_retries() -> None: session, 15_000, "quote-2", wallet.settings.primary_mint, "sat" ) raw_send.assert_awaited_once() + _, settle = _mock_lightning_payout_history + settle.assert_any_await(session, "quote-1", status="failed", amount_sats=None) @pytest.mark.asyncio @@ -579,7 +604,7 @@ async def test_fee_payout_keeps_checkpoint_when_send_outcome_is_unknown() -> Non async def send(*_args: object, **kwargs: object) -> int: checkpoint_quote = kwargs["on_melt_quote"] - await checkpoint_quote("quote-1") # type: ignore[operator] + await checkpoint_quote("quote-1", "lnbc1payout") # type: ignore[operator] raise TimeoutError("unknown outcome") with ( @@ -620,7 +645,7 @@ async def test_fee_payout_cancellation_during_send_alerts_and_propagates() -> No async def cancel_send(*_args: object, **kwargs: object) -> int: checkpoint_quote = kwargs["on_melt_quote"] - await checkpoint_quote("quote-1") # type: ignore[operator] + await checkpoint_quote("quote-1", "lnbc1payout") # type: ignore[operator] raise asyncio.CancelledError with ( @@ -666,6 +691,8 @@ async def test_fee_payout_completion_failures_use_sent_checkpoint_alert( side_effect=[ _session_context(session), _session_context(session), + _session_context(session), + RuntimeError("pool unavailable"), RuntimeError("pool unavailable"), ] ) @@ -675,7 +702,7 @@ async def test_fee_payout_completion_failures_use_sent_checkpoint_alert( async def send(*_args: object, **kwargs: object) -> int: checkpoint_quote = kwargs["on_melt_quote"] - await checkpoint_quote("quote-1") # type: ignore[operator] + await checkpoint_quote("quote-1", "lnbc1payout") # type: ignore[operator] return 5 with ( @@ -726,7 +753,7 @@ async def test_fee_payout_releases_db_connection_during_send(tmp_path: object) - async def send(*_args: object, **kwargs: object) -> int: assert engine.pool.checkedout() == 0 # type: ignore[attr-defined] checkpoint_quote = kwargs["on_melt_quote"] - await checkpoint_quote("quote-1") # type: ignore[operator] + await checkpoint_quote("quote-1", "lnbc1payout") # type: ignore[operator] assert engine.pool.checkedout() == 0 # type: ignore[attr-defined] return 5 diff --git a/tests/unit/test_lightning_settlement.py b/tests/unit/test_lightning_settlement.py index 8d7c96a9..205f3db7 100644 --- a/tests/unit/test_lightning_settlement.py +++ b/tests/unit/test_lightning_settlement.py @@ -37,6 +37,7 @@ def _invoice(**overrides: object) -> SimpleNamespace: "payment_hash": "quote-1", "amount_sats": 100, "purpose": "create", + "direction": "in", "status": "pending", "paid_at": None, "api_key_hash": None, diff --git a/tests/unit/test_lnurl_amount_and_destination.py b/tests/unit/test_lnurl_amount_and_destination.py index 72986500..3a4f033e 100644 --- a/tests/unit/test_lnurl_amount_and_destination.py +++ b/tests/unit/test_lnurl_amount_and_destination.py @@ -192,7 +192,7 @@ async def test_raw_send_to_lnurl_requotes_for_exact_input_fees_without_recursion assert paid == 485_000 assert wallet.melt_quote.await_count == 2 - checkpoint.assert_awaited_once_with("q2") + checkpoint.assert_awaited_once_with("q2", "lnbc1...") wallet.select_to_send.assert_not_called() selected = wallet.melt.await_args.kwargs["proofs"] assert sum(proof.amount for proof in selected) == 500 @@ -358,9 +358,7 @@ def _patch_getaddrinfo(ip: str) -> Any: loop = MagicMock() loop.getaddrinfo = fake_getaddrinfo - return patch.object( - lnurl_module.asyncio, "get_running_loop", return_value=loop - ) + return patch.object(lnurl_module.asyncio, "get_running_loop", return_value=loop) @pytest.mark.asyncio diff --git a/tests/unit/test_lnurl_melt_timeout.py b/tests/unit/test_lnurl_melt_timeout.py index 9c2fbdb1..4c9ea306 100644 --- a/tests/unit/test_lnurl_melt_timeout.py +++ b/tests/unit/test_lnurl_melt_timeout.py @@ -346,8 +346,9 @@ async def test_raw_send_to_lnurl_checkpoints_quote_before_melt_dispatch() -> Non wallet, proofs = _wallet() events: list[str] = [] - async def checkpoint(quote_id: str) -> None: + async def checkpoint(quote_id: str, bolt11: str) -> None: assert quote_id == "q" + assert bolt11 == "lnbc1..." events.append("checkpoint") async def melt(**_kwargs: object) -> MagicMock: diff --git a/tests/unit/test_periodic_payout.py b/tests/unit/test_periodic_payout.py index 44fdfe11..03871a82 100644 --- a/tests/unit/test_periodic_payout.py +++ b/tests/unit/test_periodic_payout.py @@ -14,11 +14,16 @@ from collections.abc import Callable, Coroutine from contextlib import asynccontextmanager from pathlib import Path from typing import Any -from unittest.mock import AsyncMock, MagicMock, patch +from unittest.mock import ANY, AsyncMock, MagicMock, patch import pytest -from routstr.wallet import _payout_units, periodic_payout +from routstr.payment.lnurl import MeltOutcomeAmbiguousError, MeltUnpaidError +from routstr.wallet import ( + _payout_units, + _reconcile_stale_payout_history, + periodic_payout, +) @pytest.fixture(autouse=True) @@ -62,11 +67,20 @@ def _one_cycle_sleep() -> Callable[[float], Coroutine[Any, Any, None]]: @pytest.mark.asyncio async def test_periodic_payout_includes_primary_mint_not_in_cashu_mints() -> None: - """primary_mint absent from cashu_mints is still paid out.""" + """primary_mint absent from cashu_mints is paid out and recorded.""" from routstr.core.settings import settings get_wallet = AsyncMock(return_value=MagicMock()) - raw_send = AsyncMock(return_value=1000) + record_payout = AsyncMock() + settle_payout = AsyncMock() + + async def send(*args: object, **kwargs: object) -> int: + await kwargs["on_melt_quote"]( # type: ignore[index,operator] + "quote-1", "lnbc1payout" + ) + return 1_000_000 + + raw_send = AsyncMock(side_effect=send) with ( patch.object(settings, "cashu_mints", []), @@ -93,6 +107,8 @@ async def test_periodic_payout_includes_primary_mint_not_in_cashu_mints() -> Non "routstr.wallet.db.total_user_liability", AsyncMock(return_value=0), ), + patch("routstr.wallet.db.record_lightning_payout", record_payout), + patch("routstr.wallet.db.settle_lightning_payout", settle_payout), patch("routstr.wallet.raw_send_to_lnurl", raw_send), ): with pytest.raises(_LoopBreak): @@ -101,6 +117,20 @@ async def test_periodic_payout_includes_primary_mint_not_in_cashu_mints() -> Non processed = {call.args[0] for call in get_wallet.await_args_list} assert processed == {"http://primary:3338"} assert raw_send.await_count >= 1 + record_payout.assert_awaited_once_with( + ANY, + quote_id="quote-1", + bolt11="lnbc1payout", + amount_sats=100_000, + mint_url="http://primary:3338", + destination="owner@ln.tld", + ) + settle_payout.assert_awaited_once_with( + ANY, + "quote-1", + status="paid", + amount_sats=1_000, + ) @pytest.mark.asyncio @@ -247,11 +277,12 @@ async def test_periodic_payout_handles_session_creation_failure() -> None: with pytest.raises(_LoopBreak): await periodic_payout() - # The liability session is opened per mint/unit (sat + msat), and each - # DB failure retains the cycle-specific alert wording while remaining - # isolated to its own iteration. - assert create_session.call_count == 2 - assert logger.error.call_count == 2 + # Per mint/unit (sat + msat) a session is opened twice: once by the stale + # payout-history sweep and once for the liability read. Each DB failure is + # logged and isolated to its own step; the liability error keeps the + # cycle-specific alert wording. + assert create_session.call_count == 4 + assert logger.error.call_count == 4 message = logger.error.call_args.args[0] extra = logger.error.call_args.kwargs["extra"] assert message == "Error in periodic payout cycle: RuntimeError" @@ -307,3 +338,224 @@ async def test_periodic_payout_caps_amount_at_max_payout_sat() -> None: assert raw_send.await_count >= 1 assert raw_send.await_args_list[0].kwargs["amount"] == 250_000 + + +@pytest.mark.asyncio +async def test_payout_history_records_the_capped_amount() -> None: + """History stores what is actually sent, not the uncapped balance.""" + from routstr.core.settings import settings + + record_payout = AsyncMock() + settle_payout = AsyncMock() + + async def send(*args: object, **kwargs: object) -> int: + await kwargs["on_melt_quote"]( # type: ignore[index,operator] + "quote-capped", "lnbc1capped" + ) + return 250_000_000 + + raw_send = AsyncMock(side_effect=send) + + with ( + patch.object(settings, "cashu_mints", ["http://mint:3338"]), + patch.object(settings, "primary_mint", "http://mint:3338"), + patch.object(settings, "receive_ln_address", "owner@ln.tld"), + patch.object(settings, "payout_interval_seconds", _INTERVAL), + patch.object(settings, "min_payout_sat", 10), + patch.object(settings, "max_payout_sat", 250_000), + patch("routstr.wallet.asyncio.sleep", _one_cycle_sleep()), + patch("routstr.wallet.db.create_session", _fake_session), + patch( + "routstr.wallet._get_supported_mint_units", + AsyncMock(return_value=["sat"]), + ), + patch("routstr.wallet.get_wallet", AsyncMock(return_value=MagicMock())), + patch( + "routstr.wallet.get_proofs_per_mint_and_unit", + MagicMock(return_value=[MagicMock(amount=1_000_000)]), + ), + patch( + "routstr.wallet.slow_filter_spend_proofs", + AsyncMock(side_effect=lambda proofs, wallet: proofs), + ), + patch("routstr.wallet.db.total_user_liability", AsyncMock(return_value=0)), + patch( + "routstr.wallet.db.list_unsettled_lightning_payouts", + AsyncMock(return_value=[]), + ), + patch("routstr.wallet.db.record_lightning_payout", record_payout), + patch("routstr.wallet.db.settle_lightning_payout", settle_payout), + patch("routstr.wallet.raw_send_to_lnurl", raw_send), + ): + with pytest.raises(_LoopBreak): + await periodic_payout() + + record_payout.assert_awaited_once_with( + ANY, + quote_id="quote-capped", + bolt11="lnbc1capped", + amount_sats=250_000, + mint_url="http://mint:3338", + destination="owner@ln.tld", + ) + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + ("error", "expected_status"), + [ + (MeltUnpaidError("mint confirmed unpaid"), "failed"), + (MeltOutcomeAmbiguousError("outcome unknown"), "reconciliation_required"), + (RuntimeError("HTTP 500 after dispatch"), "reconciliation_required"), + ], +) +async def test_payout_history_marks_failed_only_on_proven_non_payment( + error: Exception, expected_status: str +) -> None: + """Only a mint-confirmed unpaid melt is recorded as failed.""" + from routstr.core.settings import settings + + settle_payout = AsyncMock() + + async def send(*args: object, **kwargs: object) -> int: + await kwargs["on_melt_quote"]( # type: ignore[index,operator] + "quote-err", "lnbc1err" + ) + raise error + + with ( + patch.object(settings, "cashu_mints", ["http://mint:3338"]), + patch.object(settings, "primary_mint", "http://mint:3338"), + patch.object(settings, "receive_ln_address", "owner@ln.tld"), + patch.object(settings, "payout_interval_seconds", _INTERVAL), + patch.object(settings, "min_payout_sat", 10), + patch.object(settings, "max_payout_sat", 250_000), + patch("routstr.wallet.asyncio.sleep", _one_cycle_sleep()), + patch("routstr.wallet.db.create_session", _fake_session), + patch( + "routstr.wallet._get_supported_mint_units", + AsyncMock(return_value=["sat"]), + ), + patch("routstr.wallet.get_wallet", AsyncMock(return_value=MagicMock())), + patch( + "routstr.wallet.get_proofs_per_mint_and_unit", + MagicMock(return_value=[MagicMock(amount=1_000_000)]), + ), + patch( + "routstr.wallet.slow_filter_spend_proofs", + AsyncMock(side_effect=lambda proofs, wallet: proofs), + ), + patch("routstr.wallet.db.total_user_liability", AsyncMock(return_value=0)), + patch( + "routstr.wallet.db.list_unsettled_lightning_payouts", + AsyncMock(return_value=[]), + ), + patch("routstr.wallet.db.record_lightning_payout", AsyncMock()), + patch("routstr.wallet.db.settle_lightning_payout", settle_payout), + patch("routstr.wallet.raw_send_to_lnurl", AsyncMock(side_effect=send)), + ): + with pytest.raises(_LoopBreak): + await periodic_payout() + + settle_payout.assert_awaited_once_with( + ANY, "quote-err", status=expected_status, amount_sats=None + ) + + +@pytest.mark.asyncio +async def test_payout_history_write_failure_does_not_block_payout() -> None: + """A failing history insert is logged; the melt and settlement still run.""" + from routstr.core.settings import settings + + get_wallet = AsyncMock(return_value=MagicMock()) + record_payout = AsyncMock(side_effect=RuntimeError("database is locked")) + settle_payout = AsyncMock() + logger = MagicMock() + + async def send(*args: object, **kwargs: object) -> int: + await kwargs["on_melt_quote"]( # type: ignore[index,operator] + "quote-1", "lnbc1payout" + ) + return 1_000_000 + + raw_send = AsyncMock(side_effect=send) + + with ( + patch.object(settings, "cashu_mints", []), + patch.object(settings, "primary_mint", "http://primary:3338"), + patch.object(settings, "receive_ln_address", "owner@ln.tld"), + patch.object(settings, "payout_interval_seconds", _INTERVAL), + patch.object(settings, "min_payout_sat", 10), + patch.object(settings, "max_payout_sat", 250_000), + patch("routstr.wallet.asyncio.sleep", _one_cycle_sleep()), + patch("routstr.wallet.db.create_session", _fake_session), + patch( + "routstr.wallet._get_supported_mint_units", + AsyncMock(return_value=["sat"]), + ), + patch("routstr.wallet.get_wallet", get_wallet), + patch( + "routstr.wallet.get_proofs_per_mint_and_unit", + MagicMock(return_value=[MagicMock(amount=100_000)]), + ), + patch( + "routstr.wallet.slow_filter_spend_proofs", + AsyncMock(side_effect=lambda proofs, wallet: proofs), + ), + patch("routstr.wallet.db.total_user_liability", AsyncMock(return_value=0)), + patch( + "routstr.wallet.db.list_unsettled_lightning_payouts", + AsyncMock(return_value=[]), + ), + patch("routstr.wallet.db.record_lightning_payout", record_payout), + patch("routstr.wallet.db.settle_lightning_payout", settle_payout), + patch("routstr.wallet.raw_send_to_lnurl", raw_send), + patch("routstr.wallet.logger", logger), + ): + with pytest.raises(_LoopBreak): + await periodic_payout() + + record_payout.assert_awaited_once() + assert raw_send.await_count == 1 + settle_payout.assert_awaited_once_with( + ANY, "quote-1", status="paid", amount_sats=1_000 + ) + messages = [call.args[0] for call in logger.error.call_args_list] + assert "Failed to record Lightning payout history" in messages + + +@pytest.mark.asyncio +async def test_stale_payout_history_is_reconciled_from_mint_state() -> None: + """Stale out-rows follow the mint's verdict; pending/unknown are left alone.""" + stale = [ + MagicMock(payment_hash="q-paid"), + MagicMock(payment_hash="q-unpaid"), + MagicMock(payment_hash="q-pending"), + MagicMock(payment_hash="q-unknown"), + ] + states = { + "q-paid": "paid", + "q-unpaid": "unpaid", + "q-pending": "pending", + "q-unknown": "unknown", + } + settle_payout = AsyncMock() + + async def _state(_mint: str, _unit: str, quote_id: str) -> str: + return states[quote_id] + + with ( + patch("routstr.wallet.db.create_session", _fake_session), + patch( + "routstr.wallet.db.list_unsettled_lightning_payouts", + AsyncMock(return_value=stale), + ), + patch("routstr.wallet._check_bolt11_payment_status_locked", _state), + patch("routstr.wallet.db.settle_lightning_payout", settle_payout), + ): + await _reconcile_stale_payout_history("http://mint:3338", "sat") + + assert settle_payout.await_args_list == [ + ((ANY, "q-paid"), {"status": "paid", "amount_sats": None}), + ((ANY, "q-unpaid"), {"status": "failed", "amount_sats": None}), + ] diff --git a/ui/app/transactions/page.tsx b/ui/app/transactions/page.tsx index 56bc1b74..0e3eca85 100644 --- a/ui/app/transactions/page.tsx +++ b/ui/app/transactions/page.tsx @@ -264,7 +264,8 @@ function LightningInvoiceTable({ No invoices found - Lightning invoices created via /lightning/invoice will show here. + Lightning invoices created via /lightning/invoice and payouts sent + to your Lightning address will show here. @@ -290,6 +291,33 @@ function LightningInvoiceTable({ Expired ); + if (status === 'failed') + return ( + + Failed + + ); + if (status === 'settlement_pending') + return ( + + Settling + + ); + if (status === 'reconciliation_required') + return ( + + Reconciling + + ); if (status === 'cancelled') return ( + Direction Purpose Amount Status @@ -328,6 +357,11 @@ function LightningInvoiceTable({ {invoices.map((inv) => ( + + + {inv.direction === 'out' ? 'Sent' : 'Received'} + + {inv.purpose} @@ -514,15 +548,27 @@ export default function TransactionsPage() { placeholderData: keepPreviousData, }); - const LIGHTNING_STATUSES = ['pending', 'paid', 'expired', 'cancelled']; + const LIGHTNING_STATUSES = [ + 'pending', + 'settlement_pending', + 'paid', + 'failed', + 'expired', + 'cancelled', + 'reconciliation_required', + ]; const lightningStatusParam = LIGHTNING_STATUSES.includes(status) ? status : undefined; + const lightningDirectionParam = ['in', 'out'].includes(type) + ? type + : undefined; const lightningQuery = useQuery({ queryKey: [ 'lightning-invoices', lightningStatusParam, + lightningDirectionParam, searchParam, lightningPage, ], @@ -530,6 +576,7 @@ export default function TransactionsPage() { AdminService.getLightningInvoices( lightningStatusParam, undefined, + lightningDirectionParam, searchParam, PAGE_SIZE, lightningPage * PAGE_SIZE @@ -710,7 +757,9 @@ export default function TransactionsPage() { All Types Incoming (Payments) - Outgoing (Refunds) + + Outgoing (Refunds & Payouts) + @@ -727,6 +776,13 @@ export default function TransactionsPage() { Collected Swept Paid (Lightning) + + Settling (Lightning) + + Failed (Lightning) + + Reconciling (Lightning) + Expired (Lightning) Cancelled (Lightning) diff --git a/ui/lib/api/services/admin.ts b/ui/lib/api/services/admin.ts index f21e7b64..e510a571 100644 --- a/ui/lib/api/services/admin.ts +++ b/ui/lib/api/services/admin.ts @@ -938,6 +938,7 @@ export class AdminService { static async getLightningInvoices( status?: string, purpose?: string, + direction?: string, search?: string, limit: number = 50, offset: number = 0 @@ -945,6 +946,7 @@ export class AdminService { const params = new URLSearchParams(); if (status) params.append('status', status); if (purpose) params.append('purpose', purpose); + if (direction) params.append('direction', direction); if (search) params.append('search', search); params.append('limit', limit.toString()); params.append('offset', offset.toString()); @@ -1304,9 +1306,17 @@ export interface LightningInvoice { amount_sats: number; description: string; payment_hash: string; - status: 'pending' | 'paid' | 'expired' | 'cancelled'; + status: + | 'pending' + | 'settlement_pending' + | 'paid' + | 'failed' + | 'expired' + | 'cancelled' + | 'reconciliation_required'; api_key_hash: string | null; - purpose: 'create' | 'topup'; + direction: 'in' | 'out'; + purpose: 'create' | 'topup' | 'payout'; created_at: number; expires_at: number; paid_at: number | null;