Compare commits
2 commits
c50455d5f6
...
44e10caac7
| Author | SHA1 | Date | |
|---|---|---|---|
| 44e10caac7 | |||
| 4fdb358bb0 |
6 changed files with 523 additions and 66 deletions
70
crud.py
70
crud.py
|
|
@ -1696,3 +1696,73 @@ 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
|
||||
|
|
|
|||
|
|
@ -624,3 +624,30 @@ 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}
|
||||
);
|
||||
"""
|
||||
)
|
||||
|
|
|
|||
40
tasks.py
40
tasks.py
|
|
@ -179,12 +179,31 @@ 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()
|
||||
await on_invoice_paid(payment)
|
||||
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}"
|
||||
)
|
||||
|
||||
|
||||
async def on_invoice_paid(payment: Payment) -> None:
|
||||
|
|
@ -210,8 +229,20 @@ 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
|
||||
|
|
@ -245,6 +276,7 @@ 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
|
||||
|
|
@ -276,6 +308,7 @@ 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
|
||||
|
|
@ -324,6 +357,11 @@ 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
|
||||
|
|
|
|||
|
|
@ -55,6 +55,7 @@ EXPECTED_TABLES = [
|
|||
"roles",
|
||||
"role_permissions",
|
||||
"user_roles",
|
||||
"processed_payments",
|
||||
]
|
||||
|
||||
|
||||
|
|
|
|||
281
tests/test_payment_idempotency.py
Normal file
281
tests/test_payment_idempotency.py
Normal file
|
|
@ -0,0 +1,281 @@
|
|||
"""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
|
||||
176
views_api.py
176
views_api.py
|
|
@ -1850,90 +1850,130 @@ async def api_record_payment(
|
|||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=5.0) as client:
|
||||
# Get recent entries from Fava's journal endpoint
|
||||
# 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.
|
||||
response = await client.get(
|
||||
f"{fava.base_url}/api/journal",
|
||||
f"{fava.base_url}/journal",
|
||||
params={"time": ""} # Get all entries
|
||||
)
|
||||
response.raise_for_status()
|
||||
response_data = response.json()
|
||||
entries = response_data.get('entries', [])
|
||||
|
||||
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:
|
||||
# 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.
|
||||
logger.warning(f"Could not check Fava for duplicate payment: {e}")
|
||||
# Continue anyway - Fava/Beancount will catch duplicate if it exists
|
||||
|
||||
# Convert amount from millisatoshis to satoshis
|
||||
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))
|
||||
|
||||
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 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"
|
||||
status_code=HTTPStatus.SERVICE_UNAVAILABLE,
|
||||
detail="Cannot verify payment duplicate status; try again shortly",
|
||||
)
|
||||
|
||||
# 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
|
||||
# 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,
|
||||
)
|
||||
|
||||
logger.info(f"Formatted payment entry: {entry}")
|
||||
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",
|
||||
)
|
||||
|
||||
# Submit to Fava
|
||||
result = await fava.add_entry(entry)
|
||||
logger.info(f"Payment entry submitted to Fava: {result.get('data', 'Unknown')}")
|
||||
# Convert amount from millisatoshis to satoshis
|
||||
try:
|
||||
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))
|
||||
|
||||
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 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")]
|
||||
|
||||
# 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}")
|
||||
|
||||
# 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
|
||||
|
||||
# Get updated balance from Fava
|
||||
balance_data = await fava.get_user_balance_bql(target_user_id)
|
||||
|
||||
return {
|
||||
"journal_entry_id": f"fava-{datetime.now().timestamp()}",
|
||||
"journal_entry_id": entry_id,
|
||||
"new_balance": balance_data["balance"],
|
||||
"message": "Payment recorded successfully",
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue