mirror of
https://github.com/Routstr/routstr-core.git
synced 2026-08-09 02:54:37 +00:00
nostr-revenue-stats
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user