From 8135b1540497e9f7d2efc8c84d2902d9c4dd44e4 Mon Sep 17 00:00:00 2001 From: Vlad Stan Date: Tue, 30 Jun 2026 14:52:07 +0300 Subject: [PATCH] refactor: move out `dispatch_wasm_invoice_paid ` from tasks --- lnbits/app.py | 2 +- lnbits/core/extensions/events.py | 113 +++++++++++++++++++++++++++++++ lnbits/core/extensions/wasm.py | 1 - lnbits/core/tasks.py | 109 ----------------------------- 4 files changed, 114 insertions(+), 111 deletions(-) create mode 100644 lnbits/core/extensions/events.py diff --git a/lnbits/app.py b/lnbits/app.py index fd80627a8..5eb5b5b07 100644 --- a/lnbits/app.py +++ b/lnbits/app.py @@ -24,6 +24,7 @@ from lnbits.core.crud import ( update_installed_extension_state, ) from lnbits.core.crud.extensions import create_installed_extension +from lnbits.core.extensions.events import dispatch_wasm_invoice_paid from lnbits.core.extensions.loader import ( is_wasm_extension_dir, is_wasm_extension_id, @@ -37,7 +38,6 @@ from lnbits.core.services.payments import check_pending_payments from lnbits.core.tasks import ( audit_queue, collect_exchange_rates_data, - dispatch_wasm_invoice_paid, purge_audit_data, run_by_the_minute_tasks, wait_for_audit_data, diff --git a/lnbits/core/extensions/events.py b/lnbits/core/extensions/events.py new file mode 100644 index 000000000..1dd103a7b --- /dev/null +++ b/lnbits/core/extensions/events.py @@ -0,0 +1,113 @@ +import json +from typing import Any + +from loguru import logger + +from lnbits.core.db import core_app_extra + + +async def dispatch_wasm_invoice_paid(payment: Any) -> None: + extension_id = _payment_extension_id(payment) + if not extension_id: + return + + extension = core_app_extra.wasm_extension_registry.get(extension_id) + if not extension: + return + + export_name = _wasm_invoice_paid_export(extension.config) + if not export_name: + return + + if not _is_wasm_event_export(extension, export_name): + logger.warning( + f"WASM extension '{extension.id}' declares invalid onInvoicePaid " + f"export '{export_name}'." + ) + return + + try: + from lnbits.core.extensions.wasm import invoke_wasm_extension_export + + await invoke_wasm_extension_export( + extension.id, + export_name, + _wasm_invoice_paid_payload(payment), + context="event", + owner_id=await _wasm_invoice_paid_owner_id(extension, payment), + ) + except Exception as exc: + logger.warning( + f"WASM extension '{extension.id}' failed to handle paid invoice " + f"'{payment.payment_hash}': {exc!s}" + ) + + +def _payment_extension_id(payment: Any) -> str | None: + if isinstance(payment.extension, str) and payment.extension: + return payment.extension + + extra = payment.extra or {} + tag = extra.get("tag") or payment.tag + return tag if isinstance(tag, str) and tag else None + + +async def _wasm_invoice_paid_owner_id(extension: Any, payment: Any) -> str | None: + source_id = _payment_source_id(payment) + source_table = _wasm_public_invoice_source_table(extension.config) + if not source_id or not source_table: + return None + + from lnbits.core.extensions.storage import storage_get_row_owner_id + + return await storage_get_row_owner_id(extension.id, source_table, source_id) + + +def _payment_source_id(payment: Any) -> str | None: + extra = payment.extra or {} + source_id = extra.get("source_id") + return source_id if isinstance(source_id, str) and source_id else None + + +def _wasm_public_invoice_source_table(config: dict[str, Any]) -> str | None: + permissions = config.get("permissions") or [] + for permission in permissions: + if not isinstance(permission, dict): + continue + if permission.get("id") != "wallet.create_invoice_public": + continue + policy = permission.get("policy") or {} + table = policy.get("table") + return table if isinstance(table, str) and table else None + return None + + +def _wasm_invoice_paid_export(config: dict[str, Any]) -> str | None: + events = config.get("events") or {} + export_name = events.get("onInvoicePaid") + return export_name if isinstance(export_name, str) and export_name else None + + +def _is_wasm_event_export(extension: Any, export_name: str) -> bool: + for export in extension.exports: + if export.get("name") == export_name: + return export.get("visibility") == "event" + return False + + +def _wasm_invoice_paid_payload(payment: Any) -> dict[str, Any]: + return { + "checkingId": payment.checking_id, + "paymentHash": payment.payment_hash, + "walletId": payment.wallet_id, + "amount": payment.amount, + "fee": payment.fee, + "bolt11": payment.bolt11, + "memo": payment.memo, + "pending": payment.pending, + "status": payment.status, + "tag": payment.tag, + "extension": payment.extension, + "extra": payment.extra or {}, + "payment": json.loads(payment.json()), + } diff --git a/lnbits/core/extensions/wasm.py b/lnbits/core/extensions/wasm.py index 7a0b46bf7..6ec4f992a 100644 --- a/lnbits/core/extensions/wasm.py +++ b/lnbits/core/extensions/wasm.py @@ -33,7 +33,6 @@ async def invoke_wasm_extension_export( context=context, owner_id=owner_id, ) - print("### api", api) return await asyncio.to_thread( _invoke_wasm_extension_export_sync, diff --git a/lnbits/core/tasks.py b/lnbits/core/tasks.py index 756f75966..f4dda22ea 100644 --- a/lnbits/core/tasks.py +++ b/lnbits/core/tasks.py @@ -1,6 +1,4 @@ import asyncio -import json -from typing import Any from loguru import logger @@ -106,113 +104,6 @@ async def wait_for_paid_invoices(invoice_paid_queue: asyncio.Queue) -> None: await core_app_extra.dispatch_extension_invoice_paid(payment) -async def dispatch_wasm_invoice_paid(payment: Any) -> None: - extension_id = _payment_extension_id(payment) - if not extension_id: - return - - extension = core_app_extra.wasm_extension_registry.get(extension_id) - if not extension: - return - - export_name = _wasm_invoice_paid_export(extension.config) - if not export_name: - return - - if not _is_wasm_event_export(extension, export_name): - logger.warning( - f"WASM extension '{extension.id}' declares invalid onInvoicePaid " - f"export '{export_name}'." - ) - return - - try: - from lnbits.core.extensions.wasm import invoke_wasm_extension_export - - await invoke_wasm_extension_export( - extension.id, - export_name, - _wasm_invoice_paid_payload(payment), - context="event", - owner_id=await _wasm_invoice_paid_owner_id(extension, payment), - ) - except Exception as exc: - logger.warning( - f"WASM extension '{extension.id}' failed to handle paid invoice " - f"'{payment.payment_hash}': {exc!s}" - ) - - -def _payment_extension_id(payment: Any) -> str | None: - if isinstance(payment.extension, str) and payment.extension: - return payment.extension - - extra = payment.extra or {} - tag = extra.get("tag") or payment.tag - return tag if isinstance(tag, str) and tag else None - - -async def _wasm_invoice_paid_owner_id(extension: Any, payment: Any) -> str | None: - source_id = _payment_source_id(payment) - source_table = _wasm_public_invoice_source_table(extension.config) - if not source_id or not source_table: - return None - - from lnbits.core.extensions.storage import storage_get_row_owner_id - - return await storage_get_row_owner_id(extension.id, source_table, source_id) - - -def _payment_source_id(payment: Any) -> str | None: - extra = payment.extra or {} - source_id = extra.get("source_id") - return source_id if isinstance(source_id, str) and source_id else None - - -def _wasm_public_invoice_source_table(config: dict[str, Any]) -> str | None: - permissions = config.get("permissions") or [] - for permission in permissions: - if not isinstance(permission, dict): - continue - if permission.get("id") != "wallet.create_invoice_public": - continue - policy = permission.get("policy") or {} - table = policy.get("table") - return table if isinstance(table, str) and table else None - return None - - -def _wasm_invoice_paid_export(config: dict[str, Any]) -> str | None: - events = config.get("events") or {} - export_name = events.get("onInvoicePaid") - return export_name if isinstance(export_name, str) and export_name else None - - -def _is_wasm_event_export(extension: Any, export_name: str) -> bool: - for export in extension.exports: - if export.get("name") == export_name: - return export.get("visibility") == "event" - return False - - -def _wasm_invoice_paid_payload(payment: Any) -> dict[str, Any]: - return { - "checkingId": payment.checking_id, - "paymentHash": payment.payment_hash, - "walletId": payment.wallet_id, - "amount": payment.amount, - "fee": payment.fee, - "bolt11": payment.bolt11, - "memo": payment.memo, - "pending": payment.pending, - "status": payment.status, - "tag": payment.tag, - "extension": payment.extension, - "extra": payment.extra or {}, - "payment": json.loads(payment.json()), - } - - async def wait_for_audit_data() -> None: """ Waits for audit entries to be pushed to the queue.