From f9e4269d3c783a386facccd7fe40ecf96f7c1995 Mon Sep 17 00:00:00 2001 From: Shroominic Date: Mon, 1 Dec 2025 20:21:52 +0800 Subject: [PATCH] nostr-revenue-stats --- routstr/core/main.py | 7 +++ routstr/core/settings.py | 3 + routstr/revenue_stats.py | 128 +++++++++++++++++++++++++++++++++++++++ 3 files changed, 138 insertions(+) create mode 100644 routstr/revenue_stats.py diff --git a/routstr/core/main.py b/routstr/core/main.py index c3d4e104..5e8131cf 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -20,6 +20,7 @@ from ..payment.models import ( ) from ..payment.price import update_prices_periodically from ..proxy import initialize_upstreams, proxy_router, refresh_model_maps_periodically +from ..revenue_stats import publish_revenue_stats_task from ..wallet import periodic_payout from .admin import admin_router from .db import create_session, init_db, run_migrations @@ -47,6 +48,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: pricing_task = None payout_task = None nip91_task = None + revenue_stats_task = None providers_task = None models_refresh_task = None models_cleanup_task = None @@ -93,6 +95,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: model_maps_refresh_task = asyncio.create_task(refresh_model_maps_periodically()) payout_task = asyncio.create_task(periodic_payout()) nip91_task = asyncio.create_task(announce_provider()) + revenue_stats_task = asyncio.create_task(publish_revenue_stats_task()) providers_task = asyncio.create_task(providers_cache_refresher()) yield @@ -114,6 +117,8 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: payout_task.cancel() if nip91_task is not None: nip91_task.cancel() + if revenue_stats_task is not None: + revenue_stats_task.cancel() if providers_task is not None: providers_task.cancel() if models_refresh_task is not None: @@ -133,6 +138,8 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: tasks_to_wait.append(payout_task) if nip91_task is not None: tasks_to_wait.append(nip91_task) + if revenue_stats_task is not None: + tasks_to_wait.append(revenue_stats_task) if providers_task is not None: tasks_to_wait.append(providers_task) if models_refresh_task is not None: diff --git a/routstr/core/settings.py b/routstr/core/settings.py index 302df24d..10ef5c31 100644 --- a/routstr/core/settings.py +++ b/routstr/core/settings.py @@ -69,6 +69,9 @@ class Settings(BaseSettings): ) enable_pricing_refresh: bool = Field(default=True, env="ENABLE_PRICING_REFRESH") enable_models_refresh: bool = Field(default=True, env="ENABLE_MODELS_REFRESH") + enable_revenue_stats_publishing: bool = Field( + default=True, env="ENABLE_REVENUE_STATS_PUBLISHING" + ) refund_cache_ttl_seconds: int = Field(default=3600, env="REFUND_CACHE_TTL_SECONDS") # Logging diff --git a/routstr/revenue_stats.py b/routstr/revenue_stats.py new file mode 100644 index 00000000..3692088e --- /dev/null +++ b/routstr/revenue_stats.py @@ -0,0 +1,128 @@ +import asyncio +import json +from datetime import datetime, timedelta, timezone +from typing import Any + +from nostr.event import Event +from nostr.key import PrivateKey + +from .core.log_manager import log_manager +from .core.logging import get_logger +from .core.settings import settings +from .nip91 import nsec_to_keypair, publish_to_relay + +logger = get_logger(__name__) + +# Custom kind for Revenue Stats +KIND_REVENUE_STATS = 7375 + + +def _event_to_dict(ev: Event) -> dict[str, Any]: + return { + "id": ev.id, + "pubkey": ev.public_key, + "created_at": ev.created_at, + "kind": int(ev.kind) if not isinstance(ev.kind, int) else ev.kind, + "tags": ev.tags, + "content": ev.content, + "sig": ev.signature, + } + + +def create_revenue_event( + private_key_hex: str, stats: dict[str, Any], date_str: str +) -> dict[str, Any]: + pk = PrivateKey(bytes.fromhex(private_key_hex)) + + # Content is the stats JSON + content = json.dumps(stats, separators=(",", ":")) + + # Tags: + # d: revenue-YYYY-MM-DD (to make it unique per day if replaceable) + # t: revenue + # date: YYYY-MM-DD + tags = [["d", f"revenue-{date_str}"], ["t", "revenue"], ["date", date_str]] + + ev = Event(pk.public_key.hex(), content, kind=KIND_REVENUE_STATS, tags=tags) + pk.sign_event(ev) + return _event_to_dict(ev) + + +async def publish_revenue_stats_task() -> None: + """ + Background task that runs once a day (at midnight UTC) to publish revenue stats. + """ + logger.info("Revenue stats publisher task started") + + while True: + try: + # Calculate time until next midnight UTC + now = datetime.now(timezone.utc) + next_run = (now + timedelta(days=1)).replace( + hour=0, minute=0, second=0, microsecond=0 + ) + sleep_seconds = (next_run - now).total_seconds() + + logger.info( + f"Revenue stats task sleeping for {sleep_seconds:.1f}s until {next_run}" + ) + await asyncio.sleep(sleep_seconds) + + # Check if enabled + if not settings.enable_revenue_stats_publishing: + logger.info("Revenue stats publishing disabled, skipping") + continue + + # It's midnight! Fetch stats for the PREVIOUS day (last 24h) + nsec = settings.nsec + if not nsec: + logger.warning("No NSEC configured, skipping revenue stats publish") + continue + + keypair = nsec_to_keypair(nsec) + if not keypair: + logger.error("Invalid NSEC, skipping revenue stats publish") + continue + + private_key_hex, _ = keypair + + # Get date string for the day that just finished + yesterday = datetime.now(timezone.utc) - timedelta(days=1) + date_str = yesterday.strftime("%Y-%m-%d") + + logger.info(f"Publishing revenue stats for {date_str}") + + stats = log_manager.get_usage_summary(hours=24) + + # Create event + event = create_revenue_event(private_key_hex, stats, date_str) + + # Publish to relays + relay_urls = [u.strip() for u in settings.relays if u.strip()] + if not relay_urls: + relay_urls = [ + "wss://relay.nostr.band", + "wss://relay.damus.io", + "wss://relay.routstr.com", + "wss://nos.lol", + ] + + success_count = 0 + for relay in relay_urls: + if await publish_to_relay(relay, event): + success_count += 1 + + logger.info( + f"Published revenue stats to {success_count}/{len(relay_urls)} relays" + ) + + # Sleep a bit to avoid rapid loop if clock skews back + await asyncio.sleep(60) + + except asyncio.CancelledError: + logger.info("Revenue stats task cancelled") + break + except Exception as e: + logger.error(f"Error in revenue stats task: {e}") + # Sleep 1 hour before retrying loop calculation to avoid busy loop on error + await asyncio.sleep(3600)