From 9cd4ff5c2148a2c4d930c8d9cd497672831c21ad Mon Sep 17 00:00:00 2001 From: Evan Yang Date: Thu, 5 Mar 2026 15:50:54 +0800 Subject: [PATCH] Simplify analytics sharing to single snapshot with multi-window payloads --- routstr/core/main.py | 2 +- routstr/nostr/analytics.py | 508 +++++------------------------ tests/unit/test_nostr_analytics.py | 231 +++++++------ 3 files changed, 208 insertions(+), 533 deletions(-) diff --git a/routstr/core/main.py b/routstr/core/main.py index d0769706..c906204e 100644 --- a/routstr/core/main.py +++ b/routstr/core/main.py @@ -108,7 +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()) + 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()) diff --git a/routstr/nostr/analytics.py b/routstr/nostr/analytics.py index 1c33f8ff..eb148664 100644 --- a/routstr/nostr/analytics.py +++ b/routstr/nostr/analytics.py @@ -1,7 +1,7 @@ #!/usr/bin/env python3 """ Nostr usage analytics publisher. -Publishes routstr analytics snapshots for latest/day/month plus daily checkpoints. +Publishes a single replaceable analytics snapshot for each provider. """ from __future__ import annotations @@ -9,10 +9,8 @@ from __future__ import annotations import asyncio import hashlib import json -import math import time -from datetime import datetime, timezone -from typing import Any, TypedDict +from typing import Any from nostr.event import Event from nostr.key import PrivateKey @@ -25,8 +23,7 @@ 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" +ANALYTICS_SCHEMA = "routstr.analytics.snapshot.v1" DEFAULT_RELAYS = [ "wss://relay.nostr.band", "wss://relay.damus.io", @@ -37,22 +34,16 @@ PUBLISH_INTERVAL_SECONDS = 15 * 60 DISABLED_POLL_SECONDS = 60 DASHBOARD_WINDOW_HOURS = 24 DASHBOARD_INTERVAL_MINUTES = 60 -MODEL_LIMIT = 20 - +MODEL_LIMIT = 10 WINDOW_DEFINITIONS: tuple[tuple[str, int, int], ...] = ( ("24h", 24, 60), ("7d", 7 * 24, 6 * 60), ("30d", 30 * 24, 24 * 60), + ("3m", 90 * 24, 24 * 60), + ("1y", 365 * 24, 7 * 24 * 60), ) -class PayloadSpec(TypedDict): - period_type: str - period_key: str - d_tag: str - payload: dict[str, Any] - - def _event_to_dict(ev: Event) -> dict[str, Any]: return { "id": ev.id, @@ -123,31 +114,6 @@ def _to_float(value: Any) -> float: 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]]: @@ -233,80 +199,76 @@ def _build_summary_payload(summary: dict[str, Any]) -> dict[str, Any]: } -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, + interval_minutes: int, model_limit: int, ) -> dict[str, Any]: dashboard = log_manager.get_usage_dashboard( - interval=interval, + interval=interval_minutes, 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, + "interval_minutes": interval_minutes, "summary": summary_payload, - "model_revenue": model_revenue_rows, + "model_usage_mix": usage_mix_payload, "top_model_usage": top_model_usage, "others_usage": others_usage, - "model_usage_mix": usage_mix_payload, } -def build_latest_usage_analytics_payload( +def build_stats_snapshot_payload( provider_id: str, *, public_key_hex: str, generated_at: int, + window_hours: int = DASHBOARD_WINDOW_HOURS, + interval_minutes: int = DASHBOARD_INTERVAL_MINUTES, model_limit: int = MODEL_LIMIT, ) -> dict[str, Any]: + _ = (window_hours, interval_minutes) windows: dict[str, dict[str, Any]] = {} - for key, window_hours, window_interval in WINDOW_DEFINITIONS: + for key, hours, window_interval_minutes in WINDOW_DEFINITIONS: windows[key] = _build_window_payload( - hours=window_hours, - interval=window_interval, + hours=hours, + interval_minutes=window_interval_minutes, model_limit=model_limit, ) primary_window = windows.get("24h", {}) + summary_payload = ( + primary_window.get("summary", {}) + if isinstance(primary_window.get("summary", {}), dict) + else {} + ) + usage_mix_payload = ( + primary_window.get("model_usage_mix", {}) + if isinstance(primary_window.get("model_usage_mix", {}), dict) + else {} + ) + top_model_usage = ( + primary_window.get("top_model_usage", []) + if isinstance(primary_window.get("top_model_usage", []), list) + else [] + ) + others_usage = ( + primary_window.get("others_usage", {}) + if isinstance(primary_window.get("others_usage", {}), dict) + else {} + ) + return { "schema": ANALYTICS_SCHEMA, "generated_at": generated_at, @@ -314,128 +276,21 @@ def build_latest_usage_analytics_payload( "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", {}), + "window_hours": DASHBOARD_WINDOW_HOURS, + "interval_minutes": DASHBOARD_INTERVAL_MINUTES, + "summary": summary_payload, + "model_usage_mix": usage_mix_payload, + "top_model_usage": top_model_usage, + "others_usage": others_usage, "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( +def create_stats_snapshot_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)) @@ -443,13 +298,7 @@ def create_usage_analytics_event( ["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(), @@ -463,83 +312,14 @@ def create_usage_analytics_event( def _fingerprint_payload(payload: dict[str, Any]) -> str: normalized = dict(payload) - # Ignore volatile timestamps for deduping semantically identical snapshots. + # Ignore generated timestamp for semantic dedupe. 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 + last_payload_hash: str | None = None parsed_nsec: str | None = None private_key_hex: str | None = None @@ -573,11 +353,7 @@ async def publish_usage_analytics() -> None: 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 + last_payload_hash = None if private_key_hex is None or public_key_hex is None: await asyncio.sleep(DISABLED_POLL_SECONDS) @@ -591,181 +367,43 @@ async def publish_usage_analytics() -> None: 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( + payload = build_stats_snapshot_payload( resolved_provider_id, public_key_hex=public_key_hex, generated_at=now_ts, ) - payload_specs: list[PayloadSpec] = [ - { - "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, - }, - ] + payload_hash = _fingerprint_payload(payload) + if last_payload_hash == payload_hash: + await asyncio.sleep(PUBLISH_INTERVAL_SECONDS) + continue - to_publish: list[dict[str, Any]] = [] - refs: dict[str, dict[str, str]] = {} - for spec in payload_specs: - payload: dict[str, Any] = 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( + payload_json = json.dumps(payload, separators=(",", ":"), sort_keys=True) + d_tag = f"{resolved_provider_id}:stats" + event = create_stats_snapshot_event( + private_key_hex, 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, + payload_json, + d_tag=d_tag, ) - 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 + success_count = 0 + for relay_url in relay_urls: + if await publish_to_relay(relay_url, event): + 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 + if success_count > 0: + last_payload_hash = payload_hash - 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, + "Published analytics snapshot (success=%s/%s provider=%s)", + success_count, + len(relay_urls), + resolved_provider_id, 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, + "relay_success_count": success_count, + "relay_total": len(relay_urls), + "provider_id": resolved_provider_id, }, ) await asyncio.sleep(PUBLISH_INTERVAL_SECONDS) diff --git a/tests/unit/test_nostr_analytics.py b/tests/unit/test_nostr_analytics.py index adc0248d..577d5ca0 100644 --- a/tests/unit/test_nostr_analytics.py +++ b/tests/unit/test_nostr_analytics.py @@ -1,7 +1,10 @@ from __future__ import annotations +import asyncio from typing import Any +import pytest + from routstr.nostr import analytics @@ -68,7 +71,7 @@ def test_aggregate_top_model_usage_sums_metrics() -> None: } -def test_build_latest_payload_contains_windows_and_v2_schema(monkeypatch: Any) -> None: +def test_build_stats_snapshot_payload_schema_and_shape(monkeypatch: Any) -> None: seen_windows: set[tuple[int, int]] = set() def fake_usage_dashboard( @@ -76,11 +79,11 @@ def test_build_latest_payload_contains_windows_and_v2_schema(monkeypatch: Any) - ) -> dict[str, Any]: seen_windows.add((hours, interval)) assert error_limit == 1 - assert model_limit == 20 + assert model_limit == 10 return { "summary": { - "total_requests": 20, - "successful_chat_completions": 18, + "total_requests": hours, + "successful_chat_completions": max(1, hours - 1), "failed_requests": 2, "success_rate": 90.0, "unique_models_count": 2, @@ -94,27 +97,14 @@ def test_build_latest_payload_contains_windows_and_v2_schema(monkeypatch: Any) - "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}, + "model_counts": {"openai/gpt-4o": hours}, + "model_revenue_msats": {"openai/gpt-4o": float(hours * 100)}, + "model_tokens": {"openai/gpt-4o": hours * 10}, "others": 4, "others_revenue_msats": 1800.0, "others_tokens": 400, @@ -130,101 +120,148 @@ def test_build_latest_payload_contains_windows_and_v2_schema(monkeypatch: Any) - monkeypatch.setattr(analytics.settings, "http_url", "https://node.example.com") monkeypatch.setattr(analytics.settings, "onion_url", "") - payload = analytics.build_latest_usage_analytics_payload( + payload = analytics.build_stats_snapshot_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["window_hours"] == 24 + assert payload["interval_minutes"] == 60 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": []}, + assert seen_windows == { + (24, 60), + (7 * 24, 6 * 60), + (30 * 24, 24 * 60), + (90 * 24, 24 * 60), + (365 * 24, 7 * 24 * 60), + } + assert set(payload["windows"].keys()) == {"24h", "7d", "30d", "3m", "1y"} + assert payload["windows"]["1y"]["interval_minutes"] == 7 * 24 * 60 + assert payload["summary"]["total_requests"] == 24 + assert payload["top_model_usage"] == [ + { + "model": "openai/gpt-4o", + "successful_requests": 24, + "revenue_msats": 2400.0, + "total_tokens": 240, } - - 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" + ] + assert payload["others_usage"] == { + "successful_requests": 4, + "revenue_msats": 1800.0, + "total_tokens": 400, + } -def test_create_usage_analytics_event_tags() -> None: +def test_create_stats_snapshot_event_tags() -> None: private_key_hex = "11" * 32 - event = analytics.create_usage_analytics_event( + event = analytics.create_stats_snapshot_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", + payload_json='{"schema":"routstr.analytics.snapshot.v1"}', + d_tag="provider123:stats", ) tags = event["tags"] - assert ["d", "provider123:usage:day:2026-03-02"] in tags + assert ["d", "provider123:stats"] 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 + assert all(tag[0] != "period" for tag 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", - ) +def test_fingerprint_payload_ignores_generated_at() -> None: + a = {"schema": analytics.ANALYTICS_SCHEMA, "generated_at": 1000, "summary": {"x": 1}} + b = {"schema": analytics.ANALYTICS_SCHEMA, "generated_at": 2000, "summary": {"x": 1}} - 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 + assert analytics._fingerprint_payload(a) == analytics._fingerprint_payload(b) + + +@pytest.mark.asyncio +async def test_publish_usage_analytics_skips_when_disabled(monkeypatch: Any) -> None: + delays: list[int] = [] + + async def fake_sleep(seconds: int) -> None: + delays.append(seconds) + raise asyncio.CancelledError() + + def fail_build(*args: Any, **kwargs: Any) -> dict[str, Any]: + raise AssertionError("build_stats_snapshot_payload should not be called") + + monkeypatch.setattr(analytics.settings, "enable_analytics_sharing", False) + monkeypatch.setattr(analytics, "build_stats_snapshot_payload", fail_build) + monkeypatch.setattr(analytics.asyncio, "sleep", fake_sleep) + + await analytics.publish_usage_analytics() + + assert delays == [analytics.DISABLED_POLL_SECONDS] + + +@pytest.mark.asyncio +async def test_publish_usage_analytics_skips_without_nsec(monkeypatch: Any) -> None: + delays: list[int] = [] + + async def fake_sleep(seconds: int) -> None: + delays.append(seconds) + raise asyncio.CancelledError() + + def fail_build(*args: Any, **kwargs: Any) -> dict[str, Any]: + raise AssertionError("build_stats_snapshot_payload should not be called") + + monkeypatch.setattr(analytics.settings, "enable_analytics_sharing", True) + monkeypatch.setattr(analytics.settings, "nsec", "") + monkeypatch.setattr(analytics, "build_stats_snapshot_payload", fail_build) + monkeypatch.setattr(analytics.asyncio, "sleep", fake_sleep) + + await analytics.publish_usage_analytics() + + assert delays == [analytics.DISABLED_POLL_SECONDS] + + +@pytest.mark.asyncio +async def test_publish_usage_analytics_dedupes_unchanged_payload(monkeypatch: Any) -> None: + published_events: list[dict[str, Any]] = [] + sleep_calls = 0 + + async def fake_sleep(seconds: int) -> None: + nonlocal sleep_calls + sleep_calls += 1 + if sleep_calls >= 2: + raise asyncio.CancelledError() + + def fake_build_payload( + provider_id: str, + *, + public_key_hex: str, + generated_at: int, + window_hours: int = 24, + interval_minutes: int = 60, + model_limit: int = 10, + ) -> dict[str, Any]: + _ = (public_key_hex, generated_at, window_hours, interval_minutes, model_limit) + return { + "schema": analytics.ANALYTICS_SCHEMA, + "generated_at": generated_at, + "provider_id": provider_id, + "summary": {"total_requests": 1}, + } + + async def fake_publish(relay_url: str, event: dict[str, Any]) -> bool: + _ = relay_url + published_events.append(event) + return True + + monkeypatch.setattr(analytics.settings, "enable_analytics_sharing", True) + monkeypatch.setattr(analytics.settings, "nsec", "11" * 32) + monkeypatch.setattr(analytics.settings, "relays", ["wss://relay.example.com"]) + monkeypatch.setattr(analytics.settings, "provider_id", "") + monkeypatch.setattr(analytics, "build_stats_snapshot_payload", fake_build_payload) + monkeypatch.setattr(analytics, "publish_to_relay", fake_publish) + monkeypatch.setattr(analytics.asyncio, "sleep", fake_sleep) + + await analytics.publish_usage_analytics() + + assert len(published_events) == 1 + assert ["schema", analytics.ANALYTICS_SCHEMA] in published_events[0].get("tags", [])