From 7a316dede32987c909ddd85d73eb47b0a329f2b4 Mon Sep 17 00:00:00 2001 From: Shroominic Date: Thu, 19 Jun 2025 16:18:46 +0200 Subject: [PATCH] (temp) lock async calls to sixty-nuts & fix payout calling --- router/account.py | 4 +- router/auth.py | 7 +-- router/cashu.py | 108 ++++++++++++++++++++++++++++------------------ router/main.py | 8 +++- 4 files changed, 78 insertions(+), 49 deletions(-) diff --git a/router/account.py b/router/account.py index a0d11db5..aee3f558 100644 --- a/router/account.py +++ b/router/account.py @@ -5,6 +5,7 @@ from fastapi import APIRouter, Depends, Header, HTTPException from .auth import validate_bearer_key from .cashu import ( WALLET, + WALLET_LOCK, credit_balance, delete_key_if_zero_balance, refund_balance, @@ -76,7 +77,8 @@ async def refund_wallet_endpoint( status_code=400, detail="Balance too small to refund (less than 1 sat)" ) - token = await WALLET.send(remaining_balance_sats) + async with WALLET_LOCK: + token = await WALLET.send(remaining_balance_sats) result = {"msats": remaining_balance_msats, "recipient": None, "token": token} # Only after successful refund, zero out the balance diff --git a/router/auth.py b/router/auth.py index 9345761e..7fff637b 100644 --- a/router/auth.py +++ b/router/auth.py @@ -1,4 +1,3 @@ -import asyncio import hashlib import json import os @@ -7,7 +6,7 @@ from typing import Optional from fastapi import HTTPException, Request from sqlmodel import col, update -from .cashu import credit_balance, pay_out +from .cashu import credit_balance from .db import ApiKey, AsyncSession from .models import MODELS @@ -71,7 +70,7 @@ async def validate_bearer_key( key_expiry_time=key_expiry_time, ) session.add(new_key) - await session.flush() # Ensure the key is in the database before updating balance + await session.flush() msats = await credit_balance(bearer_key, new_key, session) if msats <= 0: raise Exception("Token redemption failed") @@ -295,6 +294,4 @@ async def adjust_payment_for_tokens( cost_data["total_msats"] = COST_PER_REQUEST - refund await session.refresh(key) - asyncio.create_task(pay_out()) - return cost_data diff --git a/router/cashu.py b/router/cashu.py index 2e7d79f1..8c8a8749 100644 --- a/router/cashu.py +++ b/router/cashu.py @@ -11,11 +11,14 @@ RECEIVE_LN_ADDRESS = os.environ["RECEIVE_LN_ADDRESS"] MINT = os.environ.get("MINT", "https://mint.minibits.cash/Bitcoin") MINIMUM_PAYOUT = int(os.environ.get("MINIMUM_PAYOUT", 100)) REFUND_PROCESSING_INTERVAL = int(os.environ.get("REFUND_PROCESSING_INTERVAL", 3600)) +PAYOUT_INTERVAL = int(os.environ.get("PAYOUT_INTERVAL", 300)) # Default 5 minutes DEV_LN_ADDRESS = "routstr@minibits.cash" DEVS_DONATION_RATE = float(os.environ.get("DEVS_DONATION_RATE", 0.021)) # 2.1% NSEC = os.environ["NSEC"] # Nostr private key for the wallet WALLET = Wallet(nsec=NSEC, mint_urls=[MINT]) +PAY_OUT_LOCK = asyncio.Lock() +WALLET_LOCK = asyncio.Lock() async def delete_key_if_zero_balance(key: ApiKey, session: AsyncSession) -> None: @@ -27,62 +30,83 @@ async def delete_key_if_zero_balance(key: ApiKey, session: AsyncSession) -> None async def init_wallet() -> None: global WALLET - WALLET = await Wallet.create(nsec=NSEC, mint_urls=[MINT]) + async with WALLET_LOCK: + WALLET = await Wallet.create(nsec=NSEC, mint_urls=[MINT]) async def close_wallet() -> None: global WALLET - await WALLET.aclose() + async with WALLET_LOCK: + await WALLET.aclose() async def pay_out() -> None: """ Calculates the pay-out amount based on the spent balance, profit, and donation rate. """ - try: - from .db import create_session + async with PAY_OUT_LOCK: + try: + from .db import create_session - async with create_session() as session: - result = await session.exec( - select(func.sum(col(ApiKey.balance))).where(ApiKey.balance > 0) - ) - balance = result.one_or_none() - if not balance: - # No balance to pay out - this is OK, not an error - return - - user_balance_sats = balance // 1000 - state = await WALLET.fetch_wallet_state() - wallet_balance_sats = state.balance - - # Handle edge cases more gracefully - if wallet_balance_sats < user_balance_sats: - print( - f"Warning: Wallet balance ({wallet_balance_sats} sats) is less than user balance ({user_balance_sats} sats). Skipping payout." + async with create_session() as session: + result = await session.exec( + select(func.sum(col(ApiKey.balance))).where(ApiKey.balance > 0) ) - return + balance = result.one_or_none() + if not balance: + # No balance to pay out - this is OK, not an error + return - if (revenue := wallet_balance_sats - user_balance_sats) <= MINIMUM_PAYOUT: - # Not enough revenue yet - this is OK - return + user_balance_sats = balance // 1000 + async with WALLET_LOCK: + state = await WALLET.fetch_wallet_state() + wallet_balance_sats = state.balance - devs_donation = int(revenue * DEVS_DONATION_RATE) - owners_draw = revenue - devs_donation + # Handle edge cases more gracefully + if wallet_balance_sats < user_balance_sats: + print( + f"Warning: Wallet balance ({wallet_balance_sats} sats) is less than user balance ({user_balance_sats} sats). Skipping payout." + ) + return - # Send payouts - await WALLET.send_to_lnurl(RECEIVE_LN_ADDRESS, owners_draw) - await WALLET.send_to_lnurl(DEV_LN_ADDRESS, devs_donation) + if ( + revenue := wallet_balance_sats - user_balance_sats + ) <= MINIMUM_PAYOUT: + # Not enough revenue yet - this is OK + return - except Exception as e: - # Log the error but don't crash - payouts can be retried later - print(f"Error in pay_out: {e}") + devs_donation = int(revenue * DEVS_DONATION_RATE) + owners_draw = revenue - devs_donation + + # Send payouts + async with WALLET_LOCK: + await WALLET.send_to_lnurl(RECEIVE_LN_ADDRESS, owners_draw) + await WALLET.send_to_lnurl(DEV_LN_ADDRESS, devs_donation) + + except Exception as e: + print(f"Error in pay_out: {e}") + + +# Periodic payout task +async def periodic_payout() -> None: + while True: + try: + await asyncio.sleep(300) # Run every 5 minutes + await pay_out() + except asyncio.CancelledError: + break + except Exception as e: + print(f"Error in periodic payout: {e}") + # Continue running even if payout fails async def credit_balance(cashu_token: str, key: ApiKey, session: AsyncSession) -> int: """Redeem a Cashu token and credit the amount to the API key balance.""" try: - amount_sats, _ = await WALLET.redeem(cashu_token) - except Exception: + async with WALLET_LOCK: + amount_sats, _ = await WALLET.redeem(cashu_token) + except Exception as e: + print(f"Error in credit_balance: {e}") # Ensure the balance cannot become negative if redeem fails return 0 @@ -172,13 +196,15 @@ async def refund_balance(amount_msats: int, key: ApiKey, session: AsyncSession) if key.refund_address is None: raise ValueError("Refund address not set.") - return await WALLET.send_to_lnurl( - key.refund_address, - amount=amount_sats, - ) + async with WALLET_LOCK: + return await WALLET.send_to_lnurl( + key.refund_address, + amount=amount_sats, + ) async def redeem(cashu_token: str, lnurl: str) -> int: - amount_sats, _ = await WALLET.redeem(cashu_token) - await WALLET.send_to_lnurl(lnurl, amount=amount_sats) + async with WALLET_LOCK: + amount_sats, _ = await WALLET.redeem(cashu_token) + await WALLET.send_to_lnurl(lnurl, amount=amount_sats) return amount_sats diff --git a/router/main.py b/router/main.py index f96f6624..80a646d8 100644 --- a/router/main.py +++ b/router/main.py @@ -8,7 +8,7 @@ from fastapi.middleware.cors import CORSMiddleware from .account import wallet_router from .admin import admin_router -from .cashu import check_for_refunds, close_wallet, init_wallet +from .cashu import check_for_refunds, close_wallet, init_wallet, periodic_payout from .db import init_db from .discovery import providers_router from .models import MODELS, update_sats_pricing @@ -23,13 +23,17 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: await init_wallet() pricing_task = asyncio.create_task(update_sats_pricing()) refund_task = asyncio.create_task(check_for_refunds()) + payout_task = asyncio.create_task(periodic_payout()) try: yield finally: refund_task.cancel() pricing_task.cancel() - await asyncio.gather(pricing_task, refund_task, return_exceptions=True) + payout_task.cancel() + await asyncio.gather( + pricing_task, refund_task, payout_task, return_exceptions=True + ) await close_wallet()