mirror of
https://github.com/Routstr/routstr-core.git
synced 2026-08-09 02:54:37 +00:00
Simplify analytics sharing to single snapshot with multi-window payloads
This commit is contained in:
@@ -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())
|
||||
|
||||
+73
-435
@@ -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)
|
||||
|
||||
@@ -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", [])
|
||||
|
||||
Reference in New Issue
Block a user