From 8fc1b6484c897ed00696b88e48c437eb97061e66 Mon Sep 17 00:00:00 2001 From: Evan Yang Date: Mon, 2 Mar 2026 16:52:39 +0800 Subject: [PATCH] Add Nostr analytics snapshots and expand stats model coverage --- .env.example | 1 + docs/provider/configuration.md | 2 + docs/provider/dashboard.md | 1 + routstr/core/main.py | 12 +- routstr/core/settings.py | 33 +- routstr/core/usage_analytics_store.py | 2 +- routstr/nostr/__init__.py | 3 +- routstr/nostr/analytics.py | 771 ++++++++++++++++++++++ tests/unit/test_nostr_analytics.py | 230 +++++++ tests/unit/test_settings.py | 38 +- ui/components/settings/admin-settings.tsx | 50 ++ ui/components/top-models-usage-chart.tsx | 6 +- 12 files changed, 1137 insertions(+), 12 deletions(-) create mode 100644 routstr/nostr/analytics.py create mode 100644 tests/unit/test_nostr_analytics.py diff --git a/.env.example b/.env.example index d093e006..d4c768d5 100644 --- a/.env.example +++ b/.env.example @@ -14,6 +14,7 @@ UPSTREAM_API_KEY=your-upstream-api-key # HTTP_URL=https://api.mynode.com # ONION_URL=http://mynode.onion (auto fetched from compose) # RELAYS="wss://relay.damus.io,wss://relay.nostr.band,wss://eden.nostr.land,wss://relay.routstr.com" +# ENABLE_ANALYTICS_SHARING=true # CASHU_MINTS="https://mint.minibits.cash/Bitcoin,https://mint.cubabitcoin.org,https://ecashmint.otrta.me" # RECEIVE_LN_ADDRESS= diff --git a/docs/provider/configuration.md b/docs/provider/configuration.md index 633ccd10..20c0e70c 100644 --- a/docs/provider/configuration.md +++ b/docs/provider/configuration.md @@ -97,6 +97,7 @@ Announce your node on the network: | **Npub** | Your Nostr public key | | **Nsec** | Your Nostr private key (for signing) | | **Relays** | Relays to publish announcements | +| **Share Analytics** | Publish aggregate usage stats to Nostr | See [Discovery](discovery.md) for details. @@ -122,6 +123,7 @@ Use environment variables for: | `DESCRIPTION` | Node description | `A Routstr Node` | | `NPUB` | Nostr public key (bech32) | — | | `NSEC` | Nostr private key | — | +| `ENABLE_ANALYTICS_SHARING` | Enable usage analytics sharing to Nostr | `true` | | `CASHU_MINTS` | Comma-separated mint URLs | `https://mint.minibits.cash/Bitcoin` | | `RECEIVE_LN_ADDRESS` | Lightning address for withdrawals | — | | `TOR_PROXY_URL` | SOCKS5 proxy for Tor | `socks5://127.0.0.1:9050` | diff --git a/docs/provider/dashboard.md b/docs/provider/dashboard.md index aaa860bf..66852845 100644 --- a/docs/provider/dashboard.md +++ b/docs/provider/dashboard.md @@ -152,6 +152,7 @@ Manage which mints you accept payments from: |-------|-------------| | **Nsec** | Private key for signing announcements | | **Relays** | Where to publish your node advertisement | +| **Share Analytics** | Toggle publishing aggregate usage stats to Nostr | ### Security diff --git a/routstr/core/main.py b/routstr/core/main.py index 461f8b6d..d0769706 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -12,7 +12,11 @@ from starlette.exceptions import HTTPException from ..auth import periodic_key_reset from ..balance import balance_router, deprecated_wallet_router -from ..nostr import announce_provider, providers_cache_refresher +from ..nostr import ( + announce_provider, + providers_cache_refresher, + publish_usage_analytics, +) from ..nostr.discovery import providers_router from ..payment.models import models_router, update_sats_pricing from ..payment.price import update_prices_periodically @@ -45,6 +49,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: pricing_task = None payout_task = None nip91_task = None + analytics_task = None providers_task = None models_refresh_task = None model_maps_refresh_task = None @@ -103,6 +108,7 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: payout_task = asyncio.create_task(periodic_payout()) if global_settings.nsec: nip91_task = asyncio.create_task(announce_provider()) + analytics_task = asyncio.create_task(publish_usage_analytics()) if global_settings.providers_refresh_interval_seconds > 0: providers_task = asyncio.create_task(providers_cache_refresher()) key_reset_task = asyncio.create_task(periodic_key_reset()) @@ -130,6 +136,8 @@ async def lifespan(_: FastAPI) -> AsyncGenerator[None, None]: payout_task.cancel() if nip91_task is not None: nip91_task.cancel() + if analytics_task is not None: + analytics_task.cancel() if providers_task is not None: providers_task.cancel() if models_refresh_task is not None: @@ -151,6 +159,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 analytics_task is not None: + tasks_to_wait.append(analytics_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 fa1464ac..20a3053c 100644 --- a/routstr/core/settings.py +++ b/routstr/core/settings.py @@ -92,6 +92,20 @@ class Settings(BaseSettings): # Discovery relays: list[str] = Field(default_factory=list, env="RELAYS") + enable_analytics_sharing: bool = Field( + default=True, env="ENABLE_ANALYTICS_SHARING" + ) + +def _normalize_settings_data(data: dict[str, Any]) -> dict[str, Any]: + """Discard unknown keys from persisted settings.""" + normalized: dict[str, Any] = {} + known_fields = Settings.__fields__ + + for key, value in data.items(): + if key in known_fields: + normalized[key] = value + + return normalized def _compute_primary_mint(cashu_mints: list[str]) -> str: @@ -231,16 +245,20 @@ class SettingsService: db_id, db_data, _updated_at = row try: - db_json = ( + db_json_raw = ( json.loads(db_data) if isinstance(db_data, str) else dict(db_data) ) + if not isinstance(db_json_raw, dict): + db_json_raw = {} except Exception: - db_json = {} + db_json_raw = {} + db_json = _normalize_settings_data(db_json_raw) merged_dict: dict[str, Any] = dict(env_resolved.dict()) merged_dict.update( {k: v for k, v in db_json.items() if v not in (None, "", [], {})} ) + merged_dict = Settings(**merged_dict).dict() # Ensure primary_mint is consistent with cashu_mints if not explicitly set if not merged_dict.get("primary_mint"): @@ -248,7 +266,7 @@ class SettingsService: merged_dict.get("cashu_mints", []) ) - if any(k not in db_json for k in merged_dict.keys()): + if db_json_raw != merged_dict: await db_session.exec( # type: ignore text( "UPDATE settings SET data = :data, updated_at = :updated_at WHERE id = 1" @@ -271,7 +289,7 @@ class SettingsService: ) -> Settings: async with cls._lock: current = cls.get() - candidate_dict = {**current.dict(), **partial} + candidate_dict = {**current.dict(), **_normalize_settings_data(partial)} candidate = Settings(**candidate_dict) from sqlmodel import text @@ -304,7 +322,12 @@ class SettingsService: if row is None: raise RuntimeError("Settings row missing") (data_str,) = row - data = json.loads(data_str) if isinstance(data_str, str) else dict(data_str) + data_raw = ( + json.loads(data_str) if isinstance(data_str, str) else dict(data_str) + ) + if not isinstance(data_raw, dict): + data_raw = {} + data = _normalize_settings_data(data_raw) # Update in-place for k, v in data.items(): setattr(settings, k, v) diff --git a/routstr/core/usage_analytics_store.py b/routstr/core/usage_analytics_store.py index 2df1f440..c5ce3023 100644 --- a/routstr/core/usage_analytics_store.py +++ b/routstr/core/usage_analytics_store.py @@ -1115,7 +1115,7 @@ class UsageAnalyticsStore: hours_back: int, limit: int, ) -> dict[str, Any]: - top_limit = max(1, min(int(limit), 10)) + top_limit = max(1, min(int(limit), 20)) top_rows = conn.execute( """ SELECT diff --git a/routstr/nostr/__init__.py b/routstr/nostr/__init__.py index b19039a7..afd165f5 100644 --- a/routstr/nostr/__init__.py +++ b/routstr/nostr/__init__.py @@ -1,4 +1,5 @@ +from .analytics import publish_usage_analytics from .discovery import providers_cache_refresher from .listing import announce_provider -__all__ = ["providers_cache_refresher", "announce_provider"] +__all__ = ["providers_cache_refresher", "announce_provider", "publish_usage_analytics"] diff --git a/routstr/nostr/analytics.py b/routstr/nostr/analytics.py new file mode 100644 index 00000000..137040cb --- /dev/null +++ b/routstr/nostr/analytics.py @@ -0,0 +1,771 @@ +#!/usr/bin/env python3 +""" +Nostr usage analytics publisher. +Publishes routstr analytics snapshots for latest/day/month plus daily checkpoints. +""" + +from __future__ import annotations + +import asyncio +import hashlib +import json +import math +import time +from datetime import datetime, timezone +from typing import Any + +from nostr.event import Event +from nostr.key import PrivateKey + +from ..core import get_logger +from ..core.log_manager import log_manager +from ..core.settings import settings +from .listing import nsec_to_keypair, publish_to_relay + +logger = get_logger(__name__) + +ANALYTICS_KIND = 38422 +ANALYTICS_SCHEMA = "routstr.analytics.usage.v2" +ANALYTICS_CHECKPOINT_SCHEMA = "routstr.analytics.checkpoint.v1" +DEFAULT_RELAYS = [ + "wss://relay.nostr.band", + "wss://relay.damus.io", + "wss://relay.routstr.com", + "wss://nos.lol", +] +PUBLISH_INTERVAL_SECONDS = 15 * 60 +DISABLED_POLL_SECONDS = 60 +DASHBOARD_WINDOW_HOURS = 24 +DASHBOARD_INTERVAL_MINUTES = 60 +MODEL_LIMIT = 20 + +WINDOW_DEFINITIONS: tuple[tuple[str, int, int], ...] = ( + ("24h", 24, 60), + ("7d", 7 * 24, 6 * 60), + ("30d", 30 * 24, 24 * 60), +) + + +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 _resolve_provider_id(public_key_hex: str) -> str: + explicit_provider_id = (settings.provider_id or "").strip() + if explicit_provider_id: + return explicit_provider_id + return public_key_hex[:12] + + +def _resolve_endpoint_urls() -> list[str]: + urls: list[str] = [] + http_url = (settings.http_url or "").strip() + onion_url = (settings.onion_url or "").strip() + + if http_url and http_url != "http://localhost:8000": + urls.append(http_url) + + if onion_url: + if onion_url.endswith(".onion") and not ( + onion_url.startswith("http://") or onion_url.startswith("https://") + ): + onion_url = f"http://{onion_url}" + urls.append(onion_url) + + return urls + + +def _resolve_relays() -> list[str]: + configured = [url.strip() for url in settings.relays if url.strip()] + return configured if configured else list(DEFAULT_RELAYS) + + +def _to_int(value: Any) -> int: + if isinstance(value, bool): + return int(value) + if isinstance(value, int): + return value + if isinstance(value, float): + return int(value) + if isinstance(value, str): + try: + return int(float(value)) + except ValueError: + return 0 + return 0 + + +def _to_float(value: Any) -> float: + if isinstance(value, bool): + return float(int(value)) + if isinstance(value, (int, float)): + return float(value) + if isinstance(value, str): + try: + return float(value) + except ValueError: + return 0.0 + return 0.0 + + +def _utc_day_key(unix_ts: int) -> str: + return datetime.fromtimestamp(unix_ts, tz=timezone.utc).strftime("%Y-%m-%d") + + +def _utc_month_key(unix_ts: int) -> str: + return datetime.fromtimestamp(unix_ts, tz=timezone.utc).strftime("%Y-%m") + + +def _utc_day_start_ts(unix_ts: int) -> int: + dt = datetime.fromtimestamp(unix_ts, tz=timezone.utc) + start = datetime(dt.year, dt.month, dt.day, tzinfo=timezone.utc) + return int(start.timestamp()) + + +def _utc_month_start_ts(unix_ts: int) -> int: + dt = datetime.fromtimestamp(unix_ts, tz=timezone.utc) + start = datetime(dt.year, dt.month, 1, tzinfo=timezone.utc) + return int(start.timestamp()) + + +def _hours_since(start_ts: int, end_ts: int) -> int: + elapsed = max(1, end_ts - start_ts) + return max(1, int(math.ceil(elapsed / 3600))) + + +def _aggregate_top_model_usage( + model_usage_mix: dict[str, Any], +) -> tuple[list[dict[str, Any]], dict[str, Any]]: + top_models_raw = model_usage_mix.get("top_models", []) + mix_metrics_raw = model_usage_mix.get("metrics", []) + + top_models = [model for model in top_models_raw if isinstance(model, str)] + metrics = [row for row in mix_metrics_raw if isinstance(row, dict)] + + model_totals: dict[str, dict[str, float | int]] = { + model: { + "successful_requests": 0, + "revenue_msats": 0.0, + "total_tokens": 0, + } + for model in top_models + } + others = { + "successful_requests": 0, + "revenue_msats": 0.0, + "total_tokens": 0, + } + + for metric in metrics: + model_counts = metric.get("model_counts", {}) + model_revenue = metric.get("model_revenue_msats", {}) + model_tokens = metric.get("model_tokens", {}) + + if isinstance(model_counts, dict): + for model, count in model_counts.items(): + if model in model_totals: + model_totals[model]["successful_requests"] += _to_int(count) + + if isinstance(model_revenue, dict): + for model, amount in model_revenue.items(): + if model in model_totals: + model_totals[model]["revenue_msats"] += _to_float(amount) + + if isinstance(model_tokens, dict): + for model, token_count in model_tokens.items(): + if model in model_totals: + model_totals[model]["total_tokens"] += _to_int(token_count) + + others["successful_requests"] += _to_int(metric.get("others", 0)) + others["revenue_msats"] += _to_float(metric.get("others_revenue_msats", 0.0)) + others["total_tokens"] += _to_int(metric.get("others_tokens", 0)) + + model_rows = [ + { + "model": model, + "successful_requests": int(values["successful_requests"]), + "revenue_msats": float(values["revenue_msats"]), + "total_tokens": int(values["total_tokens"]), + } + for model, values in model_totals.items() + ] + model_rows.sort(key=lambda row: row["successful_requests"], reverse=True) + + return model_rows, others + + +def _build_summary_payload(summary: dict[str, Any]) -> dict[str, Any]: + return { + "total_requests": _to_int(summary.get("total_requests", 0)), + "successful_chat_completions": _to_int( + summary.get("successful_chat_completions", 0) + ), + "failed_requests": _to_int(summary.get("failed_requests", 0)), + "success_rate": _to_float(summary.get("success_rate", 0.0)), + "unique_models_count": _to_int(summary.get("unique_models_count", 0)), + "input_tokens": _to_int(summary.get("input_tokens", 0)), + "output_tokens": _to_int(summary.get("output_tokens", 0)), + "total_tokens": _to_int(summary.get("total_tokens", 0)), + "revenue_msats": _to_float(summary.get("revenue_msats", 0.0)), + "refunds_msats": _to_float(summary.get("refunds_msats", 0.0)), + "net_revenue_msats": _to_float(summary.get("net_revenue_msats", 0.0)), + "revenue_sats": _to_float(summary.get("revenue_sats", 0.0)), + "refunds_sats": _to_float(summary.get("refunds_sats", 0.0)), + "net_revenue_sats": _to_float(summary.get("net_revenue_sats", 0.0)), + } + + +def _build_model_revenue_rows(revenue_by_model: dict[str, Any]) -> list[dict[str, Any]]: + rows: list[dict[str, Any]] = [] + models_raw = revenue_by_model.get("models", []) + if not isinstance(models_raw, list): + return rows + + for row in models_raw: + if not isinstance(row, dict): + continue + model_name = str(row.get("model", "unknown")) + rows.append( + { + "model": model_name, + "requests": _to_int(row.get("requests", 0)), + "successful": _to_int(row.get("successful", 0)), + "failed": _to_int(row.get("failed", 0)), + "revenue_sats": _to_float(row.get("revenue_sats", 0.0)), + "refunds_sats": _to_float(row.get("refunds_sats", 0.0)), + "net_revenue_sats": _to_float(row.get("net_revenue_sats", 0.0)), + } + ) + return rows + + +def _build_window_payload( + *, + hours: int, + interval: int, + model_limit: int, +) -> dict[str, Any]: + dashboard = log_manager.get_usage_dashboard( + interval=interval, + hours=hours, + error_limit=1, + model_limit=model_limit, + ) + summary = dashboard.get("summary", {}) + revenue_by_model = dashboard.get("revenue_by_model", {}) + model_usage_mix = dashboard.get("model_usage_mix", {}) + + summary_payload = _build_summary_payload(summary if isinstance(summary, dict) else {}) + model_revenue_rows = _build_model_revenue_rows( + revenue_by_model if isinstance(revenue_by_model, dict) else {} + ) + usage_mix_payload = model_usage_mix if isinstance(model_usage_mix, dict) else {} + top_model_usage, others_usage = _aggregate_top_model_usage(usage_mix_payload) + + return { + "window_hours": hours, + "interval_minutes": interval, + "summary": summary_payload, + "model_revenue": model_revenue_rows, + "top_model_usage": top_model_usage, + "others_usage": others_usage, + "model_usage_mix": usage_mix_payload, + } + + +def build_latest_usage_analytics_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + model_limit: int = MODEL_LIMIT, +) -> dict[str, Any]: + windows: dict[str, dict[str, Any]] = {} + for key, window_hours, window_interval in WINDOW_DEFINITIONS: + windows[key] = _build_window_payload( + hours=window_hours, + interval=window_interval, + model_limit=model_limit, + ) + + primary_window = windows.get("24h", {}) + return { + "schema": ANALYTICS_SCHEMA, + "generated_at": generated_at, + "provider_id": provider_id, + "pubkey": public_key_hex, + "npub": settings.npub or "", + "endpoint_urls": _resolve_endpoint_urls(), + "period_type": "latest", + "period_key": "latest", + "period_start_unix": generated_at - (24 * 3600), + "period_end_unix": generated_at, + "summary": primary_window.get("summary", {}), + "model_revenue": primary_window.get("model_revenue", []), + "top_model_usage": primary_window.get("top_model_usage", []), + "others_usage": primary_window.get("others_usage", {}), + "model_usage_mix": primary_window.get("model_usage_mix", {}), + "windows": windows, + } + + +def _build_period_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + period_type: str, + period_key: str, + period_start_unix: int, + interval_minutes: int, + model_limit: int = MODEL_LIMIT, +) -> dict[str, Any]: + hours = _hours_since(period_start_unix, generated_at) + window = _build_window_payload( + hours=hours, + interval=interval_minutes, + model_limit=model_limit, + ) + return { + "schema": ANALYTICS_SCHEMA, + "generated_at": generated_at, + "provider_id": provider_id, + "pubkey": public_key_hex, + "npub": settings.npub or "", + "endpoint_urls": _resolve_endpoint_urls(), + "period_type": period_type, + "period_key": period_key, + "period_start_unix": period_start_unix, + "period_end_unix": generated_at, + "window_hours": window.get("window_hours", hours), + "interval_minutes": window.get("interval_minutes", interval_minutes), + "summary": window.get("summary", {}), + "model_revenue": window.get("model_revenue", []), + "top_model_usage": window.get("top_model_usage", []), + "others_usage": window.get("others_usage", {}), + "model_usage_mix": window.get("model_usage_mix", {}), + } + + +def build_day_usage_analytics_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + model_limit: int = MODEL_LIMIT, +) -> dict[str, Any]: + day_key = _utc_day_key(generated_at) + day_start = _utc_day_start_ts(generated_at) + payload = _build_period_payload( + provider_id, + public_key_hex=public_key_hex, + generated_at=generated_at, + period_type="day", + period_key=day_key, + period_start_unix=day_start, + interval_minutes=60, + model_limit=model_limit, + ) + payload["day"] = day_key + return payload + + +def build_month_usage_analytics_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + model_limit: int = MODEL_LIMIT, +) -> dict[str, Any]: + month_key = _utc_month_key(generated_at) + month_start = _utc_month_start_ts(generated_at) + payload = _build_period_payload( + provider_id, + public_key_hex=public_key_hex, + generated_at=generated_at, + period_type="month", + period_key=month_key, + period_start_unix=month_start, + interval_minutes=24 * 60, + model_limit=model_limit, + ) + payload["month"] = month_key + return payload + + +def build_usage_analytics_payload( + provider_id: str, + *, + public_key_hex: str, + hours: int = DASHBOARD_WINDOW_HOURS, + interval: int = DASHBOARD_INTERVAL_MINUTES, + model_limit: int = MODEL_LIMIT, +) -> dict[str, Any]: + # Backward-compatible helper kept for existing tests/callers. + _ = (hours, interval) + return build_latest_usage_analytics_payload( + provider_id, + public_key_hex=public_key_hex, + generated_at=int(time.time()), + model_limit=model_limit, + ) + + +def create_usage_analytics_event( + private_key_hex: str, + provider_id: str, + payload_json: str, + *, + period_type: str, + period_key: str, + d_tag: str, +) -> dict[str, Any]: + private_key = PrivateKey(bytes.fromhex(private_key_hex)) + tags = [ + ["d", d_tag], + ["provider", provider_id], + ["schema", ANALYTICS_SCHEMA], + ["period", period_type], + ["period_key", period_key], + ] + if period_type == "day": + tags.append(["day", period_key]) + elif period_type == "month": + tags.append(["month", period_key]) + + event = Event( + public_key=private_key.public_key.hex(), + content=payload_json, + kind=ANALYTICS_KIND, + tags=tags, + ) + private_key.sign_event(event) + return _event_to_dict(event) + + +def _fingerprint_payload(payload: dict[str, Any]) -> str: + normalized = dict(payload) + # Ignore volatile timestamps for deduping semantically identical snapshots. + normalized.pop("generated_at", None) + normalized.pop("period_end_unix", None) + payload_json = json.dumps(normalized, separators=(",", ":"), sort_keys=True) + return hashlib.sha256(payload_json.encode("utf-8")).hexdigest() + + +def _stable_hash(data: dict[str, Any]) -> str: + encoded = json.dumps(data, separators=(",", ":"), sort_keys=True) + return hashlib.sha256(encoded.encode("utf-8")).hexdigest() + + +def build_analytics_checkpoint_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + day_utc: str, + refs: dict[str, dict[str, str]], + previous_checkpoint_hash: str | None, +) -> dict[str, Any]: + base = { + "schema": ANALYTICS_CHECKPOINT_SCHEMA, + "generated_at": generated_at, + "provider_id": provider_id, + "pubkey": public_key_hex, + "npub": settings.npub or "", + "day_utc": day_utc, + "refs": refs, + "previous_checkpoint_hash": previous_checkpoint_hash or "", + } + checkpoint_hash = _stable_hash( + { + "provider_id": provider_id, + "day_utc": day_utc, + "refs": refs, + "previous_checkpoint_hash": previous_checkpoint_hash or "", + } + ) + base["checkpoint_hash"] = checkpoint_hash + return base + + +def create_analytics_checkpoint_event( + private_key_hex: str, + provider_id: str, + payload_json: str, + *, + day_utc: str, + previous_checkpoint_hash: str | None, +) -> dict[str, Any]: + private_key = PrivateKey(bytes.fromhex(private_key_hex)) + tags = [ + ["d", f"{provider_id}:usage:checkpoint:{day_utc}"], + ["provider", provider_id], + ["schema", ANALYTICS_CHECKPOINT_SCHEMA], + ["day", day_utc], + ] + if previous_checkpoint_hash: + tags.append(["prev", previous_checkpoint_hash]) + + event = Event( + public_key=private_key.public_key.hex(), + content=payload_json, + kind=ANALYTICS_KIND, + tags=tags, + ) + private_key.sign_event(event) + return _event_to_dict(event) + + +async def publish_usage_analytics() -> None: + last_period_state: dict[str, tuple[str, str]] = {} + last_checkpoint_state: tuple[str, str] | None = None + checkpoint_day: str | None = None + checkpoint_hash_for_day: str | None = None + previous_checkpoint_hash: str | None = None + + parsed_nsec: str | None = None + private_key_hex: str | None = None + public_key_hex: str | None = None + provider_id: str | None = None + warned_missing_nsec = False + + logger.info("Usage analytics sharing task started") + + while True: + try: + if not settings.enable_analytics_sharing: + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + nsec = (settings.nsec or "").strip() + if not nsec: + if not warned_missing_nsec: + logger.info("NSEC is not configured; skipping analytics sharing to Nostr") + warned_missing_nsec = True + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + warned_missing_nsec = False + if nsec != parsed_nsec or private_key_hex is None or public_key_hex is None: + keypair = nsec_to_keypair(nsec) + if not keypair: + logger.error("Invalid NSEC; analytics sharing is paused") + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + private_key_hex, public_key_hex = keypair + parsed_nsec = nsec + provider_id = _resolve_provider_id(public_key_hex) + last_period_state = {} + last_checkpoint_state = None + checkpoint_day = None + checkpoint_hash_for_day = None + previous_checkpoint_hash = None + + if private_key_hex is None or public_key_hex is None: + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + relay_urls = _resolve_relays() + if not relay_urls: + logger.warning("No Nostr relays configured; analytics sharing skipped") + await asyncio.sleep(DISABLED_POLL_SECONDS) + continue + + resolved_provider_id = provider_id or _resolve_provider_id(public_key_hex) + now_ts = int(time.time()) + day_key = _utc_day_key(now_ts) + month_key = _utc_month_key(now_ts) + + latest_payload = build_latest_usage_analytics_payload( + resolved_provider_id, + public_key_hex=public_key_hex, + generated_at=now_ts, + ) + day_payload = build_day_usage_analytics_payload( + resolved_provider_id, + public_key_hex=public_key_hex, + generated_at=now_ts, + ) + month_payload = build_month_usage_analytics_payload( + resolved_provider_id, + public_key_hex=public_key_hex, + generated_at=now_ts, + ) + + payload_specs = [ + { + "period_type": "latest", + "period_key": "latest", + "d_tag": f"{resolved_provider_id}:usage:latest", + "payload": latest_payload, + }, + { + "period_type": "day", + "period_key": day_key, + "d_tag": f"{resolved_provider_id}:usage:day:{day_key}", + "payload": day_payload, + }, + { + "period_type": "month", + "period_key": month_key, + "d_tag": f"{resolved_provider_id}:usage:month:{month_key}", + "payload": month_payload, + }, + ] + + to_publish: list[dict[str, Any]] = [] + refs: dict[str, dict[str, str]] = {} + for spec in payload_specs: + payload = spec["payload"] + payload_hash = _fingerprint_payload(payload) + d_tag = str(spec["d_tag"]) + period_type = str(spec["period_type"]) + + refs[period_type] = {"d": d_tag, "payload_hash": payload_hash} + last_state = last_period_state.get(period_type) + if last_state is not None and last_state[0] == d_tag and last_state[1] == payload_hash: + continue + + payload_json = json.dumps(payload, separators=(",", ":"), sort_keys=True) + event = create_usage_analytics_event( + private_key_hex, + resolved_provider_id, + payload_json, + period_type=period_type, + period_key=str(spec["period_key"]), + d_tag=d_tag, + ) + to_publish.append( + { + "period_type": period_type, + "d_tag": d_tag, + "payload_hash": payload_hash, + "event": event, + } + ) + + period_attempted = {str(item["period_type"]) for item in to_publish} + period_successes = {period_type: 0 for period_type in period_attempted} + if to_publish: + for relay_url in relay_urls: + for item in to_publish: + if await publish_to_relay(relay_url, item["event"]): + period_successes[item["period_type"]] += 1 + + for item in to_publish: + period_type = item["period_type"] + if period_successes.get(period_type, 0) > 0: + last_period_state[period_type] = ( + item["d_tag"], + item["payload_hash"], + ) + + if checkpoint_day is None: + checkpoint_day = day_key + elif checkpoint_day != day_key: + if checkpoint_hash_for_day: + previous_checkpoint_hash = checkpoint_hash_for_day + checkpoint_day = day_key + checkpoint_hash_for_day = None + last_checkpoint_state = None + + checkpoint_payload = build_analytics_checkpoint_payload( + resolved_provider_id, + public_key_hex=public_key_hex, + generated_at=now_ts, + day_utc=day_key, + refs=refs, + previous_checkpoint_hash=previous_checkpoint_hash, + ) + checkpoint_d = f"{resolved_provider_id}:usage:checkpoint:{day_key}" + checkpoint_hash = _fingerprint_payload(checkpoint_payload) + + checkpoint_attempted = False + checkpoint_success_count = 0 + if ( + last_checkpoint_state is None + or last_checkpoint_state[0] != checkpoint_d + or last_checkpoint_state[1] != checkpoint_hash + ): + checkpoint_attempted = True + checkpoint_payload_json = json.dumps( + checkpoint_payload, + separators=(",", ":"), + sort_keys=True, + ) + checkpoint_event = create_analytics_checkpoint_event( + private_key_hex, + resolved_provider_id, + checkpoint_payload_json, + day_utc=day_key, + previous_checkpoint_hash=previous_checkpoint_hash, + ) + for relay_url in relay_urls: + if await publish_to_relay(relay_url, checkpoint_event): + checkpoint_success_count += 1 + + if checkpoint_success_count > 0: + last_checkpoint_state = (checkpoint_d, checkpoint_hash) + checkpoint_hash_for_day = str( + checkpoint_payload.get("checkpoint_hash", "") + ) or None + + relay_total = len(relay_urls) + latest_result = ( + f"{period_successes.get('latest', 0)}/{relay_total}" + if "latest" in period_attempted + else "skip" + ) + day_result = ( + f"{period_successes.get('day', 0)}/{relay_total}" + if "day" in period_attempted + else "skip" + ) + month_result = ( + f"{period_successes.get('month', 0)}/{relay_total}" + if "month" in period_attempted + else "skip" + ) + checkpoint_result = ( + f"{checkpoint_success_count}/{relay_total}" + if checkpoint_attempted + else "skip" + ) + logger.info( + "Published analytics snapshots " + "(latest=%s day=%s month=%s checkpoint=%s day_utc=%s month_utc=%s)", + latest_result, + day_result, + month_result, + checkpoint_result, + day_key, + month_key, + extra={ + "latest_relays": period_successes.get("latest", 0), + "day_relays": period_successes.get("day", 0), + "month_relays": period_successes.get("month", 0), + "checkpoint_relays": checkpoint_success_count, + "relay_total": relay_total, + "day": day_key, + "month": month_key, + }, + ) + await asyncio.sleep(PUBLISH_INTERVAL_SECONDS) + + except asyncio.CancelledError: + logger.info("Usage analytics sharing task cancelled") + break + except Exception as e: + logger.error( + "Usage analytics sharing error", + extra={"error": str(e), "error_type": type(e).__name__}, + ) + await asyncio.sleep(DISABLED_POLL_SECONDS) diff --git a/tests/unit/test_nostr_analytics.py b/tests/unit/test_nostr_analytics.py new file mode 100644 index 00000000..adc0248d --- /dev/null +++ b/tests/unit/test_nostr_analytics.py @@ -0,0 +1,230 @@ +from __future__ import annotations + +from typing import Any + +from routstr.nostr import analytics + + +def test_aggregate_top_model_usage_sums_metrics() -> None: + model_usage_mix = { + "top_models": ["openai/gpt-4o", "anthropic/claude-3.5-sonnet"], + "metrics": [ + { + "model_counts": { + "openai/gpt-4o": 4, + "anthropic/claude-3.5-sonnet": 2, + }, + "model_revenue_msats": { + "openai/gpt-4o": 1500, + "anthropic/claude-3.5-sonnet": 700, + }, + "model_tokens": { + "openai/gpt-4o": 1200, + "anthropic/claude-3.5-sonnet": 600, + }, + "others": 1, + "others_revenue_msats": 300, + "others_tokens": 200, + }, + { + "model_counts": { + "openai/gpt-4o": 3, + "anthropic/claude-3.5-sonnet": 1, + }, + "model_revenue_msats": { + "openai/gpt-4o": 1000, + "anthropic/claude-3.5-sonnet": 500, + }, + "model_tokens": { + "openai/gpt-4o": 800, + "anthropic/claude-3.5-sonnet": 300, + }, + "others": 2, + "others_revenue_msats": 450, + "others_tokens": 350, + }, + ], + } + + rows, others = analytics._aggregate_top_model_usage(model_usage_mix) + assert rows == [ + { + "model": "openai/gpt-4o", + "successful_requests": 7, + "revenue_msats": 2500.0, + "total_tokens": 2000, + }, + { + "model": "anthropic/claude-3.5-sonnet", + "successful_requests": 3, + "revenue_msats": 1200.0, + "total_tokens": 900, + }, + ] + assert others == { + "successful_requests": 3, + "revenue_msats": 750.0, + "total_tokens": 550, + } + + +def test_build_latest_payload_contains_windows_and_v2_schema(monkeypatch: Any) -> None: + seen_windows: set[tuple[int, int]] = set() + + def fake_usage_dashboard( + *, interval: int, hours: int, error_limit: int, model_limit: int + ) -> dict[str, Any]: + seen_windows.add((hours, interval)) + assert error_limit == 1 + assert model_limit == 20 + return { + "summary": { + "total_requests": 20, + "successful_chat_completions": 18, + "failed_requests": 2, + "success_rate": 90.0, + "unique_models_count": 2, + "input_tokens": 2000, + "output_tokens": 1000, + "total_tokens": 3000, + "revenue_msats": 9000.0, + "refunds_msats": 1000.0, + "net_revenue_msats": 8000.0, + "revenue_sats": 9.0, + "refunds_sats": 1.0, + "net_revenue_sats": 8.0, + }, + "revenue_by_model": { + "models": [ + { + "model": "openai/gpt-4o", + "requests": 15, + "successful": 14, + "failed": 1, + "revenue_sats": 7.2, + "refunds_sats": 0.3, + "net_revenue_sats": 6.9, + } + ] + }, + "model_usage_mix": { + "top_models": ["openai/gpt-4o"], + "metrics": [ + { + "timestamp": "2026-03-02 10:00:00", + "model_counts": {"openai/gpt-4o": 14}, + "model_revenue_msats": {"openai/gpt-4o": 7200.0}, + "model_tokens": {"openai/gpt-4o": 2600}, + "others": 4, + "others_revenue_msats": 1800.0, + "others_tokens": 400, + } + ], + }, + } + + monkeypatch.setattr( + analytics.log_manager, "get_usage_dashboard", fake_usage_dashboard + ) + monkeypatch.setattr(analytics.settings, "npub", "npub1example") + monkeypatch.setattr(analytics.settings, "http_url", "https://node.example.com") + monkeypatch.setattr(analytics.settings, "onion_url", "") + + payload = analytics.build_latest_usage_analytics_payload( + "provider123", + public_key_hex="ab" * 32, + generated_at=1772451600, + model_limit=20, + ) + assert seen_windows == {(24, 60), (7 * 24, 6 * 60), (30 * 24, 24 * 60)} + assert payload["schema"] == analytics.ANALYTICS_SCHEMA + assert payload["provider_id"] == "provider123" + assert payload["period_type"] == "latest" + assert payload["period_key"] == "latest" + assert payload["endpoint_urls"] == ["https://node.example.com"] + assert set(payload["windows"].keys()) == {"24h", "7d", "30d"} + + +def test_day_and_month_payload_keys(monkeypatch: Any) -> None: + def fake_usage_dashboard( + *, interval: int, hours: int, error_limit: int, model_limit: int + ) -> dict[str, Any]: + _ = (error_limit, model_limit) + return { + "summary": { + "total_requests": max(1, hours), + "successful_chat_completions": max(1, hours), + "failed_requests": 0, + "total_tokens": max(1, hours) * 100, + "revenue_sats": float(max(1, hours)), + }, + "revenue_by_model": {"models": []}, + "model_usage_mix": {"top_models": [], "metrics": []}, + } + + monkeypatch.setattr( + analytics.log_manager, "get_usage_dashboard", fake_usage_dashboard + ) + monkeypatch.setattr(analytics.settings, "npub", "npub1example") + monkeypatch.setattr(analytics.settings, "http_url", "https://node.example.com") + monkeypatch.setattr(analytics.settings, "onion_url", "") + + generated_at = 1772451600 # 2026-03-02 + day_payload = analytics.build_day_usage_analytics_payload( + "provider123", + public_key_hex="ab" * 32, + generated_at=generated_at, + ) + month_payload = analytics.build_month_usage_analytics_payload( + "provider123", + public_key_hex="ab" * 32, + generated_at=generated_at, + ) + + assert day_payload["period_type"] == "day" + assert day_payload["period_key"] == "2026-03-02" + assert day_payload["day"] == "2026-03-02" + assert month_payload["period_type"] == "month" + assert month_payload["period_key"] == "2026-03" + assert month_payload["month"] == "2026-03" + + +def test_create_usage_analytics_event_tags() -> None: + private_key_hex = "11" * 32 + event = analytics.create_usage_analytics_event( + private_key_hex, + "provider123", + payload_json='{"schema":"routstr.analytics.usage.v2"}', + period_type="day", + period_key="2026-03-02", + d_tag="provider123:usage:day:2026-03-02", + ) + + tags = event["tags"] + assert ["d", "provider123:usage:day:2026-03-02"] in tags + assert ["provider", "provider123"] in tags + assert ["schema", analytics.ANALYTICS_SCHEMA] in tags + assert ["period", "day"] in tags + assert ["period_key", "2026-03-02"] in tags + assert ["day", "2026-03-02"] in tags + + +def test_checkpoint_payload_contains_chain_hash() -> None: + payload = analytics.build_analytics_checkpoint_payload( + "provider123", + public_key_hex="ab" * 32, + generated_at=1772451600, + day_utc="2026-03-02", + refs={ + "latest": {"d": "provider123:usage:latest", "payload_hash": "a"}, + "day": {"d": "provider123:usage:day:2026-03-02", "payload_hash": "b"}, + "month": {"d": "provider123:usage:month:2026-03", "payload_hash": "c"}, + }, + previous_checkpoint_hash="prev-hash", + ) + + assert payload["schema"] == analytics.ANALYTICS_CHECKPOINT_SCHEMA + assert payload["day_utc"] == "2026-03-02" + assert payload["previous_checkpoint_hash"] == "prev-hash" + assert isinstance(payload["checkpoint_hash"], str) + assert len(payload["checkpoint_hash"]) == 64 diff --git a/tests/unit/test_settings.py b/tests/unit/test_settings.py index 5c5cb048..770af976 100644 --- a/tests/unit/test_settings.py +++ b/tests/unit/test_settings.py @@ -2,6 +2,7 @@ import os import pytest from sqlalchemy.ext.asyncio import create_async_engine +from sqlmodel import text from sqlmodel.ext.asyncio.session import AsyncSession from routstr.core.settings import SettingsService @@ -11,6 +12,7 @@ from routstr.core.settings import SettingsService async def test_settings_seed_from_env_and_persist() -> None: os.environ["UPSTREAM_BASE_URL"] = "https://api.test/v1" os.environ.pop("ONION_URL", None) + os.environ.pop("ENABLE_ANALYTICS_SHARING", None) engine = create_async_engine("sqlite+aiosqlite:///:memory:") async with AsyncSession(engine, expire_on_commit=False) as session: @@ -19,19 +21,53 @@ async def test_settings_seed_from_env_and_persist() -> None: assert settings.upstream_base_url == "https://api.test/v1" # ONION_URL may be empty if not discoverable assert isinstance(settings.onion_url, str) + assert settings.enable_analytics_sharing is True @pytest.mark.asyncio async def test_settings_db_precedence_over_env() -> None: os.environ["UPSTREAM_BASE_URL"] = "https://api.env/v1" + os.environ["ENABLE_ANALYTICS_SHARING"] = "true" engine = create_async_engine("sqlite+aiosqlite:///:memory:") async with AsyncSession(engine, expire_on_commit=False) as session: _ = await SettingsService.initialize(session) - updated = await SettingsService.update({"name": "DBName"}, session) + updated = await SettingsService.update( + {"name": "DBName", "enable_analytics_sharing": False}, session + ) assert updated.name == "DBName" + assert updated.enable_analytics_sharing is False # Change env and re-initialize; DB should still win os.environ["NAME"] = "EnvName" + os.environ["ENABLE_ANALYTICS_SHARING"] = "true" again = await SettingsService.initialize(session) assert again.name == "DBName" + assert again.enable_analytics_sharing is False + + +@pytest.mark.asyncio +async def test_settings_initialize_discards_unknown_keys() -> None: + engine = create_async_engine("sqlite+aiosqlite:///:memory:") + async with AsyncSession(engine, expire_on_commit=False) as session: + _ = await SettingsService.initialize(session) + + # Simulate older persisted key name and an unknown key. + await session.exec( # type: ignore + text( + "UPDATE settings SET data = :data WHERE id = 1" + ).bindparams( + data='{"name":"LegacyNode","nostr_analytics_enabled":false,"unknown_key":123}' + ) + ) + await session.commit() + + reloaded = await SettingsService.initialize(session) + assert reloaded.name == "LegacyNode" + assert reloaded.enable_analytics_sharing is True + + row = await session.exec(text("SELECT data FROM settings WHERE id = 1")) # type: ignore + stored_data = row.first()[0] + assert '"enable_analytics_sharing": true' in stored_data + assert "nostr_analytics_enabled" not in stored_data + assert "unknown_key" not in stored_data diff --git a/ui/components/settings/admin-settings.tsx b/ui/components/settings/admin-settings.tsx index ceab6128..30e0f79d 100644 --- a/ui/components/settings/admin-settings.tsx +++ b/ui/components/settings/admin-settings.tsx @@ -26,6 +26,7 @@ interface SettingsData { description?: string; npub?: string; nsec?: string; + enable_analytics_sharing?: boolean; upstream_api_key?: string; http_url?: string; onion_url?: string; @@ -43,6 +44,7 @@ const HANDLED_KEYS = [ 'nsec', 'cashu_mints', 'relays', + 'enable_analytics_sharing', 'admin_password', 'id', 'updated_at', @@ -366,6 +368,7 @@ export function AdminSettings() { const nostrChanged = ['npub', 'nsec'].some(hasFieldChanged); const cashuMintsChanged = hasFieldChanged('cashu_mints'); const relaysChanged = hasFieldChanged('relays'); + const analyticsSharingChanged = hasFieldChanged('enable_analytics_sharing'); const advancedKeys = Object.keys(settings).filter( (key) => !HANDLED_KEYS.includes(key) && !IGNORED_KEYS.includes(key) ); @@ -397,6 +400,7 @@ export function AdminSettings() { resetFields(['relays']); setNewRelay(''); }; + const resetAnalyticsSharing = () => resetFields(['enable_analytics_sharing']); const resetAdvanced = () => resetFields(advancedKeys); if (loading) { @@ -686,6 +690,52 @@ export function AdminSettings() { ) : null} + {/* Analytics Sharing */} + + + Analytics Sharing + + Publish aggregate usage stats to Nostr for external dashboards + + + +
+
+ +

+ When enabled, Routstr periodically publishes aggregate model + usage and revenue stats. +

+
+ + handleInputChange('enable_analytics_sharing', checked) + } + /> +
+
+ {analyticsSharingChanged ? ( + +
+ + +
+
+ ) : null} +
+ {/* Other Settings */} diff --git a/ui/components/top-models-usage-chart.tsx b/ui/components/top-models-usage-chart.tsx index 87474803..26eb5330 100644 --- a/ui/components/top-models-usage-chart.tsx +++ b/ui/components/top-models-usage-chart.tsx @@ -218,11 +218,11 @@ export function TopModelsUsageChart({ ); const chartModels = useMemo( - () => mixTopModels.slice(0, 10), + () => mixTopModels.slice(0, 20), [mixTopModels] ); const leaderboardModels = useMemo( - () => mixTopModels.slice(0, 10), + () => mixTopModels.slice(0, 20), [mixTopModels] ); const revenueDisplayUnit: DisplayUnit = useMemo(() => { @@ -496,7 +496,7 @@ export function TopModelsUsageChart({ }) .filter((row) => row.totalRaw > 0) .sort((a, b) => b.totalRaw - a.totalRaw) - .slice(0, 10) + .slice(0, 20) .map((row, index) => ({ ...row, rank: index + 1,