diff --git a/crud.py b/crud.py index b49bb9b..0692806 100644 --- a/crud.py +++ b/crud.py @@ -1696,73 +1696,3 @@ async def check_user_has_role_permission( return True return False - - -# ============================================================================= -# PROCESSED PAYMENTS (Lightning payment idempotency gate) -# ============================================================================= -# The Fava-side duplicate checks are read-then-write races; this table's -# primary key on payment_hash makes exactly one claimant win. Shared by the -# background invoice listener (tasks.on_invoice_paid) and the client-driven -# /record-payment endpoint. - - -async def claim_payment(payment_hash: str) -> bool: - """Atomically claim a Lightning payment for recording. - - Returns True when this caller owns the claim; False when the payment - is already recorded or another coroutine is recording it right now. - """ - result = await db.execute( - """ - INSERT INTO processed_payments (payment_hash, status) - VALUES (:payment_hash, 'processing') - ON CONFLICT (payment_hash) DO NOTHING - """, - {"payment_hash": payment_hash}, - ) - return result.rowcount == 1 - - -async def get_processed_payment(payment_hash: str) -> Optional[dict]: - row = await db.fetchone( - "SELECT payment_hash, status, entry_id FROM processed_payments" - " WHERE payment_hash = :payment_hash", - {"payment_hash": payment_hash}, - ) - return dict(row) if row else None - - -async def mark_payment_done(payment_hash: str, entry_id: Optional[str] = None) -> None: - await db.execute( - """ - UPDATE processed_payments SET status = 'done', entry_id = :entry_id - WHERE payment_hash = :payment_hash - """, - {"payment_hash": payment_hash, "entry_id": entry_id}, - ) - - -async def release_payment_claim(payment_hash: str) -> None: - """Compensating delete after a failed recording, so redelivery retries. - - Only removes an in-flight claim — a 'done' row is permanent. - """ - await db.execute( - "DELETE FROM processed_payments" - " WHERE payment_hash = :payment_hash AND status = 'processing'", - {"payment_hash": payment_hash}, - ) - - -async def clear_stale_payment_claims() -> int: - """Drop 'processing' claims left behind by a previous process life. - - A live claim only exists inside a running coroutine, so anything - still 'processing' at listener startup belongs to a crashed or - restarted process and would otherwise block that payment forever. - """ - result = await db.execute( - "DELETE FROM processed_payments WHERE status = 'processing'" - ) - return result.rowcount diff --git a/migrations.py b/migrations.py index d8a4bba..2c39507 100644 --- a/migrations.py +++ b/migrations.py @@ -624,30 +624,3 @@ async def m004_add_rbac_tables(db): "created_by": "system", # System-created default roles }, ) - - -async def m005_add_processed_payments(db): - """ - Local idempotency gate for Lightning payment recording. - - The Fava-side duplicate check (`add_entry_idempotent`, journal-link - scan) is a read-then-write race: the background invoice listener and - the client-driven /record-payment endpoint can both pass the "not - present" check for the same payment_hash and both insert. The - primary key on payment_hash makes exactly one claimant win. - - status lifecycle: 'processing' (claimed, write in flight) → 'done' - (entry recorded). Failed claims are deleted so redelivery retries; - 'processing' rows from a crashed process are cleared at listener - startup. - """ - await db.execute( - f""" - CREATE TABLE IF NOT EXISTS processed_payments ( - payment_hash TEXT PRIMARY KEY, - status TEXT NOT NULL DEFAULT 'processing', - entry_id TEXT, - created_at TIMESTAMP NOT NULL DEFAULT {db.timestamp_now} - ); - """ - ) diff --git a/tasks.py b/tasks.py index 158d913..8ed5a33 100644 --- a/tasks.py +++ b/tasks.py @@ -179,31 +179,12 @@ async def wait_for_paid_invoices(): This ensures payments are recorded even if the user closes their browser before the payment is detected by client-side polling. """ - from .crud import clear_stale_payment_claims - invoice_queue = Queue() register_invoice_listener(invoice_queue, "ext_libra") - # Claims from a previous process life can't be live anymore — clear - # them so those payments aren't blocked forever. - cleared = await clear_stale_payment_claims() - if cleared: - logger.warning( - f"[LIBRA] Cleared {cleared} stale in-flight payment claim(s) " - "from a previous run" - ) - while True: payment = await invoice_queue.get() - try: - await on_invoice_paid(payment) - except Exception: - # One bad payment must not kill the listener for the rest of - # the process lifetime; its claim was released, so redelivery - # can retry it. - logger.exception( - f"[LIBRA] Failed to record payment {payment.payment_hash}" - ) + await on_invoice_paid(payment) async def on_invoice_paid(payment: Payment) -> None: @@ -229,20 +210,8 @@ async def on_invoice_paid(payment: Payment) -> None: logger.warning(f"Libra invoice {payment.payment_hash} missing user_id in metadata") return - from .crud import claim_payment, mark_payment_done, release_payment_claim from .fava_client import get_fava_client - # Local idempotency gate: exactly one claimant (this listener or the - # /record-payment endpoint) gets to record a given payment_hash. The - # Fava-side idempotent write below stays as a second layer for - # entries recorded before this table existed. - if not await claim_payment(payment.payment_hash): - logger.info( - f"Payment {payment.payment_hash} already recorded or being " - "recorded; skipping" - ) - return - fava = get_fava_client() # Use idempotency key based on payment hash - this ensures duplicate @@ -276,7 +245,6 @@ async def on_invoice_paid(payment: Payment) -> None: if not fiat_currency or not fiat_amount: logger.error(f"Payment {payment.payment_hash} missing fiat currency/amount metadata") - await release_payment_claim(payment.payment_hash) return # Get user's current balance to determine receivables and payables @@ -308,7 +276,6 @@ async def on_invoice_paid(payment: Payment) -> None: lightning_account = await get_account_by_name("Assets:Bitcoin:Lightning") if not lightning_account: logger.error("Lightning account 'Assets:Bitcoin:Lightning' not found") - await release_payment_claim(payment.payment_hash) return # Query for unsettled entries to link this settlement back to them @@ -357,11 +324,6 @@ async def on_invoice_paid(payment: Payment) -> None: f"{result.get('data', 'Unknown')}" ) - await mark_payment_done(payment.payment_hash, idempotency_key) - except Exception as e: logger.error(f"Error recording Libra payment {payment.payment_hash}: {e}") - # Release the claim so redelivery (or the /record-payment - # endpoint) can retry this payment. - await release_payment_claim(payment.payment_hash) raise diff --git a/tests/test_migrations.py b/tests/test_migrations.py index 1586473..b66b595 100644 --- a/tests/test_migrations.py +++ b/tests/test_migrations.py @@ -55,7 +55,6 @@ EXPECTED_TABLES = [ "roles", "role_permissions", "user_roles", - "processed_payments", ] diff --git a/tests/test_payment_idempotency.py b/tests/test_payment_idempotency.py deleted file mode 100644 index 969107b..0000000 --- a/tests/test_payment_idempotency.py +++ /dev/null @@ -1,281 +0,0 @@ -"""Lightning payment idempotency — the `processed_payments` claim gate. - -The background invoice listener (`tasks.on_invoice_paid`) and the -client-driven `POST /record-payment` endpoint can both fire for the -same `payment_hash` (queue redelivery after restart, webhook + poller). -The Fava-side duplicate checks are read-then-write races; the local -`processed_payments` primary key makes exactly one claimant win. - -These tests bypass invoice generation (blocked by libra/issues/40) by -delivering synthetic paid `Payment` objects straight to -`on_invoice_paid` and by inserting paid payment rows via the LNbits -core crud for the endpoint tests. -""" -import asyncio -import importlib -from uuid import uuid4 - -import pytest - -from lnbits.core.crud.payments import create_payment -from lnbits.core.models.payments import CreatePayment, Payment, PaymentState - -from .helpers import list_user_entries, post_receivable - -pytestmark = pytest.mark.anyio - - -def _module(name: str): - """Import a libra submodule under whichever path the active LNbits layout - uses (default `lnbits.extensions.libra` or bare `libra`).""" - for prefix in ("lnbits.extensions.libra", "libra"): - try: - return importlib.import_module(f"{prefix}.{name}") - except ModuleNotFoundError: - continue - raise ModuleNotFoundError(f"libra.{name}: tried both import paths") - - -tasks = _module("tasks") -libra_crud = _module("crud") - - -def _paid_payment( - wallet_id: str, - user_id: str, - *, - fiat_amount: str = "100.00", - fiat_currency: str = "EUR", - sats: int = 100_000, -) -> Payment: - payment_hash = uuid4().hex + uuid4().hex[:32] - return Payment( - checking_id=payment_hash, - payment_hash=payment_hash, - wallet_id=wallet_id, - amount=sats * 1000, - fee=0, - bolt11="lnbcfake", - status=PaymentState.SUCCESS, - extra={ - "tag": "libra", - "user_id": user_id, - "fiat_currency": fiat_currency, - "fiat_amount": fiat_amount, - }, - ) - - -async def _setup_receivable( - client, super_user_headers, configured_user, standard_accounts, - amount: str = "100.00", -): - user, wallet = configured_user - await post_receivable( - client, - super_user_headers=super_user_headers, - user_id=user.id, - amount=amount, - description=f"Idempotency setup {uuid4().hex[:6]}", - revenue_account=standard_accounts["revenue_rent"]["name"], - ) - # Force a Fava reload before downstream balance reads (see #37). - await list_user_entries(client, wallet_inkey=wallet.inkey) - return user, wallet - - -async def _entries_with_link(client, wallet_inkey: str, link: str) -> list: - payload = await list_user_entries(client, wallet_inkey=wallet_inkey) - return [ - e for e in payload["entries"] if link in (e.get("links") or []) - ] - - -# --------------------------------------------------------------------------- -# on_invoice_paid — the background listener path -# --------------------------------------------------------------------------- - - -async def test_double_delivery_records_exactly_once( - client, super_user_headers, configured_user, standard_accounts -): - """Same payment delivered twice (queue redelivery) → one ledger entry.""" - user, wallet = await _setup_receivable( - client, super_user_headers, configured_user, standard_accounts - ) - payment = _paid_payment(wallet.id, user.id) - - await tasks.on_invoice_paid(payment) - await tasks.on_invoice_paid(payment) - - link = f"ln-{payment.payment_hash[:16]}" - assert len(await _entries_with_link(client, wallet.inkey, link)) == 1 - - row = await libra_crud.get_processed_payment(payment.payment_hash) - assert row is not None and row["status"] == "done" - - -async def test_failed_recording_releases_claim_and_retry_succeeds( - client, super_user_headers, configured_user, standard_accounts, monkeypatch -): - """A Fava failure mid-write must not permanently block the payment.""" - user, wallet = await _setup_receivable( - client, super_user_headers, configured_user, standard_accounts - ) - payment = _paid_payment(wallet.id, user.id) - - fava_client = _module("fava_client") - fava = fava_client.get_fava_client() - - async def _boom(*args, **kwargs): - raise RuntimeError("fava down") - - monkeypatch.setattr(fava, "add_entry_idempotent", _boom) - with pytest.raises(RuntimeError): - await tasks.on_invoice_paid(payment) - monkeypatch.undo() - - # Claim released → nothing recorded, retry allowed. - assert await libra_crud.get_processed_payment(payment.payment_hash) is None - - await tasks.on_invoice_paid(payment) - row = await libra_crud.get_processed_payment(payment.payment_hash) - assert row is not None and row["status"] == "done" - link = f"ln-{payment.payment_hash[:16]}" - assert len(await _entries_with_link(client, wallet.inkey, link)) == 1 - - -async def test_listener_survives_poison_payment_and_clears_stale_claims( - client, super_user_headers, configured_user, standard_accounts, monkeypatch -): - """One bad payment must not kill the listener; stale 'processing' - claims from a previous process life are cleared at startup.""" - user, wallet = await _setup_receivable( - client, super_user_headers, configured_user, standard_accounts - ) - - # A claim left behind by a "crashed" previous run. - stale_hash = uuid4().hex + uuid4().hex[:32] - assert await libra_crud.claim_payment(stale_hash) - - captured: dict = {} - monkeypatch.setattr( - tasks, - "register_invoice_listener", - lambda queue, name: captured.update(queue=queue), - ) - - listener = asyncio.create_task(tasks.wait_for_paid_invoices()) - try: - for _ in range(50): - if "queue" in captured: - break - await asyncio.sleep(0.05) - assert "queue" in captured, "listener never registered its queue" - - poison = _paid_payment(wallet.id, user.id, fiat_amount="not-a-number") - good = _paid_payment(wallet.id, user.id) - captured["queue"].put_nowait(poison) - captured["queue"].put_nowait(good) - - row = None - for _ in range(100): - row = await libra_crud.get_processed_payment(good.payment_hash) - if row and row["status"] == "done": - break - await asyncio.sleep(0.1) - assert row is not None and row["status"] == "done", ( - "good payment was not recorded after the poison payment" - ) - finally: - listener.cancel() - - # Startup cleared the stale claim; the poison payment's claim was - # released on failure so redelivery could retry it. - assert await libra_crud.get_processed_payment(stale_hash) is None - assert await libra_crud.get_processed_payment(poison.payment_hash) is None - - -# --------------------------------------------------------------------------- -# POST /record-payment — the client-driven path -# --------------------------------------------------------------------------- - - -async def _insert_paid_payment_row(wallet_id: str, user_id: str) -> str: - payment_hash = uuid4().hex + uuid4().hex[:32] - await create_payment( - checking_id=payment_hash, - data=CreatePayment( - wallet_id=wallet_id, - payment_hash=payment_hash, - bolt11="lnbcfake", - amount_msat=100_000_000, - memo="idempotency test", - extra={ - "tag": "libra", - "user_id": user_id, - "fiat_currency": "EUR", - "fiat_amount": "100.00", - }, - ), - status=PaymentState.SUCCESS, - ) - return payment_hash - - -async def test_record_payment_conflicts_while_in_flight( - client, super_user_headers, configured_user, standard_accounts -): - user, wallet = await _setup_receivable( - client, super_user_headers, configured_user, standard_accounts - ) - payment_hash = await _insert_paid_payment_row(wallet.id, user.id) - - # Another claimant (e.g. the background listener) is mid-recording. - assert await libra_crud.claim_payment(payment_hash) - - r = await client.post( - "/libra/api/v1/record-payment", - headers={"X-Api-Key": wallet.inkey}, - json={"payment_hash": payment_hash}, - ) - assert r.status_code == 409, r.text - - # Once that claimant finishes, a replay reports "already recorded" - # instead of writing a second entry. - await libra_crud.mark_payment_done(payment_hash, f"ln-{payment_hash[:16]}") - r = await client.post( - "/libra/api/v1/record-payment", - headers={"X-Api-Key": wallet.inkey}, - json={"payment_hash": payment_hash}, - ) - assert r.status_code == 200, r.text - assert "already recorded" in r.json()["message"].lower() - - -async def test_record_payment_records_once_then_replays_safely( - client, super_user_headers, configured_user, standard_accounts -): - user, wallet = await _setup_receivable( - client, super_user_headers, configured_user, standard_accounts - ) - payment_hash = await _insert_paid_payment_row(wallet.id, user.id) - - r = await client.post( - "/libra/api/v1/record-payment", - headers={"X-Api-Key": wallet.inkey}, - json={"payment_hash": payment_hash}, - ) - assert r.status_code == 200, r.text - assert r.json()["message"] == "Payment recorded successfully" - - r = await client.post( - "/libra/api/v1/record-payment", - headers={"X-Api-Key": wallet.inkey}, - json={"payment_hash": payment_hash}, - ) - assert r.status_code == 200, r.text - assert "already recorded" in r.json()["message"].lower() - - link = f"ln-{payment_hash[:16]}" - assert len(await _entries_with_link(client, wallet.inkey, link)) == 1 diff --git a/views_api.py b/views_api.py index 4c36427..1b3149a 100644 --- a/views_api.py +++ b/views_api.py @@ -1850,130 +1850,90 @@ async def api_record_payment( try: async with httpx.AsyncClient(timeout=5.0) as client: - # Get recent entries from Fava's journal endpoint. base_url - # already ends in /api — the previous "/api/journal" path - # 404'd, so this duplicate check silently never ran. + # Get recent entries from Fava's journal endpoint response = await client.get( - f"{fava.base_url}/journal", + f"{fava.base_url}/api/journal", params={"time": ""} # Get all entries ) - response.raise_for_status() - response_data = response.json() - entries = response_data.get('entries', []) - # Check if any entry has our payment link - for entry in entries: - entry_links = entry.get('links', []) - if link_to_find in entry_links: - # Payment already recorded, return existing entry - balance_data = await fava.get_user_balance_bql(target_user_id) - return { - "journal_entry_id": f"fava-exists-{data.payment_hash[:16]}", - "new_balance": balance_data["balance"], - "message": "Payment already recorded", - } - except httpx.HTTPError as e: - # Fail CLOSED: if Fava can't confirm the payment isn't already - # recorded, refuse to write — proceeding on a transient blip is - # how double entries happen. The client can simply retry. + if response.status_code == 200: + response_data = response.json() + entries = response_data.get('entries', []) + + # Check if any entry has our payment link + for entry in entries: + entry_links = entry.get('links', []) + if link_to_find in entry_links: + # Payment already recorded, return existing entry + balance_data = await fava.get_user_balance_bql(target_user_id) + return { + "journal_entry_id": f"fava-exists-{data.payment_hash[:16]}", + "new_balance": balance_data["balance"], + "message": "Payment already recorded", + } + except Exception as e: logger.warning(f"Could not check Fava for duplicate payment: {e}") - raise HTTPException( - status_code=HTTPStatus.SERVICE_UNAVAILABLE, - detail="Cannot verify payment duplicate status; try again shortly", - ) - - # Local idempotency gate shared with the background invoice listener - # (tasks.on_invoice_paid): exactly one claimant records a payment_hash. - from .crud import ( - claim_payment, - get_processed_payment, - mark_payment_done, - release_payment_claim, - ) - - if not await claim_payment(data.payment_hash): - existing = await get_processed_payment(data.payment_hash) - if existing and existing["status"] == "done": - balance_data = await fava.get_user_balance_bql(target_user_id) - return { - "journal_entry_id": existing.get("entry_id") - or f"fava-exists-{data.payment_hash[:16]}", - "new_balance": balance_data["balance"], - "message": "Payment already recorded", - } - raise HTTPException( - status_code=HTTPStatus.CONFLICT, - detail="Payment is being recorded; check balance shortly", - ) + # Continue anyway - Fava/Beancount will catch duplicate if it exists # Convert amount from millisatoshis to satoshis - try: - amount_sats = payment.amount // 1000 + amount_sats = payment.amount // 1000 - # Extract fiat metadata from invoice (if present) - fiat_currency = None - fiat_amount = None - if payment.extra and isinstance(payment.extra, dict): - logger.info(f"Payment.extra contents: {payment.extra}") - fiat_currency = payment.extra.get("fiat_currency") - fiat_amount_str = payment.extra.get("fiat_amount") - if fiat_amount_str: - from decimal import Decimal - fiat_amount = Decimal(str(fiat_amount_str)) + # Extract fiat metadata from invoice (if present) + fiat_currency = None + fiat_amount = None + if payment.extra and isinstance(payment.extra, dict): + logger.info(f"Payment.extra contents: {payment.extra}") + fiat_currency = payment.extra.get("fiat_currency") + fiat_amount_str = payment.extra.get("fiat_amount") + if fiat_amount_str: + from decimal import Decimal + fiat_amount = Decimal(str(fiat_amount_str)) - logger.info(f"Extracted fiat metadata - currency: {fiat_currency}, amount: {fiat_amount}") + logger.info(f"Extracted fiat metadata - currency: {fiat_currency}, amount: {fiat_amount}") - # Get user's receivable account (what user owes) - user_receivable = await get_or_create_user_account( - target_user_id, AccountType.ASSET, "Accounts Receivable" + # Get user's receivable account (what user owes) + user_receivable = await get_or_create_user_account( + target_user_id, AccountType.ASSET, "Accounts Receivable" + ) + + # Get lightning account + lightning_account = await get_account_by_name("Assets:Bitcoin:Lightning") + if not lightning_account: + raise HTTPException( + status_code=HTTPStatus.NOT_FOUND, detail="Lightning account not found" ) - # Get lightning account - lightning_account = await get_account_by_name("Assets:Bitcoin:Lightning") - if not lightning_account: - raise HTTPException( - status_code=HTTPStatus.NOT_FOUND, detail="Lightning account not found" - ) + # Get unsettled receivable entries to link to this settlement + unsettled = await fava.get_unsettled_entries_bql(target_user_id, "receivable") + settled_links = [e["link"] for e in unsettled if e.get("link")] - # Get unsettled receivable entries to link to this settlement - unsettled = await fava.get_unsettled_entries_bql(target_user_id, "receivable") - settled_links = [e["link"] for e in unsettled if e.get("link")] + # Format payment entry and submit to Fava + entry = format_payment_entry( + user_id=target_user_id, + payment_account=lightning_account.name, + payable_or_receivable_account=user_receivable.name, + amount_sats=amount_sats, + description=f"Lightning payment from user {target_user_id[:8]}", + entry_date=datetime.now().date(), + is_payable=False, # User paying libra (receivable settlement) + fiat_currency=fiat_currency, + fiat_amount=fiat_amount, + payment_hash=data.payment_hash, + reference=data.payment_hash, + settled_entry_links=settled_links + ) - # Format payment entry and submit to Fava - entry = format_payment_entry( - user_id=target_user_id, - payment_account=lightning_account.name, - payable_or_receivable_account=user_receivable.name, - amount_sats=amount_sats, - description=f"Lightning payment from user {target_user_id[:8]}", - entry_date=datetime.now().date(), - is_payable=False, # User paying libra (receivable settlement) - fiat_currency=fiat_currency, - fiat_amount=fiat_amount, - payment_hash=data.payment_hash, - reference=data.payment_hash, - settled_entry_links=settled_links - ) + logger.info(f"Formatted payment entry: {entry}") - logger.info(f"Formatted payment entry: {entry}") - - # Submit to Fava - result = await fava.add_entry(entry) - logger.info(f"Payment entry submitted to Fava: {result.get('data', 'Unknown')}") - - entry_id = f"ln-{data.payment_hash[:16]}" - await mark_payment_done(data.payment_hash, entry_id) - except BaseException: - # Release the claim so a retry (client or background listener) - # can record this payment. - await release_payment_claim(data.payment_hash) - raise + # Submit to Fava + result = await fava.add_entry(entry) + logger.info(f"Payment entry submitted to Fava: {result.get('data', 'Unknown')}") # Get updated balance from Fava balance_data = await fava.get_user_balance_bql(target_user_id) return { - "journal_entry_id": entry_id, + "journal_entry_id": f"fava-{datetime.now().timestamp()}", "new_balance": balance_data["balance"], "message": "Payment recorded successfully", }