Merge pull request #758 from Routstr/feat/lightning-payout-history

feat: show Lightning payouts in transaction history
This commit is contained in:
9qeklajc
2026-09-22 23:30:14 +02:00
committed by GitHub
15 changed files with 749 additions and 56 deletions
@@ -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")
+3
View File
@@ -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(
+79 -2
View File
@@ -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.
+12 -3
View File
@@ -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,
)
+2 -2
View File
@@ -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
+170 -11
View File
@@ -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
@@ -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
# ---------------------------------------------------------------------------
+47 -9
View File
@@ -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"
+37 -10
View File
@@ -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
+1
View File
@@ -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,
@@ -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
+2 -1
View File
@@ -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:
+261 -9
View File
@@ -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}),
]
+59 -3
View File
@@ -264,7 +264,8 @@ function LightningInvoiceTable({
</EmptyMedia>
<EmptyTitle>No invoices found</EmptyTitle>
<EmptyDescription>
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.
</EmptyDescription>
</EmptyHeader>
</Empty>
@@ -290,6 +291,33 @@ function LightningInvoiceTable({
Expired
</Badge>
);
if (status === 'failed')
return (
<Badge
variant='outline'
className='border-red-500/20 bg-red-500/10 text-red-500'
>
Failed
</Badge>
);
if (status === 'settlement_pending')
return (
<Badge
variant='outline'
className='border-amber-500/20 bg-amber-500/10 text-amber-500'
>
Settling
</Badge>
);
if (status === 'reconciliation_required')
return (
<Badge
variant='outline'
className='border-amber-500/20 bg-amber-500/10 text-amber-500'
>
Reconciling
</Badge>
);
if (status === 'cancelled')
return (
<Badge
@@ -315,6 +343,7 @@ function LightningInvoiceTable({
<Table>
<TableHeader>
<TableRow>
<TableHead>Direction</TableHead>
<TableHead>Purpose</TableHead>
<TableHead>Amount</TableHead>
<TableHead>Status</TableHead>
@@ -328,6 +357,11 @@ function LightningInvoiceTable({
<TableBody>
{invoices.map((inv) => (
<TableRow key={inv.id}>
<TableCell>
<Badge variant='outline' className='capitalize'>
{inv.direction === 'out' ? 'Sent' : 'Received'}
</Badge>
</TableCell>
<TableCell>
<span className='capitalize'>{inv.purpose}</span>
</TableCell>
@@ -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() {
<SelectContent>
<SelectItem value='all'>All Types</SelectItem>
<SelectItem value='in'>Incoming (Payments)</SelectItem>
<SelectItem value='out'>Outgoing (Refunds)</SelectItem>
<SelectItem value='out'>
Outgoing (Refunds & Payouts)
</SelectItem>
</SelectContent>
</Select>
</div>
@@ -727,6 +776,13 @@ export default function TransactionsPage() {
<SelectItem value='collected'>Collected</SelectItem>
<SelectItem value='swept'>Swept</SelectItem>
<SelectItem value='paid'>Paid (Lightning)</SelectItem>
<SelectItem value='settlement_pending'>
Settling (Lightning)
</SelectItem>
<SelectItem value='failed'>Failed (Lightning)</SelectItem>
<SelectItem value='reconciliation_required'>
Reconciling (Lightning)
</SelectItem>
<SelectItem value='expired'>Expired (Lightning)</SelectItem>
<SelectItem value='cancelled'>
Cancelled (Lightning)
+12 -2
View File
@@ -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;