refactor: move out dispatch_wasm_invoice_paid from tasks
This commit is contained in:
+1
-1
@@ -24,6 +24,7 @@ from lnbits.core.crud import (
|
|||||||
update_installed_extension_state,
|
update_installed_extension_state,
|
||||||
)
|
)
|
||||||
from lnbits.core.crud.extensions import create_installed_extension
|
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 (
|
from lnbits.core.extensions.loader import (
|
||||||
is_wasm_extension_dir,
|
is_wasm_extension_dir,
|
||||||
is_wasm_extension_id,
|
is_wasm_extension_id,
|
||||||
@@ -37,7 +38,6 @@ from lnbits.core.services.payments import check_pending_payments
|
|||||||
from lnbits.core.tasks import (
|
from lnbits.core.tasks import (
|
||||||
audit_queue,
|
audit_queue,
|
||||||
collect_exchange_rates_data,
|
collect_exchange_rates_data,
|
||||||
dispatch_wasm_invoice_paid,
|
|
||||||
purge_audit_data,
|
purge_audit_data,
|
||||||
run_by_the_minute_tasks,
|
run_by_the_minute_tasks,
|
||||||
wait_for_audit_data,
|
wait_for_audit_data,
|
||||||
|
|||||||
@@ -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()),
|
||||||
|
}
|
||||||
@@ -33,7 +33,6 @@ async def invoke_wasm_extension_export(
|
|||||||
context=context,
|
context=context,
|
||||||
owner_id=owner_id,
|
owner_id=owner_id,
|
||||||
)
|
)
|
||||||
print("### api", api)
|
|
||||||
|
|
||||||
return await asyncio.to_thread(
|
return await asyncio.to_thread(
|
||||||
_invoke_wasm_extension_export_sync,
|
_invoke_wasm_extension_export_sync,
|
||||||
|
|||||||
@@ -1,6 +1,4 @@
|
|||||||
import asyncio
|
import asyncio
|
||||||
import json
|
|
||||||
from typing import Any
|
|
||||||
|
|
||||||
from loguru import logger
|
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)
|
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:
|
async def wait_for_audit_data() -> None:
|
||||||
"""
|
"""
|
||||||
Waits for audit entries to be pushed to the queue.
|
Waits for audit entries to be pushed to the queue.
|
||||||
|
|||||||
Reference in New Issue
Block a user